1 /* $NetBSD: work.c,v 1.5 2026/08/29 14:55:18 christos Exp $ */ 2 3 /* 4 * Copyright (C) Internet Systems Consortium, Inc. ("ISC") 5 * 6 * SPDX-License-Identifier: MPL-2.0 7 * 8 * This Source Code Form is subject to the terms of the Mozilla Public 9 * License, v. 2.0. If a copy of the MPL was not distributed with this 10 * file, you can obtain one at https://mozilla.org/MPL/2.0/. 11 * 12 * See the COPYRIGHT file distributed with this work for additional 13 * information regarding copyright ownership. 14 */ 15 16 #include <limits.h> 17 #include <stddef.h> 18 #include <stdint.h> 19 20 #include <isc/async.h> 21 #include <isc/job.h> 22 #include <isc/loop.h> 23 #include <isc/magic.h> 24 #include <isc/queue.h> 25 #include <isc/thread.h> 26 #include <isc/urcu.h> 27 #include <isc/util.h> 28 #include <isc/uv.h> 29 #include <isc/work.h> 30 31 #include "loop_p.h" 32 33 #define WORK_MAGIC ISC_MAGIC('W', 'o', 'r', 'k') 34 #define VALID_WORK(t) ISC_MAGIC_VALID(t, WORK_MAGIC) 35 #define WORKTHREAD_MAGIC ISC_MAGIC('W', 'k', 'T', 'h') 36 #define VALID_WORKTHREAD(t) ISC_MAGIC_VALID(t, WORKTHREAD_MAGIC) 37 38 enum waitstate { 39 /* The value a sleeping worker blocks on in FUTEX_WAIT. */ 40 THREAD_WAITING = 0, 41 /* Any non-zero bit keeps FUTEX_WAIT from blocking. */ 42 THREAD_WAKEUP = (1 << 0), 43 THREAD_RUNNING = (1 << 1), 44 THREAD_SHUTDOWN = (1 << 2), 45 THREAD_PAUSE = (1 << 3), /* request from the owning loop */ 46 THREAD_PAUSED = (1 << 4), /* ack from the worker */ 47 }; 48 49 /* Sticky bits a paused worker must not drop. */ 50 #define THREAD_STICKY (THREAD_SHUTDOWN | THREAD_PAUSE | THREAD_PAUSED) 51 52 enum workstate { 53 WORK_QUEUED = 0, 54 WORK_RUNNING, 55 WORK_CANCELED, 56 }; 57 58 struct isc_work { 59 unsigned int magic; 60 uint32_t state; /* enum workstate */ 61 isc_result_t result; 62 isc_work_cb cb; /* runs on a worker thread */ 63 isc_work_done_cb done_cb; /* runs on the origin loop */ 64 void *cbarg; 65 isc_loop_t *loop; /* origin loop, referenced */ 66 struct cds_wfcq_node node; /* dispatch queue linkage */ 67 }; 68 69 typedef struct isc__workthread { 70 union { 71 struct { 72 unsigned int magic; 73 isc_worklane_t lane; 74 isc_loop_t *loop; 75 isc_thread_t thread; 76 struct __cds_wfcq_head qhead; 77 int32_t state; /* enum waitstate */ 78 }; 79 uint8_t __padding0[ISC_OS_CACHELINE_SIZE]; 80 }; 81 union { 82 struct cds_wfcq_tail qtail; 83 uint8_t __padding1[ISC_OS_CACHELINE_SIZE]; 84 }; 85 } isc__workthread_t; 86 87 STATIC_ASSERT(ISC_OS_CACHELINE_SIZE >= sizeof(struct cds_wfcq_tail), 88 "ISC_OS_CACHELINE_SIZE smaller than sizeof(struct " 89 "cds_wfcq_tail)"); 90 STATIC_ASSERT(offsetof(isc__workthread_t, qtail) == ISC_OS_CACHELINE_SIZE, 91 "isc__workthread_t.qtail not on second cacheline"); 92 STATIC_ASSERT(sizeof(isc__workthread_t) == 2 * ISC_OS_CACHELINE_SIZE, 93 "isc__workthread_t is not two cachelines"); 94 95 static void 96 workthread_wake(isc__workthread_t *thread) { 97 cmm_smp_mb(); 98 if ((uatomic_load(&thread->state, CMM_RELAXED) & THREAD_RUNNING) != 0) { 99 /* Actively running; it will notice the queue on its own. */ 100 return; 101 } 102 103 uatomic_or(&thread->state, THREAD_WAKEUP); 104 if (futex_noasync(&thread->state, FUTEX_WAKE, 1, NULL, NULL, 0) < 0) { 105 FATAL_ERROR("futex_noasync(FUTEX_WAKE): %s", strerror(errno)); 106 } 107 } 108 109 static void 110 workthread_slumber(isc__workthread_t *thread) { 111 rcu_thread_offline(); 112 while (futex_noasync(&thread->state, FUTEX_WAIT, THREAD_WAITING, NULL, 113 NULL, 0) != 0) 114 { 115 if (errno == EWOULDBLOCK) { 116 break; 117 } else if (errno != EINTR) { 118 FATAL_ERROR("futex_noasync(FUTEX_WAIT): %s", 119 strerror(errno)); 120 } 121 /* Or retry if interrupted by signal. */ 122 } 123 rcu_thread_online(); 124 } 125 126 static void 127 workthread_sleep(isc__workthread_t *thread) { 128 /* 129 * Drop to WAITING while keeping a pending SHUTDOWN/PAUSE sticky, so the 130 * FUTEX_WAIT below refuses to block once either is signalled. 131 */ 132 uatomic_and(&thread->state, THREAD_STICKY); 133 cmm_smp_mb(); 134 135 /* 136 * The queue is the one wake condition that can't live in 'state', so 137 * recheck it under the fence; SHUTDOWN and WAKEUP are handled by 138 * FUTEX_WAIT's own value check. 139 */ 140 if (cds_wfcq_empty(&thread->qhead, &thread->qtail)) { 141 workthread_slumber(thread); 142 } 143 144 /* Tell the waker we are running (keeping any sticky SHUTDOWN/PAUSE). */ 145 uatomic_or(&thread->state, THREAD_RUNNING); 146 } 147 148 /* 149 * Acknowledge a pause request: publish PAUSED (dropping RUNNING/WAKEUP) and 150 * wake the waiting pauser. A new pause clears PAUSED, so the worker re-acks 151 * and the pauser only ever observes an ack set for its own request, never a 152 * stale one from the previous pause generation. 153 */ 154 static void 155 workthread_ack_pause(isc__workthread_t *thread) { 156 int32_t old, next; 157 do { 158 old = uatomic_load(&thread->state, CMM_RELAXED); 159 next = (old & THREAD_STICKY) | THREAD_PAUSED; 160 } while (uatomic_cmpxchg(&thread->state, old, next) != old); 161 162 (void)futex_noasync(&thread->state, FUTEX_WAKE, INT_MAX, NULL, NULL, 0); 163 } 164 165 /* 166 * Honour a pause: pause until the owning loop clears PAUSE (resume). A fresh 167 * pause clears PAUSED (see isc__workthread_pause), so (re-)ack whenever PAUSED 168 * is gone the pauser only proceeds on an ack set for *its* request, never a 169 * stale one from the previous generation. Stays RCU-offline while paused so 170 * it can't hold up an exclusive-mode grace period. 171 */ 172 static void 173 workthread_pause(isc__workthread_t *thread) { 174 rcu_thread_offline(); 175 176 while (true) { 177 int32_t old = uatomic_load(&thread->state, CMM_ACQUIRE); 178 if ((old & (THREAD_PAUSE | THREAD_SHUTDOWN)) != THREAD_PAUSE) { 179 break; 180 } 181 if ((old & THREAD_PAUSED) == 0) { 182 workthread_ack_pause(thread); 183 continue; 184 } 185 (void)futex_noasync(&thread->state, FUTEX_WAIT, old, NULL, NULL, 186 0); 187 } 188 189 uatomic_and(&thread->state, ~THREAD_PAUSED); 190 rcu_thread_online(); 191 } 192 193 static void 194 work_done(void *arg) { 195 isc_work_t *work = arg; 196 isc_loop_t *loop = work->loop; 197 198 /* work_run() has settled work->result before scheduling us. */ 199 INSIST(work->result != ISC_R_UNSET); 200 201 work->done_cb(work->cbarg, work->result); 202 203 work->magic = 0; 204 isc_mem_put(work->loop->mctx, work, sizeof(*work)); 205 isc_loop_unref(loop); 206 } 207 208 static void 209 work_run(void *arg) { 210 isc_work_t *work = arg; 211 /* 212 * The CAS *is* the tombstone check: whoever moves the item out 213 * of WORK_QUEUED first this worker or isc_work_cancel() 214 * decides whether the callback runs. uatomic_cmpxchg returns the 215 * prior state, so WORK_QUEUED means we won the race. 216 */ 217 uint32_t prev = uatomic_cmpxchg(&work->state, WORK_QUEUED, 218 WORK_RUNNING); 219 switch (prev) { 220 case WORK_QUEUED: 221 work->result = work->cb(work->cbarg); 222 break; 223 case WORK_CANCELED: 224 work->result = ISC_R_CANCELED; 225 break; 226 default: 227 UNREACHABLE(); 228 } 229 230 /* Completion always routes back to the origin loop. */ 231 isc_async_run(work->loop, work_done, work); 232 } 233 234 static void * 235 workthread_thread(void *arg) { 236 isc__workthread_t *thread = arg; 237 238 isc__loopmgr_starting(thread->loop); 239 240 while (true) { 241 /* 242 * Honour a pause before touching the queue (gated on !SHUTDOWN 243 * so a shutting-down worker exits instead of pausing). 244 */ 245 int32_t state = uatomic_load(&thread->state, CMM_ACQUIRE); 246 if ((state & (THREAD_PAUSE | THREAD_SHUTDOWN)) == THREAD_PAUSE) 247 { 248 workthread_pause(thread); 249 continue; 250 } 251 252 struct cds_wfcq_node *node; 253 node = __cds_wfcq_dequeue_blocking(&thread->qhead, 254 &thread->qtail); 255 256 if (node == NULL) { 257 /* 258 * Only exit the loop if there's nothing to do. 259 */ 260 if ((uatomic_load(&thread->state, CMM_ACQUIRE) & 261 THREAD_SHUTDOWN) != 0) 262 { 263 synchronize_rcu(); 264 if (!cds_wfcq_empty(&thread->qhead, 265 &thread->qtail)) 266 { 267 continue; 268 } 269 break; 270 } 271 272 workthread_sleep(thread); 273 274 continue; 275 } 276 277 isc_work_t *work = caa_container_of(node, isc_work_t, node); 278 work_run(work); 279 } 280 281 isc__loopmgr_stopping(thread->loop); 282 283 return NULL; 284 } 285 286 isc_work_t * 287 isc_work_enqueue(isc_loop_t *loop, isc_worklane_t lane, isc_work_cb cb, 288 isc_work_done_cb done_cb, void *cbarg) { 289 REQUIRE(loop == isc_loop()); 290 291 isc__workthread_t *thread = isc__loopmgr_workthread(loop, lane); 292 293 isc_work_t *work = isc_mem_get(loop->mctx, sizeof(*work)); 294 *work = (isc_work_t){ 295 .magic = WORK_MAGIC, 296 .result = ISC_R_UNSET, 297 .cb = cb, 298 .done_cb = done_cb, 299 .cbarg = cbarg, 300 .loop = isc_loop_ref(loop), 301 .state = WORK_QUEUED, 302 }; 303 304 rcu_read_lock(); 305 if ((uatomic_load(&thread->state, CMM_ACQUIRE) & THREAD_SHUTDOWN) != 0) 306 { 307 rcu_read_unlock(); 308 309 /* 310 * We are shutting down, so immedaitely run task instead of 311 * adding more in the queue. (The worker is running the 312 * remaining enqueue tasks and shutdown after, see 313 * workthread_thread().) 314 */ 315 isc_async_run(loop, work_run, work); 316 } else { 317 (void)cds_wfcq_enqueue(&thread->qhead, &thread->qtail, 318 &work->node); 319 rcu_read_unlock(); 320 321 if ((uatomic_load(&thread->state, CMM_ACQUIRE) & 322 THREAD_RUNNING) == 0) 323 { 324 workthread_wake(thread); 325 } 326 } 327 328 return work; 329 } 330 331 bool 332 isc_work_cancel(isc_work_t *work) { 333 REQUIRE(VALID_WORK(work)); 334 335 /* 336 * Tombstone: QUEUED -> CANCELED. The node stays in the queue 337 * (no interior unlink in a singly-linked lock-free queue) and 338 * is discarded by whichever worker dequeues it; done_cb still 339 * fires with ISC_R_CANCELED. Nothing is freed here. False 340 * means the callback is running or done uv_cancel semantics. 341 */ 342 return uatomic_cmpxchg(&work->state, WORK_QUEUED, WORK_CANCELED) == 343 WORK_QUEUED; 344 } 345 346 isc__workthread_t * 347 isc__workthread_create(isc_loop_t *loop, isc_worklane_t lane) { 348 isc__workthread_t *thread = isc_mem_get(loop->mctx, sizeof(*thread)); 349 350 *thread = (isc__workthread_t){ 351 .lane = lane, 352 .magic = WORKTHREAD_MAGIC, 353 .state = THREAD_WAITING, 354 .loop = loop, 355 }; 356 357 __cds_wfcq_init(&thread->qhead, &thread->qtail); 358 359 isc_thread_create(workthread_thread, thread, &thread->thread); 360 361 return thread; 362 } 363 364 void 365 isc__workthread_shutdown(isc__workthread_t *thread) { 366 REQUIRE(VALID_WORKTHREAD(thread)); 367 368 /* 369 * Not called while the worker is paused by isc__workthread_pause(): 370 * shutdown callbacks run from uv loops, and loopmgr pause keeps every 371 * loop out of uv_run() until resume, so PAUSE and SHUTDOWN never 372 * coexist on a worker (the SHUTDOWN checks in the pause path are only 373 * a belt-and-braces exit if that ever changed). 374 */ 375 376 /* Set the sticky SHUTDOWN bit once; bail if already shutting down. */ 377 int32_t old; 378 do { 379 old = uatomic_load(&thread->state, CMM_RELAXED); 380 if ((old & THREAD_SHUTDOWN) != 0) { 381 return; 382 } 383 } while (uatomic_cmpxchg(&thread->state, old, old | THREAD_SHUTDOWN) != 384 old); 385 386 /* Fence in-flight enqueues (which touch the queue) before draining. */ 387 synchronize_rcu(); 388 389 workthread_wake(thread); 390 } 391 392 void 393 isc__workthread_destroy(isc__workthread_t **threadp) { 394 REQUIRE(threadp != NULL && VALID_WORKTHREAD(*threadp)); 395 isc__workthread_t *thread = MOVE_OWNERSHIP(*threadp); 396 397 isc_thread_join(thread->thread, NULL); 398 399 INSIST(cds_wfcq_empty(&thread->qhead, &thread->qtail)); 400 401 thread->magic = 0; 402 isc_mem_put(thread->loop->mctx, thread, sizeof(*thread)); 403 } 404 405 void 406 isc__workthread_pause(isc__workthread_t *thread) { 407 REQUIRE(VALID_WORKTHREAD(thread)); 408 409 /* 410 * Request a pause, but only if not already shutting down a 411 * shutting-down worker heads for the stopping barrier and must never 412 * be waited on here (that'd be a deadlock). Clearing PAUSED as we set 413 * PAUSE invalidates any ack left over from the previous generation, so 414 * the wait below can only succeed on an ack for this request. 415 */ 416 int32_t old; 417 do { 418 old = uatomic_load(&thread->state, CMM_RELAXED); 419 if ((old & THREAD_SHUTDOWN) != 0) { 420 return; 421 } 422 } while (uatomic_cmpxchg(&thread->state, old, 423 (old | THREAD_PAUSE) & ~THREAD_PAUSED) != old); 424 425 workthread_wake(thread); 426 427 /* 428 * Wait for the worker to acknowledge (PAUSED, form workthread_thread() 429 * calling workthread_pause()) or for shutdown. 430 */ 431 while (true) { 432 old = uatomic_load(&thread->state, CMM_ACQUIRE); 433 if ((old & (THREAD_PAUSED | THREAD_SHUTDOWN)) != 0) { 434 return; 435 } 436 (void)futex_noasync(&thread->state, FUTEX_WAIT, old, NULL, NULL, 437 0); 438 } 439 } 440 441 void 442 isc__workthread_resume(isc__workthread_t *thread) { 443 REQUIRE(VALID_WORKTHREAD(thread)); 444 445 /* Clear the request and wake the paused worker. */ 446 uatomic_and(&thread->state, ~THREAD_PAUSE); 447 (void)futex_noasync(&thread->state, FUTEX_WAKE, INT_MAX, NULL, NULL, 0); 448 } 449