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.
Independent broadcast consumers
Section titled “Independent broadcast consumers”This example registers two handlers. Audit verifies sequence order and metrics computes a count and sum; each owns separate mutable state.
zig build example-broadcastconst 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", .{});}Concurrent producers
Section titled “Concurrent producers”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.
zig build example-multi_producerconst 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.
