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