Home | History | Annotate | Line # | Download | only in isc
      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