1*22649d4dSMartin Matuska // SPDX-License-Identifier: CDDL-1.0
2*22649d4dSMartin Matuska /*
3*22649d4dSMartin Matuska * This file and its contents are supplied under the terms of the
4*22649d4dSMartin Matuska * Common Development and Distribution License ("CDDL"), version 1.0.
5*22649d4dSMartin Matuska * You may only use this file in accordance with the terms of version
6*22649d4dSMartin Matuska * 1.0 of the CDDL.
7*22649d4dSMartin Matuska *
8*22649d4dSMartin Matuska * A full copy of the text of the CDDL should have accompanied this
9*22649d4dSMartin Matuska * source. A copy of the CDDL is also available via the Internet at
10*22649d4dSMartin Matuska * https://opensource.org/license/CDDL-1.0.
11*22649d4dSMartin Matuska */
12*22649d4dSMartin Matuska
13*22649d4dSMartin Matuska /*
14*22649d4dSMartin Matuska * Copyright (c) 2026 by Garth Snyder. All rights reserved.
15*22649d4dSMartin Matuska */
16*22649d4dSMartin Matuska
17*22649d4dSMartin Matuska #include <assert.h>
18*22649d4dSMartin Matuska #include <atomic.h>
19*22649d4dSMartin Matuska #include <err.h>
20*22649d4dSMartin Matuska #include <errno.h>
21*22649d4dSMartin Matuska #include <pthread.h>
22*22649d4dSMartin Matuska #include <sched.h>
23*22649d4dSMartin Matuska #include <stddef.h>
24*22649d4dSMartin Matuska #include <stdint.h>
25*22649d4dSMartin Matuska #include <stdio.h>
26*22649d4dSMartin Matuska #include <stdlib.h>
27*22649d4dSMartin Matuska #include <string.h>
28*22649d4dSMartin Matuska #include <sys/param.h>
29*22649d4dSMartin Matuska #include <sys/random.h>
30*22649d4dSMartin Matuska #include <sys/stdtypes.h>
31*22649d4dSMartin Matuska #include <sys/sysmacros.h>
32*22649d4dSMartin Matuska #include <sys/time.h>
33*22649d4dSMartin Matuska #include <time.h>
34*22649d4dSMartin Matuska #include <unistd.h>
35*22649d4dSMartin Matuska
36*22649d4dSMartin Matuska #include "zstream_queue.h"
37*22649d4dSMartin Matuska #include "zstream_util.h"
38*22649d4dSMartin Matuska
39*22649d4dSMartin Matuska #define ENQUEUE_DELAY_NSEC (100 * 1000) /* 100us */
40*22649d4dSMartin Matuska #define DISPATCH_BACKUP_NSEC (1000 * 1000) /* 1ms */
41*22649d4dSMartin Matuska
42*22649d4dSMartin Matuska #define PLENTY_OF_WORK 6 /* "Many" items to claim */
43*22649d4dSMartin Matuska #define NO_WORK 1.0E-6 /* No-work score threshold */
44*22649d4dSMartin Matuska #define DEQUEUE_SCORE_WEIGHT 0.3 /* Dequeue score relative weight */
45*22649d4dSMartin Matuska
46*22649d4dSMartin Matuska #define Q_MOD(queue, index) ((index) % (queue)->zq_params.qp_queue_length)
47*22649d4dSMartin Matuska #define Q_SLOT(queue, index) ((queue)->zq_slots[Q_MOD((queue), (index))])
48*22649d4dSMartin Matuska
49*22649d4dSMartin Matuska #define Q_FULL(queue) ((queue)->zq_ix.enqueue - (queue)->zq_ix.dequeue >= \
50*22649d4dSMartin Matuska (queue)->zq_params.qp_queue_length)
51*22649d4dSMartin Matuska
52*22649d4dSMartin Matuska /*
53*22649d4dSMartin Matuska * A zstream_queue is a ring buffer with four indexes: enqueue, claim,
54*22649d4dSMartin Matuska * complete, and dequeue, in that order. No index can move beyond its
55*22649d4dSMartin Matuska * preceding index. Every interval between indexes contains work items in a
56*22649d4dSMartin Matuska * particular state: enqueued, claimed for work, or completed. Items never
57*22649d4dSMartin Matuska * leave the ring buffer, so FIFO order is guaranteed on dequeueing.
58*22649d4dSMartin Matuska *
59*22649d4dSMartin Matuska * In concept, every index has a corresponding condition that threads can
60*22649d4dSMartin Matuska * wait on if they are interested in knowing when that index moves:
61*22649d4dSMartin Matuska * enqueued, claimed, completed, dequeued. However, the reality deviates
62*22649d4dSMartin Matuska * from this model in two ways:
63*22649d4dSMartin Matuska *
64*22649d4dSMartin Matuska * - There is no "claimed" condition, because no thread would wait on it.
65*22649d4dSMartin Matuska * Claiming and processing are one unified operation. Dequeuers await the
66*22649d4dSMartin Matuska * "completed" condition.
67*22649d4dSMartin Matuska *
68*22649d4dSMartin Matuska * - All queues share one thread pool, so idle threads are not bound to any
69*22649d4dSMartin Matuska * particular queue. Instead of having queue-specific "enqueued" conditions,
70*22649d4dSMartin Matuska * queues share a centralized dispatch system. On being awakened, worker
71*22649d4dSMartin Matuska * threads assign themselves to a queue through a scoring mechanism
72*22649d4dSMartin Matuska * described in the comments at score_queue().
73*22649d4dSMartin Matuska *
74*22649d4dSMartin Matuska * LOCKING
75*22649d4dSMartin Matuska *
76*22649d4dSMartin Matuska * There are three types of lock:
77*22649d4dSMartin Matuska *
78*22649d4dSMartin Matuska * - One global lock that gates changes to the thread pool and queue cohort.
79*22649d4dSMartin Matuska * This lock also acts as the mutex for the tp_wake_worker condition.
80*22649d4dSMartin Matuska *
81*22649d4dSMartin Matuska * - A second, low-contention global lock that protects the dispatch system
82*22649d4dSMartin Matuska *
83*22649d4dSMartin Matuska * - One lock for each queue
84*22649d4dSMartin Matuska *
85*22649d4dSMartin Matuska * Any operation that adds or removes queues or threads should hold the pool
86*22649d4dSMartin Matuska * lock. Any operation that moves a queue's indexes should hold the queue
87*22649d4dSMartin Matuska * lock. Any thread waiting for work waits on tp_wake_worker.
88*22649d4dSMartin Matuska *
89*22649d4dSMartin Matuska * Worker threads hold no locks while they are actually processing items.
90*22649d4dSMartin Matuska *
91*22649d4dSMartin Matuska * The global locking order is dispatch -> pool -> queue.
92*22649d4dSMartin Matuska *
93*22649d4dSMartin Matuska * DISPATCH
94*22649d4dSMartin Matuska *
95*22649d4dSMartin Matuska * Four events trigger dispatch loops:
96*22649d4dSMartin Matuska *
97*22649d4dSMartin Matuska * 1) A worker thread completing its batch. Threads always check to see if
98*22649d4dSMartin Matuska * there's more claimable work before going to sleep.
99*22649d4dSMartin Matuska *
100*22649d4dSMartin Matuska * 2) A worker thread discovering more work than it can handle on its own.
101*22649d4dSMartin Matuska * Before starting work on its own batch, the worker attempts to signal
102*22649d4dSMartin Matuska * another thread to wake up and assess the current state.
103*22649d4dSMartin Matuska *
104*22649d4dSMartin Matuska * 3) Enqueues. These go through the dispatch system and are batched.
105*22649d4dSMartin Matuska * Roughly ENQUEUE_DELAY_NSEC after an enqueue (on any queue), the dispatch
106*22649d4dSMartin Matuska * thread attempts to awaken a worker.
107*22649d4dSMartin Matuska *
108*22649d4dSMartin Matuska * 4) The expiration of a backup timer. The atomic value tp_unclaimed tracks
109*22649d4dSMartin Matuska * the total number of enqueued-but-unclaimed items across all queues. When
110*22649d4dSMartin Matuska * zero, it indicates that no worker dispatch is currently necessary. This
111*22649d4dSMartin Matuska * value is rigorously maintained and the increments and decrements are
112*22649d4dSMartin Matuska * sequentially consistent. However, the reads are relaxed, so a reader may
113*22649d4dSMartin Matuska * see a stale value. In the event that a critical worker wakeup is dropped,
114*22649d4dSMartin Matuska * the backup timer intervenes to keep dispatches running.
115*22649d4dSMartin Matuska */
116*22649d4dSMartin Matuska
117*22649d4dSMartin Matuska typedef struct {
118*22649d4dSMartin Matuska queue_item_t *qs_item;
119*22649d4dSMartin Matuska size_t qs_cost;
120*22649d4dSMartin Matuska boolean_t qs_completed;
121*22649d4dSMartin Matuska boolean_t qs_end_of_stream;
122*22649d4dSMartin Matuska } queue_slot_t;
123*22649d4dSMartin Matuska
124*22649d4dSMartin Matuska typedef struct {
125*22649d4dSMartin Matuska uint64_t enqueue;
126*22649d4dSMartin Matuska uint64_t claim;
127*22649d4dSMartin Matuska uint64_t complete;
128*22649d4dSMartin Matuska uint64_t dequeue;
129*22649d4dSMartin Matuska } zq_indexes_t;
130*22649d4dSMartin Matuska
131*22649d4dSMartin Matuska typedef struct {
132*22649d4dSMartin Matuska pthread_cond_t completed;
133*22649d4dSMartin Matuska pthread_cond_t dequeued;
134*22649d4dSMartin Matuska } zq_conditions_t;
135*22649d4dSMartin Matuska
136*22649d4dSMartin Matuska typedef struct {
137*22649d4dSMartin Matuska int min_depth;
138*22649d4dSMartin Matuska int max_depth;
139*22649d4dSMartin Matuska } zq_stats_t;
140*22649d4dSMartin Matuska
141*22649d4dSMartin Matuska struct zstream_queue {
142*22649d4dSMartin Matuska int zq_id;
143*22649d4dSMartin Matuska queue_slot_t *zq_slots;
144*22649d4dSMartin Matuska pthread_mutex_t zq_mutex;
145*22649d4dSMartin Matuska zq_indexes_t zq_ix;
146*22649d4dSMartin Matuska zq_conditions_t zq_cond;
147*22649d4dSMartin Matuska zq_params_t zq_params;
148*22649d4dSMartin Matuska boolean_t zq_disallow_enqueue;
149*22649d4dSMartin Matuska #ifdef MONITOR_QUEUES
150*22649d4dSMartin Matuska zq_stats_t zq_stats;
151*22649d4dSMartin Matuska uint64_t zq_histogram[ZQ_MAX_BATCH+1]; /* Batch sizes */
152*22649d4dSMartin Matuska #endif
153*22649d4dSMartin Matuska };
154*22649d4dSMartin Matuska
155*22649d4dSMartin Matuska typedef struct {
156*22649d4dSMartin Matuska pthread_mutex_t tp_pool_mutex;
157*22649d4dSMartin Matuska pthread_cond_t tp_wake_worker; /* Awaited by workers */
158*22649d4dSMartin Matuska
159*22649d4dSMartin Matuska pthread_mutex_t tp_dispatch_mutex;
160*22649d4dSMartin Matuska pthread_cond_t tp_request_dispatch; /* By dispatch thread */
161*22649d4dSMartin Matuska boolean_t tp_dispatch_requested;
162*22649d4dSMartin Matuska
163*22649d4dSMartin Matuska zstream_queue_t *tp_queues[ZQ_MAX_QUEUES];
164*22649d4dSMartin Matuska int tp_num_queues;
165*22649d4dSMartin Matuska
166*22649d4dSMartin Matuska boolean_t tp_threads_created;
167*22649d4dSMartin Matuska int tp_num_threads;
168*22649d4dSMartin Matuska
169*22649d4dSMartin Matuska uint64_t tp_unclaimed; /* Atomic, all queues */
170*22649d4dSMartin Matuska } thread_pool_t;
171*22649d4dSMartin Matuska
172*22649d4dSMartin Matuska typedef union {
173*22649d4dSMartin Matuska long long ll;
174*22649d4dSMartin Matuska long double ld;
175*22649d4dSMartin Matuska void *p;
176*22649d4dSMartin Matuska void (*fp)(void);
177*22649d4dSMartin Matuska } worst_case_alignment_t;
178*22649d4dSMartin Matuska
179*22649d4dSMartin Matuska static void *queue_worker(void *);
180*22649d4dSMartin Matuska static void *dispatch_worker(void *);
181*22649d4dSMartin Matuska
182*22649d4dSMartin Matuska #ifdef MONITOR_QUEUES
183*22649d4dSMartin Matuska static void *cpu_and_queue_monitor(void *);
184*22649d4dSMartin Matuska static void print_batch_size_histogram(zstream_queue_t *);
185*22649d4dSMartin Matuska #endif
186*22649d4dSMartin Matuska
187*22649d4dSMartin Matuska static thread_pool_t pool = {0};
188*22649d4dSMartin Matuska static pthread_once_t once_control = PTHREAD_ONCE_INIT;
189*22649d4dSMartin Matuska
190*22649d4dSMartin Matuska /*
191*22649d4dSMartin Matuska * The dispatch timer needs sub-millisecond accuracy. POSIX timers on
192*22649d4dSMartin Matuska * FreeBSD don't implement that, but nanosleep() works fine.
193*22649d4dSMartin Matuska */
194*22649d4dSMartin Matuska static void
sleep_nsec(uint64_t nsec)195*22649d4dSMartin Matuska sleep_nsec(uint64_t nsec)
196*22649d4dSMartin Matuska {
197*22649d4dSMartin Matuska struct timespec ts = {
198*22649d4dSMartin Matuska .tv_sec = nsec / NANOSEC,
199*22649d4dSMartin Matuska .tv_nsec = nsec % NANOSEC
200*22649d4dSMartin Matuska };
201*22649d4dSMartin Matuska while (nanosleep(&ts, &ts) != 0) {
202*22649d4dSMartin Matuska if (errno != EINTR)
203*22649d4dSMartin Matuska err(1, "nanosleep failed");
204*22649d4dSMartin Matuska }
205*22649d4dSMartin Matuska }
206*22649d4dSMartin Matuska
207*22649d4dSMartin Matuska static void
thread_pool_init(void)208*22649d4dSMartin Matuska thread_pool_init(void)
209*22649d4dSMartin Matuska {
210*22649d4dSMartin Matuska pthread_mutex_init(&pool.tp_pool_mutex, NULL);
211*22649d4dSMartin Matuska pthread_cond_init(&pool.tp_wake_worker, NULL);
212*22649d4dSMartin Matuska
213*22649d4dSMartin Matuska pthread_mutex_init(&pool.tp_dispatch_mutex, NULL);
214*22649d4dSMartin Matuska pthread_cond_init(&pool.tp_request_dispatch, NULL);
215*22649d4dSMartin Matuska
216*22649d4dSMartin Matuska safe_create_thread(dispatch_worker, NULL, "dispatch", B_TRUE);
217*22649d4dSMartin Matuska }
218*22649d4dSMartin Matuska
219*22649d4dSMartin Matuska /*
220*22649d4dSMartin Matuska * If this function is to be called at all, it must be called before any
221*22649d4dSMartin Matuska * queues have been created.
222*22649d4dSMartin Matuska */
223*22649d4dSMartin Matuska void
zstream_queue_set_num_threads(int n)224*22649d4dSMartin Matuska zstream_queue_set_num_threads(int n)
225*22649d4dSMartin Matuska {
226*22649d4dSMartin Matuska pthread_once(&once_control, thread_pool_init);
227*22649d4dSMartin Matuska pthread_mutex_lock(&pool.tp_pool_mutex);
228*22649d4dSMartin Matuska if (pool.tp_threads_created) {
229*22649d4dSMartin Matuska errx(1, "thread pool size must be set before creating queues");
230*22649d4dSMartin Matuska } else if (n < 1) {
231*22649d4dSMartin Matuska errx(1, "number of threads must be at least 1");
232*22649d4dSMartin Matuska } else if (n < ZQ_MIN_THREADS) {
233*22649d4dSMartin Matuska warnx("using only %d threads may limit performance, setting "
234*22649d4dSMartin Matuska "anyway...", n);
235*22649d4dSMartin Matuska } else if (n > 256) {
236*22649d4dSMartin Matuska warnx("num_threads = %d seems suspiciously high, setting "
237*22649d4dSMartin Matuska "anyway...", n);
238*22649d4dSMartin Matuska }
239*22649d4dSMartin Matuska pool.tp_num_threads = n;
240*22649d4dSMartin Matuska pthread_mutex_unlock(&pool.tp_pool_mutex);
241*22649d4dSMartin Matuska }
242*22649d4dSMartin Matuska
243*22649d4dSMartin Matuska /*
244*22649d4dSMartin Matuska * Locking: the caller must hold the pool mutex.
245*22649d4dSMartin Matuska *
246*22649d4dSMartin Matuska * If tp_num_threads is nonzero, it sets the number of threads to spawn.
247*22649d4dSMartin Matuska * Otherwise, one thread is spawned per core, with a minimum of 6 threads.
248*22649d4dSMartin Matuska *
249*22649d4dSMartin Matuska * sched_getaffinity() is a better estimate of available threads than
250*22649d4dSMartin Matuska * sysconf because sysconf doesn't account for limits that might be set on,
251*22649d4dSMartin Matuska * e.g., a container.
252*22649d4dSMartin Matuska */
253*22649d4dSMartin Matuska static void
thread_pool_spinup(void)254*22649d4dSMartin Matuska thread_pool_spinup(void)
255*22649d4dSMartin Matuska {
256*22649d4dSMartin Matuska if (pool.tp_num_threads == 0) {
257*22649d4dSMartin Matuska #ifdef CPU_COUNT
258*22649d4dSMartin Matuska cpu_set_t cpu_set;
259*22649d4dSMartin Matuska if (sched_getaffinity(0, sizeof (cpu_set_t), &cpu_set) != 0) {
260*22649d4dSMartin Matuska warn("sched_getaffinity failed, using sysconf");
261*22649d4dSMartin Matuska pool.tp_num_threads = sysconf(_SC_NPROCESSORS_ONLN);
262*22649d4dSMartin Matuska } else {
263*22649d4dSMartin Matuska pool.tp_num_threads = CPU_COUNT(&cpu_set);
264*22649d4dSMartin Matuska }
265*22649d4dSMartin Matuska #else
266*22649d4dSMartin Matuska pool.tp_num_threads = sysconf(_SC_NPROCESSORS_ONLN);
267*22649d4dSMartin Matuska #endif
268*22649d4dSMartin Matuska pool.tp_num_threads = MAX(pool.tp_num_threads, ZQ_MIN_THREADS);
269*22649d4dSMartin Matuska }
270*22649d4dSMartin Matuska for (int i = 0; i < pool.tp_num_threads; i++) {
271*22649d4dSMartin Matuska char name[32];
272*22649d4dSMartin Matuska snprintf(name, sizeof (name), "queue-%d", i);
273*22649d4dSMartin Matuska safe_create_thread(queue_worker, NULL, name, B_TRUE);
274*22649d4dSMartin Matuska }
275*22649d4dSMartin Matuska #ifdef MONITOR_QUEUES
276*22649d4dSMartin Matuska safe_create_thread(cpu_and_queue_monitor, NULL, "monitor", B_TRUE);
277*22649d4dSMartin Matuska #endif
278*22649d4dSMartin Matuska }
279*22649d4dSMartin Matuska
280*22649d4dSMartin Matuska zstream_queue_t *
zstream_queue_create(zq_params_t * params)281*22649d4dSMartin Matuska zstream_queue_create(zq_params_t *params)
282*22649d4dSMartin Matuska {
283*22649d4dSMartin Matuska static int next_queue_id = 0;
284*22649d4dSMartin Matuska
285*22649d4dSMartin Matuska VERIFY3P(params->qp_process, !=, NULL);
286*22649d4dSMartin Matuska VERIFY3P(params->qp_cost, !=, NULL);
287*22649d4dSMartin Matuska VERIFY3U(params->qp_item_size, >, 0);
288*22649d4dSMartin Matuska VERIFY3U(params->qp_queue_length, >, 0);
289*22649d4dSMartin Matuska VERIFY3U(params->qp_queue_length, <, 1 << 18);
290*22649d4dSMartin Matuska
291*22649d4dSMartin Matuska pthread_once(&once_control, thread_pool_init);
292*22649d4dSMartin Matuska pthread_mutex_lock(&pool.tp_pool_mutex);
293*22649d4dSMartin Matuska VERIFY3S(pool.tp_num_queues, <, ZQ_MAX_QUEUES);
294*22649d4dSMartin Matuska
295*22649d4dSMartin Matuska if (!pool.tp_threads_created) {
296*22649d4dSMartin Matuska thread_pool_spinup();
297*22649d4dSMartin Matuska pool.tp_threads_created = B_TRUE;
298*22649d4dSMartin Matuska }
299*22649d4dSMartin Matuska
300*22649d4dSMartin Matuska zstream_queue_t *queue = safe_malloc(sizeof (zstream_queue_t));
301*22649d4dSMartin Matuska *queue = (zstream_queue_t) {
302*22649d4dSMartin Matuska .zq_id = next_queue_id++,
303*22649d4dSMartin Matuska .zq_params = *params,
304*22649d4dSMartin Matuska .zq_slots = safe_malloc(params->qp_queue_length *
305*22649d4dSMartin Matuska (sizeof (queue_slot_t))),
306*22649d4dSMartin Matuska #ifdef MONITOR_QUEUES
307*22649d4dSMartin Matuska .zq_stats.min_depth = INT_MAX
308*22649d4dSMartin Matuska #endif
309*22649d4dSMartin Matuska };
310*22649d4dSMartin Matuska pool.tp_queues[pool.tp_num_queues] = queue;
311*22649d4dSMartin Matuska
312*22649d4dSMartin Matuska size_t qpis_rounded = P2ROUNDUP(params->qp_item_size,
313*22649d4dSMartin Matuska _Alignof(worst_case_alignment_t));
314*22649d4dSMartin Matuska uint8_t *items = safe_malloc(params->qp_queue_length * qpis_rounded);
315*22649d4dSMartin Matuska for (size_t i = 0; i < params->qp_queue_length; i++) {
316*22649d4dSMartin Matuska queue->zq_slots[i].qs_item =
317*22649d4dSMartin Matuska (queue_item_t *)(items + i * qpis_rounded);
318*22649d4dSMartin Matuska }
319*22649d4dSMartin Matuska
320*22649d4dSMartin Matuska pthread_mutex_init(&queue->zq_mutex, NULL);
321*22649d4dSMartin Matuska pthread_cond_init(&queue->zq_cond.completed, NULL);
322*22649d4dSMartin Matuska pthread_cond_init(&queue->zq_cond.dequeued, NULL);
323*22649d4dSMartin Matuska
324*22649d4dSMartin Matuska pool.tp_num_queues++;
325*22649d4dSMartin Matuska pthread_mutex_unlock(&pool.tp_pool_mutex);
326*22649d4dSMartin Matuska return (queue);
327*22649d4dSMartin Matuska }
328*22649d4dSMartin Matuska
329*22649d4dSMartin Matuska /*
330*22649d4dSMartin Matuska * Try to advance the "claim" and "complete" indexes as far as possible by
331*22649d4dSMartin Matuska * examining the qs_completed flag on each item. This can't be done directly
332*22649d4dSMartin Matuska * by the threads that complete work, for a couple of reasons:
333*22649d4dSMartin Matuska *
334*22649d4dSMartin Matuska * - Items can be completed in any order. Just because you (a thread) have
335*22649d4dSMartin Matuska * finished your batch doesn't mean that all prior batches have completed.
336*22649d4dSMartin Matuska * If there are uncompleted items ahead of you in the ring buffer, you can't
337*22649d4dSMartin Matuska * advance the completion index past them on your way out.
338*22649d4dSMartin Matuska *
339*22649d4dSMartin Matuska * - Items for which the cost function returns 0 are marked as qs_completed
340*22649d4dSMartin Matuska * on enqueue and are never seen by a worker thread. So, there needs to be
341*22649d4dSMartin Matuska * an independent mechanism to sweep the completion index past these items
342*22649d4dSMartin Matuska * whenever that becomes possible.
343*22649d4dSMartin Matuska *
344*22649d4dSMartin Matuska * This function is called:
345*22649d4dSMartin Matuska *
346*22649d4dSMartin Matuska * - Whenever a thread completes a batch
347*22649d4dSMartin Matuska * - Whenever a thread claims a batch
348*22649d4dSMartin Matuska * - Whenever an item of cost 0 is enqueued
349*22649d4dSMartin Matuska *
350*22649d4dSMartin Matuska * Strictly speaking, advancing on claiming a batch is not logically
351*22649d4dSMartin Matuska * necessary. However, the claimer already holds the queue mutex, and it's
352*22649d4dSMartin Matuska * in our interest to make completed items available for dequeueing as
353*22649d4dSMartin Matuska * expeditiously as possible.
354*22649d4dSMartin Matuska *
355*22649d4dSMartin Matuska * Sweeping of the "claim" index is also an optimization. It is not
356*22649d4dSMartin Matuska * necessary for correctness. However, if we don't do it here, it can only
357*22649d4dSMartin Matuska * be done by threads as they claim jobs to work on. In some cases, not
358*22649d4dSMartin Matuska * advancing the "claim" index here can result in an empty batch and a
359*22649d4dSMartin Matuska * wasted claim cycle.
360*22649d4dSMartin Matuska *
361*22649d4dSMartin Matuska * Locking: the caller must hold the queue mutex.
362*22649d4dSMartin Matuska */
363*22649d4dSMartin Matuska static inline void
advance_indexes(zstream_queue_t * queue)364*22649d4dSMartin Matuska advance_indexes(zstream_queue_t *queue)
365*22649d4dSMartin Matuska {
366*22649d4dSMartin Matuska boolean_t any_completed = B_FALSE;
367*22649d4dSMartin Matuska uint64_t claimed = 0;
368*22649d4dSMartin Matuska
369*22649d4dSMartin Matuska while (queue->zq_ix.claim < queue->zq_ix.enqueue &&
370*22649d4dSMartin Matuska Q_SLOT(queue, queue->zq_ix.claim).qs_completed) {
371*22649d4dSMartin Matuska queue->zq_ix.claim++;
372*22649d4dSMartin Matuska claimed++;
373*22649d4dSMartin Matuska }
374*22649d4dSMartin Matuska if (claimed > 0) {
375*22649d4dSMartin Matuska /*
376*22649d4dSMartin Matuska * tp_unclaimed is decremented both here and in
377*22649d4dSMartin Matuska * claim_batch(). The conditions are mutually exclusive, so
378*22649d4dSMartin Matuska * double counting will not occur.
379*22649d4dSMartin Matuska */
380*22649d4dSMartin Matuska atomic_sub_64(&pool.tp_unclaimed, claimed);
381*22649d4dSMartin Matuska }
382*22649d4dSMartin Matuska while (queue->zq_ix.complete < queue->zq_ix.claim &&
383*22649d4dSMartin Matuska Q_SLOT(queue, queue->zq_ix.complete).qs_completed) {
384*22649d4dSMartin Matuska queue->zq_ix.complete++;
385*22649d4dSMartin Matuska any_completed = B_TRUE;
386*22649d4dSMartin Matuska }
387*22649d4dSMartin Matuska if (any_completed) {
388*22649d4dSMartin Matuska pthread_cond_signal(&queue->zq_cond.completed);
389*22649d4dSMartin Matuska }
390*22649d4dSMartin Matuska }
391*22649d4dSMartin Matuska
392*22649d4dSMartin Matuska /*
393*22649d4dSMartin Matuska * Score a queue according to its need for workers. Higher is better. The
394*22649d4dSMartin Matuska * scoring tries to assign threads to queues that are running out of space
395*22649d4dSMartin Matuska * for new enqueuements or that have little completed work available to
396*22649d4dSMartin Matuska * dequeue. The broader goal is to try to avoid pipeline stalls.
397*22649d4dSMartin Matuska *
398*22649d4dSMartin Matuska * Two measures are used for scoring. The "open score" is 1/M where M is the
399*22649d4dSMartin Matuska * number of slots available to receive new items. The "dequeue score" is
400*22649d4dSMartin Matuska * 1/N where N is the number of completed items available to dequeue. These
401*22649d4dSMartin Matuska * two measures are added together with the dequeue score scaled by
402*22649d4dSMartin Matuska * DEQUEUE_SCORE_WEIGHT.
403*22649d4dSMartin Matuska *
404*22649d4dSMartin Matuska * The composite score is scaled by a factor that reflects how much work is
405*22649d4dSMartin Matuska * actually available to be claimed on the queue; there's no point assigning
406*22649d4dSMartin Matuska * threads to queues that have no work.
407*22649d4dSMartin Matuska *
408*22649d4dSMartin Matuska * Locking: the caller must hold the thread pool mutex and the queue mutex.
409*22649d4dSMartin Matuska */
410*22649d4dSMartin Matuska static inline double
score_queue(zstream_queue_t * queue)411*22649d4dSMartin Matuska score_queue(zstream_queue_t *queue)
412*22649d4dSMartin Matuska {
413*22649d4dSMartin Matuska uint64_t claimable = queue->zq_ix.enqueue - queue->zq_ix.claim;
414*22649d4dSMartin Matuska uint64_t dequeueable = queue->zq_ix.complete - queue->zq_ix.dequeue;
415*22649d4dSMartin Matuska uint64_t in_queue = queue->zq_ix.enqueue - queue->zq_ix.dequeue;
416*22649d4dSMartin Matuska uint64_t open_slots = queue->zq_params.qp_queue_length - in_queue;
417*22649d4dSMartin Matuska
418*22649d4dSMartin Matuska double open_score = (open_slots > 0) ? (1.0 / open_slots) : 2.0;
419*22649d4dSMartin Matuska double dq_score = (dequeueable > 0) ? (1.0 / dequeueable) : 2.0;
420*22649d4dSMartin Matuska double claim_factor = MIN(claimable, (uint64_t)PLENTY_OF_WORK) /
421*22649d4dSMartin Matuska (double)PLENTY_OF_WORK;
422*22649d4dSMartin Matuska double need = open_score + dq_score * DEQUEUE_SCORE_WEIGHT;
423*22649d4dSMartin Matuska return (need * claim_factor);
424*22649d4dSMartin Matuska }
425*22649d4dSMartin Matuska
426*22649d4dSMartin Matuska /*
427*22649d4dSMartin Matuska * Return a random index from an array of doubles, with the likelihood of
428*22649d4dSMartin Matuska * index i being selected equal to weights[i] / sum(weights). Returns index
429*22649d4dSMartin Matuska * 0 if no weight is greater than 0.
430*22649d4dSMartin Matuska */
431*22649d4dSMartin Matuska static inline int
select_stochastic(double weights[],int num_values)432*22649d4dSMartin Matuska select_stochastic(double weights[], int num_values)
433*22649d4dSMartin Matuska {
434*22649d4dSMartin Matuska const double denominator = (double)UINT64_MAX;
435*22649d4dSMartin Matuska uint64_t numerator;
436*22649d4dSMartin Matuska double total = 0.0;
437*22649d4dSMartin Matuska
438*22649d4dSMartin Matuska for (int i = 0; i < num_values; i++) {
439*22649d4dSMartin Matuska total += weights[i];
440*22649d4dSMartin Matuska }
441*22649d4dSMartin Matuska random_get_pseudo_bytes((uint8_t *)&numerator, sizeof (numerator));
442*22649d4dSMartin Matuska double select_val = total * numerator / denominator;
443*22649d4dSMartin Matuska for (int i = 0; i < num_values; i++) {
444*22649d4dSMartin Matuska if (select_val < weights[i])
445*22649d4dSMartin Matuska return (i);
446*22649d4dSMartin Matuska select_val -= weights[i];
447*22649d4dSMartin Matuska }
448*22649d4dSMartin Matuska /* Fallback in case of FP rounding not producing a winner */
449*22649d4dSMartin Matuska for (int i = num_values - 1; i >= 0; i--) {
450*22649d4dSMartin Matuska if (weights[i] != 0.0)
451*22649d4dSMartin Matuska return (i);
452*22649d4dSMartin Matuska }
453*22649d4dSMartin Matuska return (0);
454*22649d4dSMartin Matuska }
455*22649d4dSMartin Matuska
456*22649d4dSMartin Matuska /*
457*22649d4dSMartin Matuska * Claim up to ZQ_MAX_BATCH work items from the given queue, trying to
458*22649d4dSMartin Matuska * accumulate at least qp_batch_budget worth of work data (== "cost"). All
459*22649d4dSMartin Matuska * items in a batch will be drawn from the same queue.
460*22649d4dSMartin Matuska *
461*22649d4dSMartin Matuska * Does not block waiting to fill the budget; returns whatever is available.
462*22649d4dSMartin Matuska *
463*22649d4dSMartin Matuska * Locking: this function must be called with both the queue mutex and the
464*22649d4dSMartin Matuska * thread pool mutex held. zstream_queue_destroy() can't hold a queue's
465*22649d4dSMartin Matuska * mutex while destroying it (because destruction entails destroying the
466*22649d4dSMartin Matuska * queue mutex, which must be unlocked), so holding the queue mutex while
467*22649d4dSMartin Matuska * attempting to claim work is not a sufficient guarantee of correctness.
468*22649d4dSMartin Matuska *
469*22649d4dSMartin Matuska * In other contexts, we have more certainty about whether a queue still has
470*22649d4dSMartin Matuska * work to do. If it does, it can't be destroyed while we hold the queue
471*22649d4dSMartin Matuska * mutex alone. But here, we merely suspect that there's work available
472*22649d4dSMartin Matuska * based on possibly outdated queue scoring information. By the time we get
473*22649d4dSMartin Matuska * here, the queue might already have been finalized. Holding the thread
474*22649d4dSMartin Matuska * pool mutex guarantees that the queue won't have been destroyed out from
475*22649d4dSMartin Matuska * under us.
476*22649d4dSMartin Matuska */
477*22649d4dSMartin Matuska static int
claim_batch(zstream_queue_t * queue,queue_slot_t ** batch)478*22649d4dSMartin Matuska claim_batch(zstream_queue_t *queue, queue_slot_t **batch)
479*22649d4dSMartin Matuska {
480*22649d4dSMartin Matuska size_t cost_claimed = 0;
481*22649d4dSMartin Matuska int count = 0;
482*22649d4dSMartin Matuska uint64_t passed = 0;
483*22649d4dSMartin Matuska boolean_t more_to_claim, more_slots, more_budget;
484*22649d4dSMartin Matuska boolean_t first_and_only, ok_to_claim;
485*22649d4dSMartin Matuska
486*22649d4dSMartin Matuska while (B_TRUE) {
487*22649d4dSMartin Matuska more_to_claim = queue->zq_ix.claim < queue->zq_ix.enqueue;
488*22649d4dSMartin Matuska more_slots = count < ZQ_MAX_BATCH;
489*22649d4dSMartin Matuska more_budget = cost_claimed < queue->zq_params.qp_batch_budget;
490*22649d4dSMartin Matuska first_and_only = queue->zq_params.qp_batch_budget == 0 &&
491*22649d4dSMartin Matuska count == 0;
492*22649d4dSMartin Matuska ok_to_claim = first_and_only || more_budget;
493*22649d4dSMartin Matuska
494*22649d4dSMartin Matuska if (!more_to_claim || !more_slots || !ok_to_claim) {
495*22649d4dSMartin Matuska break;
496*22649d4dSMartin Matuska }
497*22649d4dSMartin Matuska queue_slot_t *slot = &Q_SLOT(queue, queue->zq_ix.claim);
498*22649d4dSMartin Matuska if (!slot->qs_completed) {
499*22649d4dSMartin Matuska cost_claimed += slot->qs_cost;
500*22649d4dSMartin Matuska batch[count++] = slot;
501*22649d4dSMartin Matuska }
502*22649d4dSMartin Matuska queue->zq_ix.claim++;
503*22649d4dSMartin Matuska passed++;
504*22649d4dSMartin Matuska }
505*22649d4dSMartin Matuska
506*22649d4dSMartin Matuska /*
507*22649d4dSMartin Matuska * Every slot the claim index moved over leaves the unclaimed pool,
508*22649d4dSMartin Matuska * whether we took it for the batch or skipped it as already complete.
509*22649d4dSMartin Matuska */
510*22649d4dSMartin Matuska if (passed > 0) {
511*22649d4dSMartin Matuska atomic_sub_64(&pool.tp_unclaimed, passed);
512*22649d4dSMartin Matuska }
513*22649d4dSMartin Matuska advance_indexes(queue);
514*22649d4dSMartin Matuska #ifdef MONITOR_QUEUES
515*22649d4dSMartin Matuska queue->zq_histogram[count]++;
516*22649d4dSMartin Matuska #endif
517*22649d4dSMartin Matuska return (count);
518*22649d4dSMartin Matuska }
519*22649d4dSMartin Matuska
520*22649d4dSMartin Matuska /*
521*22649d4dSMartin Matuska * Threads are assigned to a queue on each loop so they can be shifted
522*22649d4dSMartin Matuska * dynamically to follow available work. Idle threads will typically be
523*22649d4dSMartin Matuska * waiting on the tp_wake_worker condition within this function.
524*22649d4dSMartin Matuska *
525*22649d4dSMartin Matuska * Locking: we hold the pool mutex throughout, both to keep a queue from
526*22649d4dSMartin Matuska * being destroyed out from under us while we score it or claim from it, and
527*22649d4dSMartin Matuska * because it is the mutex for tp_wake_worker. Individual queues are locked
528*22649d4dSMartin Matuska * for only as long as it takes to score or claim from them.
529*22649d4dSMartin Matuska */
530*22649d4dSMartin Matuska static int
assign_queue_and_get_work(zstream_queue_t ** queue,queue_slot_t ** batch)531*22649d4dSMartin Matuska assign_queue_and_get_work(zstream_queue_t **queue, queue_slot_t **batch)
532*22649d4dSMartin Matuska {
533*22649d4dSMartin Matuska pthread_mutex_lock(&pool.tp_pool_mutex);
534*22649d4dSMartin Matuska
535*22649d4dSMartin Matuska while (B_TRUE) {
536*22649d4dSMartin Matuska int num_queues = pool.tp_num_queues;
537*22649d4dSMartin Matuska double weights[ZQ_MAX_QUEUES];
538*22649d4dSMartin Matuska int queues_with_work = 0;
539*22649d4dSMartin Matuska
540*22649d4dSMartin Matuska for (int i = 0; i < num_queues; i++) {
541*22649d4dSMartin Matuska zstream_queue_t *to_score = pool.tp_queues[i];
542*22649d4dSMartin Matuska pthread_mutex_lock(&to_score->zq_mutex);
543*22649d4dSMartin Matuska weights[i] = score_queue(to_score);
544*22649d4dSMartin Matuska pthread_mutex_unlock(&to_score->zq_mutex);
545*22649d4dSMartin Matuska if (weights[i] > NO_WORK)
546*22649d4dSMartin Matuska queues_with_work++;
547*22649d4dSMartin Matuska }
548*22649d4dSMartin Matuska if (!queues_with_work) {
549*22649d4dSMartin Matuska pthread_cond_wait(&pool.tp_wake_worker,
550*22649d4dSMartin Matuska &pool.tp_pool_mutex);
551*22649d4dSMartin Matuska } else {
552*22649d4dSMartin Matuska int q = select_stochastic(weights, num_queues);
553*22649d4dSMartin Matuska *queue = pool.tp_queues[q];
554*22649d4dSMartin Matuska pthread_mutex_lock(&(*queue)->zq_mutex);
555*22649d4dSMartin Matuska int count = claim_batch(*queue, batch);
556*22649d4dSMartin Matuska pthread_mutex_unlock(&(*queue)->zq_mutex);
557*22649d4dSMartin Matuska /*
558*22649d4dSMartin Matuska * Try to wake up another worker thread if there
559*22649d4dSMartin Matuska * still seems to be work available (on any queue).
560*22649d4dSMartin Matuska */
561*22649d4dSMartin Matuska if (atomic_load_64(&pool.tp_unclaimed) > 0) {
562*22649d4dSMartin Matuska pthread_cond_signal(&pool.tp_wake_worker);
563*22649d4dSMartin Matuska }
564*22649d4dSMartin Matuska pthread_mutex_unlock(&pool.tp_pool_mutex);
565*22649d4dSMartin Matuska return (count);
566*22649d4dSMartin Matuska }
567*22649d4dSMartin Matuska }
568*22649d4dSMartin Matuska }
569*22649d4dSMartin Matuska
570*22649d4dSMartin Matuska /*
571*22649d4dSMartin Matuska * Batches are processed without holding any locks. The existence of the
572*22649d4dSMartin Matuska * items we're working on guarantees that the queue can't be destroyed out
573*22649d4dSMartin Matuska * from under us.
574*22649d4dSMartin Matuska *
575*22649d4dSMartin Matuska * However, we can't mark items completed without holding the queue lock
576*22649d4dSMartin Matuska * because that creates a potential race condition with advance_indexes()
577*22649d4dSMartin Matuska * being called on another thread.
578*22649d4dSMartin Matuska */
579*22649d4dSMartin Matuska static void *
queue_worker(void * dummy)580*22649d4dSMartin Matuska queue_worker(void *dummy)
581*22649d4dSMartin Matuska {
582*22649d4dSMartin Matuska (void) dummy;
583*22649d4dSMartin Matuska zstream_queue_t *queue;
584*22649d4dSMartin Matuska queue_slot_t *batch[ZQ_MAX_BATCH];
585*22649d4dSMartin Matuska int count;
586*22649d4dSMartin Matuska
587*22649d4dSMartin Matuska while (B_TRUE) {
588*22649d4dSMartin Matuska count = assign_queue_and_get_work(&queue, batch);
589*22649d4dSMartin Matuska if (count) {
590*22649d4dSMartin Matuska zq_process_item_f *process =
591*22649d4dSMartin Matuska queue->zq_params.qp_process;
592*22649d4dSMartin Matuska void *context = queue->zq_params.qp_context;
593*22649d4dSMartin Matuska for (int i = 0; i < count; i++) {
594*22649d4dSMartin Matuska process(batch[i]->qs_item, context);
595*22649d4dSMartin Matuska }
596*22649d4dSMartin Matuska pthread_mutex_lock(&queue->zq_mutex);
597*22649d4dSMartin Matuska for (int i = 0; i < count; i++) {
598*22649d4dSMartin Matuska batch[i]->qs_completed = B_TRUE;
599*22649d4dSMartin Matuska }
600*22649d4dSMartin Matuska advance_indexes(queue);
601*22649d4dSMartin Matuska pthread_mutex_unlock(&queue->zq_mutex);
602*22649d4dSMartin Matuska }
603*22649d4dSMartin Matuska }
604*22649d4dSMartin Matuska return (NULL);
605*22649d4dSMartin Matuska }
606*22649d4dSMartin Matuska
607*22649d4dSMartin Matuska /*
608*22649d4dSMartin Matuska * Locking: must be called with the dispatch mutex held
609*22649d4dSMartin Matuska *
610*22649d4dSMartin Matuska * Skips the wakeup if tp_unclaimed == 0.
611*22649d4dSMartin Matuska */
612*22649d4dSMartin Matuska static inline void
maybe_wake_worker(void)613*22649d4dSMartin Matuska maybe_wake_worker(void)
614*22649d4dSMartin Matuska {
615*22649d4dSMartin Matuska pool.tp_dispatch_requested = B_FALSE;
616*22649d4dSMartin Matuska if (atomic_load_64(&pool.tp_unclaimed) > 0) {
617*22649d4dSMartin Matuska pthread_mutex_lock(&pool.tp_pool_mutex);
618*22649d4dSMartin Matuska pthread_cond_signal(&pool.tp_wake_worker);
619*22649d4dSMartin Matuska pthread_mutex_unlock(&pool.tp_pool_mutex);
620*22649d4dSMartin Matuska }
621*22649d4dSMartin Matuska }
622*22649d4dSMartin Matuska
623*22649d4dSMartin Matuska static inline struct timespec
timeout_timespec(void)624*22649d4dSMartin Matuska timeout_timespec(void)
625*22649d4dSMartin Matuska {
626*22649d4dSMartin Matuska struct timespec expire;
627*22649d4dSMartin Matuska struct timeval tv;
628*22649d4dSMartin Matuska
629*22649d4dSMartin Matuska if (gettimeofday(&tv, NULL) != 0)
630*22649d4dSMartin Matuska err(1, "couldn't gettimeofday()");
631*22649d4dSMartin Matuska uint64_t nsec = tv.tv_usec * 1000 + DISPATCH_BACKUP_NSEC;
632*22649d4dSMartin Matuska expire.tv_sec = tv.tv_sec + nsec / NANOSEC;
633*22649d4dSMartin Matuska expire.tv_nsec = nsec % NANOSEC;
634*22649d4dSMartin Matuska return (expire);
635*22649d4dSMartin Matuska }
636*22649d4dSMartin Matuska
637*22649d4dSMartin Matuska /*
638*22649d4dSMartin Matuska * The enqueue notification pacing thread, which converts a notification
639*22649d4dSMartin Matuska * from an enqueuer into a possible worker wakeup roughly ENQUEUE_DELAY_NSEC
640*22649d4dSMartin Matuska * later. The delay facilitates larger batch sizes and keeps enqueuers on a
641*22649d4dSMartin Matuska * less-contested mutex.
642*22649d4dSMartin Matuska *
643*22649d4dSMartin Matuska * The condwait timeout is necessary because the tp_unclaimed count is not
644*22649d4dSMartin Matuska * the final word on whether there is actually any work to claim. It is
645*22649d4dSMartin Matuska * calculated rigorously. However, it's a bare atomic and therefore
646*22649d4dSMartin Matuska * potentially out of date at any given moment. A backup strategy is
647*22649d4dSMartin Matuska * necessary to restart processing in the event of a race.
648*22649d4dSMartin Matuska */
649*22649d4dSMartin Matuska static void *
dispatch_worker(void * nope)650*22649d4dSMartin Matuska dispatch_worker(void *nope)
651*22649d4dSMartin Matuska {
652*22649d4dSMartin Matuska (void) nope;
653*22649d4dSMartin Matuska pthread_mutex_lock(&pool.tp_dispatch_mutex);
654*22649d4dSMartin Matuska while (B_TRUE) {
655*22649d4dSMartin Matuska while (!pool.tp_dispatch_requested) {
656*22649d4dSMartin Matuska int rc;
657*22649d4dSMartin Matuska struct timespec expire = timeout_timespec();
658*22649d4dSMartin Matuska rc = pthread_cond_timedwait(&pool.tp_request_dispatch,
659*22649d4dSMartin Matuska &pool.tp_dispatch_mutex, &expire);
660*22649d4dSMartin Matuska if (rc == ETIMEDOUT) {
661*22649d4dSMartin Matuska maybe_wake_worker();
662*22649d4dSMartin Matuska } else if (rc != 0) {
663*22649d4dSMartin Matuska errx(1, "pthread_cond_timedwait() failed: %s",
664*22649d4dSMartin Matuska strerror(rc));
665*22649d4dSMartin Matuska }
666*22649d4dSMartin Matuska }
667*22649d4dSMartin Matuska pthread_mutex_unlock(&pool.tp_dispatch_mutex);
668*22649d4dSMartin Matuska sleep_nsec(ENQUEUE_DELAY_NSEC);
669*22649d4dSMartin Matuska pthread_mutex_lock(&pool.tp_dispatch_mutex);
670*22649d4dSMartin Matuska maybe_wake_worker();
671*22649d4dSMartin Matuska }
672*22649d4dSMartin Matuska return (NULL);
673*22649d4dSMartin Matuska }
674*22649d4dSMartin Matuska
675*22649d4dSMartin Matuska /*
676*22649d4dSMartin Matuska * Implements both _enqueue and _fini. item == NULL for fini.
677*22649d4dSMartin Matuska */
678*22649d4dSMartin Matuska void
zstream_enqueue(zstream_queue_t * queue,queue_item_t * item)679*22649d4dSMartin Matuska zstream_enqueue(zstream_queue_t *queue, queue_item_t *item)
680*22649d4dSMartin Matuska {
681*22649d4dSMartin Matuska VERIFY3P(queue, !=, NULL);
682*22649d4dSMartin Matuska pthread_mutex_lock(&queue->zq_mutex);
683*22649d4dSMartin Matuska
684*22649d4dSMartin Matuska VERIFY3B(queue->zq_disallow_enqueue, ==, B_FALSE);
685*22649d4dSMartin Matuska while (Q_FULL(queue)) {
686*22649d4dSMartin Matuska pthread_cond_wait(&queue->zq_cond.dequeued, &queue->zq_mutex);
687*22649d4dSMartin Matuska }
688*22649d4dSMartin Matuska VERIFY3B(queue->zq_disallow_enqueue, ==, B_FALSE);
689*22649d4dSMartin Matuska queue_slot_t *slot = &Q_SLOT(queue, queue->zq_ix.enqueue);
690*22649d4dSMartin Matuska if (item) {
691*22649d4dSMartin Matuska slot->qs_cost =
692*22649d4dSMartin Matuska queue->zq_params.qp_cost(item, queue->zq_params.qp_context);
693*22649d4dSMartin Matuska slot->qs_completed = slot->qs_cost == 0;
694*22649d4dSMartin Matuska slot->qs_end_of_stream = B_FALSE;
695*22649d4dSMartin Matuska memcpy(slot->qs_item, item, queue->zq_params.qp_item_size);
696*22649d4dSMartin Matuska } else {
697*22649d4dSMartin Matuska slot->qs_cost = 0;
698*22649d4dSMartin Matuska slot->qs_completed = B_TRUE;
699*22649d4dSMartin Matuska slot->qs_end_of_stream = B_TRUE;
700*22649d4dSMartin Matuska queue->zq_disallow_enqueue = B_TRUE;
701*22649d4dSMartin Matuska }
702*22649d4dSMartin Matuska queue->zq_ix.enqueue++;
703*22649d4dSMartin Matuska atomic_inc_64(&pool.tp_unclaimed);
704*22649d4dSMartin Matuska if (slot->qs_cost == 0)
705*22649d4dSMartin Matuska advance_indexes(queue);
706*22649d4dSMartin Matuska
707*22649d4dSMartin Matuska #ifdef MONITOR_QUEUES
708*22649d4dSMartin Matuska /* Maintain queue usage data per monitor interval */
709*22649d4dSMartin Matuska uint64_t depth = queue->zq_ix.enqueue - queue->zq_ix.dequeue;
710*22649d4dSMartin Matuska queue->zq_stats.max_depth = MAX(queue->zq_stats.max_depth, depth);
711*22649d4dSMartin Matuska queue->zq_stats.min_depth = MIN(queue->zq_stats.min_depth, depth);
712*22649d4dSMartin Matuska #endif
713*22649d4dSMartin Matuska
714*22649d4dSMartin Matuska pthread_mutex_unlock(&queue->zq_mutex);
715*22649d4dSMartin Matuska
716*22649d4dSMartin Matuska pthread_mutex_lock(&pool.tp_dispatch_mutex);
717*22649d4dSMartin Matuska pool.tp_dispatch_requested = B_TRUE;
718*22649d4dSMartin Matuska pthread_cond_signal(&pool.tp_request_dispatch);
719*22649d4dSMartin Matuska pthread_mutex_unlock(&pool.tp_dispatch_mutex);
720*22649d4dSMartin Matuska }
721*22649d4dSMartin Matuska
722*22649d4dSMartin Matuska void
zstream_queue_fini(zstream_queue_t * queue)723*22649d4dSMartin Matuska zstream_queue_fini(zstream_queue_t *queue)
724*22649d4dSMartin Matuska {
725*22649d4dSMartin Matuska zstream_enqueue(queue, NULL);
726*22649d4dSMartin Matuska }
727*22649d4dSMartin Matuska
728*22649d4dSMartin Matuska /*
729*22649d4dSMartin Matuska * This function is not public. The only way to destroy a queue through the
730*22649d4dSMartin Matuska * public API is to call zstream_queue_fini(), wait for all items to be
731*22649d4dSMartin Matuska * processed, and then dequeue all items. As a consequence, threads are
732*22649d4dSMartin Matuska * entitled to assume that any queue with unprocessed work will not be
733*22649d4dSMartin Matuska * removed without locking the pool mutex.
734*22649d4dSMartin Matuska *
735*22649d4dSMartin Matuska * Locking: the caller must NOT hold the queue lock. The pool mutex is held
736*22649d4dSMartin Matuska * while destroying the queue.
737*22649d4dSMartin Matuska */
738*22649d4dSMartin Matuska static void
zstream_queue_destroy(zstream_queue_t * queue)739*22649d4dSMartin Matuska zstream_queue_destroy(zstream_queue_t *queue)
740*22649d4dSMartin Matuska {
741*22649d4dSMartin Matuska pthread_mutex_lock(&pool.tp_pool_mutex);
742*22649d4dSMartin Matuska
743*22649d4dSMartin Matuska #ifdef MONITOR_QUEUES
744*22649d4dSMartin Matuska print_batch_size_histogram(queue);
745*22649d4dSMartin Matuska #endif
746*22649d4dSMartin Matuska
747*22649d4dSMartin Matuska VERIFY0(pthread_mutex_destroy(&queue->zq_mutex));
748*22649d4dSMartin Matuska VERIFY0(pthread_cond_destroy(&queue->zq_cond.dequeued));
749*22649d4dSMartin Matuska if (pthread_cond_destroy(&queue->zq_cond.completed) != 0) {
750*22649d4dSMartin Matuska errx(1, "cannot destroy zstream_queue completed condition - "
751*22649d4dSMartin Matuska "are you attempting to dequeue from multiple threads "
752*22649d4dSMartin Matuska "simultaneously?");
753*22649d4dSMartin Matuska }
754*22649d4dSMartin Matuska pool.tp_num_queues--;
755*22649d4dSMartin Matuska if (pool.tp_num_queues > 0) {
756*22649d4dSMartin Matuska /* Gaps are not allowed in the tp_queues array */
757*22649d4dSMartin Matuska zstream_queue_t **qscan = &pool.tp_queues[0];
758*22649d4dSMartin Matuska int i = pool.tp_num_queues;
759*22649d4dSMartin Matuska while (*qscan != queue) { qscan++; i--; }
760*22649d4dSMartin Matuska if (i > 0)
761*22649d4dSMartin Matuska memmove(qscan, qscan + 1, i * sizeof (*qscan));
762*22649d4dSMartin Matuska }
763*22649d4dSMartin Matuska /*
764*22649d4dSMartin Matuska * Items are allocated as a single block. The address of the first
765*22649d4dSMartin Matuska * item field is in fact the start of the block.
766*22649d4dSMartin Matuska */
767*22649d4dSMartin Matuska free(queue->zq_slots[0].qs_item);
768*22649d4dSMartin Matuska free(queue->zq_slots);
769*22649d4dSMartin Matuska queue->zq_slots = NULL;
770*22649d4dSMartin Matuska free(queue);
771*22649d4dSMartin Matuska
772*22649d4dSMartin Matuska pthread_mutex_unlock(&pool.tp_pool_mutex);
773*22649d4dSMartin Matuska }
774*22649d4dSMartin Matuska
775*22649d4dSMartin Matuska /*
776*22649d4dSMartin Matuska * Locking: if more than one thread attempts to dequeue items
777*22649d4dSMartin Matuska * simultaneously, disaster is likely. It will work fine until the end of
778*22649d4dSMartin Matuska * the stream, at which point it becomes a tossup between a race condition
779*22649d4dSMartin Matuska * with multiple attempts to destroy the whole queue vs. an attempt to
780*22649d4dSMartin Matuska * delete a condition that another thread is waiting on. Hence the warning
781*22649d4dSMartin Matuska * not to do multithreaded dequeues in zstream_queue.h.
782*22649d4dSMartin Matuska *
783*22649d4dSMartin Matuska * Returns B_TRUE if real data is returned, B_FALSE if the end of the queue
784*22649d4dSMartin Matuska * has been reached.
785*22649d4dSMartin Matuska */
786*22649d4dSMartin Matuska boolean_t
zstream_dequeue(zstream_queue_t * queue,queue_item_t * item)787*22649d4dSMartin Matuska zstream_dequeue(zstream_queue_t *queue, queue_item_t *item)
788*22649d4dSMartin Matuska {
789*22649d4dSMartin Matuska pthread_mutex_lock(&queue->zq_mutex);
790*22649d4dSMartin Matuska while (queue->zq_ix.dequeue >= queue->zq_ix.complete) {
791*22649d4dSMartin Matuska pthread_cond_wait(&queue->zq_cond.completed, &queue->zq_mutex);
792*22649d4dSMartin Matuska }
793*22649d4dSMartin Matuska queue_slot_t *slot = &Q_SLOT(queue, queue->zq_ix.dequeue);
794*22649d4dSMartin Matuska queue->zq_ix.dequeue++;
795*22649d4dSMartin Matuska if (slot->qs_end_of_stream) {
796*22649d4dSMartin Matuska pthread_mutex_unlock(&queue->zq_mutex);
797*22649d4dSMartin Matuska /* Potential multi-dequeuer race point */
798*22649d4dSMartin Matuska zstream_queue_destroy(queue);
799*22649d4dSMartin Matuska return (B_FALSE);
800*22649d4dSMartin Matuska } else {
801*22649d4dSMartin Matuska memcpy(item, slot->qs_item, queue->zq_params.qp_item_size);
802*22649d4dSMartin Matuska pthread_cond_signal(&queue->zq_cond.dequeued);
803*22649d4dSMartin Matuska pthread_mutex_unlock(&queue->zq_mutex);
804*22649d4dSMartin Matuska return (B_TRUE);
805*22649d4dSMartin Matuska }
806*22649d4dSMartin Matuska }
807*22649d4dSMartin Matuska
808*22649d4dSMartin Matuska #ifdef MONITOR_QUEUES
809*22649d4dSMartin Matuska
810*22649d4dSMartin Matuska #define USEC_PER_JIFFY 10000
811*22649d4dSMartin Matuska #define SAMPLE_DURATION_USEC 1000000
812*22649d4dSMartin Matuska #define CPU_FIELD_WIDTH 14
813*22649d4dSMartin Matuska
814*22649d4dSMartin Matuska /*
815*22649d4dSMartin Matuska * Called only during zstream_queue_destroy(), under the pool mutex
816*22649d4dSMartin Matuska */
817*22649d4dSMartin Matuska static void
print_batch_size_histogram(zstream_queue_t * queue)818*22649d4dSMartin Matuska print_batch_size_histogram(zstream_queue_t *queue)
819*22649d4dSMartin Matuska {
820*22649d4dSMartin Matuska int last_nonzero = 0;
821*22649d4dSMartin Matuska static int lines_printed = 0;
822*22649d4dSMartin Matuska
823*22649d4dSMartin Matuska if (lines_printed++ == 0)
824*22649d4dSMartin Matuska fprintf(stderr, "\nBatch size histograms:\n");
825*22649d4dSMartin Matuska for (last_nonzero = ZQ_MAX_BATCH; last_nonzero >= 0; last_nonzero--) {
826*22649d4dSMartin Matuska if (queue->zq_histogram[last_nonzero] > 0)
827*22649d4dSMartin Matuska break;
828*22649d4dSMartin Matuska }
829*22649d4dSMartin Matuska fprintf(stderr, "Queue %d: ", queue->zq_id);
830*22649d4dSMartin Matuska const char *sep = "";
831*22649d4dSMartin Matuska for (int i = 0; i <= last_nonzero; i++) {
832*22649d4dSMartin Matuska fprintf(stderr, "%s%llu", sep,
833*22649d4dSMartin Matuska (u_longlong_t)queue->zq_histogram[i]);
834*22649d4dSMartin Matuska sep = ", ";
835*22649d4dSMartin Matuska }
836*22649d4dSMartin Matuska fprintf(stderr, "\n");
837*22649d4dSMartin Matuska fflush(stderr);
838*22649d4dSMartin Matuska }
839*22649d4dSMartin Matuska
840*22649d4dSMartin Matuska /*
841*22649d4dSMartin Matuska * Monitor queue and CPU usage from a separate thread. This is all
842*22649d4dSMartin Matuska * Linux-specific, but it's needed only while tuning queue lengths and
843*22649d4dSMartin Matuska * batch sizes. Prints the minimum and maximum queue depth observed
844*22649d4dSMartin Matuska * during each period.
845*22649d4dSMartin Matuska *
846*22649d4dSMartin Matuska * Example output:
847*22649d4dSMartin Matuska *
848*22649d4dSMartin Matuska * CPU: 99.85% Queue 0: 745-1024 Queue 1: 183-256
849*22649d4dSMartin Matuska */
850*22649d4dSMartin Matuska static void *
cpu_and_queue_monitor(void * dummy)851*22649d4dSMartin Matuska cpu_and_queue_monitor(void *dummy)
852*22649d4dSMartin Matuska {
853*22649d4dSMartin Matuska (void) dummy;
854*22649d4dSMartin Matuska uint64_t period = SAMPLE_DURATION_USEC;
855*22649d4dSMartin Matuska struct timespec clock = {0};
856*22649d4dSMartin Matuska uint64_t start_us, end_us;
857*22649d4dSMartin Matuska uint64_t cpu_jif_prior = 0;
858*22649d4dSMartin Matuska uint64_t delta_jif, delta_cpu_jif;
859*22649d4dSMartin Matuska long unsigned int utime, stime;
860*22649d4dSMartin Matuska char buff[1024];
861*22649d4dSMartin Matuska FILE *fp;
862*22649d4dSMartin Matuska
863*22649d4dSMartin Matuska fprintf(stderr, "Queue depths:\n");
864*22649d4dSMartin Matuska
865*22649d4dSMartin Matuska while (B_TRUE) {
866*22649d4dSMartin Matuska
867*22649d4dSMartin Matuska usleep(period);
868*22649d4dSMartin Matuska
869*22649d4dSMartin Matuska fp = fopen("/proc/self/stat", "r");
870*22649d4dSMartin Matuska VERIFY3P(fp, !=, NULL);
871*22649d4dSMartin Matuska VERIFY3P(fgets(buff, sizeof (buff), fp), !=, NULL);
872*22649d4dSMartin Matuska fclose(fp);
873*22649d4dSMartin Matuska char *p = strrchr(buff, ')');
874*22649d4dSMartin Matuska VERIFY3P(p, !=, NULL);
875*22649d4dSMartin Matuska p += 2; /* skip ") " and fields 3-13 */
876*22649d4dSMartin Matuska for (int i = 0; i < 11; i++) {
877*22649d4dSMartin Matuska p = strchr(p, ' ');
878*22649d4dSMartin Matuska VERIFY3P(p, !=, NULL);
879*22649d4dSMartin Matuska p++;
880*22649d4dSMartin Matuska }
881*22649d4dSMartin Matuska VERIFY3U(sscanf(p, "%lu %lu", &utime, &stime), ==, 2);
882*22649d4dSMartin Matuska
883*22649d4dSMartin Matuska pthread_mutex_lock(&pool.tp_pool_mutex);
884*22649d4dSMartin Matuska
885*22649d4dSMartin Matuska clock_gettime(CLOCK_MONOTONIC, &clock);
886*22649d4dSMartin Matuska end_us = clock.tv_sec * 1000000 + clock.tv_nsec / 1000;
887*22649d4dSMartin Matuska
888*22649d4dSMartin Matuska if (cpu_jif_prior > 0) {
889*22649d4dSMartin Matuska delta_cpu_jif = utime + stime - cpu_jif_prior;
890*22649d4dSMartin Matuska delta_jif = (end_us - start_us) / USEC_PER_JIFFY;
891*22649d4dSMartin Matuska double cpu_pct = (double)delta_cpu_jif /
892*22649d4dSMartin Matuska (pool.tp_num_threads * delta_jif);
893*22649d4dSMartin Matuska cpu_pct = MIN(cpu_pct, 0.9999); /* Don't print 100% */
894*22649d4dSMartin Matuska fprintf(stderr, "CPU: %5.2f%% ", 100 * cpu_pct);
895*22649d4dSMartin Matuska } else {
896*22649d4dSMartin Matuska /* No CPU data available for the first interval */
897*22649d4dSMartin Matuska fprintf(stderr, "%*s", CPU_FIELD_WIDTH, "");
898*22649d4dSMartin Matuska }
899*22649d4dSMartin Matuska
900*22649d4dSMartin Matuska for (int i = 0; i < pool.tp_num_queues; i++) {
901*22649d4dSMartin Matuska zstream_queue_t *q = pool.tp_queues[i];
902*22649d4dSMartin Matuska pthread_mutex_lock(&q->zq_mutex);
903*22649d4dSMartin Matuska int min = q->zq_stats.min_depth;
904*22649d4dSMartin Matuska int max = q->zq_stats.max_depth;
905*22649d4dSMartin Matuska if (min > max)
906*22649d4dSMartin Matuska min = max = 0;
907*22649d4dSMartin Matuska fprintf(stderr, "Queue %d: %4d-%-4d ",
908*22649d4dSMartin Matuska q->zq_id, min, max);
909*22649d4dSMartin Matuska q->zq_stats.min_depth = INT_MAX;
910*22649d4dSMartin Matuska q->zq_stats.max_depth = 0;
911*22649d4dSMartin Matuska pthread_mutex_unlock(&q->zq_mutex);
912*22649d4dSMartin Matuska }
913*22649d4dSMartin Matuska
914*22649d4dSMartin Matuska pthread_mutex_unlock(&pool.tp_pool_mutex);
915*22649d4dSMartin Matuska
916*22649d4dSMartin Matuska fprintf(stderr, "\n");
917*22649d4dSMartin Matuska fflush(stderr);
918*22649d4dSMartin Matuska
919*22649d4dSMartin Matuska cpu_jif_prior = utime + stime;
920*22649d4dSMartin Matuska start_us = end_us;
921*22649d4dSMartin Matuska }
922*22649d4dSMartin Matuska return (NULL);
923*22649d4dSMartin Matuska }
924*22649d4dSMartin Matuska
925*22649d4dSMartin Matuska #endif /* MONITOR_QUEUES */
926