diff --git a/docs/incident-2026-05-31-seq-rewind.md b/docs/incident-2026-05-31-seq-rewind.md new file mode 100644 index 0000000..00ecc94 --- /dev/null +++ b/docs/incident-2026-05-31-seq-rewind.md @@ -0,0 +1,91 @@ +# seq rewind on unclean restart — relay_seq goes backwards + +**Punchline:** zlay's event log is a Zig port of indigo's `diskpersist.go`, but it broke +indigo's one load-bearing invariant: **an event is broadcast to consumers only after it is +durably on disk.** zlay moved the broadcast out of the flush and into the persist callers, so +it broadcast `relay_seq` up to one flush interval *before* writing it. An unclean kill (OOM / +crashloop) then lost the unflushed tail, and on restart the seq counter resumed *below* seq +values already broadcast — consumers saw `seq` go backwards. **Fixed** by restoring indigo's +broadcast-after-durable-flush ordering. + +## evidence (downstream consumer, 00:51:22–23Z 2026-05-31) + +`prev` = consumer's last-seen seq, `seq` = incoming. `prev` pinned at the pre-crash +high-water (~…8994xx) while incoming climbs monotonically from ~2,953 behind: + +``` +prev 7161747611348899444 seq 7161747611348896501 +... +prev 7161747611348899454 seq 7161747611348899032 +``` + +One-time rewind then monotonic climb (a counter rewind, not an ordering race), surfaced on a +cursor reconnect after one of several overnight restarts. The ~2,953-frame gap is ~7× the +400-event flush threshold — see "why the gap exceeded the flush bound" below. + +## root cause: a divergence from indigo + +Indigo broadcasts from inside the flush, after the durable write, and nowhere else +(`cmd/relay/stream/persist/diskpersist/diskpersist.go:371`): + +```go +func (dp *DiskPersistence) flushLog(ctx) error { + io.Copy(dp.logfi, dp.outbuf) // 1. write batch to disk + for _, ej := range dp.evtbuf { + dp.broadcast(ej.Evt) // 2. THEN broadcast + } +} +``` + +`broadcast` is wired once via `SetEventBroadcaster` (event_manager.go:36). So broadcast +high-water ≤ durable boundary, always; on restart `resumeLog` sets `curSeq = lastDurable+1 = +lastBroadcast+1`, and a rewind is structurally impossible. + +zlay had relocated the broadcast into the persist callers (`frame_worker`, `subscriber`, +`host_ops`): each called `persist()`, got the seq, resequenced, and pushed to the broadcast +queue *before* `flushLocked` wrote to disk. `flushLocked` did not broadcast at all. That single +relocation is the entire defect. + +## the fix (implemented) + +Restored indigo's contract: + +- `DiskPersist` gained `setBroadcaster(ctx, fn)` and broadcasts each event **from + `flushLocked`, after the durable write**, in seq order (`evtbuf` is append-ordered). +- The three producers no longer broadcast or resequence; they just `persist()`. seq order is + preserved by `persist()`'s own mutex, so the per-producer `persist_order` spinlock that + guarded broadcast ordering is no longer needed in the hot path (field kept for now; see + design notes). +- `broadcaster.broadcastPersisted` is the registered callback: it resequences to `relay_seq` + and enqueues for fan-out, invoked only post-flush. + +Regression test: `event_log.zig` "broadcast is gated on durable flush" — persists below the +flush threshold, asserts **zero** broadcasts, flushes, then asserts all events broadcast once +in seq order. (Also fixed the DB test harness to give each test an isolated database, since the +suite's runner executes tests concurrently against one server — they were contaminating each +other through the shared relay tables.) + +## still to harden: why the gap exceeded the flush bound + +The 2,953-frame gap is larger than the 400-event/100ms flush window should allow, which points +at a *second*, recovery-side weakness that the fix above does not fully close: + +- **`scanForLastSeq` under-recovers silently.** zlay's scan advances `pos += header_size + + hdr.len` with **no magic, no checksum, no length validation** (`event_log.zig`), and on a + short read just `break`s and returns whatever it reached. A torn record (what an interrupted + flush leaves) desyncs the walk and truncates the recovered tail by an arbitrary amount, + compounding across crashloops. Indigo, by contrast, **fails loud** here + (`return fmt.Errorf("did not seek to next event properly")`) and `readHeader` errors on a + partial record — so a bad recovery refuses to start instead of silently rewinding. + +Follow-ups (separate change): +1. add a per-record magic + length bound (ideally CRC) so `scanForLastSeq` detects a torn + record and truncates to the last valid one, or fails loud like indigo. +2. fsync + ordering on flush/rotation so the recovery boundary can't name a non-durable tail. + +## status + +The broadcast-ordering fix is committed and the suite is green (incl. DB tests against an +isolated postgres). The downstream consumer (tap) is on `relay1.us-east.bsky.network` and +healthy, so redeploying zlay is not urgent — but with this fix plus the recovery-hardening +follow-ups, zlay can be a trustworthy relay again. diff --git a/src/broadcaster.zig b/src/broadcaster.zig index a81ed54..1ba6bf1 100644 --- a/src/broadcaster.zig +++ b/src/broadcaster.zig @@ -280,6 +280,22 @@ pub fn resequenceFrame(allocator: Allocator, data: []const u8, relay_seq: u64) ? return result; } +/// callback registered with DiskPersist.setBroadcaster. resequences a persisted +/// frame to its relay seq and enqueues it for fan-out. invoked from flushLocked, +/// i.e. ONLY after the event is durably on disk — so an enqueued seq can never +/// exceed the on-disk seq. `payload` is owned by the persister and freed right +/// after this returns, so anything enqueued must be a fresh allocation. +pub fn broadcastPersisted(ctx: *anyopaque, seq: u64, payload: []const u8) void { + const self: *Broadcaster = @ptrCast(@alignCast(ctx)); + if (resequenceFrame(self.allocator, payload, seq)) |reseq| { + self.broadcast_queue.push(seq, reseq, &self.stats); + } else { + // decode failed — broadcast the raw frame verbatim (must be owned). + const owned = self.allocator.dupe(u8, payload) catch return; + self.broadcast_queue.push(seq, owned, &self.stats); + } +} + // --- broadcast queue (worker → fiber handoff) --- /// item produced by frame workers, consumed by the broadcaster fiber. diff --git a/src/event_log.zig b/src/event_log.zig index d9bdec6..c55c613 100644 --- a/src/event_log.zig +++ b/src/event_log.zig @@ -234,6 +234,13 @@ pub const DiskPersist = struct { flush_future: ?Io.Future(void) = null, alive: std.atomic.Value(bool) = .{ .raw = true }, + // live broadcaster, invoked from flushLocked AFTER the durable write (indigo's + // SetEventBroadcaster contract). null until wired in main. broadcasting only + // post-flush guarantees a broadcast seq never exceeds the on-disk seq, so the + // seq recovered on restart can't trail one already sent to consumers. + broadcast_ctx: ?*anyopaque = null, + broadcast_fn: ?*const fn (ctx: *anyopaque, seq: u64, payload: []const u8) void = null, + io: Io, /// last successful DB interaction (epoch seconds, set by Threaded workers). @@ -274,6 +281,14 @@ pub const DiskPersist = struct { return self.outbuf.capacity; } + /// register the live broadcaster. mirrors indigo's SetEventBroadcaster: + /// events are handed to it ONLY from flushLocked, after the durable write. + /// must be called before any producers start (no events in flight). + pub fn setBroadcaster(self: *DiskPersist, ctx: *anyopaque, f: *const fn (ctx: *anyopaque, seq: u64, payload: []const u8) void) void { + self.broadcast_ctx = ctx; + self.broadcast_fn = f; + } + pub fn init(allocator: Allocator, dir_path: []const u8, database_url: []const u8, db_pool_size: u16, io: Io) !DiskPersist { // ensure directory exists try Io.Dir.cwd().createDirPath(io, dir_path); @@ -1211,11 +1226,18 @@ pub const DiskPersist = struct { }; self.current_file_pos += self.outbuf.items.len; - // clear buffers + // clear write buffer (bytes are now durable) self.outbuf.clearRetainingCapacity(); - // free job data + // broadcast each event AFTER its durable write, in seq order (evtbuf is + // append-ordered = seq-ordered). this is the invariant that prevents the + // post-crash seq rewind: consumers never see a seq that isn't on disk. + // payload is the original frame (data after the fixed header); the + // broadcaster resequences it to job.seq. for (self.evtbuf.items) |job| { + if (self.broadcast_fn) |bf| { + bf(self.broadcast_ctx.?, job.seq, job.data[header_size..]); + } self.allocator.free(job.data); } self.evtbuf.clearRetainingCapacity(); @@ -1428,8 +1450,57 @@ fn requireDatabaseUrl() ![]const u8 { return getenv("DATABASE_URL") orelse return error.SkipZigTest; } +/// per-test database isolation. the test binary's runner executes tests +/// concurrently against one server, and the relay tables (log_file_refs / +/// account / host) are global — so sharing one database lets tests see each +/// other's rows (playback(0) returns another test's events). each test gets its +/// own freshly-created database instead, dropped on teardown. unique across +/// processes (pid) and within one (counter). +var test_db_counter: std.atomic.Value(u32) = .{ .raw = 0 }; + +const TestDb = struct { + base_url: []const u8, + name: []u8, + url: []u8, + + fn create(base_url: []const u8) !TestDb { + const pid: u64 = @intCast(std.c.getpid()); + const n = test_db_counter.fetchAdd(1, .monotonic); + const name = try std.fmt.allocPrint(std.testing.allocator, "zlay_test_{d}_{d}", .{ pid, n }); + errdefer std.testing.allocator.free(name); + + { + const uri = std.Uri.parse(base_url) catch return error.InvalidDatabaseUrl; + const pool = try pg.Pool.initUri(std.testing.allocator, std.testing.io, uri, .{ .size = 1 }); + defer pool.deinit(); + var buf: [128]u8 = undefined; + // name is generated from pid/counter — no injection risk + const stmt = try std.fmt.bufPrint(&buf, "CREATE DATABASE \"{s}\"", .{name}); + _ = try pool.exec(stmt, .{}); + } + + const slash = std.mem.lastIndexOfScalar(u8, base_url, '/') orelse return error.InvalidDatabaseUrl; + const url = try std.fmt.allocPrint(std.testing.allocator, "{s}/{s}", .{ base_url[0..slash], name }); + return .{ .base_url = base_url, .name = name, .url = url }; + } + + fn destroy(self: *TestDb) void { + const uri = std.Uri.parse(self.base_url) catch return; + const pool = pg.Pool.initUri(std.testing.allocator, std.testing.io, uri, .{ .size = 1 }) catch return; + defer pool.deinit(); + var buf: [128]u8 = undefined; + const stmt = std.fmt.bufPrint(&buf, "DROP DATABASE IF EXISTS \"{s}\" WITH (FORCE)", .{self.name}) catch return; + _ = pool.exec(stmt, .{}) catch {}; + std.testing.allocator.free(self.name); + std.testing.allocator.free(self.url); + } +}; + test "persist and playback" { - const database_url = try requireDatabaseUrl(); + 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(); @@ -1471,8 +1542,62 @@ test "persist and playback" { try std.testing.expectEqualStrings("payload-three", entries.items[2].data); } +test "broadcast is gated on durable flush" { + // regression for the seq-rewind incident (docs/incident-2026-05-31-seq-rewind.md): + // an event must NOT reach the broadcaster until it is durably written, so the + // seq recovered on restart can never trail an already-broadcast seq. matches + // indigo, which broadcasts only from flushLog. + 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(); + + const Recorder = struct { + seqs: std.ArrayListUnmanaged(u64) = .empty, + fn cb(ctx: *anyopaque, seq: u64, payload: []const u8) void { + _ = payload; + const self: *@This() = @ptrCast(@alignCast(ctx)); + self.seqs.append(std.testing.allocator, seq) catch {}; + } + }; + var rec = Recorder{}; + defer rec.seqs.deinit(std.testing.allocator); + dp.setBroadcaster(&rec, Recorder.cb); + // unwire before rec is destroyed so deinit's final flush can't call back + defer dp.broadcast_fn = null; + + // persist below the flush threshold — nothing broadcast yet + _ = try dp.persist(.commit, 1, "a"); + _ = try dp.persist(.commit, 2, "b"); + _ = try dp.persist(.commit, 3, "c"); + try std.testing.expectEqual(@as(usize, 0), rec.seqs.items.len); + + // flush → all three broadcast exactly once, in seq order + { + dp.mutex.lockUncancelable(dp.io); + defer dp.mutex.unlock(dp.io); + try dp.flushLocked(); + } + try std.testing.expectEqual(@as(usize, 3), rec.seqs.items.len); + try std.testing.expectEqual(@as(u64, 1), rec.seqs.items[0]); + try std.testing.expectEqual(@as(u64, 2), rec.seqs.items[1]); + try std.testing.expectEqual(@as(u64, 3), rec.seqs.items[2]); +} + test "playback with cursor" { - const database_url = try requireDatabaseUrl(); + 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(); @@ -1506,7 +1631,10 @@ test "playback with cursor" { } test "seq recovery after reinit" { - const database_url = try requireDatabaseUrl(); + 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(); @@ -1537,7 +1665,10 @@ test "seq recovery after reinit" { } test "takedown zeros payload" { - const database_url = try requireDatabaseUrl(); + 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(); @@ -1572,7 +1703,10 @@ test "takedown zeros payload" { } test "uidForDid assigns and caches UIDs" { - const database_url = try requireDatabaseUrl(); + 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(); @@ -1598,7 +1732,10 @@ test "uidForDid assigns and caches UIDs" { } test "uidForDid survives reinit" { - const database_url = try requireDatabaseUrl(); + 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(); @@ -1623,7 +1760,10 @@ test "uidForDid survives reinit" { } test "takedown with real UIDs" { - const database_url = try requireDatabaseUrl(); + 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(); diff --git a/src/frame_worker.zig b/src/frame_worker.zig index 76f44a1..ee1d328 100644 --- a/src/frame_worker.zig +++ b/src/frame_worker.zig @@ -274,34 +274,18 @@ pub fn processFrame(work: *FrameWork) void { else .identity; - // persist + resequence + enqueue under ordering lock to guarantee - // broadcast_queue insertion order matches seq assignment order. + // persist the event. seq assignment and the live broadcast both happen inside + // DiskPersist now: persist() assigns seq under its mutex (so order is kept), and + // the broadcast fires from flushLocked AFTER the durable write — so a broadcast + // seq can never outrun the on-disk seq. see docs/incident-2026-05-31-seq-rewind.md. if (work.persist) |dp| { - 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 relay_seq = dp.persist(kind, uid, data) catch |err| { - work.bc.persist_order.store(0, .release); log.warn("persist failed: {s}", .{@errorName(err)}); return; }; work.bc.stats.relay_seq.store(relay_seq, .release); - const broadcast_data = broadcaster.resequenceFrame(alloc, data, relay_seq) orelse data; - const owned = work.allocator.dupe(u8, broadcast_data) catch { - work.bc.persist_order.store(0, .release); - return; - }; - work.bc.broadcast_queue.push(relay_seq, owned, &work.bc.stats); - work.bc.persist_order.store(0, .release); - - // update per-DID state outside the ordering lock (Postgres round-trip) + // update per-DID state (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| diff --git a/src/host_ops.zig b/src/host_ops.zig index 0dd1448..f2ef796 100644 --- a/src/host_ops.zig +++ b/src/host_ops.zig @@ -313,25 +313,12 @@ pub const HostOpsQueue = struct { const frame_bytes = td.frameSlice(); const bc = self.bc orelse return; - // persist the #account event under ordering lock - while (bc.persist_order.cmpxchgWeak(0, 1, .acquire, .monotonic) != null) { - std.atomic.spinLoopHint(); - } - + // persist the #account event. broadcast fires from flushLocked after the + // durable write (see docs/incident-2026-05-31-seq-rewind.md). if (self.persist.persist(.account, td.uid, frame_bytes)) |relay_seq| { bc.stats.relay_seq.store(relay_seq, .release); - - const broadcast_data = broadcaster.resequenceFrame(self.persist.allocator, frame_bytes, relay_seq) orelse frame_bytes; - const owned = self.persist.allocator.dupe(u8, broadcast_data) catch { - bc.persist_order.store(0, .release); - log.warn("host_ops: failed to alloc broadcast data for takedown uid={d}", .{td.uid}); - return; - }; - bc.broadcast_queue.push(relay_seq, owned, &bc.stats); - bc.persist_order.store(0, .release); log.info("host_ops: emitted #account takedown for uid={d} (seq={d})", .{ td.uid, relay_seq }); } else |err| { - bc.persist_order.store(0, .release); log.warn("host_ops: persist #account takedown failed: {s}", .{@errorName(err)}); } } diff --git a/src/main.zig b/src/main.zig index b971aac..a3f0436 100644 --- a/src/main.zig +++ b/src/main.zig @@ -259,6 +259,12 @@ pub fn main() !void { bc.db_queue = &db_queue; val.persist = &dp; + // wire the live broadcaster into the persister: events are broadcast from + // flushLocked AFTER their durable write (indigo's contract), not by the + // producers pre-flush. set before any producer starts. see + // docs/incident-2026-05-31-seq-rewind.md. + dp.setBroadcaster(&bc, broadcaster.broadcastPersisted); + // init collection index (RocksDB — inspired by lightrail/microcosm.blue) const ci_dir = getenv("COLLECTION_INDEX_DIR") orelse "data/collection-index"; var ci = collection_index_mod.CollectionIndex.open(allocator, ci_dir) catch |err| { diff --git a/src/subscriber.zig b/src/subscriber.zig index 36b7a7b..af9fa8f 100644 --- a/src/subscriber.zig +++ b/src/subscriber.zig @@ -738,34 +738,18 @@ const FrameHandler = struct { else // is_identity (unknown types already filtered above) .identity; - // persist + resequence + enqueue under ordering lock to guarantee - // broadcast_queue insertion order matches seq assignment order. + // persist the event. seq assignment + live broadcast both happen inside + // DiskPersist: persist() assigns seq under its mutex (order preserved), the + // broadcast fires from flushLocked AFTER the durable write — so a broadcast + // seq can never outrun the on-disk seq. see docs/incident-2026-05-31-seq-rewind.md. if (sub.persist) |dp| { - 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 relay_seq = dp.persist(kind, uid, data) catch |err| { - sub.bc.persist_order.store(0, .release); log.warn("persist failed: {s}", .{@errorName(err)}); return; }; sub.bc.stats.relay_seq.store(relay_seq, .release); - const broadcast_data = broadcaster.resequenceFrame(alloc, data, relay_seq) orelse data; - const owned = sub.allocator.dupe(u8, broadcast_data) catch { - sub.bc.persist_order.store(0, .release); - return; - }; - sub.bc.broadcast_queue.push(relay_seq, owned, &sub.bc.stats); - sub.bc.persist_order.store(0, .release); - - // update per-DID state outside the ordering lock (Postgres round-trip) + // update per-DID state (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|