diff --git a/src/internal/broadcaster.zig b/src/internal/broadcaster.zig index d3d029a..f8c6e01 100644 --- a/src/internal/broadcaster.zig +++ b/src/internal/broadcaster.zig @@ -417,11 +417,14 @@ pub const Consumer = struct { io: Io, /// push a shared frame to this consumer's send buffer. - /// acquires a reference. returns false if full (consumer too slow). + /// acquires a reference. returns false if full (consumer too slow) or + /// dead — a disconnected consumer must not keep accumulating frames + /// (replay for a flapping consumer compounded into an OOM otherwise). pub fn enqueue(self: *Consumer, frame: *SharedFrame) bool { self.mutex.lockUncancelable(self.io); defer self.mutex.unlock(self.io); + if (!self.alive.load(.acquire)) return false; if (self.buf_len == BUFFER_CAP) return false; frame.acquire(); @@ -522,6 +525,16 @@ pub const Consumer = struct { // --- broadcaster --- +/// per-chunk bound for disk replay: caps the transient allocation of one +/// consumer connect at ~16 MiB regardless of how much backlog its cursor +/// spans. worst-case per-frame size (5 MiB max_message_size) is still +/// covered: max_bytes is checked before each read, so a chunk holds at most +/// max_bytes + one frame. +const replay_chunk_limits: event_log_mod.PlaybackLimits = .{ + .max_entries = 4096, + .max_bytes = 16 * 1024 * 1024, +}; + const history_capacity = 50_000; const FrameHistory = ring_buffer.RingBuffer(history_capacity); @@ -694,49 +707,61 @@ pub const Broadcaster = struct { /// disk playback runs on pool_io (Threaded) via request/reply — playback() /// holds the DiskPersist mutex and reads files, all of which require Threaded Io. pub fn replayTo(self: *Broadcaster, consumer: *Consumer, cursor: u64) void { - // phase 1: disk replay from diskpersist via request/reply + // phase 1: disk replay from diskpersist via request/reply, in bounded + // chunks. materializing the whole backlog per connect is an + // O(retention) allocation — a flapping consumer stacking those + // replays OOMed the canary (2026-08-06). each chunk is capped, the + // consumer's liveness is re-checked between chunks, and a dead or + // slow consumer aborts the replay entirely. if (self.persist) |dp| { - var req: event_log_mod.PlaybackRequest = .{ .since = cursor, .allocator = self.allocator }; - dp.enqueuePlayback(&req); - - // poll until pool_io worker completes the request. - // SAFETY: req is stack-local — we MUST wait for the worker to finish before - // returning, otherwise the stack frame unwinds while the worker still holds &req. - // if sleep fails (shutdown/io error), fall back to spin-wait. the host_ops worker - // drains all pending playback requests before exiting, so this is bounded. - while (!req.done.load(.acquire)) { - self.io.sleep(Io.Duration.fromMicroseconds(100), .awake) catch { - while (!req.done.load(.acquire)) std.atomic.spinLoopHint(); - break; + var since = cursor; + var replayed: usize = 0; + while (consumer.alive.load(.acquire)) { + var req: event_log_mod.PlaybackRequest = .{ + .since = since, + .allocator = self.allocator, + .limits = replay_chunk_limits, }; - } + dp.enqueuePlayback(&req); + + // poll until pool_io worker completes the request. + // SAFETY: req is stack-local — we MUST wait for the worker to finish before + // returning, otherwise the stack frame unwinds while the worker still holds &req. + // if sleep fails (shutdown/io error), fall back to spin-wait. the host_ops worker + // drains all pending playback requests before exiting, so this is bounded. + while (!req.done.load(.acquire)) { + self.io.sleep(Io.Duration.fromMicroseconds(100), .awake) catch { + while (!req.done.load(.acquire)) std.atomic.spinLoopHint(); + break; + }; + } - if (req.err) |err| { - log.warn("disk replay failed: {s}, falling back to memory", .{@errorName(err)}); - // clean up any partial entries - for (req.entries.items) |e| self.allocator.free(e.data); - req.entries.deinit(self.allocator); - self.replayFromMemory(consumer, cursor); - return; - } + defer { + for (req.entries.items) |e| self.allocator.free(e.data); + req.entries.deinit(self.allocator); + } - defer { - for (req.entries.items) |e| self.allocator.free(e.data); - req.entries.deinit(self.allocator); - } + if (req.err) |err| { + log.warn("disk replay failed: {s}, falling back to memory", .{@errorName(err)}); + self.replayFromMemory(consumer, since); + return; + } - var replayed: usize = 0; - var reseq_arena = std.heap.ArenaAllocator.init(self.allocator); - defer reseq_arena.deinit(); - - for (req.entries.items) |entry| { - // resequence: replace upstream seq in CBOR with relay-assigned seq - const frame_data = resequenceFrame(reseq_arena.allocator(), entry.data, entry.seq) orelse entry.data; - if (!consumer.enqueueRaw(frame_data)) { - log.warn("replay buffer full after {d} frames", .{replayed}); - break; + if (req.entries.items.len == 0) break; // caught up + + var reseq_arena = std.heap.ArenaAllocator.init(self.allocator); + defer reseq_arena.deinit(); + + for (req.entries.items) |entry| { + // resequence: replace upstream seq in CBOR with relay-assigned seq + const frame_data = resequenceFrame(reseq_arena.allocator(), entry.data, entry.seq) orelse entry.data; + if (!consumer.enqueueRaw(frame_data)) { + log.warn("replay aborted after {d} frames (consumer dead or buffer full, cursor={d})", .{ replayed, cursor }); + return; + } + replayed += 1; + since = entry.seq; } - replayed += 1; } if (replayed > 0) { log.info("replayed {d} frames from disk (cursor={d})", .{ replayed, cursor }); @@ -1583,6 +1608,26 @@ test "frame history supports cursor replay" { try std.testing.expectEqual(@as(u64, 5), frames[1].seq); } +test "enqueue refuses frames for a dead consumer" { + // regression: replay for an abruptly-disconnected consumer kept + // accumulating frames, compounding into the 2026-08-06 canary OOM. + var consumer = Consumer{ + .conn = undefined, // enqueue never touches the connection + .allocator = std.testing.allocator, + .io = std.testing.io, + }; + + try std.testing.expect(consumer.enqueueRaw("frame1")); + + consumer.alive.store(false, .release); + try std.testing.expect(!consumer.enqueueRaw("frame2")); + + // drain the accepted frame + consumer.mutex.lockUncancelable(consumer.io); + while (consumer.dequeue()) |f| f.release(); + consumer.mutex.unlock(consumer.io); +} + test "shared frame ref counting" { const frame = try SharedFrame.create(std.testing.allocator, "hello"); frame.acquire(); // ref=2 diff --git a/src/internal/event_log.zig b/src/internal/event_log.zig index 42e1fc1..2f8656f 100644 --- a/src/internal/event_log.zig +++ b/src/internal/event_log.zig @@ -83,10 +83,21 @@ const PersistJob = struct { // --- disk persistence --- +/// bounds on a single playback call. a consumer connect with an old cursor +/// must never materialize the whole retained backlog in one shot — that is +/// an O(retention) allocation per connect, and a flapping consumer turns it +/// into an OOM (see docs/handoffs/HANDOFF-2026-08-06-consumer-flap-oom.md). +/// callers replay in bounded chunks instead, advancing `since` between calls. +pub const PlaybackLimits = struct { + max_entries: usize = std.math.maxInt(usize), + max_bytes: usize = std.math.maxInt(usize), +}; + /// playback request — caller posts, pool_io worker executes pub const PlaybackRequest = struct { since: u64, allocator: Allocator, + limits: PlaybackLimits = .{}, entries: std.ArrayListUnmanaged(PlaybackEntry) = .empty, err: ?anyerror = null, done: std.atomic.Value(bool) = .{ .raw = false }, @@ -961,8 +972,10 @@ pub const DiskPersist = struct { return seq; } - /// playback events with seq > since. calls cb for each event. - pub fn playback(self: *DiskPersist, since: u64, allocator: Allocator, entries: *std.ArrayListUnmanaged(PlaybackEntry)) !void { + /// playback events with seq > since, bounded by `limits`. a truncated + /// result is detected by the caller via entries hitting the limits; + /// resume by calling again with since = last returned seq. + pub fn playback(self: *DiskPersist, since: u64, allocator: Allocator, entries: *std.ArrayListUnmanaged(PlaybackEntry), limits: PlaybackLimits) !void { self.mutex.lockUncancelable(self.io); defer self.mutex.unlock(self.io); @@ -1004,11 +1017,13 @@ pub const DiskPersist = struct { defer for (start_files.items) |f| allocator.free(f.path); - // read events from each file + // read events from each file until the limits are hit + var bytes_used: usize = 0; for (start_files.items) |ref| { + if (entries.items.len >= limits.max_entries or bytes_used >= limits.max_bytes) break; var file = self.dir.openFile(self.io, ref.path, .{}) catch continue; defer file.close(self.io); - try readEventsFrom(allocator, file, self.io, since, entries); + try readEventsFrom(allocator, file, self.io, since, entries, limits, &bytes_used); } } @@ -1340,7 +1355,7 @@ const GcSizeRef = struct { // --- file-level operations --- -fn readEventsFrom(allocator: Allocator, file: Io.File, io: Io, since: u64, result: *std.ArrayListUnmanaged(PlaybackEntry)) !void { +fn readEventsFrom(allocator: Allocator, file: Io.File, io: Io, since: u64, result: *std.ArrayListUnmanaged(PlaybackEntry), limits: PlaybackLimits, bytes_used: *usize) !void { const file_size = (try file.stat(io)).size; // if since > 0, scan to the right position @@ -1351,6 +1366,7 @@ fn readEventsFrom(allocator: Allocator, file: Io.File, io: Io, since: u64, resul // read events while (pos + header_size <= file_size) { + if (result.items.len >= limits.max_entries or bytes_used.* >= limits.max_bytes) return; var hdr_buf: [header_size]u8 = undefined; const n = file.readPositionalAll(io, &hdr_buf, pos) catch break; if (n < header_size) break; @@ -1390,6 +1406,7 @@ fn readEventsFrom(allocator: Allocator, file: Io.File, io: Io, since: u64, resul allocator.free(data); break; }; + bytes_used.* += data.len; } } @@ -1705,7 +1722,7 @@ test "persist and playback" { for (entries.items) |e| std.testing.allocator.free(e.data); entries.deinit(std.testing.allocator); } - try dp.playback(0, std.testing.allocator, &entries); + try dp.playback(0, std.testing.allocator, &entries, .{}); try std.testing.expectEqual(@as(usize, 3), entries.items.len); try std.testing.expectEqualStrings("payload-one", entries.items[0].data); @@ -1783,13 +1800,69 @@ test "playback with cursor" { for (entries.items) |e| std.testing.allocator.free(e.data); entries.deinit(std.testing.allocator); } - try dp.playback(2, std.testing.allocator, &entries); + try dp.playback(2, std.testing.allocator, &entries, .{}); try std.testing.expectEqual(@as(usize, 1), entries.items.len); try std.testing.expectEqual(@as(u64, 3), entries.items[0].seq); try std.testing.expectEqualStrings("c", entries.items[0].data); } +test "playback respects entry and byte limits, resumable by cursor" { + // regression for the 2026-08-06 consumer-flap OOM: an old cursor must + // not materialize the whole backlog in one playback call. + const base_url = try requireDatabaseUrl(); + var tdb = try TestDb.create(base_url); + defer tdb.destroy(); + const database_url = tdb.url; + + var tmp = std.testing.tmpDir(.{}); + defer tmp.cleanup(); + + const dir_path = try tmpDirRealPath(std.testing.allocator, tmp); + defer std.testing.allocator.free(dir_path); + + var dp = try DiskPersist.init(std.testing.allocator, dir_path, database_url, 5, std.testing.io); + defer dp.deinit(); + + for (0..10) |_| _ = try dp.persist(.commit, 1, "0123456789"); + { + dp.mutex.lockUncancelable(dp.io); + defer dp.mutex.unlock(dp.io); + try dp.flushLocked(); + } + + // max_entries caps the chunk; resuming from the last seq gets the rest + var chunk: std.ArrayListUnmanaged(PlaybackEntry) = .empty; + defer { + for (chunk.items) |e| std.testing.allocator.free(e.data); + chunk.deinit(std.testing.allocator); + } + try dp.playback(0, std.testing.allocator, &chunk, .{ .max_entries = 4 }); + try std.testing.expectEqual(@as(usize, 4), chunk.items.len); + try std.testing.expectEqual(@as(u64, 4), chunk.items[3].seq); + + var rest: std.ArrayListUnmanaged(PlaybackEntry) = .empty; + defer { + for (rest.items) |e| std.testing.allocator.free(e.data); + rest.deinit(std.testing.allocator); + } + try dp.playback(chunk.items[3].seq, std.testing.allocator, &rest, .{ .max_entries = 100 }); + try std.testing.expectEqual(@as(usize, 6), rest.items.len); + try std.testing.expectEqual(@as(u64, 5), rest.items[0].seq); + try std.testing.expectEqual(@as(u64, 10), rest.items[5].seq); + + // max_bytes caps the chunk: 10-byte payloads, 25-byte budget → the + // check runs before each read, so the 3rd entry crosses the budget + // and the 4th is never read + var byte_chunk: std.ArrayListUnmanaged(PlaybackEntry) = .empty; + defer { + for (byte_chunk.items) |e| std.testing.allocator.free(e.data); + byte_chunk.deinit(std.testing.allocator); + } + try dp.playback(0, std.testing.allocator, &byte_chunk, .{ .max_bytes = 25 }); + try std.testing.expectEqual(@as(usize, 3), byte_chunk.items.len); +} + test "seq recovery after reinit" { const base_url = try requireDatabaseUrl(); var tdb = try TestDb.create(base_url); @@ -1856,7 +1929,7 @@ test "takedown zeros payload" { for (entries.items) |e| std.testing.allocator.free(e.data); entries.deinit(std.testing.allocator); } - try dp.playback(0, std.testing.allocator, &entries); + try dp.playback(0, std.testing.allocator, &entries, .{}); try std.testing.expectEqual(@as(usize, 1), entries.items.len); try std.testing.expectEqualStrings("other-data", entries.items[0].data); @@ -2016,7 +2089,7 @@ test "takedown with real UIDs" { for (entries.items) |e| std.testing.allocator.free(e.data); entries.deinit(std.testing.allocator); } - try dp.playback(0, std.testing.allocator, &entries); + try dp.playback(0, std.testing.allocator, &entries, .{}); try std.testing.expectEqual(@as(usize, 1), entries.items.len); try std.testing.expectEqualStrings("bob-post", entries.items[0].data); diff --git a/src/internal/host_ops.zig b/src/internal/host_ops.zig index c3ad961..2363a4b 100644 --- a/src/internal/host_ops.zig +++ b/src/internal/host_ops.zig @@ -269,7 +269,7 @@ pub const HostOpsQueue = struct { // Treiber stack pops in LIFO order — fine for playback (each request is independent) while (maybe_batch) |req| { maybe_batch = req.next.load(.acquire); - self.persist.playback(req.since, req.allocator, &req.entries) catch |e| { + self.persist.playback(req.since, req.allocator, &req.entries, req.limits) catch |e| { req.err = e; }; req.done.store(true, .release);