From e42ba3b9d297df7f934e0af3bb89c6fe463923ff Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Mon, 2 Mar 2026 13:42:20 -0600 Subject: [PATCH] fix: serialize persist + broadcast for monotonic firehose sequences MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit zlay's thread-per-PDS architecture (2,700 concurrent subscribers) allowed interleaving between seq assignment (under persist lock) and broadcast (unlocked), delivering frames out of order to consumers. adds broadcast_order mutex to Broadcaster — subscriber threads now hold it across persist → resequence → broadcast, matching Indigo's pattern of serializing the entire pipeline through one lock. also fixes the error fallback that broadcast with upstream seq on persist failure (mixing sequence domains). now drops the frame instead. includes regression test that spawns 8 threads doing concurrent broadcast and verifies monotonic output. confirmed: test catches the bug when the lock is removed (seq 3 after 4). Co-Authored-By: Claude Opus 4.6 --- src/broadcaster.zig | 56 +++++++++++++++++++++++++++++++++++++++++++++ src/subscriber.zig | 28 +++++++++++++++-------- 2 files changed, 74 insertions(+), 10 deletions(-) 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); } -- 2.51.2