From 1d2e18d1970063e7700f749f87c86f60320f018b Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Sat, 4 Apr 2026 13:41:18 -0500 Subject: [PATCH] add pipeline contention metrics + zero-consumer fast path instrumentation for the ~1,300 host CPU cliff: - relay_persist_order_spins_total: spin iterations on the ordering lock - relay_broadcast_queue_full_total: spin iterations on full broadcast queue - relay_broadcast_queue_depth_hwm: high-water mark of queue depth - relay_broadcast_no_consumers_total: frames that skipped SharedFrame alloc zero-consumer fast path: when no consumers are connected, broadcast() returns after history.push() without allocating SharedFrame or taking consumers_mutex. saves one heap alloc + one mutex per frame. also includes the cursor coalesce fix (CursorMap) and slot reuse (free list with unregister on subscriber exit) from previous commits. Co-Authored-By: Claude Opus 4.6 --- src/api/admin.zig | 2 +- src/broadcaster.zig | 59 ++++++++++++++++++++++++++++++++++++++++++-- src/frame_worker.zig | 9 +++++-- src/subscriber.zig | 9 +++++-- 4 files changed, 72 insertions(+), 7 deletions(-) diff --git a/src/api/admin.zig b/src/api/admin.zig index 5ae9446..415f8a0 100644 --- a/src/api/admin.zig +++ b/src/api/admin.zig @@ -91,7 +91,7 @@ pub fn handleBan(conn: *h.Conn, body: []const u8, headers: *const websocket.Hand h.respondJson(conn, .ok, "{\"success\":true}"); return; }; - ctx.bc.broadcast_queue.push(relay_seq, owned); + ctx.bc.broadcast_queue.push(relay_seq, owned, &ctx.bc.stats); ctx.bc.persist_order.store(0, .release); log.info("admin: emitted #account takedown event for {s} (seq={d})", .{ did, relay_seq }); } else |err| { diff --git a/src/broadcaster.zig b/src/broadcaster.zig index ba0cbab..c294cf9 100644 --- a/src/broadcaster.zig +++ b/src/broadcaster.zig @@ -57,6 +57,11 @@ pub const Stats = struct { host_authority_time_us: std.atomic.Value(u64) = .{ .raw = 0 }, // frame pool memory pressure pool_queued_bytes: std.atomic.Value(u64) = .{ .raw = 0 }, + // persist/broadcast pipeline contention + persist_order_spins: std.atomic.Value(u64) = .{ .raw = 0 }, + broadcast_queue_full: std.atomic.Value(u64) = .{ .raw = 0 }, + broadcast_queue_depth_hwm: std.atomic.Value(u32) = .{ .raw = 0 }, + broadcast_no_consumers: std.atomic.Value(u64) = .{ .raw = 0 }, start_time: i64 = 0, }; @@ -256,10 +261,18 @@ pub const BroadcastQueue = struct { return .{ .allocator = allocator }; } + /// approximate queue depth (may be slightly stale — for metrics only) + pub fn depth(self: *const BroadcastQueue) u32 { + const tail = self.tail.load(.monotonic); + const head = self.head.load(.monotonic); + return if (tail >= head) tail - head else CAPACITY - head + tail; + } + /// push an item (called by worker threads). spins until space is available. /// matches Indigo semantics: every persisted event reaches live broadcast. /// the broadcaster fiber drains at memory speed (no I/O), so spin is brief. - pub fn push(self: *BroadcastQueue, seq: u64, data: []const u8) void { + pub fn push(self: *BroadcastQueue, seq: u64, data: []const u8, stats: *Stats) void { + var full_spins: u32 = 0; while (true) { // acquire spinlock while (self.push_lock.cmpxchgWeak(0, 1, .acquire, .monotonic) != null) { @@ -271,6 +284,7 @@ pub const BroadcastQueue = struct { if (next_tail == self.head.load(.acquire)) { // full — release lock, yield, retry self.push_lock.store(0, .release); + full_spins += 1; std.atomic.spinLoopHint(); continue; } @@ -278,6 +292,16 @@ pub const BroadcastQueue = struct { self.items[tail] = .{ .seq = seq, .data = data }; self.tail.store(next_tail, .release); self.push_lock.store(0, .release); + + if (full_spins > 0) { + _ = stats.broadcast_queue_full.fetchAdd(full_spins, .monotonic); + } + // update high-water mark (racy but directionally correct) + const d = self.depth(); + const hwm = stats.broadcast_queue_depth_hwm.load(.monotonic); + if (d > hwm) { + _ = stats.broadcast_queue_depth_hwm.cmpxchgWeak(hwm, d, .monotonic, .monotonic); + } return; } } @@ -522,6 +546,12 @@ pub const Broadcaster = struct { // add to history for cursor replay _ = self.history.push(seq, data); + // fast path: skip SharedFrame allocation when nobody is listening + if (self.stats.consumer_count.load(.acquire) == 0) { + _ = self.stats.broadcast_no_consumers.fetchAdd(1, .monotonic); + return; + } + // create one shared frame for all consumers const frame = SharedFrame.create(self.allocator, data) catch return; defer frame.release(); // release broadcaster's reference @@ -934,6 +964,31 @@ pub fn formatPrometheusMetrics(stats: *const Stats, cache_entries: usize, attrib stats.pool_queued_bytes.load(.acquire), }) catch return w.buffered(); + // pipeline contention metrics (separate print to stay under 32-arg limit) + w.print( + \\# TYPE relay_persist_order_spins_total counter + \\# HELP relay_persist_order_spins_total spin iterations waiting for persist ordering lock + \\relay_persist_order_spins_total {d} + \\ + \\# TYPE relay_broadcast_queue_full_total counter + \\# HELP relay_broadcast_queue_full_total spin iterations on full broadcast queue + \\relay_broadcast_queue_full_total {d} + \\ + \\# TYPE relay_broadcast_queue_depth_hwm gauge + \\# HELP relay_broadcast_queue_depth_hwm high-water mark of broadcast queue depth + \\relay_broadcast_queue_depth_hwm {d} + \\ + \\# TYPE relay_broadcast_no_consumers_total counter + \\# HELP relay_broadcast_no_consumers_total frames skipped broadcast (no consumers) + \\relay_broadcast_no_consumers_total {d} + \\ + , .{ + stats.persist_order_spins.load(.acquire), + stats.broadcast_queue_full.load(.acquire), + stats.broadcast_queue_depth_hwm.load(.acquire), + stats.broadcast_no_consumers.load(.acquire), + }) catch {}; + // validation failure breakdown by reason w.print( \\# TYPE relay_validation_failed counter @@ -1424,7 +1479,7 @@ test "concurrent broadcast queue + drain produces monotonic sequences" { bc_ptr.persist_order.store(0, .release); continue; }; - bc_ptr.broadcast_queue.push(seq, data); + bc_ptr.broadcast_queue.push(seq, data, &bc_ptr.stats); bc_ptr.persist_order.store(0, .release); } } diff --git a/src/frame_worker.zig b/src/frame_worker.zig index 87cdd9e..79d369c 100644 --- a/src/frame_worker.zig +++ b/src/frame_worker.zig @@ -278,9 +278,14 @@ pub fn processFrame(work: *FrameWork) void { // worker threads never touch consumer state directly. if (work.persist) |dp| { const relay_seq = blk: { + var spins: u64 = 0; while (work.bc.persist_order.cmpxchgWeak(0, 1, .acquire, .monotonic) != null) { + spins += 1; std.atomic.spinLoopHint(); } + if (spins > 0) { + _ = work.bc.stats.persist_order_spins.fetchAdd(spins, .monotonic); + } const seq = dp.persist(kind, uid, data) catch |err| { work.bc.persist_order.store(0, .release); @@ -294,7 +299,7 @@ pub fn processFrame(work: *FrameWork) void { work.bc.persist_order.store(0, .release); return; }; - work.bc.broadcast_queue.push(seq, owned); + work.bc.broadcast_queue.push(seq, owned, &work.bc.stats); work.bc.persist_order.store(0, .release); break :blk seq; }; @@ -321,7 +326,7 @@ pub fn processFrame(work: *FrameWork) void { } else { const upstream_seq = payload.getUint("seq") orelse 0; const owned = work.allocator.dupe(u8, data) catch return; - work.bc.broadcast_queue.push(upstream_seq, owned); + work.bc.broadcast_queue.push(upstream_seq, owned, &work.bc.stats); } } diff --git a/src/subscriber.zig b/src/subscriber.zig index d239ce7..b2a0194 100644 --- a/src/subscriber.zig +++ b/src/subscriber.zig @@ -713,9 +713,14 @@ const FrameHandler = struct { // ordering spinlock ensures frames are persisted + enqueued in seq order. if (sub.persist) |dp| { const relay_seq = blk: { + var spins: u64 = 0; while (sub.bc.persist_order.cmpxchgWeak(0, 1, .acquire, .monotonic) != null) { + spins += 1; std.atomic.spinLoopHint(); } + if (spins > 0) { + _ = sub.bc.stats.persist_order_spins.fetchAdd(spins, .monotonic); + } const seq = dp.persist(kind, uid, data) catch |err| { sub.bc.persist_order.store(0, .release); @@ -728,7 +733,7 @@ const FrameHandler = struct { sub.bc.persist_order.store(0, .release); return; }; - sub.bc.broadcast_queue.push(seq, owned); + sub.bc.broadcast_queue.push(seq, owned, &sub.bc.stats); sub.bc.persist_order.store(0, .release); break :blk seq; }; @@ -755,7 +760,7 @@ const FrameHandler = struct { } else { const seq = upstream_seq orelse 0; const owned = sub.allocator.dupe(u8, data) catch return; - sub.bc.broadcast_queue.push(seq, owned); + sub.bc.broadcast_queue.push(seq, owned, &sub.bc.stats); } } -- 2.51.2