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 #ifndef _ZSTREAM_QUEUE_H 18*22649d4dSMartin Matuska #define _ZSTREAM_QUEUE_H 19*22649d4dSMartin Matuska 20*22649d4dSMartin Matuska #ifdef __cplusplus 21*22649d4dSMartin Matuska extern "C" { 22*22649d4dSMartin Matuska #endif 23*22649d4dSMartin Matuska 24*22649d4dSMartin Matuska #include <stddef.h> 25*22649d4dSMartin Matuska #include <sys/stdtypes.h> 26*22649d4dSMartin Matuska 27*22649d4dSMartin Matuska /* 28*22649d4dSMartin Matuska * This is a generalized implementation of multithreaded FIFO work queues. 29*22649d4dSMartin Matuska * 30*22649d4dSMartin Matuska * Callers define a fixed item size to be used by each queue and supply two 31*22649d4dSMartin Matuska * thread-safe functions that 1) estimate individual items' processing costs 32*22649d4dSMartin Matuska * and 2) perform the actual processing. The queue never inspects or 33*22649d4dSMartin Matuska * interprets items in the queue, so processing functions can modify them as 34*22649d4dSMartin Matuska * they wish. 35*22649d4dSMartin Matuska * 36*22649d4dSMartin Matuska * The cost function assigns a size_t cost that estimates the amount of work 37*22649d4dSMartin Matuska * needed to process an item. For operations like hashing and data 38*22649d4dSMartin Matuska * compression, the natural cost is typically payload length. 39*22649d4dSMartin Matuska * 40*22649d4dSMartin Matuska * It's expected that only a subset of input items will require processing. 41*22649d4dSMartin Matuska * If an item's cost is 0, it is fast-tracked and never presented to the 42*22649d4dSMartin Matuska * processing function. 43*22649d4dSMartin Matuska * 44*22649d4dSMartin Matuska * The cost function is run as items enter the queue, with the queue mutex 45*22649d4dSMartin Matuska * held, so it should return a value promptly. If cost estimation is 46*22649d4dSMartin Matuska * expensive and important, use a separate queue to implement it. 47*22649d4dSMartin Matuska * 48*22649d4dSMartin Matuska * Dispatch granularity is specified as a per-batch budget that is set for 49*22649d4dSMartin Matuska * each queue in the same (arbitrary) units used for item costs. Threads 50*22649d4dSMartin Matuska * claim items until the budget is met, there are no more items available, 51*22649d4dSMartin Matuska * or ZQ_MAX_BATCH items have been claimed. When claiming items to work on, 52*22649d4dSMartin Matuska * threads never block waiting for additional work to arrive. They start 53*22649d4dSMartin Matuska * work as quickly as possible even if the budget has not been reached. 54*22649d4dSMartin Matuska * 55*22649d4dSMartin Matuska * A batch budget of 0 means that all batches will have a size of 1. 56*22649d4dSMartin Matuska * 57*22649d4dSMartin Matuska * All queues share a single thread pool that is managed to avoid 58*22649d4dSMartin Matuska * contention. Threads are assigned to queues dynamically according to where 59*22649d4dSMartin Matuska * work is available. When multiple queues have work, threads are allocated 60*22649d4dSMartin Matuska * among them stochastically with an eye toward preventing pipeline stalls. 61*22649d4dSMartin Matuska * 62*22649d4dSMartin Matuska * The shared thread pool persists until the process exits. 63*22649d4dSMartin Matuska */ 64*22649d4dSMartin Matuska 65*22649d4dSMartin Matuska #define ZQ_MAX_BATCH 32 /* The most items that can be claimed at once */ 66*22649d4dSMartin Matuska #define ZQ_MAX_QUEUES 16 /* The maximum number of simultaneous queues */ 67*22649d4dSMartin Matuska #define ZQ_MIN_THREADS 6 68*22649d4dSMartin Matuska 69*22649d4dSMartin Matuska typedef void queue_item_t; 70*22649d4dSMartin Matuska 71*22649d4dSMartin Matuska typedef struct zstream_queue zstream_queue_t; 72*22649d4dSMartin Matuska 73*22649d4dSMartin Matuska /* 74*22649d4dSMartin Matuska * Signatures that cost and processing functions must conform to. 75*22649d4dSMartin Matuska */ 76*22649d4dSMartin Matuska typedef size_t 77*22649d4dSMartin Matuska zq_estimate_cost_f(queue_item_t *item, void *context); 78*22649d4dSMartin Matuska 79*22649d4dSMartin Matuska typedef void 80*22649d4dSMartin Matuska zq_process_item_f(queue_item_t *item, void *context); 81*22649d4dSMartin Matuska 82*22649d4dSMartin Matuska /* 83*22649d4dSMartin Matuska * Set the number of threads to be spawned for queue work. Since all queues 84*22649d4dSMartin Matuska * share a thread pool, this value affects all queues. The value must be set 85*22649d4dSMartin Matuska * before any queues are created. By default, one thread is spawned for 86*22649d4dSMartin Matuska * every CPU core, but always at least ZQ_MIN_THREADS threads. 87*22649d4dSMartin Matuska */ 88*22649d4dSMartin Matuska void 89*22649d4dSMartin Matuska zstream_queue_set_num_threads(int num_threads); 90*22649d4dSMartin Matuska 91*22649d4dSMartin Matuska /* 92*22649d4dSMartin Matuska * Create a queue. The qp_context field is passed to the cost and processing 93*22649d4dSMartin Matuska * functions and is not examined by the queue itself. 94*22649d4dSMartin Matuska */ 95*22649d4dSMartin Matuska 96*22649d4dSMartin Matuska typedef struct { 97*22649d4dSMartin Matuska zq_process_item_f *qp_process; 98*22649d4dSMartin Matuska zq_estimate_cost_f *qp_cost; 99*22649d4dSMartin Matuska void *qp_context; 100*22649d4dSMartin Matuska size_t qp_item_size; 101*22649d4dSMartin Matuska size_t qp_batch_budget; 102*22649d4dSMartin Matuska size_t qp_queue_length; 103*22649d4dSMartin Matuska } zq_params_t; 104*22649d4dSMartin Matuska 105*22649d4dSMartin Matuska zstream_queue_t * 106*22649d4dSMartin Matuska zstream_queue_create(zq_params_t *params); 107*22649d4dSMartin Matuska 108*22649d4dSMartin Matuska /* 109*22649d4dSMartin Matuska * Submit a work item. Blocks if the queue is full. The work item is 110*22649d4dSMartin Matuska * shallow-copied into the queue. Multiple threads may enqueue at once. 111*22649d4dSMartin Matuska */ 112*22649d4dSMartin Matuska void 113*22649d4dSMartin Matuska zstream_enqueue(zstream_queue_t *queue, queue_item_t *item); 114*22649d4dSMartin Matuska 115*22649d4dSMartin Matuska /* 116*22649d4dSMartin Matuska * Retrieve a completed work item. The caller must provide a buffer into 117*22649d4dSMartin Matuska * which the dequeued item is shallow-copied. If the next item is not yet 118*22649d4dSMartin Matuska * ready, this call will block. 119*22649d4dSMartin Matuska * 120*22649d4dSMartin Matuska * If zstream_dequeue returns B_FALSE, the stream is complete. The returned 121*22649d4dSMartin Matuska * item is not valid and no further calls may be made on the queue. 122*22649d4dSMartin Matuska * 123*22649d4dSMartin Matuska * Only one thread may dequeue at once. 124*22649d4dSMartin Matuska */ 125*22649d4dSMartin Matuska boolean_t 126*22649d4dSMartin Matuska zstream_dequeue(zstream_queue_t *queue, queue_item_t *item); 127*22649d4dSMartin Matuska 128*22649d4dSMartin Matuska /* 129*22649d4dSMartin Matuska * Declare that all items have been submitted. The queue will continue to 130*22649d4dSMartin Matuska * function normally for dequeuers and worker threads until zstream_dequeue() 131*22649d4dSMartin Matuska * returns B_FALSE, at which point the queue will be destroyed. 132*22649d4dSMartin Matuska */ 133*22649d4dSMartin Matuska void 134*22649d4dSMartin Matuska zstream_queue_fini(zstream_queue_t *queue); 135*22649d4dSMartin Matuska 136*22649d4dSMartin Matuska #ifdef __cplusplus 137*22649d4dSMartin Matuska } 138*22649d4dSMartin Matuska #endif 139*22649d4dSMartin Matuska 140*22649d4dSMartin Matuska #endif /* _ZSTREAM_QUEUE_H */ 141