Coverage Report

Created: 2026-09-03 06:30

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/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
}