diff --git a/scripts/archive_backfill_smoke.zig b/scripts/archive_backfill_smoke.zig index fb4c4fe..4e44312 100644 --- a/scripts/archive_backfill_smoke.zig +++ b/scripts/archive_backfill_smoke.zig @@ -8,20 +8,67 @@ pub fn main() !void { std.debug.print("archive backfill smoke: pollz collections from stream.waow.tech\n", .{}); + const io = std.Options.debug_io; var handler = Handler{}; - const result = try zat.ArchiveBackfill.run(std.Options.debug_io, allocator, .{ + const result = try zat.ArchiveBackfill.run(io, allocator, .{ .host = "https://stream.waow.tech", .collections = &.{ "tech.waow.pollz.poll", "tech.waow.pollz.vote" }, // a full plan spans tens of GB; the smoke only proves plan/fetch/decode - .max_blocks = 25, + .max_blocks = 15, }, &handler); std.debug.print( - "done: {d} blocks (truncated={}), {d} events delivered ({d} polls, {d} votes, {d} deletes), planned_through_seq={d}, sealed_tip_seq={d}, last_time_us={?d}\n", - .{ result.blocks_decoded, result.truncated, result.events_delivered, handler.polls, handler.votes, handler.deletes, result.planned_through_seq, result.sealed_tip_seq, result.last_time_us }, + "pass 1: {d} blocks (truncated={}), {d} events delivered ({d} polls, {d} votes, {d} deletes), position={?any}, last_time_us={?d}\n", + .{ result.blocks_decoded, result.truncated, result.events_delivered, handler.polls, handler.votes, handler.deletes, result.position, result.last_time_us }, ); if (result.events_delivered == 0) return error.NoEvents; if (handler.missing_record != 0) return error.MissingRecords; + + // resume from pass 1's position: must pick up strictly after it + const resumed = try zat.ArchiveBackfill.run(io, allocator, .{ + .host = "https://stream.waow.tech", + .collections = &.{ "tech.waow.pollz.poll", "tech.waow.pollz.vote" }, + .max_blocks = 10, + .start_after = result.position, + }, &handler); + std.debug.print("pass 2 (resumed): {d} blocks, {d} events, position={?any}\n", .{ resumed.blocks_decoded, resumed.events_delivered, resumed.position }); + const p1 = result.position.?; + const p2 = resumed.position orelse return error.NoResumeProgress; + if (p2.segment_index < p1.segment_index or + (p2.segment_index == p1.segment_index and p2.block_index <= p1.block_index)) + return error.ResumeWentBackwards; + + // checksum verification on a real sealed segment (smallest in the archive) + try verifySmallestSegment(io, allocator); +} + +fn verifySmallestSegment(io: std.Io, allocator: std.mem.Allocator) !void { + var transport = zat.HttpTransport.init(io, allocator); + defer transport.deinit(); + + var listed = try transport.fetch(.{ .url = "https://stream.waow.tech/xrpc/network.bsky.jetstream.listSegments?limit=1000" }); + defer listed.deinit(allocator); + if (listed.status != .ok) return error.ListSegmentsFailed; + + var arena_state = std.heap.ArenaAllocator.init(allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + const parsed = try std.json.parseFromSliceLeaky(std.json.Value, arena, listed.body, .{}); + const segments = parsed.object.get("segments").?.array.items; + var smallest = segments[0].object; + for (segments[1..]) |s| { + if (s.object.get("sizeBytes").?.integer < smallest.get("sizeBytes").?.integer) smallest = s.object; + } + const name = smallest.get("name").?.string; + const checksum = try std.fmt.parseInt(u64, smallest.get("checksum").?.string, 16); + + const url = try std.fmt.allocPrint(arena, "https://stream.waow.tech/xrpc/network.bsky.jetstream.getSegment?name={s}", .{name}); + var fetched = try transport.fetch(.{ .url = url }); + defer fetched.deinit(allocator); + if (fetched.status != .ok) return error.SegmentFetchFailed; + + try zat.ArchiveBackfill.verifySegment(fetched.body, checksum); + std.debug.print("checksum verified: {s} ({d} bytes, xxh3={x:0>16})\n", .{ name, fetched.body.len, checksum }); } const Handler = struct { diff --git a/src/internal/streaming/archive_backfill.zig b/src/internal/streaming/archive_backfill.zig index 07d9c1b..a89c400 100644 --- a/src/internal/streaming/archive_backfill.zig +++ b/src/internal/streaming/archive_backfill.zig @@ -10,6 +10,17 @@ //! the returned `last_time_us` (or the plan boundary) — overlap is fine for //! idempotent consumers; delivery is at-least-once. //! +//! scope (measured 2026-08-07): retrieval is block-granular, so a SPARSE +//! collection scattered across the archive reads ~1,300x more than it keeps +//! (streamplace chat: 40 GB fetched for ~30 MB of rows; one block per ~4096 +//! events regardless of how few match). the honest fit is dense collections, +//! bulk drains, per-DID history, and deleted/historical records. for a +//! sparse collection's current records, use a collection directory +//! (lightrail.microcosm.blue) + per-repo listRecords, then the live tail — +//! unless you need deletes or event ordering, which listRecords cannot +//! express. whole-segment mode buffers one full segment in memory (~277 MB); +//! block mode buffers one block. +//! //! jss v1 format: https://tangled.org/zat.dev/stream → docs/jss-format-v1.md const std = @import("std"); @@ -39,6 +50,20 @@ pub const Options = struct { /// the run early, the plan is not fully covered — do not treat /// `planned_through_seq` as a completed-backfill boundary. max_blocks: ?u64 = null, + /// resume a prior run: skip everything up to and including this position + /// (take it from the prior Result.position). segment indices are stable + /// for a given archive, so a re-plan covers at least the same range. + start_after: ?Position = null, + /// verify each whole-mode segment's jss self-checksum (xxh3 over header + /// and footer, cross-checked against the plan's checksum) before decoding. + /// block-mode fetches have no verifiable checksum here. off by default. + verify_checksums: bool = false, +}; + +/// a block's place in the plan; Result.position reports the last decoded one +pub const Position = struct { + segment_index: u32, + block_index: u32, }; pub const Result = struct { @@ -50,6 +75,8 @@ pub const Result = struct { blocks_decoded: u64 = 0, /// true when Options.max_blocks stopped the run before the plan was exhausted truncated: bool = false, + /// last decoded block; feed to Options.start_after to resume + position: ?Position = null, /// witnessed_at of the last delivered event, in the same µs domain as the /// live client's cursor. null if nothing matched. last_time_us: ?i64 = null, @@ -77,6 +104,7 @@ pub fn run(io: Io, allocator: Allocator, options: Options, handler: anytype) !Re for (ranges) |range| { var bi = range.first; while (bi <= range.last) : (bi += 1) { + if (skippedByResume(options.start_after, segment.index, bi)) continue; if (limitReached(options, &result)) return result; const url = try std.fmt.allocPrint( allocator, @@ -84,13 +112,16 @@ pub fn run(io: Io, allocator: Allocator, options: Options, handler: anytype) !Re .{ options.host, segment.name, bi }, ); defer allocator.free(url); - var fetched = try transport.fetch(.{ .url = url }); + var fetched = try fetchWithRetry(io, &transport, url); defer fetched.deinit(allocator); - if (fetched.status != .ok) return error.BlockFetchFailed; - try deliverFrame(allocator, fetched.body, options, handler, &result); + try deliverFrame(allocator, fetched.body, options, .{ + .segment_index = segment.index, + .block_index = bi, + }, handler, &result); } } } else { + if (options.start_after) |sa| if (segment.index < sa.segment_index) continue; if (limitReached(options, &result)) return result; const url = try std.fmt.allocPrint( allocator, @@ -98,16 +129,51 @@ pub fn run(io: Io, allocator: Allocator, options: Options, handler: anytype) !Re .{ options.host, segment.name }, ); defer allocator.free(url); - var fetched = try transport.fetch(.{ .url = url }); + var fetched = try fetchWithRetry(io, &transport, url); defer fetched.deinit(allocator); - if (fetched.status != .ok) return error.SegmentFetchFailed; - try deliverSegment(allocator, fetched.body, options, handler, &result); + try deliverSegment(allocator, fetched.body, options, segment, handler, &result); } } return result; } +/// resume skip: everything at or before start_after is already delivered +fn skippedByResume(start_after: ?Position, segment_index: u32, block_index: u32) bool { + const sa = start_after orelse return false; + if (segment_index < sa.segment_index) return true; + return segment_index == sa.segment_index and block_index <= sa.block_index; +} + +const fetch_attempts = 4; + +/// GET with bounded retries: transport errors and retryable statuses +/// (5xx, 429) back off 1s/2s/4s; other non-200s fail immediately. a +/// multi-GB run should not be discarded over one transient failure. +fn fetchWithRetry(io: Io, transport: *HttpTransport, url: []const u8) !HttpTransport.FetchResult { + var attempt: u32 = 0; + var backoff_s: u64 = 1; + while (true) { + attempt += 1; + if (transport.fetch(.{ .url = url })) |result| { + if (result.status == .ok) return result; + var discarded = result; + discarded.deinit(transport.allocator); + const code = @intFromEnum(result.status); + const retryable = code >= 500 or result.status == .too_many_requests; + if (!retryable) return error.FetchFailed; + if (attempt >= fetch_attempts) return error.FetchRetriesExhausted; + log.warn("fetch {d} (attempt {d}/{d}), retrying in {d}s: {s}", .{ code, attempt, fetch_attempts, backoff_s, url }); + } else |err| { + if (err == error.Canceled) return err; + if (attempt >= fetch_attempts) return err; + log.warn("fetch error {s} (attempt {d}/{d}), retrying in {d}s: {s}", .{ @errorName(err), attempt, fetch_attempts, backoff_s, url }); + } + try io.sleep(Io.Duration.fromSeconds(@intCast(backoff_s)), .awake); + backoff_s *= 2; + } +} + fn limitReached(options: Options, result: *Result) bool { const max = options.max_blocks orelse return false; if (result.blocks_decoded < max) return false; @@ -121,6 +187,9 @@ const BlockRange = struct { first: u32, last: u32 }; const PlannedSegment = struct { name: []const u8, + index: u32, + /// jss self-checksum from the plan (16 hex chars on the wire); null if absent + checksum: ?u64, /// null = whole-segment mode blocks: ?[]const BlockRange, }; @@ -201,7 +270,16 @@ fn parsePlan(arena: Allocator, body: []const u8) !Plan { } blocks = ranges; } - out.* = .{ .name = name, .blocks = blocks }; + const checksum: ?u64 = if (seg.get("checksum")) |v| switch (v) { + .string => |s| std.fmt.parseInt(u64, s, 16) catch null, + else => null, + } else null; + out.* = .{ + .name = name, + .index = try getJsonU32(seg_val, "index"), + .checksum = checksum, + .blocks = blocks, + }; } return .{ @@ -233,14 +311,17 @@ const block_index_entry_size = 52; const max_block_event_count = 1 << 18; const max_block_count = 1 << 20; -fn deliverSegment(allocator: Allocator, body: []const u8, options: Options, handler: anytype, result: *Result) !void { +fn deliverSegment(allocator: Allocator, body: []const u8, options: Options, segment: PlannedSegment, handler: anytype, result: *Result) !void { if (body.len < header_size) return error.MalformedSegment; if (!mem.eql(u8, body[0..4], "jss0")) return error.MalformedSegment; const block_count = mem.readInt(u32, body[14..18], .little); const block_index_offset = mem.readInt(u64, body[90..98], .little); if (block_count > max_block_count) return error.MalformedSegment; + if (options.verify_checksums) try verifySegment(body, segment.checksum); + for (0..block_count) |bi| { + if (skippedByResume(options.start_after, segment.index, @intCast(bi))) continue; if (limitReached(options, result)) return; const entry_off = block_index_offset + bi * block_index_entry_size; if (entry_off + block_index_entry_size > body.len) return error.MalformedSegment; @@ -248,11 +329,33 @@ fn deliverSegment(allocator: Allocator, body: []const u8, options: Options, hand const compressed_size = mem.readInt(u32, body[entry_off + 8 ..][0..4], .little); const frame_start = offset + 8; if (frame_start + compressed_size > body.len) return error.MalformedSegment; - try deliverFrame(allocator, body[frame_start..][0..compressed_size], options, handler, result); + try deliverFrame(allocator, body[frame_start..][0..compressed_size], options, .{ + .segment_index = segment.index, + .block_index = @intCast(bi), + }, handler, result); + } +} + +/// verify a sealed segment's self-checksum: xxh3_64(seed 0) over +/// header[12..256] ++ file[footer_offset..EOF], stored at header offset 4. +/// also cross-checks the plan's checksum for the segment when present. +pub fn verifySegment(body: []const u8, plan_checksum: ?u64) !void { + if (body.len < header_size) return error.MalformedSegment; + if (!mem.eql(u8, body[0..4], "jss0")) return error.MalformedSegment; + const stored = mem.readInt(u64, body[4..12], .little); + if (stored == 0) return error.SegmentUnsealed; + if (plan_checksum) |expected| { + if (expected != stored) return error.ChecksumMismatch; } + const footer_offset = mem.readInt(u64, body[58..66], .little); + if (footer_offset < header_size or footer_offset > body.len) return error.MalformedSegment; + var hasher = std.hash.XxHash3.init(0); + hasher.update(body[12..header_size]); + hasher.update(body[footer_offset..]); + if (hasher.final() != stored) return error.ChecksumMismatch; } -fn deliverFrame(allocator: Allocator, compressed: []const u8, options: Options, handler: anytype, result: *Result) !void { +fn deliverFrame(allocator: Allocator, compressed: []const u8, options: Options, position: Position, handler: anytype, result: *Result) !void { var arena_state = std.heap.ArenaAllocator.init(allocator); defer arena_state.deinit(); const arena = arena_state.allocator(); @@ -264,6 +367,7 @@ fn deliverFrame(allocator: Allocator, compressed: []const u8, options: Options, try deliverBlock(arena, out.written(), options, handler, result); result.blocks_decoded += 1; + result.position = position; } // jss event kinds (1..7) @@ -599,8 +703,11 @@ test "deliverFrame decompresses a plain zstd frame" { }; var handler = NoopHandler{}; var result: Result = .{ .planned_through_seq = 0, .sealed_tip_seq = 0 }; - try deliverFrame(arena, &frame, .{ .host = "https://example.test" }, &handler, &result); + try deliverFrame(arena, &frame, .{ .host = "https://example.test" }, .{ .segment_index = 7, .block_index = 3 }, &handler, &result); try testing.expectEqual(@as(u64, 0), result.events_delivered); + try testing.expectEqual(@as(u64, 1), result.blocks_decoded); + try testing.expectEqual(@as(u32, 7), result.position.?.segment_index); + try testing.expectEqual(@as(u32, 3), result.position.?.block_index); } test "cborToJson converts bytes and cids to lex-JSON conventions" { @@ -640,6 +747,9 @@ test "parsePlan handles blocks and whole-segment modes" { try testing.expectEqual(@as(u64, 12000), plan.sealed_tip_seq); try testing.expectEqual(@as(usize, 2), plan.segments.len); try testing.expectEqualStrings("seg_0000000001", plan.segments[0].name); + try testing.expectEqual(@as(u32, 0), plan.segments[0].index); + try testing.expectEqual(@as(u64, 0xabc), plan.segments[0].checksum.?); + try testing.expectEqual(@as(u32, 1), plan.segments[1].index); const ranges = plan.segments[0].blocks.?; try testing.expectEqual(@as(u32, 0), ranges[0].first); try testing.expectEqual(@as(u32, 2), ranges[0].last); @@ -679,8 +789,53 @@ test "deliverSegment walks header and block index" { }; var handler = NoopHandler{}; var result: Result = .{ .planned_through_seq = 0, .sealed_tip_seq = 0 }; - try deliverSegment(arena, out.written(), .{ .host = "https://example.test" }, &handler, &result); + const seg = PlannedSegment{ .name = "seg_test", .index = 0, .checksum = null, .blocks = null }; + try deliverSegment(arena, out.written(), .{ .host = "https://example.test" }, seg, &handler, &result); try testing.expectEqual(@as(u64, 0), result.events_delivered); + try testing.expectEqual(@as(u64, 1), result.blocks_decoded); +} + +test "skippedByResume skips through the resume position, nothing after" { + const sa = Position{ .segment_index = 5, .block_index = 10 }; + try testing.expect(skippedByResume(sa, 4, 999)); + try testing.expect(skippedByResume(sa, 5, 9)); + try testing.expect(skippedByResume(sa, 5, 10)); + try testing.expect(!skippedByResume(sa, 5, 11)); + try testing.expect(!skippedByResume(sa, 6, 0)); + try testing.expect(!skippedByResume(null, 0, 0)); +} + +test "verifySegment checks the jss self-checksum and the plan's" { + var arena_state = std.heap.ArenaAllocator.init(testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + + // minimal sealed segment: 256B header + 4-byte footer + var out: std.Io.Writer.Allocating = .init(arena); + const w = &out.writer; + try w.writeAll("jss0"); + try w.splatByteAll(0, 54); // checksum + version..offset 58 + try w.writeInt(u64, header_size, .little); // footer_offset @58 + try w.splatByteAll(0, 190); // rest of header + try w.writeAll("foot"); + const body = out.written(); + + var hasher = std.hash.XxHash3.init(0); + hasher.update(body[12..header_size]); + hasher.update(body[header_size..]); + const checksum = hasher.final(); + mem.writeInt(u64, body[4..12], checksum, .little); + + try verifySegment(body, checksum); + try verifySegment(body, null); + try testing.expectError(error.ChecksumMismatch, verifySegment(body, checksum + 1)); + + body[body.len - 1] ^= 1; // corrupt the footer + try testing.expectError(error.ChecksumMismatch, verifySegment(body, checksum)); + body[body.len - 1] ^= 1; + + mem.writeInt(u64, body[4..12], 0, .little); // unsealed + try testing.expectError(error.SegmentUnsealed, verifySegment(body, null)); } test "deliverSegment rejects bad magic" { @@ -694,6 +849,7 @@ test "deliverSegment rejects bad magic" { testing.allocator, &bogus, .{ .host = "https://example.test" }, + .{ .name = "seg_test", .index = 0, .checksum = null, .blocks = null }, &handler, &result, ));