Skip to content

Producers and broadcast

The second EventBus argument selects the producer mode. Use false for one publishing thread or true when multiple threads may call produce() concurrently. A worker can publish while the initializing thread owns startup and shutdown.

A publication copies one event value into the ring and may wait for capacity. Calls from separate producers can return in a different order from the stream’s claim order. Consumers observe the same sequence order and wait at gaps. Never use observed callback scheduling to infer cross-handler ordering.

This example registers two handlers. Audit verifies sequence order and metrics computes a count and sum; each owns separate mutable state.

Terminal window
zig build example-broadcast
const std = @import("std");
const zigrupt = @import("zigrupt");
// Separate storage for each handler. Main reads only after both threads join.
const Audit = struct {
var next: usize = 0;
fn handle(events: []const usize) usize {
for (events) |event| {
if (event != next) @panic("audit received an unexpected sequence");
next += 1;
}
return events.len;
}
};
const Metrics = struct {
var count: usize = 0;
var sum: usize = 0;
fn handle(events: []const usize) usize {
for (events) |event| {
count += 1;
sum += event;
}
return events.len;
}
};
const Bus = zigrupt.EventBus(usize, false, &.{ Audit.handle, Metrics.handle }, zigrupt.waiting_strategy.BusySpin, zigrupt.waiting_strategy.BusySpin);
pub fn main() !void {
var bus = try Bus.init(std.heap.page_allocator, 8, 4);
defer bus.deinit();
try bus.start();
errdefer bus.stop() catch {};
for (0..100) |value| try bus.produce(value);
try bus.stop();
if (Audit.next != 100 or Metrics.count != 100 or Metrics.sum != 4950)
return error.UnexpectedDelivery;
std.debug.print("broadcast: both handlers received 100 events\n", .{});
}

Create the bus and start its handlers before spawning producers. Join every producer before calling stop(). The scoped joins also cover a partial thread-spawn failure. Producer errors are reported after the threads join.

Terminal window
zig build example-multi_producer
const std = @import("std");
const zigrupt = @import("zigrupt");
const producer_count = 3;
const events_per_producer = 100;
const Event = struct { producer: usize, offset: usize };
var next: [producer_count]usize = @splat(0);
fn handle(events: []const Event) usize {
for (events) |event| {
// Cross-producer interleaving is unspecified; each producer stays ordered.
if (event.offset != next[event.producer]) @panic("unexpected producer order");
next[event.producer] += 1;
}
return events.len;
}
const Bus = zigrupt.EventBus(Event, true, &.{handle}, zigrupt.waiting_strategy.BusySpin, zigrupt.waiting_strategy.BusySpin);
const Publisher = struct {
bus: *Bus,
id: usize,
failure: ?anyerror = null, // Read only after this publisher is joined.
fn run(self: *Publisher) void {
for (0..events_per_producer) |offset| {
self.bus.produce(.{ .producer = self.id, .offset = offset }) catch |err| {
self.failure = err;
return;
};
}
}
};
pub fn main() !void {
var bus = try Bus.init(std.heap.page_allocator, 16, 8);
defer bus.deinit();
try bus.start();
errdefer bus.stop() catch {};
var publishers: [producer_count]Publisher = undefined;
{
var threads: [producer_count]std.Thread = undefined;
var started: usize = 0;
// Also join already-spawned producers if a later spawn fails.
// This scope exits before the bus's error cleanup calls stop.
defer for (threads[0..started]) |thread| thread.join();
for (&publishers, 0..) |*publisher, id| {
publisher.* = .{ .bus = &bus, .id = id };
threads[started] = try std.Thread.spawn(.{}, Publisher.run, .{publisher});
started += 1;
}
}
try bus.stop(); // Every producer has joined, so draining is safe.
for (publishers) |publisher| if (publisher.failure) |err| return err;
for (next) |count| if (count != events_per_producer) return error.MissingEvents;
std.debug.print("multi_producer: 300 events, each producer ordered\n", .{});
}

Avoid synchronous publishing to the same bus from a handler: when the ring fills, the publication can wait for that handler’s own acknowledgement and deadlock. Multi-producer mode does not remove this dependency. Structure application pipelines so their progress does not form a cycle.