From f202c9b8fcc8caff038f45a08619930161703903 Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Sun, 16 Aug 2026 22:59:27 -0500 Subject: [PATCH] sweep: enforce the exact (after_seq, before_seq] window per row upstream matcher parity (filter.go: the plan prunes whole blocks, the exact seq window is applied client-side per row). found live by the new cutover smoke against production: the sweep delivered ~2 whole blocks of rows below the requested resume point. adds the cutover smoke example (archive sweep -> gapless cutover -> live tail, contiguity asserted through one handler) and a deliverBlock regression test. Co-Authored-By: Claude Fable 5 --- build.zig | 18 +++++++++ examples/cutover_smoke.zig | 76 ++++++++++++++++++++++++++++++++++++++ src/archive_backfill.zig | 39 +++++++++++++++++++ 3 files changed, 133 insertions(+) create mode 100644 examples/cutover_smoke.zig diff --git a/build.zig b/build.zig index 9e06c36..0c63c32 100644 --- a/build.zig +++ b/build.zig @@ -97,4 +97,22 @@ pub fn build(b: *std.Build) void { if (b.args) |args| run_smoke.addArgs(args); const smoke_step = b.step("example-failover-smoke", "live smoke: dead primary rotates to a real fallback"); smoke_step.dependOn(&run_smoke.step); + + const cutover = b.addExecutable(.{ + .name = "example-cutover-smoke", + .root_module = b.createModule(.{ + .root_source_file = b.path("examples/cutover_smoke.zig"), + .target = target, + .optimize = optimize, + .link_libc = true, + .imports = &.{ + .{ .name = "jetstream", .module = mod }, + }, + }), + }); + b.installArtifact(cutover); + const run_cutover = b.addRunArtifact(cutover); + if (b.args) |args| run_cutover.addArgs(args); + const cutover_step = b.step("example-cutover-smoke", "live smoke: archive sweep -> gapless cutover -> live tail"); + cutover_step.dependOn(&run_cutover.step); } diff --git a/examples/cutover_smoke.zig b/examples/cutover_smoke.zig new file mode 100644 index 0000000..4d54767 --- /dev/null +++ b/examples/cutover_smoke.zig @@ -0,0 +1,76 @@ +//! live backfill→cutover smoke: sweep the REAL sealed archive from just +//! below the sealed tip, cut over to the live tail, and verify the seam. +//! +//! this is the production proof the loopback e2e can't give: planSnapshot +//! → getBlock → jss decode against the real (token-gated) archive, then +//! the live websocket continuing PAST the sealed tip, delivered through +//! one handler with every seq contiguous across the seam (unfiltered +//! subscribe sees every seq, so any gap or duplicate is loud). +//! +//! run (from a box that may consume the firehose; needs the archive key): +//! zig build example-cutover-smoke -- https://stream.waow.tech $API_KEY + +const std = @import("std"); +const jetstream_sdk = @import("jetstream"); + +const sweep_back: u64 = 2000; // rows of sealed archive to replay +const live_target: usize = 25; // live rows to observe past the sealed tip + +const SeamHandler = struct { + sealed_tip: u64, + first_seq: u64 = 0, + last_seq: u64 = 0, + total: usize = 0, + archive_rows: usize = 0, + live_rows: usize = 0, + gaps: usize = 0, + + pub fn onEvent(self: *@This(), event: jetstream_sdk.Event) bool { + if (self.total == 0) { + self.first_seq = event.seq; + } else if (event.seq != self.last_seq + 1) { + self.gaps += 1; + std.debug.print(" GAP/DUP: {d} -> {d}\n", .{ self.last_seq, event.seq }); + } + self.last_seq = event.seq; + self.total += 1; + if (event.seq <= self.sealed_tip) self.archive_rows += 1 else self.live_rows += 1; + if (self.live_rows == 1 and event.seq == self.sealed_tip + 1) + std.debug.print(" cutover seam crossed contiguously at seq {d}\n", .{event.seq}); + return self.live_rows < live_target; + } +}; + +pub fn main(init: std.process.Init) !void { + const allocator = init.gpa; + const args = try init.minimal.args.toSlice(init.arena.allocator()); + if (args.len != 3) { + std.debug.print("usage: example-cutover-smoke \n", .{}); + return error.MissingArgs; + } + const host = args[1]; + const api_key = args[2]; + + var arena_state = std.heap.ArenaAllocator.init(allocator); + defer arena_state.deinit(); + const spans = try jetstream_sdk.ArchiveBackfill.fetchSegmentSpans(init.io, allocator, arena_state.allocator(), host, api_key); + var sealed_tip: u64 = 0; + for (spans) |span| { + if (span.max_seq > sealed_tip) sealed_tip = span.max_seq; + } + if (sealed_tip == 0) return error.NoSealedArchive; + const after = sealed_tip - sweep_back; + std.debug.print("sealed tip {d}; sweeping (after_seq {d}, {d} sealed rows] then live\n\n", .{ sealed_tip, after, sweep_back }); + + var handler = SeamHandler{ .sealed_tip = sealed_tip }; + try jetstream_sdk.subscribe(init.io, allocator, .{ + .hosts = &.{host}, + .after_seq = after, + .api_key = api_key, + }, &handler); + + std.debug.print("\nseqs {d}..{d}: {d} archive rows + {d} live rows, {d} gaps\n", .{ handler.first_seq, handler.last_seq, handler.archive_rows, handler.live_rows, handler.gaps }); + if (handler.gaps != 0) return error.SeamNotContiguous; + if (handler.archive_rows == 0) return error.NoArchiveRows; + std.debug.print("smoke PASS: archive sweep -> gapless cutover -> live tail, one handler\n", .{}); +} diff --git a/src/archive_backfill.zig b/src/archive_backfill.zig index 8601d4a..42cba62 100644 --- a/src/archive_backfill.zig +++ b/src/archive_backfill.zig @@ -792,6 +792,12 @@ fn deliverBlock(arena: Allocator, raw: []const u8, options: Options, handler: an const witnessed_at_row = mem.readInt(i64, raw[wit_base + 8 * i ..][0..8], .little); const seq_row = mem.readInt(u64, raw[seq_base + 8 * i ..][0..8], .little); + // the plan prunes whole segments/blocks; the exact (after_seq, + // before_seq] window is the client's to enforce per row, exactly + // like upstream's matcher (filter.go: "applied on the client") + if (options.after_seq) |a| if (seq_row <= a) continue; + if (options.before_seq) |b| if (seq_row > b) continue; + // extended row handler (the unified client's archive half): every // matching row is delivered as a v2 event WITH its seq — including // identity/account/sync marker rows, which ride inline exactly like @@ -1315,3 +1321,36 @@ test "deliverSegment rejects bad magic" { &result, )); } + +test "deliverBlock enforces the exact (after_seq, before_seq] window per row" { + // upstream matcher parity (filter.go): the plan prunes whole blocks, + // but the exact seq window is the client's to apply on each row — + // found live by the cutover smoke, which received ~2 whole blocks of + // rows below the requested resume point + var arena_state = std.heap.ArenaAllocator.init(testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + + const raw = try buildTestBlock(arena, &.{ + .{ .seq = 10, .witnessed_at = 1000, .kind = kind_create, .collection = "c", .did = "did:plc:a", .rkey = "r1", .rev = "v1", .payload = &test_record_cbor }, + .{ .seq = 11, .witnessed_at = 1001, .kind = kind_create, .collection = "c", .did = "did:plc:a", .rkey = "r2", .rev = "v2", .payload = &test_record_cbor }, + .{ .seq = 12, .witnessed_at = 1002, .kind = kind_create, .collection = "c", .did = "did:plc:a", .rkey = "r3", .rev = "v3", .payload = &test_record_cbor }, + .{ .seq = 13, .witnessed_at = 1003, .kind = kind_create, .collection = "c", .did = "did:plc:a", .rkey = "r4", .rev = "v4", .payload = &test_record_cbor }, + }); + + const RowCollector = struct { + seqs: std.array_list.Managed(u64), + pub fn onRow(self: *@This(), event: livedecode.Event) bool { + self.seqs.append(event.seq) catch unreachable; + return true; + } + }; + var handler = RowCollector{ .seqs = .init(arena) }; + var result: Result = .{ .planned_through_seq = 0, .sealed_tip_seq = 0 }; + try deliverBlock(arena, raw, .{ + .host = "https://example.test", + .after_seq = 10, + .before_seq = 12, + }, &handler, &result); + try testing.expectEqualSlices(u64, &.{ 11, 12 }, handler.seqs.items); +} -- 2.51.2