diff --git a/src/host_ops.zig b/src/host_ops.zig index 45d28c8..e976c77 100644 --- a/src/host_ops.zig +++ b/src/host_ops.zig @@ -1,12 +1,15 @@ -//! host ops queue — MPSC ring buffer for cross-Io host bookkeeping +//! host ops — cross-Io host bookkeeping without pg.Pool from Evented fibers //! -//! subscriber fibers (Evented) push host operations into this queue. -//! a single background thread (std.Thread on pool_io / Threaded) pops and -//! executes them against DiskPersist (pg.Pool). this avoids Evented fibers -//! calling pg.Pool.acquire(), which invokes a Threaded futex → NULL -//! Thread.current() → heap corruption. +//! two mechanisms: //! -//! same atomic spinlock MPSC pattern as BroadcastQueue (broadcaster.zig). +//! 1. CursorMap — coalescing cursor store. subscribers atomically write their +//! latest seq (nanoseconds, no lock). the worker thread sweeps every 5s and +//! batch-flushes only changed cursors to the DB. this replaces the old FIFO +//! cursor queue that caused ~350 individual UPDATEs/sec and busy-spin on full. +//! +//! 2. HostOpsQueue — MPSC ring for rare ops (increment_failures, reset_failures, +//! update_status). same atomic spinlock pattern as BroadcastQueue. these ops +//! are infrequent (connect/disconnect/error), so the queue never fills. const std = @import("std"); const event_log_mod = @import("event_log.zig"); @@ -14,20 +17,124 @@ const event_log_mod = @import("event_log.zig"); const Io = std.Io; const log = std.log.scoped(.relay); +// --------------------------------------------------------------------------- +// cursor coalescing map +// --------------------------------------------------------------------------- + +pub const CursorMap = struct { + const MAX_SLOTS = 4096; + /// sentinel: slot is inactive (freed or never used) — worker skips these + const EMPTY: u64 = 0; + + host_ids: [MAX_SLOTS]u64 = @splat(0), + seqs: [MAX_SLOTS]std.atomic.Value(u64) = @splat(.{ .raw = EMPTY }), + /// last value successfully written to DB (worker-thread only, no atomics) + last_flushed: [MAX_SLOTS]u64 = @splat(0), + /// high-water mark — slots [0..count) have been used at least once + count: std.atomic.Value(u32) = .{ .raw = 0 }, + + /// free list for slot reuse (protected by free_lock spinlock) + free_stack: [MAX_SLOTS]u32 = undefined, + free_top: u32 = 0, + free_lock: std.atomic.Value(u32) = .{ .raw = 0 }, + + /// register a host and return its slot index. called at spawn time (slurper). + pub fn register(self: *CursorMap, host_id: u64, initial_seq: u64) ?u32 { + // try reusing a freed slot first + const slot = self.popFree() orelse blk: { + const s = self.count.fetchAdd(1, .acq_rel); + if (s >= MAX_SLOTS) { + _ = self.count.fetchSub(1, .monotonic); + log.warn("cursor_map: out of slots (max {d})", .{MAX_SLOTS}); + return null; + } + break :blk s; + }; + // write host_id and last_flushed before making the slot visible via seqs + self.host_ids[slot] = host_id; + self.last_flushed[slot] = initial_seq; + // release: makes host_ids write visible to worker thread's acquire load + self.seqs[slot].store(initial_seq, .release); + return slot; + } + + /// return a slot to the free list. called when a subscriber exits (slurper). + pub fn unregister(self: *CursorMap, slot: u32) void { + // mark inactive first — worker sees EMPTY on next sweep and skips + self.seqs[slot].store(EMPTY, .release); + self.host_ids[slot] = 0; + self.pushFree(slot); + } + + /// update cursor (called from subscriber — any Io context). single atomic store. + pub fn update(self: *CursorMap, slot: u32, seq: u64) void { + self.seqs[slot].store(seq, .release); + } + + /// sweep all slots, flush changed cursors to DB. called from worker thread only. + fn flush(self: *CursorMap, persist: *event_log_mod.DiskPersist) void { + const n = self.count.load(.acquire); + var flushed: u32 = 0; + for (0..n) |i| { + const current = self.seqs[i].load(.acquire); + // EMPTY means freed or never-used slot — skip + if (current == EMPTY) continue; + if (current != self.last_flushed[i]) { + persist.updateHostSeq(self.host_ids[i], current) catch |err| { + log.debug("cursor flush failed for host_id={d}: {s}", .{ self.host_ids[i], @errorName(err) }); + continue; + }; + self.last_flushed[i] = current; + flushed += 1; + } + } + if (flushed > 0) { + log.debug("cursor_map: flushed {d}/{d} cursors", .{ flushed, n }); + } + } + + fn popFree(self: *CursorMap) ?u32 { + self.lockFree(); + defer self.unlockFree(); + if (self.free_top == 0) return null; + self.free_top -= 1; + return self.free_stack[self.free_top]; + } + + fn pushFree(self: *CursorMap, slot: u32) void { + self.lockFree(); + defer self.unlockFree(); + self.free_stack[self.free_top] = slot; + self.free_top += 1; + } + + fn lockFree(self: *CursorMap) void { + while (self.free_lock.cmpxchgWeak(0, 1, .acquire, .monotonic) != null) { + std.atomic.spinLoopHint(); + } + } + + fn unlockFree(self: *CursorMap) void { + self.free_lock.store(0, .release); + } +}; + +// --------------------------------------------------------------------------- +// rare ops queue (failures, status) +// --------------------------------------------------------------------------- + pub const HostOp = struct { host_id: u64, kind: Kind, payload: Payload, pub const Kind = enum { - flush_cursor, increment_failures, reset_failures, update_status, }; pub const Payload = union { - seq: u64, /// pointer to subscriber's host_shutdown atomic — set by worker on exhaustion host_shutdown: *std.atomic.Value(bool), none: void, @@ -54,6 +161,7 @@ pub const HostOp = struct { pub const HostOpsQueue = struct { const CAPACITY = 4096; + const CURSOR_FLUSH_INTERVAL_SEC: i64 = 5; items: [CAPACITY]HostOp = undefined, head: std.atomic.Value(u32) = .{ .raw = 0 }, @@ -61,11 +169,12 @@ pub const HostOpsQueue = struct { push_lock: std.atomic.Value(u32) = .{ .raw = 0 }, persist: *event_log_mod.DiskPersist, + cursor_map: *CursorMap, shutdown: *std.atomic.Value(bool), max_consecutive_failures: u32, - /// push an op (called from any Io context — Evented fibers or Threaded threads). - /// spins until space is available. + /// push a rare op (called from any Io context). spins until space is available. + /// only used for failures/status — cursors go through CursorMap. pub fn push(self: *HostOpsQueue, op: HostOp) void { while (true) { // acquire spinlock @@ -100,22 +209,33 @@ pub const HostOpsQueue = struct { } /// worker thread entry point. runs on pool_io (Threaded). - /// pops ops and executes them against DiskPersist. + /// drains rare ops immediately, sweeps cursor map every 5s. pub fn run(self: *HostOpsQueue, pool_io: Io) void { + var last_cursor_flush: i64 = timestamp(pool_io); + while (!self.shutdown.load(.acquire)) { + // drain rare ops (failures, status) var drained: u32 = 0; while (self.pop()) |op| { self.execute(op); drained += 1; } + // periodic cursor sweep + const now = timestamp(pool_io); + if (now - last_cursor_flush >= CURSOR_FLUSH_INTERVAL_SEC) { + self.cursor_map.flush(self.persist); + last_cursor_flush = now; + } + if (drained == 0) { - // nothing to do — brief sleep via pool_io (Threaded, safe from std.Thread) - pool_io.sleep(Io.Duration.fromMilliseconds(10), .awake) catch return; + // brief sleep — rare ops are infrequent, cursor sweep is timer-driven + pool_io.sleep(Io.Duration.fromMilliseconds(100), .awake) catch return; } } - // drain remaining ops on shutdown + // final cursor flush + drain remaining ops on shutdown + self.cursor_map.flush(self.persist); while (self.pop()) |op| { self.execute(op); } @@ -123,11 +243,6 @@ pub const HostOpsQueue = struct { fn execute(self: *HostOpsQueue, op: HostOp) void { switch (op.kind) { - .flush_cursor => { - self.persist.updateHostSeq(op.host_id, op.payload.seq) catch |err| { - log.debug("host_ops: cursor flush failed for host_id={d}: {s}", .{ op.host_id, @errorName(err) }); - }; - }, .increment_failures => { const failures = self.persist.incrementHostFailures(op.host_id) catch 0; if (failures >= self.max_consecutive_failures) { @@ -148,4 +263,8 @@ pub const HostOpsQueue = struct { }, } } + + fn timestamp(io: Io) i64 { + return @intCast(@divFloor(Io.Timestamp.now(io, .real).nanoseconds, std.time.ns_per_s)); + } }; diff --git a/src/main.zig b/src/main.zig index 8418e32..50cc62d 100644 --- a/src/main.zig +++ b/src/main.zig @@ -261,10 +261,14 @@ pub fn main() !void { try resyncer.start(); defer resyncer.deinit(); - // init host ops queue — subscriber fibers (Evented) push DB ops here, - // a background thread (Threaded) executes them. avoids cross-Io pg.Pool access. + // init cursor map + host ops queue — subscriber fibers (Evented) write + // cursor seqs to the coalescing map (atomic store, no lock) and push rare + // ops (failures, status) to the MPSC queue. a background thread (Threaded) + // sweeps cursors every 5s and drains rare ops immediately. + var cursor_map: host_ops_mod.CursorMap = .{}; var host_ops_queue: host_ops_mod.HostOpsQueue = .{ .persist = &dp, + .cursor_map = &cursor_map, .shutdown = &shutdown_flag, .max_consecutive_failures = 15, }; @@ -294,6 +298,7 @@ pub fn main() !void { slurper.collection_index = &ci; slurper.resyncer = &resyncer; slurper.host_ops = &host_ops_queue; + slurper.cursor_map = &cursor_map; // start: loads active hosts from DB, spawns subscriber threads try slurper.start(); diff --git a/src/slurper.zig b/src/slurper.zig index 9d99f39..aa49666 100644 --- a/src/slurper.zig +++ b/src/slurper.zig @@ -230,6 +230,7 @@ pub const Slurper = struct { collection_index: ?*collection_index_mod.CollectionIndex = null, resyncer: ?*resync_mod.Resyncer = null, host_ops: ?*host_ops_mod.HostOpsQueue = null, + cursor_map: ?*host_ops_mod.CursorMap = null, shutdown: *std.atomic.Value(bool), options: Options, @@ -489,6 +490,10 @@ pub const Slurper = struct { if (self.frame_pool) |*fp| sub.pool = fp; sub.pool_io = self.pool_io; sub.host_ops = self.host_ops; + if (self.cursor_map) |cm| { + sub.cursor_map = cm; + sub.cursor_slot = cm.register(host_id, last_seq); + } if (last_seq > 0) sub.last_upstream_seq = last_seq; const future = try self.io.concurrent(runWorker, .{ self, host_id, sub }); @@ -506,6 +511,13 @@ pub const Slurper = struct { fn runWorker(self: *Slurper, host_id: u64, sub: *subscriber_mod.Subscriber) void { sub.run(); + // return cursor slot to free list before destroying subscriber + if (self.cursor_map) |cm| { + if (sub.cursor_slot) |slot| { + cm.unregister(slot); + } + } + // subscriber returned — remove from active workers self.workers_mutex.lockUncancelable(self.io); defer self.workers_mutex.unlock(self.io); diff --git a/src/subscriber.zig b/src/subscriber.zig index 9153db2..d239ce7 100644 --- a/src/subscriber.zig +++ b/src/subscriber.zig @@ -211,8 +211,11 @@ pub const Subscriber = struct { pool: ?*frame_worker_mod.FramePool = null, /// dedicated Threaded io for frame workers — safe from plain OS threads pool_io: ?Io = null, - /// host ops queue — pushes DB ops to a background thread (avoids cross-Io pg.Pool access) + /// host ops queue — pushes rare DB ops to a background thread (avoids cross-Io pg.Pool access) host_ops: ?*host_ops_mod.HostOpsQueue = null, + /// coalescing cursor map — subscriber writes latest seq atomically, worker sweeps every 5s + cursor_map: ?*host_ops_mod.CursorMap = null, + cursor_slot: ?u32 = null, shutdown: *std.atomic.Value(bool), last_upstream_seq: ?u64 = null, last_cursor_flush: i64 = 0, @@ -296,16 +299,14 @@ pub const Subscriber = struct { } } - /// flush cursor position to the host table (via host_ops queue — avoids cross-Io pg.Pool) + /// flush cursor position (atomic store to coalescing map — worker sweeps every 5s) fn flushCursor(self: *Subscriber) void { if (self.options.host_id == 0) return; const seq = self.last_upstream_seq orelse return; - if (self.host_ops) |hq| { - hq.push(.{ - .host_id = self.options.host_id, - .kind = .flush_cursor, - .payload = .{ .seq = seq }, - }); + if (self.cursor_map) |cm| { + if (self.cursor_slot) |slot| { + cm.update(slot, seq); + } } }