xref: /freebsd/sys/contrib/openzfs/cmd/zstream/zstream_selftest_queue.c (revision 22649d4dba730d46244fd2dff4fd174903c8379f)
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 /*
18*22649d4dSMartin Matuska  * Selftests for the zstream_queue multithreaded FIFO queue API.
19*22649d4dSMartin Matuska  *
20*22649d4dSMartin Matuska  * All tests are built on one generic workload runner. A workload is
21*22649d4dSMartin Matuska  * described by a qtest_config_t: some number of producer threads each
22*22649d4dSMartin Matuska  * enqueue a stream of self-describing items with randomized costs,
23*22649d4dSMartin Matuska  * payloads, and processing delays, while one consumer thread per queue
24*22649d4dSMartin Matuska  * dequeues and verifies. Several workloads can run concurrently on separate
25*22649d4dSMartin Matuska  * queues to exercise the shared thread pool.
26*22649d4dSMartin Matuska  *
27*22649d4dSMartin Matuska  * Every item carries enough information to be verified independently:
28*22649d4dSMartin Matuska  *
29*22649d4dSMartin Matuska  * - The tuple (qi_producer, qi_seq) identifies each item; the consumer
30*22649d4dSMartin Matuska  *   checks that each producer's items arrive in the same order they
31*22649d4dSMartin Matuska  *   were enqueued.
32*22649d4dSMartin Matuska  *
33*22649d4dSMartin Matuska  * - qi_check is a hash of (qi_seed, qi_producer, qi_seq). The processing
34*22649d4dSMartin Matuska  *   function verifies it and then XORs in TRANSFORM_MAGIC. The consumer
35*22649d4dSMartin Matuska  *   checks that the transform happened iff cost > 0.
36*22649d4dSMartin Matuska  *
37*22649d4dSMartin Matuska  * - qi_pattern[] is filled from qi_check and verified both by the process
38*22649d4dSMartin Matuska  *   function and the consumer, to catch any corruption of the shallow
39*22649d4dSMartin Matuska  *   copies in and out of the ring buffer.
40*22649d4dSMartin Matuska  *
41*22649d4dSMartin Matuska  * - qi_process_count counts invocations of the process function, which
42*22649d4dSMartin Matuska  *   must be exactly one for cost > 0 items and zero for cost == 0 items.
43*22649d4dSMartin Matuska  *
44*22649d4dSMartin Matuska  * Global conservation checks: the number of items dequeued must equal the
45*22649d4dSMartin Matuska  * number enqueued, and the total number of process-function invocations
46*22649d4dSMartin Matuska  * must equal the number of nonzero-cost items enqueued.
47*22649d4dSMartin Matuska  */
48*22649d4dSMartin Matuska 
49*22649d4dSMartin Matuska #include <assert.h>
50*22649d4dSMartin Matuska #include <atomic.h>
51*22649d4dSMartin Matuska #include <err.h>
52*22649d4dSMartin Matuska #include <pthread.h>
53*22649d4dSMartin Matuska #include <stdalign.h>
54*22649d4dSMartin Matuska #include <stdint.h>
55*22649d4dSMartin Matuska #include <stdio.h>
56*22649d4dSMartin Matuska #include <string.h>
57*22649d4dSMartin Matuska #include <unistd.h>
58*22649d4dSMartin Matuska 
59*22649d4dSMartin Matuska #include "zstream_queue.h"
60*22649d4dSMartin Matuska #include "zstream_selftest.h"
61*22649d4dSMartin Matuska 
62*22649d4dSMartin Matuska #define	TRANSFORM_MAGIC	0xf00dfeedbeefcafeULL
63*22649d4dSMartin Matuska 
64*22649d4dSMartin Matuska /*
65*22649d4dSMartin Matuska  * Number of times per 1000 processing function invocations to use an
66*22649d4dSMartin Matuska  * extra-long "outlier" processing delay to force overtly out-of-order
67*22649d4dSMartin Matuska  * completion.
68*22649d4dSMartin Matuska  */
69*22649d4dSMartin Matuska #define	LONG_DELAYS_PER_THOUSAND	3
70*22649d4dSMartin Matuska #define	LONG_DELAY_MULTIPLIER		20
71*22649d4dSMartin Matuska 
72*22649d4dSMartin Matuska typedef struct {
73*22649d4dSMartin Matuska 	uint32_t	qi_producer;
74*22649d4dSMartin Matuska 	uint32_t	qi_delay_us;
75*22649d4dSMartin Matuska 	uint64_t	qi_seq;
76*22649d4dSMartin Matuska 	uint64_t	qi_check;
77*22649d4dSMartin Matuska 	size_t		qi_cost;
78*22649d4dSMartin Matuska 	uint32_t	qi_process_count;
79*22649d4dSMartin Matuska 	uint8_t		qi_pattern[];
80*22649d4dSMartin Matuska } qtest_item_t;
81*22649d4dSMartin Matuska 
82*22649d4dSMartin Matuska typedef struct {
83*22649d4dSMartin Matuska 	uint32_t	qc_producers;		/* Number of producers */
84*22649d4dSMartin Matuska 	uint64_t	qc_items;		/* Items per producer */
85*22649d4dSMartin Matuska 	size_t		qc_queue_length;
86*22649d4dSMartin Matuska 	size_t		qc_batch_budget;
87*22649d4dSMartin Matuska 	size_t		qc_pattern_len;		/* Extra payload bytes */
88*22649d4dSMartin Matuska 	uint32_t	qc_zero_cost_pct;	/* % of items fast-tracked */
89*22649d4dSMartin Matuska 	size_t		qc_max_cost;		/* Nonzero costs are 1..max */
90*22649d4dSMartin Matuska 	uint32_t	qc_delay_pct;		/* % of items slept on */
91*22649d4dSMartin Matuska 	uint32_t	qc_max_delay_us;
92*22649d4dSMartin Matuska 	uint32_t	qc_producer_stall_pct;	/* % chance producer naps */
93*22649d4dSMartin Matuska 	uint32_t	qc_consumer_stall_pct;	/* % chance consumer naps */
94*22649d4dSMartin Matuska 	uint32_t	qc_stall_max_us;
95*22649d4dSMartin Matuska 	uint64_t	qc_rng_stream;		/* base PRNG stream number */
96*22649d4dSMartin Matuska } qtest_config_t;
97*22649d4dSMartin Matuska 
98*22649d4dSMartin Matuska typedef struct {
99*22649d4dSMartin Matuska 	const qtest_config_t	*qr_cfg;
100*22649d4dSMartin Matuska 	zstream_queue_t		*qr_queue;
101*22649d4dSMartin Matuska 	uint32_t		qr_producers_left;
102*22649d4dSMartin Matuska 	uint64_t		qr_expect_processed;	/* Atomic */
103*22649d4dSMartin Matuska 	uint64_t		qr_processed;		/* Atomic */
104*22649d4dSMartin Matuska 	uint64_t		qr_dequeued;
105*22649d4dSMartin Matuska } qtest_run_t;
106*22649d4dSMartin Matuska 
107*22649d4dSMartin Matuska typedef struct {
108*22649d4dSMartin Matuska 	qtest_run_t	*qp_run;
109*22649d4dSMartin Matuska 	uint32_t	qp_id;
110*22649d4dSMartin Matuska } qtest_producer_arg_t;
111*22649d4dSMartin Matuska 
112*22649d4dSMartin Matuska static uint64_t
item_check_value(uint32_t producer,uint64_t seq)113*22649d4dSMartin Matuska item_check_value(uint32_t producer, uint64_t seq)
114*22649d4dSMartin Matuska {
115*22649d4dSMartin Matuska 	return (selftest_mix64(selftest_seed ^
116*22649d4dSMartin Matuska 	    (((uint64_t)producer << 40) + seq)));
117*22649d4dSMartin Matuska }
118*22649d4dSMartin Matuska 
119*22649d4dSMartin Matuska static void
fill_pattern(uint8_t * pattern,size_t len,uint64_t check)120*22649d4dSMartin Matuska fill_pattern(uint8_t *pattern, size_t len, uint64_t check)
121*22649d4dSMartin Matuska {
122*22649d4dSMartin Matuska 	for (size_t i = 0; i < len; i++)
123*22649d4dSMartin Matuska 		pattern[i] = (uint8_t)(check >> ((i & 7) << 3)) ^ (uint8_t)i;
124*22649d4dSMartin Matuska }
125*22649d4dSMartin Matuska 
126*22649d4dSMartin Matuska static void
verify_pattern(const uint8_t * pattern,size_t len,uint64_t check,const char * who)127*22649d4dSMartin Matuska verify_pattern(const uint8_t *pattern, size_t len, uint64_t check,
128*22649d4dSMartin Matuska     const char *who)
129*22649d4dSMartin Matuska {
130*22649d4dSMartin Matuska 	for (size_t i = 0; i < len; i++) {
131*22649d4dSMartin Matuska 		uint8_t expect =
132*22649d4dSMartin Matuska 		    (uint8_t)(check >> ((i & 7) << 3)) ^ (uint8_t)i;
133*22649d4dSMartin Matuska 		if (pattern[i] != expect) {
134*22649d4dSMartin Matuska 			errx(1, "%s: payload corrupted at byte %zu "
135*22649d4dSMartin Matuska 			    "(0x%02x != 0x%02x)", who, i, pattern[i], expect);
136*22649d4dSMartin Matuska 		}
137*22649d4dSMartin Matuska 	}
138*22649d4dSMartin Matuska }
139*22649d4dSMartin Matuska 
140*22649d4dSMartin Matuska static size_t
qtest_cost(void * item_in,void * context)141*22649d4dSMartin Matuska qtest_cost(void *item_in, void *context)
142*22649d4dSMartin Matuska {
143*22649d4dSMartin Matuska 	(void) context;
144*22649d4dSMartin Matuska 	qtest_item_t *item = item_in;
145*22649d4dSMartin Matuska 	return (item->qi_cost);
146*22649d4dSMartin Matuska }
147*22649d4dSMartin Matuska 
148*22649d4dSMartin Matuska static void
qtest_process(void * item_in,void * context)149*22649d4dSMartin Matuska qtest_process(void *item_in, void *context)
150*22649d4dSMartin Matuska {
151*22649d4dSMartin Matuska 	qtest_run_t *run = context;
152*22649d4dSMartin Matuska 	qtest_item_t *item = item_in;
153*22649d4dSMartin Matuska 
154*22649d4dSMartin Matuska 	/* Cost-0 items should never reach the process function */
155*22649d4dSMartin Matuska 	VERIFY3U(item->qi_cost, >, 0);
156*22649d4dSMartin Matuska 	VERIFY3U(item->qi_check, ==,
157*22649d4dSMartin Matuska 	    item_check_value(item->qi_producer, item->qi_seq));
158*22649d4dSMartin Matuska 	verify_pattern(item->qi_pattern, run->qr_cfg->qc_pattern_len,
159*22649d4dSMartin Matuska 	    item->qi_check, "process");
160*22649d4dSMartin Matuska 	VERIFY3U(atomic_add_32_nv(&item->qi_process_count, 1), ==, 1);
161*22649d4dSMartin Matuska 
162*22649d4dSMartin Matuska 	if (item->qi_delay_us > 0)
163*22649d4dSMartin Matuska 		(void) usleep(item->qi_delay_us);
164*22649d4dSMartin Matuska 
165*22649d4dSMartin Matuska 	item->qi_check ^= TRANSFORM_MAGIC;
166*22649d4dSMartin Matuska 	atomic_add_64(&run->qr_processed, 1);
167*22649d4dSMartin Matuska }
168*22649d4dSMartin Matuska 
169*22649d4dSMartin Matuska /*
170*22649d4dSMartin Matuska  * Pthreads worker function for enqueuers
171*22649d4dSMartin Matuska  */
172*22649d4dSMartin Matuska static void *
qtest_producer(void * arg)173*22649d4dSMartin Matuska qtest_producer(void *arg)
174*22649d4dSMartin Matuska {
175*22649d4dSMartin Matuska 	qtest_producer_arg_t *pa = arg;
176*22649d4dSMartin Matuska 	qtest_run_t *run = pa->qp_run;
177*22649d4dSMartin Matuska 	const qtest_config_t *cfg = run->qr_cfg;
178*22649d4dSMartin Matuska 	uint64_t local_expect = 0;
179*22649d4dSMartin Matuska 	selftest_rng_t rng;
180*22649d4dSMartin Matuska 	alignas(uint64_t) uint8_t item_buffer[sizeof (qtest_item_t) +
181*22649d4dSMartin Matuska 	    cfg->qc_pattern_len];
182*22649d4dSMartin Matuska 	qtest_item_t *item = (qtest_item_t *)item_buffer;
183*22649d4dSMartin Matuska 
184*22649d4dSMartin Matuska 	selftest_rng_init(&rng, cfg->qc_rng_stream + 1000 + pa->qp_id);
185*22649d4dSMartin Matuska 
186*22649d4dSMartin Matuska 	for (uint64_t seq = 0; seq < cfg->qc_items; seq++) {
187*22649d4dSMartin Matuska 
188*22649d4dSMartin Matuska 		qtest_item_t item_xfer = {
189*22649d4dSMartin Matuska 			.qi_producer = pa->qp_id,
190*22649d4dSMartin Matuska 			.qi_seq = seq,
191*22649d4dSMartin Matuska 			.qi_process_count = 0,
192*22649d4dSMartin Matuska 			.qi_check = item_check_value(pa->qp_id, seq)
193*22649d4dSMartin Matuska 		};
194*22649d4dSMartin Matuska 		*item = item_xfer;
195*22649d4dSMartin Matuska 		fill_pattern(item->qi_pattern, cfg->qc_pattern_len,
196*22649d4dSMartin Matuska 		    item->qi_check);
197*22649d4dSMartin Matuska 
198*22649d4dSMartin Matuska 		if (selftest_rng_below(&rng, 100) < cfg->qc_zero_cost_pct) {
199*22649d4dSMartin Matuska 			item->qi_cost = 0;
200*22649d4dSMartin Matuska 		} else {
201*22649d4dSMartin Matuska 			item->qi_cost =
202*22649d4dSMartin Matuska 			    1 + selftest_rng_below(&rng, cfg->qc_max_cost);
203*22649d4dSMartin Matuska 			local_expect++;
204*22649d4dSMartin Matuska 		}
205*22649d4dSMartin Matuska 
206*22649d4dSMartin Matuska 		if (item->qi_cost > 0 && cfg->qc_max_delay_us > 0) {
207*22649d4dSMartin Matuska 			if (selftest_rng_below(&rng, 1000) <
208*22649d4dSMartin Matuska 			    LONG_DELAYS_PER_THOUSAND) {
209*22649d4dSMartin Matuska 				item->qi_delay_us = cfg->qc_max_delay_us *
210*22649d4dSMartin Matuska 				    LONG_DELAY_MULTIPLIER;
211*22649d4dSMartin Matuska 			} else if (selftest_rng_below(&rng, 100) <
212*22649d4dSMartin Matuska 			    cfg->qc_delay_pct) {
213*22649d4dSMartin Matuska 				item->qi_delay_us = selftest_rng_below(&rng,
214*22649d4dSMartin Matuska 				    cfg->qc_max_delay_us);
215*22649d4dSMartin Matuska 			}
216*22649d4dSMartin Matuska 		}
217*22649d4dSMartin Matuska 
218*22649d4dSMartin Matuska 		if (cfg->qc_producer_stall_pct > 0 &&
219*22649d4dSMartin Matuska 		    selftest_rng_below(&rng, 100) < cfg->qc_producer_stall_pct)
220*22649d4dSMartin Matuska 			(void) usleep(selftest_rng_below(&rng,
221*22649d4dSMartin Matuska 			    cfg->qc_stall_max_us));
222*22649d4dSMartin Matuska 
223*22649d4dSMartin Matuska 		zstream_enqueue(run->qr_queue, item);
224*22649d4dSMartin Matuska 	}
225*22649d4dSMartin Matuska 
226*22649d4dSMartin Matuska 	atomic_add_64(&run->qr_expect_processed, local_expect);
227*22649d4dSMartin Matuska 	if (atomic_add_32_nv(&run->qr_producers_left, -1) == 0)
228*22649d4dSMartin Matuska 		zstream_queue_fini(run->qr_queue);
229*22649d4dSMartin Matuska 	return (NULL);
230*22649d4dSMartin Matuska }
231*22649d4dSMartin Matuska 
232*22649d4dSMartin Matuska /*
233*22649d4dSMartin Matuska  * Pthreads worker function for dequeuers
234*22649d4dSMartin Matuska  */
235*22649d4dSMartin Matuska static void *
qtest_consumer(void * arg)236*22649d4dSMartin Matuska qtest_consumer(void *arg)
237*22649d4dSMartin Matuska {
238*22649d4dSMartin Matuska 	qtest_run_t *run = arg;
239*22649d4dSMartin Matuska 	const qtest_config_t *cfg = run->qr_cfg;
240*22649d4dSMartin Matuska 	selftest_rng_t rng;
241*22649d4dSMartin Matuska 	uint64_t expected_seq[cfg->qc_producers];
242*22649d4dSMartin Matuska 	alignas(uint64_t) uint8_t item_buffer[sizeof (qtest_item_t) +
243*22649d4dSMartin Matuska 	    cfg->qc_pattern_len];
244*22649d4dSMartin Matuska 	qtest_item_t *item = (qtest_item_t *)item_buffer;
245*22649d4dSMartin Matuska 
246*22649d4dSMartin Matuska 	memset(expected_seq, 0, sizeof (expected_seq));
247*22649d4dSMartin Matuska 	selftest_rng_init(&rng, cfg->qc_rng_stream + 999);
248*22649d4dSMartin Matuska 
249*22649d4dSMartin Matuska 	while (zstream_dequeue(run->qr_queue, item)) {
250*22649d4dSMartin Matuska 		VERIFY3U(item->qi_producer, <, cfg->qc_producers);
251*22649d4dSMartin Matuska 		if (item->qi_seq != expected_seq[item->qi_producer]) {
252*22649d4dSMartin Matuska 			errx(1, "consumer: FIFO order violated: got "
253*22649d4dSMartin Matuska 			    "producer %u seq %ju, expected seq %ju",
254*22649d4dSMartin Matuska 			    item->qi_producer, (uintmax_t)item->qi_seq,
255*22649d4dSMartin Matuska 			    (uintmax_t)expected_seq[item->qi_producer]);
256*22649d4dSMartin Matuska 		}
257*22649d4dSMartin Matuska 		expected_seq[item->qi_producer]++;
258*22649d4dSMartin Matuska 
259*22649d4dSMartin Matuska 		uint64_t check =
260*22649d4dSMartin Matuska 		    item_check_value(item->qi_producer, item->qi_seq);
261*22649d4dSMartin Matuska 		if (item->qi_cost > 0) {
262*22649d4dSMartin Matuska 			VERIFY3U(item->qi_process_count, ==, 1);
263*22649d4dSMartin Matuska 			VERIFY3U(item->qi_check, ==, check ^ TRANSFORM_MAGIC);
264*22649d4dSMartin Matuska 		} else {
265*22649d4dSMartin Matuska 			VERIFY3U(item->qi_process_count, ==, 0);
266*22649d4dSMartin Matuska 			VERIFY3U(item->qi_check, ==, check);
267*22649d4dSMartin Matuska 		}
268*22649d4dSMartin Matuska 		verify_pattern(item->qi_pattern, cfg->qc_pattern_len, check,
269*22649d4dSMartin Matuska 		    "consumer");
270*22649d4dSMartin Matuska 		run->qr_dequeued++;
271*22649d4dSMartin Matuska 
272*22649d4dSMartin Matuska 		if (cfg->qc_consumer_stall_pct > 0 &&
273*22649d4dSMartin Matuska 		    selftest_rng_below(&rng, 100) < cfg->qc_consumer_stall_pct)
274*22649d4dSMartin Matuska 			(void) usleep(selftest_rng_below(&rng,
275*22649d4dSMartin Matuska 			    cfg->qc_stall_max_us));
276*22649d4dSMartin Matuska 	}
277*22649d4dSMartin Matuska 
278*22649d4dSMartin Matuska 	for (uint32_t p = 0; p < cfg->qc_producers; p++)
279*22649d4dSMartin Matuska 		VERIFY3U(expected_seq[p], ==, cfg->qc_items);
280*22649d4dSMartin Matuska 	VERIFY3U(run->qr_dequeued, ==,
281*22649d4dSMartin Matuska 	    (uint64_t)cfg->qc_producers * cfg->qc_items);
282*22649d4dSMartin Matuska 
283*22649d4dSMartin Matuska 	return (NULL);
284*22649d4dSMartin Matuska }
285*22649d4dSMartin Matuska 
286*22649d4dSMartin Matuska /*
287*22649d4dSMartin Matuska  * Run several workloads at once, one queue per config, with a dedicated
288*22649d4dSMartin Matuska  * consumer thread and qc_producers producer threads per queue. Returns
289*22649d4dSMartin Matuska  * after every queue has been drained to end-of-stream (and therefore
290*22649d4dSMartin Matuska  * destroyed) and all verification checks have passed.
291*22649d4dSMartin Matuska  */
292*22649d4dSMartin Matuska static void
run_queue_workloads(const qtest_config_t * cfgs,int ncfg)293*22649d4dSMartin Matuska run_queue_workloads(const qtest_config_t *cfgs, int ncfg)
294*22649d4dSMartin Matuska {
295*22649d4dSMartin Matuska 	qtest_run_t runs[ncfg];
296*22649d4dSMartin Matuska 	pthread_t consumers[ncfg];
297*22649d4dSMartin Matuska 	uint32_t total_producers = 0;
298*22649d4dSMartin Matuska 
299*22649d4dSMartin Matuska 	for (int i = 0; i < ncfg; i++)
300*22649d4dSMartin Matuska 		total_producers += cfgs[i].qc_producers;
301*22649d4dSMartin Matuska 
302*22649d4dSMartin Matuska 	pthread_t producers[total_producers];
303*22649d4dSMartin Matuska 	qtest_producer_arg_t pargs[total_producers];
304*22649d4dSMartin Matuska 	memset(runs, 0, sizeof (runs));
305*22649d4dSMartin Matuska 	memset(pargs, 0, sizeof (pargs));
306*22649d4dSMartin Matuska 
307*22649d4dSMartin Matuska 	for (int i = 0; i < ncfg; i++) {
308*22649d4dSMartin Matuska 		runs[i].qr_cfg = &cfgs[i];
309*22649d4dSMartin Matuska 		runs[i].qr_producers_left = cfgs[i].qc_producers;
310*22649d4dSMartin Matuska 		zq_params_t params = {
311*22649d4dSMartin Matuska 			.qp_process = qtest_process,
312*22649d4dSMartin Matuska 			.qp_cost = qtest_cost,
313*22649d4dSMartin Matuska 			.qp_context = &runs[i],
314*22649d4dSMartin Matuska 			.qp_item_size =
315*22649d4dSMartin Matuska 			    sizeof (qtest_item_t) + cfgs[i].qc_pattern_len,
316*22649d4dSMartin Matuska 			.qp_batch_budget = cfgs[i].qc_batch_budget,
317*22649d4dSMartin Matuska 			.qp_queue_length = cfgs[i].qc_queue_length,
318*22649d4dSMartin Matuska 		};
319*22649d4dSMartin Matuska 		runs[i].qr_queue = zstream_queue_create(&params);
320*22649d4dSMartin Matuska 	}
321*22649d4dSMartin Matuska 
322*22649d4dSMartin Matuska 	int p = 0;
323*22649d4dSMartin Matuska 	for (int i = 0; i < ncfg; i++) {
324*22649d4dSMartin Matuska 		VERIFY3S(pthread_create(&consumers[i], NULL, qtest_consumer,
325*22649d4dSMartin Matuska 		    &runs[i]), ==, 0);
326*22649d4dSMartin Matuska 		for (uint32_t j = 0; j < cfgs[i].qc_producers; j++, p++) {
327*22649d4dSMartin Matuska 			pargs[p].qp_run = &runs[i];
328*22649d4dSMartin Matuska 			pargs[p].qp_id = j;
329*22649d4dSMartin Matuska 			VERIFY3S(pthread_create(&producers[p], NULL,
330*22649d4dSMartin Matuska 			    qtest_producer, &pargs[p]), ==, 0);
331*22649d4dSMartin Matuska 		}
332*22649d4dSMartin Matuska 	}
333*22649d4dSMartin Matuska 
334*22649d4dSMartin Matuska 	for (uint32_t i = 0; i < total_producers; i++)
335*22649d4dSMartin Matuska 		VERIFY3S(pthread_join(producers[i], NULL), ==, 0);
336*22649d4dSMartin Matuska 	for (int i = 0; i < ncfg; i++)
337*22649d4dSMartin Matuska 		VERIFY3S(pthread_join(consumers[i], NULL), ==, 0);
338*22649d4dSMartin Matuska 
339*22649d4dSMartin Matuska 	for (int i = 0; i < ncfg; i++)
340*22649d4dSMartin Matuska 		VERIFY3U(runs[i].qr_processed, ==, runs[i].qr_expect_processed);
341*22649d4dSMartin Matuska }
342*22649d4dSMartin Matuska 
343*22649d4dSMartin Matuska static void
run_queue_workload(const qtest_config_t * cfg)344*22649d4dSMartin Matuska run_queue_workload(const qtest_config_t *cfg)
345*22649d4dSMartin Matuska {
346*22649d4dSMartin Matuska 	run_queue_workloads(cfg, 1);
347*22649d4dSMartin Matuska }
348*22649d4dSMartin Matuska 
349*22649d4dSMartin Matuska /*
350*22649d4dSMartin Matuska  * Basic single-producer smoke test: deterministic-ish costs, no delays.
351*22649d4dSMartin Matuska  */
352*22649d4dSMartin Matuska static void
queue_basic(void)353*22649d4dSMartin Matuska queue_basic(void)
354*22649d4dSMartin Matuska {
355*22649d4dSMartin Matuska 	qtest_config_t cfg = {
356*22649d4dSMartin Matuska 		.qc_producers = 1,
357*22649d4dSMartin Matuska 		.qc_items = 5000,
358*22649d4dSMartin Matuska 		.qc_queue_length = 64,
359*22649d4dSMartin Matuska 		.qc_batch_budget = 256,
360*22649d4dSMartin Matuska 		.qc_pattern_len = 32,
361*22649d4dSMartin Matuska 		.qc_zero_cost_pct = 20,
362*22649d4dSMartin Matuska 		.qc_max_cost = 64,
363*22649d4dSMartin Matuska 	};
364*22649d4dSMartin Matuska 	run_queue_workload(&cfg);
365*22649d4dSMartin Matuska }
366*22649d4dSMartin Matuska 
367*22649d4dSMartin Matuska /*
368*22649d4dSMartin Matuska  * A long, randomized stream with heavy-tailed processing delays, a large
369*22649d4dSMartin Matuska  * fraction of fast-tracked items, costs that exceed the batch budget, and a
370*22649d4dSMartin Matuska  * consumer that periodically stalls so the queue backs up and enqueue
371*22649d4dSMartin Matuska  * blocks on Q_FULL. The ring indices wrap hundreds of times.
372*22649d4dSMartin Matuska  */
373*22649d4dSMartin Matuska static void
queue_torture(void)374*22649d4dSMartin Matuska queue_torture(void)
375*22649d4dSMartin Matuska {
376*22649d4dSMartin Matuska 	qtest_config_t cfg = {
377*22649d4dSMartin Matuska 		.qc_producers = 1,
378*22649d4dSMartin Matuska 		.qc_items = 100000,
379*22649d4dSMartin Matuska 		.qc_queue_length = 512,
380*22649d4dSMartin Matuska 		.qc_batch_budget = 2048,
381*22649d4dSMartin Matuska 		.qc_pattern_len = 64,
382*22649d4dSMartin Matuska 		.qc_zero_cost_pct = 30,
383*22649d4dSMartin Matuska 		.qc_max_cost = 4096,
384*22649d4dSMartin Matuska 		.qc_delay_pct = 5,
385*22649d4dSMartin Matuska 		.qc_max_delay_us = 100,
386*22649d4dSMartin Matuska 		.qc_consumer_stall_pct = 1,
387*22649d4dSMartin Matuska 		.qc_stall_max_us = 500,
388*22649d4dSMartin Matuska 		.qc_rng_stream = 100,
389*22649d4dSMartin Matuska 	};
390*22649d4dSMartin Matuska 	run_queue_workload(&cfg);
391*22649d4dSMartin Matuska }
392*22649d4dSMartin Matuska 
393*22649d4dSMartin Matuska /*
394*22649d4dSMartin Matuska  * Off-by-one hunting: sweep the degenerate corners of queue length,
395*22649d4dSMartin Matuska  * batch budget, and stream length, including a zero-item stream and
396*22649d4dSMartin Matuska  * zero-length payloads.
397*22649d4dSMartin Matuska  */
398*22649d4dSMartin Matuska static void
queue_edge_cases(void)399*22649d4dSMartin Matuska queue_edge_cases(void)
400*22649d4dSMartin Matuska {
401*22649d4dSMartin Matuska 	static const size_t lengths[] =
402*22649d4dSMartin Matuska 	    { 1, 2, ZQ_MAX_BATCH - 1, ZQ_MAX_BATCH, ZQ_MAX_BATCH + 1, 64 };
403*22649d4dSMartin Matuska 	static const size_t budgets[] = { 0, 1, 16, SIZE_MAX / 2 };
404*22649d4dSMartin Matuska 	uint64_t stream = 200;
405*22649d4dSMartin Matuska 
406*22649d4dSMartin Matuska 	for (int l = 0; l < 6; l++) {
407*22649d4dSMartin Matuska 		for (int b = 0; b < 4; b++) {
408*22649d4dSMartin Matuska 			uint64_t counts[] = { 0, lengths[l], lengths[l] + 1,
409*22649d4dSMartin Matuska 			    4 * lengths[l] + 3 };
410*22649d4dSMartin Matuska 			for (int n = 0; n < 4; n++) {
411*22649d4dSMartin Matuska 				qtest_config_t cfg = {
412*22649d4dSMartin Matuska 					.qc_producers = 1,
413*22649d4dSMartin Matuska 					.qc_items = counts[n],
414*22649d4dSMartin Matuska 					.qc_queue_length = lengths[l],
415*22649d4dSMartin Matuska 					.qc_batch_budget = budgets[b],
416*22649d4dSMartin Matuska 					.qc_pattern_len =
417*22649d4dSMartin Matuska 					    (lengths[l] & 1) ? 0 : 24,
418*22649d4dSMartin Matuska 					.qc_zero_cost_pct = 25,
419*22649d4dSMartin Matuska 					.qc_max_cost = 8,
420*22649d4dSMartin Matuska 					.qc_rng_stream = stream++,
421*22649d4dSMartin Matuska 				};
422*22649d4dSMartin Matuska 				run_queue_workload(&cfg);
423*22649d4dSMartin Matuska 			}
424*22649d4dSMartin Matuska 		}
425*22649d4dSMartin Matuska 	}
426*22649d4dSMartin Matuska }
427*22649d4dSMartin Matuska 
428*22649d4dSMartin Matuska /*
429*22649d4dSMartin Matuska  * All items cost 0, so every item takes the fast track and the process
430*22649d4dSMartin Matuska  * function must never run (qtest_process VERIFYs cost > 0, and the
431*22649d4dSMartin Matuska  * conservation check at the end of the run confirms zero invocations).
432*22649d4dSMartin Matuska  * This exercises the completion-index sweep for items no worker ever
433*22649d4dSMartin Matuska  * touches.
434*22649d4dSMartin Matuska  */
435*22649d4dSMartin Matuska static void
queue_zero_cost(void)436*22649d4dSMartin Matuska queue_zero_cost(void)
437*22649d4dSMartin Matuska {
438*22649d4dSMartin Matuska 	qtest_config_t cfg = {
439*22649d4dSMartin Matuska 		.qc_producers = 1,
440*22649d4dSMartin Matuska 		.qc_items = 20000,
441*22649d4dSMartin Matuska 		.qc_queue_length = 128,
442*22649d4dSMartin Matuska 		.qc_batch_budget = 1024,
443*22649d4dSMartin Matuska 		.qc_pattern_len = 16,
444*22649d4dSMartin Matuska 		.qc_zero_cost_pct = 100,
445*22649d4dSMartin Matuska 		.qc_max_cost = 8,
446*22649d4dSMartin Matuska 		.qc_rng_stream = 300,
447*22649d4dSMartin Matuska 	};
448*22649d4dSMartin Matuska 	run_queue_workload(&cfg);
449*22649d4dSMartin Matuska }
450*22649d4dSMartin Matuska 
451*22649d4dSMartin Matuska /*
452*22649d4dSMartin Matuska  * Eight producer threads hammering one queue with random pacing. The
453*22649d4dSMartin Matuska  * consumer verifies per-producer FIFO order and exact counts.
454*22649d4dSMartin Matuska  */
455*22649d4dSMartin Matuska static void
queue_multi_producer(void)456*22649d4dSMartin Matuska queue_multi_producer(void)
457*22649d4dSMartin Matuska {
458*22649d4dSMartin Matuska 	qtest_config_t cfg = {
459*22649d4dSMartin Matuska 		.qc_producers = 8,
460*22649d4dSMartin Matuska 		.qc_items = 15000,
461*22649d4dSMartin Matuska 		.qc_queue_length = 256,
462*22649d4dSMartin Matuska 		.qc_batch_budget = 512,
463*22649d4dSMartin Matuska 		.qc_pattern_len = 24,
464*22649d4dSMartin Matuska 		.qc_zero_cost_pct = 25,
465*22649d4dSMartin Matuska 		.qc_max_cost = 512,
466*22649d4dSMartin Matuska 		.qc_delay_pct = 2,
467*22649d4dSMartin Matuska 		.qc_max_delay_us = 50,
468*22649d4dSMartin Matuska 		.qc_producer_stall_pct = 1,
469*22649d4dSMartin Matuska 		.qc_stall_max_us = 200,
470*22649d4dSMartin Matuska 		.qc_rng_stream = 400,
471*22649d4dSMartin Matuska 	};
472*22649d4dSMartin Matuska 	run_queue_workload(&cfg);
473*22649d4dSMartin Matuska }
474*22649d4dSMartin Matuska 
475*22649d4dSMartin Matuska /*
476*22649d4dSMartin Matuska  * Many dissimilar queues live at once, stressing worker scoring and
477*22649d4dSMartin Matuska  * assignment, per-queue index isolation, and destruction of queues while
478*22649d4dSMartin Matuska  * others remain active (which compacts the pool's queue array).
479*22649d4dSMartin Matuska  */
480*22649d4dSMartin Matuska static void
queue_multi_queue(void)481*22649d4dSMartin Matuska queue_multi_queue(void)
482*22649d4dSMartin Matuska {
483*22649d4dSMartin Matuska 	qtest_config_t cfgs[12];
484*22649d4dSMartin Matuska 	for (int i = 0; i < 12; i++) {
485*22649d4dSMartin Matuska 		uint32_t producers = 1 + i % 3;
486*22649d4dSMartin Matuska 		qtest_config_t cfg = {
487*22649d4dSMartin Matuska 			.qc_producers = producers,
488*22649d4dSMartin Matuska 			.qc_items = 4000 / producers,
489*22649d4dSMartin Matuska 			.qc_queue_length = (size_t)4 << (i % 6),
490*22649d4dSMartin Matuska 			.qc_batch_budget =
491*22649d4dSMartin Matuska 			    (i % 4 == 0) ? 0 : (size_t)64 << (i % 5),
492*22649d4dSMartin Matuska 			.qc_pattern_len = 8 * (i % 5),
493*22649d4dSMartin Matuska 			.qc_zero_cost_pct = 10 * (i % 6),
494*22649d4dSMartin Matuska 			.qc_max_cost = (size_t)16 << (i % 8),
495*22649d4dSMartin Matuska 			.qc_delay_pct = i % 3,
496*22649d4dSMartin Matuska 			.qc_max_delay_us = 60,
497*22649d4dSMartin Matuska 			.qc_rng_stream = 500 + i * 10000,
498*22649d4dSMartin Matuska 		};
499*22649d4dSMartin Matuska 		cfgs[i] = cfg;
500*22649d4dSMartin Matuska 	}
501*22649d4dSMartin Matuska 	run_queue_workloads(cfgs, 12);
502*22649d4dSMartin Matuska }
503*22649d4dSMartin Matuska 
504*22649d4dSMartin Matuska /*
505*22649d4dSMartin Matuska  * Seeded chaos: randomize every workload parameter within sane bounds
506*22649d4dSMartin Matuska  * and run a few rounds of concurrent queues. Whatever the targeted tests
507*22649d4dSMartin Matuska  * miss, this net catches over many CI runs; failures replay with -s.
508*22649d4dSMartin Matuska  */
509*22649d4dSMartin Matuska static void
queue_stress(void)510*22649d4dSMartin Matuska queue_stress(void)
511*22649d4dSMartin Matuska {
512*22649d4dSMartin Matuska 	selftest_rng_t rng;
513*22649d4dSMartin Matuska 	selftest_rng_init(&rng, 900);
514*22649d4dSMartin Matuska 
515*22649d4dSMartin Matuska 	for (int iter = 0; iter < 8; iter++) {
516*22649d4dSMartin Matuska 		int nqueues = 1 + selftest_rng_below(&rng, 4);
517*22649d4dSMartin Matuska 		qtest_config_t cfgs[4];
518*22649d4dSMartin Matuska 
519*22649d4dSMartin Matuska 		for (int i = 0; i < nqueues; i++) {
520*22649d4dSMartin Matuska 			uint32_t producers = 1 + selftest_rng_below(&rng, 4);
521*22649d4dSMartin Matuska 			qtest_config_t cfg = {
522*22649d4dSMartin Matuska 				.qc_producers = producers,
523*22649d4dSMartin Matuska 				.qc_items = (2000 +
524*22649d4dSMartin Matuska 				    selftest_rng_below(&rng, 8000)) /
525*22649d4dSMartin Matuska 				    producers,
526*22649d4dSMartin Matuska 				.qc_queue_length = (size_t)1 <<
527*22649d4dSMartin Matuska 				    selftest_rng_below(&rng, 10),
528*22649d4dSMartin Matuska 				.qc_batch_budget =
529*22649d4dSMartin Matuska 				    (selftest_rng_below(&rng, 3) == 0) ? 0 :
530*22649d4dSMartin Matuska 				    selftest_rng_below(&rng, 4096),
531*22649d4dSMartin Matuska 				.qc_pattern_len =
532*22649d4dSMartin Matuska 				    selftest_rng_below(&rng, 64),
533*22649d4dSMartin Matuska 				.qc_zero_cost_pct =
534*22649d4dSMartin Matuska 				    selftest_rng_below(&rng, 101),
535*22649d4dSMartin Matuska 				.qc_max_cost =
536*22649d4dSMartin Matuska 				    1 + selftest_rng_below(&rng, 2048),
537*22649d4dSMartin Matuska 				.qc_delay_pct = selftest_rng_below(&rng, 4),
538*22649d4dSMartin Matuska 				.qc_max_delay_us =
539*22649d4dSMartin Matuska 				    selftest_rng_below(&rng, 120),
540*22649d4dSMartin Matuska 				.qc_producer_stall_pct =
541*22649d4dSMartin Matuska 				    selftest_rng_below(&rng, 2),
542*22649d4dSMartin Matuska 				.qc_consumer_stall_pct =
543*22649d4dSMartin Matuska 				    selftest_rng_below(&rng, 2),
544*22649d4dSMartin Matuska 				.qc_stall_max_us =
545*22649d4dSMartin Matuska 				    selftest_rng_below(&rng, 400),
546*22649d4dSMartin Matuska 				.qc_rng_stream = 1000000 + iter * 1000 +
547*22649d4dSMartin Matuska 				    i * 100,
548*22649d4dSMartin Matuska 			};
549*22649d4dSMartin Matuska 			cfgs[i] = cfg;
550*22649d4dSMartin Matuska 		}
551*22649d4dSMartin Matuska 		run_queue_workloads(cfgs, nqueues);
552*22649d4dSMartin Matuska 	}
553*22649d4dSMartin Matuska }
554*22649d4dSMartin Matuska 
555*22649d4dSMartin Matuska const test_case_t selftest_queue_cases[] = {
556*22649d4dSMartin Matuska 	{ "queue_basic",		queue_basic },
557*22649d4dSMartin Matuska 	{ "queue_edge_cases",		queue_edge_cases },
558*22649d4dSMartin Matuska 	{ "queue_zero_cost",		queue_zero_cost },
559*22649d4dSMartin Matuska 	{ "queue_torture",		queue_torture },
560*22649d4dSMartin Matuska 	{ "queue_multi_producer",	queue_multi_producer },
561*22649d4dSMartin Matuska 	{ "queue_multi_queue",		queue_multi_queue },
562*22649d4dSMartin Matuska 	{ "queue_stress",		queue_stress },
563*22649d4dSMartin Matuska 	{ NULL,				NULL },
564*22649d4dSMartin Matuska };
565