diff --git a/CHANGELOG.md b/CHANGELOG.md index d5e152a..fea9112 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,32 @@ # changelog +## unreleased + +`zat.ArchiveBackfill` works against post-proposal-0015 jetstream archives, +and can bound a plan to a time window instead of draining the world. + +- **fix**: the plan endpoint is `network.bsky.jetstream.planSnapshot` — + upstream renamed it from `planBackfill` (bluesky-social/jetstream PR #330) + and removed the old name, so every `ArchiveBackfill.run` against a current + archive failed with `PlanFailed` until this rename. +- **feat**: `Options.after_seq` / `Options.before_seq` bound the plan + (exclusive lower / inclusive upper, upstream's semantics). New + `fetchSeqBounds(io, allocator, host, start_us, end_us)` maps a + witnessed-time window to those bounds via `listSegments` metadata; its + `covered=false` result is the honest signal that a window predates the + archive's live capture. Measured against stream.waow.tech: a two-hour + recent window plans ~25 blocks where the unbounded plan reads the whole + collection's scatter. +- **feat**: `run` logs one progress line per 256 decoded blocks — a bounded + but large plan (a session inside the merged bootstrap region) must not + look like a hang. +- **fix**: `zig build` compiles again: the `cbor-bench` target imported + `cbor.zig` as a file, whose `../crypto/multibase.zig` import escapes the + bench module's path since 0acbedf; it now imports the `zat` module. +- example 07 (streamplace chat) replays historical sessions from the + archive instead of relying on Jetstream cursor retention, then finishes + the unsealed tail over the live socket, deduping the seam by rkey. + ## 0.3.29 Idle-connection failover works on both streaming clients: a connection that diff --git a/src/internal/streaming/archive_backfill.zig b/src/internal/streaming/archive_backfill.zig index 0bbd626..5cf1904 100644 --- a/src/internal/streaming/archive_backfill.zig +++ b/src/internal/streaming/archive_backfill.zig @@ -1,24 +1,29 @@ //! archive backfill — replay a jetstream archive through the live-tail handler //! -//! consumes the network.bsky.jetstream archive API (planBackfill / getBlock / -//! getSegment) served by archiving jetstream instances such as +//! consumes the network.bsky.jetstream archive API (planSnapshot / getBlock / +//! getSegment; planSnapshot was named planBackfill before jetstream's +//! proposal-0015 rename) served by archiving jetstream instances such as //! stream.waow.tech: plan → fetch sealed jss v1 segment blocks → decompress //! (plain zstd, no dictionary) → columnar decode → deliver each matching row //! as a jetstream.Event through the same onEvent handler the live client uses. //! //! after `run` returns, connect the live client with a cursor at or before //! the returned `last_time_us` (or the plan boundary) — overlap is fine for -//! idempotent consumers; delivery is at-least-once. +//! idempotent consumers; delivery is at-least-once. the plan covers SEALED +//! segments only: the newest ~one segment of events is unsealed, so a +//! complete drain is always archive pass + live replay from the boundary. //! -//! scope (measured 2026-08-07): retrieval is block-granular, so a SPARSE -//! collection scattered across the archive reads ~1,300x more than it keeps +//! scope (measured 2026-08-07): retrieval is block-granular, so an UNBOUNDED +//! plan over a sparse collection 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); +//! events regardless of how few match). two ways that cost disappears: +//! - bound the plan: `after_seq`/`before_seq` (see fetchSeqBounds for +//! mapping a witnessed-time window to bounds — a two-hour window planned +//! ~25 blocks when measured 2026-08-12) +//! - or, for a sparse collection's CURRENT records only, use a collection +//! directory (lightrail.microcosm.blue) + per-repo listRecords — though +//! that cannot express deletes or event ordering. +//! 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 @@ -46,6 +51,11 @@ pub const Options = struct { collections: []const []const u8 = &.{}, /// dids to replay; empty = all dids: []const []const u8 = &.{}, + /// exclusive lower seq bound for the plan; null = from the beginning. + /// derive from a time window with fetchSeqBounds. + after_seq: ?u64 = null, + /// inclusive upper seq bound for the plan; null = to the sealed tip + before_seq: ?u64 = null, /// stop after decoding this many blocks. a full plan can span tens of GB; /// use this to bound a partial or exploratory run. when the limit stops /// the run early, the plan is not fully covered — do not treat @@ -339,7 +349,7 @@ const Plan = struct { segments: []const PlannedSegment, }; -fn fetchPlan(arena: Allocator, transport: *HttpTransport, options: Options) !Plan { +fn planRequestBody(arena: Allocator, options: Options) ![]const u8 { var body: std.Io.Writer.Allocating = .init(arena); var stringify: json.Stringify = .{ .writer = &body.writer }; try stringify.beginObject(); @@ -351,22 +361,127 @@ fn fetchPlan(arena: Allocator, transport: *HttpTransport, options: Options) !Pla try stringify.beginArray(); for (options.dids) |d| try stringify.write(d); try stringify.endArray(); + if (options.after_seq) |seq| { + try stringify.objectField("afterSeq"); + try stringify.write(seq); + } + if (options.before_seq) |seq| { + try stringify.objectField("beforeSeq"); + try stringify.write(seq); + } try stringify.endObject(); + return body.written(); +} - const url = try std.fmt.allocPrint(arena, "{s}/xrpc/network.bsky.jetstream.planBackfill", .{options.host}); +fn fetchPlan(arena: Allocator, transport: *HttpTransport, options: Options) !Plan { + const body = try planRequestBody(arena, options); + const url = try std.fmt.allocPrint(arena, "{s}/xrpc/network.bsky.jetstream.planSnapshot", .{options.host}); var fetched = try transport.fetch(.{ .url = url, .method = .POST, - .payload = body.written(), + .payload = body, }); defer fetched.deinit(transport.allocator); if (fetched.status != .ok) { - log.err("planBackfill failed: {d} {s}", .{ @intFromEnum(fetched.status), fetched.body }); + log.err("planSnapshot failed: {d} {s}", .{ @intFromEnum(fetched.status), fetched.body }); return error.PlanFailed; } return parsePlan(arena, fetched.body); } +// === time window → seq bounds === + +/// one sealed segment's metadata row from listSegments, reduced to what +/// bounds derivation needs +pub const SegmentSpan = struct { + min_seq: u64, + max_seq: u64, + min_witnessed_us: i64, + max_witnessed_us: i64, +}; + +pub const SeqBounds = struct { + after_seq: ?u64, + before_seq: ?u64, + /// false when no sealed segment's witnessed window intersects the + /// requested window — the window predates the archive's live capture or + /// lies entirely in the unsealed tail. plan nothing; go straight to the + /// live client (or drop the bounds for a full historical scan). + covered: bool, +}; + +/// reduce a witnessed-time window to plan seq bounds: the bounds admit every +/// segment whose witnessed span intersects [start_us, end_us]. sound for +/// live-captured rows (an event is witnessed at or after it happened); +/// bootstrap-merged rows carry bootstrap-era witnessed times, so windows +/// older than the archive's live capture report covered=false rather than +/// pretending. +pub fn seqBoundsForWindow(spans: []const SegmentSpan, start_us: i64, end_us: ?i64) SeqBounds { + var after: ?u64 = null; + var before: ?u64 = null; + var covered = false; + for (spans) |span| { + const ends_before = span.max_witnessed_us < start_us; + const starts_after = if (end_us) |e| span.min_witnessed_us > e else false; + if (ends_before or starts_after) continue; + covered = true; + if (after == null or span.min_seq -| 1 < after.?) after = span.min_seq -| 1; + if (before == null or span.max_seq > before.?) before = span.max_seq; + } + // an unbounded end means "through the sealed tip": drop the upper bound + if (end_us == null) before = null; + return .{ .after_seq = after, .before_seq = before, .covered = covered }; +} + +/// fetch every sealed segment's metadata from `host` and derive the seq +/// bounds covering the witnessed window [start_us, end_us]. end_us == null +/// means "through the sealed tip". +pub fn fetchSeqBounds(io: Io, allocator: Allocator, host: []const u8, start_us: i64, end_us: ?i64) !SeqBounds { + var transport = HttpTransport.init(io, allocator); + defer transport.deinit(); + var arena_state = std.heap.ArenaAllocator.init(allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + + var spans: std.ArrayList(SegmentSpan) = .empty; + var cursor: ?[]const u8 = null; + while (true) { + const url = if (cursor) |c| + try std.fmt.allocPrint(arena, "{s}/xrpc/network.bsky.jetstream.listSegments?limit=1000&cursor={s}", .{ host, c }) + else + try std.fmt.allocPrint(arena, "{s}/xrpc/network.bsky.jetstream.listSegments?limit=1000", .{host}); + var fetched = try fetchWithRetry(io, &transport, url); + defer fetched.deinit(allocator); + const parsed = try json.parseFromSliceLeaky(json.Value, arena, fetched.body, .{}); + const root = switch (parsed) { + .object => |o| o, + else => return error.MalformedSegmentList, + }; + const segments = switch (root.get("segments") orelse return error.MalformedSegmentList) { + .array => |a| a.items, + else => return error.MalformedSegmentList, + }; + if (segments.len == 0) break; + for (segments) |seg| { + try spans.append(arena, .{ + .min_seq = try getJsonU64(seg, "minSeq"), + .max_seq = try getJsonU64(seg, "maxSeq"), + .min_witnessed_us = @intCast(try getJsonU64(seg, "minWitnessedAt")), + .max_witnessed_us = @intCast(try getJsonU64(seg, "maxWitnessedAt")), + }); + } + // a short page is the end even if a cursor is present + if (segments.len < 1000) break; + const next = root.get("cursor") orelse break; + cursor = switch (next) { + .string => |s| try arena.dupe(u8, s), + .integer => |i| try std.fmt.allocPrint(arena, "{d}", .{i}), + else => break, + }; + } + return seqBoundsForWindow(spans.items, start_us, end_us); +} + fn parsePlan(arena: Allocator, body: []const u8) !Plan { const parsed = try json.parseFromSliceLeaky(json.Value, arena, body, .{}); const root = switch (parsed) { @@ -507,6 +622,10 @@ fn deliverFrame(allocator: Allocator, compressed: []const u8, options: Options, try deliverBlock(arena, raw, options, handler, result); result.blocks_decoded += 1; result.position = position; + // a bounded-but-large plan (a session inside the merged bootstrap + // region can touch thousands of blocks) must not look like a hang + if (result.blocks_decoded % 256 == 0) + log.info("archive backfill: {d} blocks decoded, {d} events delivered", .{ result.blocks_decoded, result.events_delivered }); } // jss event kinds (1..7) @@ -871,6 +990,57 @@ test "cborToJson converts bytes and cids to lex-JSON conventions" { try testing.expectEqual(true, nested.array.items[2].bool); } +test "plan request carries seq bounds only when set" { + var arena_state = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + + const unbounded = try planRequestBody(arena, .{ .host = "", .collections = &.{"a.b.c"} }); + try std.testing.expect(mem.indexOf(u8, unbounded, "afterSeq") == null); + try std.testing.expect(mem.indexOf(u8, unbounded, "beforeSeq") == null); + + const bounded = try planRequestBody(arena, .{ + .host = "", + .collections = &.{"a.b.c"}, + .after_seq = 100, + .before_seq = 200, + }); + try std.testing.expect(mem.indexOf(u8, bounded, "\"afterSeq\":100") != null); + try std.testing.expect(mem.indexOf(u8, bounded, "\"beforeSeq\":200") != null); +} + +test "seqBoundsForWindow admits intersecting segments, reports uncovered windows" { + // three sealed segments; witnessed spans overlap at the edges the way + // real segments do (per-connection skew) + const spans = [_]SegmentSpan{ + .{ .min_seq = 1, .max_seq = 100, .min_witnessed_us = 1000, .max_witnessed_us = 2000 }, + .{ .min_seq = 101, .max_seq = 200, .min_witnessed_us = 1900, .max_witnessed_us = 3000 }, + .{ .min_seq = 201, .max_seq = 300, .min_witnessed_us = 2900, .max_witnessed_us = 4000 }, + }; + + // a window inside the middle segment admits it plus the overlapping edge + const mid = seqBoundsForWindow(&spans, 2500, 2600); + try std.testing.expect(mid.covered); + try std.testing.expectEqual(@as(?u64, 100), mid.after_seq); + try std.testing.expectEqual(@as(?u64, 200), mid.before_seq); + + // an open end means "through the sealed tip": no upper bound + const open = seqBoundsForWindow(&spans, 3500, null); + try std.testing.expect(open.covered); + try std.testing.expectEqual(@as(?u64, 200), open.after_seq); + try std.testing.expectEqual(@as(?u64, null), open.before_seq); + + // a window before all witnessed spans (bootstrap-era session) is honest + const before_archive = seqBoundsForWindow(&spans, 10, 20); + try std.testing.expect(!before_archive.covered); + + // straddling windows admit everything they touch + const straddle = seqBoundsForWindow(&spans, 1500, 3500); + try std.testing.expect(straddle.covered); + try std.testing.expectEqual(@as(?u64, 0), straddle.after_seq); + try std.testing.expectEqual(@as(?u64, 300), straddle.before_seq); +} + test "parsePlan handles blocks and whole-segment modes" { var arena_state = std.heap.ArenaAllocator.init(testing.allocator); defer arena_state.deinit();