diff --git a/src/broadcaster.zig b/src/broadcaster.zig index caff55a..109a3fd 100644 --- a/src/broadcaster.zig +++ b/src/broadcaster.zig @@ -288,6 +288,7 @@ pub const Broadcaster = struct { allocator: Allocator, consumers: std.ArrayListUnmanaged(*Consumer) = .{}, consumers_mutex: std.Thread.Mutex = .{}, + broadcast_order: std.Thread.Mutex = .{}, history: FrameHistory, persist: ?*event_log_mod.DiskPersist = null, stats: Stats = .{}, @@ -912,3 +913,58 @@ test "resequenceFrame preserves identity frame fields" { try std.testing.expectEqualStrings("did:plc:alice", p.getString("did").?); try std.testing.expectEqualStrings("2024-01-15T10:30:00Z", p.getString("time").?); } + +test "concurrent broadcast through ordering mutex produces monotonic sequences" { + // regression test: without broadcast_order serialization, concurrent + // subscriber threads can interleave persist (seq assignment) and broadcast, + // delivering frames out of order to consumers. + // + // this simulates the subscriber pattern: N threads each acquire the + // ordering lock, assign a seq (atomic increment, like persist), and + // broadcast. the ring buffer history must be strictly monotonic. + + var bc = Broadcaster.init(std.testing.allocator); + defer bc.deinit(); + + const num_threads = 8; + const frames_per_thread = 500; + const total_frames = num_threads * frames_per_thread; + var seq_counter = std.atomic.Value(u64).init(0); + + const Worker = struct { + fn run(broadcaster: *Broadcaster, counter: *std.atomic.Value(u64)) void { + for (0..frames_per_thread) |_| { + broadcaster.broadcast_order.lock(); + defer broadcaster.broadcast_order.unlock(); + + const seq = counter.fetchAdd(1, .monotonic) + 1; + broadcaster.broadcast(seq, "x"); + } + } + }; + + var threads: [num_threads]std.Thread = undefined; + for (&threads) |*t| { + t.* = std.Thread.spawn(.{}, Worker.run, .{ &bc, &seq_counter }) catch unreachable; + } + for (&threads) |*t| t.join(); + + // verify: ring buffer history has strictly monotonic sequences + const frames = try bc.history.framesSince(std.testing.allocator, 0); + defer { + for (frames) |f| std.testing.allocator.free(f.data); + std.testing.allocator.free(frames); + } + + // history is capped at 50K, we pushed 4K — all should be present + try std.testing.expectEqual(total_frames, frames.len); + + var prev_seq: u64 = 0; + for (frames) |f| { + if (f.seq <= prev_seq) { + std.debug.print("monotonicity violation: seq {d} after {d}\n", .{ f.seq, prev_seq }); + return error.NonMonotonicSequence; + } + prev_seq = f.seq; + } +} diff --git a/src/subscriber.zig b/src/subscriber.zig index 8a58fee..b39d39a 100644 --- a/src/subscriber.zig +++ b/src/subscriber.zig @@ -425,15 +425,27 @@ const FrameHandler = struct { else .identity; - // persist and get relay-assigned seq, broadcast raw bytes + // persist and get relay-assigned seq, broadcast raw bytes. + // ordering mutex ensures frames are broadcast in seq order — + // without it, concurrent subscriber threads can interleave + // persist (seq assignment) and broadcast, delivering out-of-order. if (sub.persist) |dp| { - const relay_seq = dp.persist(kind, uid, data) catch |err| { - log.warn("persist failed: {s}", .{@errorName(err)}); - sub.bc.broadcast(upstream_seq orelse 0, data); - return; + const relay_seq = blk: { + sub.bc.broadcast_order.lock(); + defer sub.bc.broadcast_order.unlock(); + + const seq = dp.persist(kind, uid, data) catch |err| { + log.warn("persist failed: {s}", .{@errorName(err)}); + return; + }; + sub.bc.stats.relay_seq.store(seq, .release); + const broadcast_data = broadcaster.resequenceFrame(alloc, data, seq) orelse data; + sub.bc.broadcast(seq, broadcast_data); + break :blk seq; }; + _ = relay_seq; - // update per-DID state after successful commit/sync validation + // update per-DID state outside the ordering lock (Postgres round-trip) if ((is_commit or is_sync) and uid > 0) { if (commit_rev) |rev| { const cid_str: []const u8 = if (commit_data_cid) |cid_raw| @@ -445,10 +457,6 @@ const FrameHandler = struct { }; } } - - sub.bc.stats.relay_seq.store(relay_seq, .release); - const broadcast_data = broadcaster.resequenceFrame(alloc, data, relay_seq) orelse data; - sub.bc.broadcast(relay_seq, broadcast_data); } else { sub.bc.broadcast(upstream_seq orelse 0, data); }