Recipe: A Producer/Consumer Queue
The problem
One part of your program discovers work; another part does it. The classic answer is a queue between them, but hand-rolled versions keep growing parts: a mutex, a condition variable, a “no more work” flag, and a subtle bug where workers sleep forever on shutdown.
Current Zig master ships the whole thing: std.Io.Queue(T) is a bounded,
thread-safe channel with blocking put and get and a close protocol.
The plan
- Create
std.Io.Queue(u64)over a fixed buffer you provide. The buffer length is the queue’s capacity, and capacity is backpressure: a producer more than that far ahead blocks until workers catch up. - Spawn workers that loop on
getOne, which blocks until an item arrives. - Produce, then call
close. Workers drain the remaining items, thengetOnereturnserror.Closedand each worker exits its loop. jointhe workers.
const std = @import("std");
const Queue = std.Io.Queue(u64);
fn consumer(q: *Queue, io: std.Io, total: *std.atomic.Value(u64)) void {
while (true) {
// Blocks until an item arrives. After close(), remaining items
// are still delivered; only then does getOne return Closed.
const job = q.getOne(io) catch return;
_ = total.fetchAdd(job * job, .monotonic);
}
}
pub fn main(init: std.process.Init) !void {
const io = init.io;
var buf: [256]u8 = undefined;
var file_writer = std.Io.File.stdout().writerStreaming(io, &buf);
const out = &file_writer.interface;
// The queue borrows its storage; eight slots of backpressure. A
// producer that runs ahead by more than eight items blocks, which is
// the mechanism that keeps a fast producer from drowning slow workers.
var storage: [8]u64 = undefined;
var queue: Queue = .init(&storage);
var total = std.atomic.Value(u64).init(0);
var workers: [3]std.Thread = undefined;
for (&workers) |*t| {
t.* = try std.Thread.spawn(.{}, consumer, .{ &queue, io, &total });
}
// Produce twenty jobs, then close. Close is the shutdown protocol:
// no sentinel values, no stop flag, no leaked worker.
for (1..21) |job| try queue.putAll(io, &.{job});
queue.close(io);
for (workers) |t| t.join();
try out.print("processed total: {d}\n", .{total.load(.monotonic)});
try out.flush();
}Close is the shutdown protocol
The traditional failure mode of worker queues is shutdown: a sentinel
value per worker, or a flag plus a broadcast, each easy to get wrong.
close folds it into the queue’s contract: puts after close fail,
items already queued are still delivered, and only then do getters see
error.Closed. The worker loop is therefore two lines with no special
cases, and no worker can sleep through shutdown.
Why the queue takes io
Every blocking operation goes through the std.Io interface, queues
included. That is the same design the rest of the standard library
follows now: whoever provides the io implementation decides how blocking
is implemented (OS threads here; an event loop elsewhere) and blocking
operations become cancelation points. Passing io into the worker is not
noise; it is what makes the same worker code portable across execution
models.
Variations
- Batching:
putandgetmove slices, not just single items, and return how many they moved;putAllblocks until everything is in. Batches amortize the synchronization per item. - Multiple producers: put is threadsafe; several producers can feed one queue. Close once, from the coordinator that knows production is finished.
- Results channel: for a full pipeline, give workers a second queue
to put results into, and close it after the worker
join.