//! archive backfill — replay a jetstream archive through the live-tail handler //! //! 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. 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 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). 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 const std = @import("std"); const zat = @import("zat"); const jetstream = zat.jetstream; const sync = zat.firehose; // re-exports sync.CommitAction const cbor = zat.cbor; const zstd = @import("zstd.zig"); const livedecode = @import("livedecode.zig"); const filter = @import("filter.zig"); const multibase = zat.multibase; const HttpTransport = zat.HttpTransport; const mem = std.mem; const json = std.json; const Allocator = mem.Allocator; const Io = std.Io; const log = std.log.scoped(.zat); pub const Event = jetstream.Event; pub const Options = struct { /// base URL of an archiving jetstream instance, e.g. "https://stream.waow.tech" host: []const u8, /// collections to replay; empty = all collections: []const []const u8 = &.{}, /// plan-level kind pruning (upstream WithKinds -> plan request /// "kinds"). without it the plan admits marker/sentinel blocks — /// measured 83x more blocks for a network-wide collection slice. /// empty = all kinds. kinds: []const livedecode.Kind = &.{}, /// dids to replay; empty = all dids: []const []const u8 = &.{}, /// bearer credential for token-gated archives (Bluesky's hosted /// instances 401 archive endpoints without one; the live tail never /// needs it). mirrors the Go client's WithAPIKey. sent as /// "Authorization: Bearer " on plan/list/segment/block requests. api_key: ?[]const u8 = null, /// 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 /// `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, /// Decode each commit's record CBOR into Event.record. Like upstream's /// raw-record mode, false skips decoding and leaves record null — for consumers /// that only want (seq, did, collection, rkey, op), it avoids the cost /// and means one undecodable stored payload cannot stop the sweep. decode_records: bool = true, /// verify each whole-mode segment's jss self-checksum (xxh3 over header /// and footer, cross-checked against the plan's checksum) before decoding. /// off by default. per-block zstd content checksums are always verified /// by libzstd regardless of this flag. verify_checksums: bool = false, /// block fetches kept in flight. serial fetch is round-trip-bound (~0.2 /// MB/s at 900ms RTT measured 2026-08-08; 8 in flight was 7.2x faster), /// so WAN consumers want 4-8. delivery order is unaffected: events /// always arrive in plan order. stay modest — each slot holds its own /// connection, and politeness is the caller's job. clamped to [1, 16]. concurrency: usize = 1, }; /// 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 { /// archive coverage boundary from the plan. resume the live client with a /// cursor at or before the timestamp of the last delivered event. planned_through_seq: u64, sealed_tip_seq: u64, /// true when an extended row handler (onRow) returned false mid-sweep; /// the plan was abandoned cleanly at that point stopped: bool = false, events_delivered: u64 = 0, blocks_decoded: u64 = 0, /// compressed bytes fetched over the wire (getBlock frames and /// getSegment bodies), for IO accounting bytes_fetched: 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, }; /// replay the archive through `handler` (same contract as JetstreamClient: /// `fn onEvent(*H, Event) void`). event slices are only valid during the /// onEvent call. Pages through the sealed tip observed on the first plan, /// unless the handler stops or max_blocks is reached. PlanStalled means /// the server returned a continuation that cannot advance the sweep. /// Optional onError(anyerror) bool accepts a recoverable archive error with /// true or stops cleanly with false. Without it, the error is returned. pub fn run(io: Io, allocator: Allocator, options: Options, handler: anytype) !Result { var page_options = options; var pinned_tip: ?u64 = null; var result: Result = .{ .planned_through_seq = options.after_seq orelse 0, .sealed_tip_seq = 0 }; while (true) { if (options.max_blocks) |max| { if (result.blocks_decoded >= max) { result.truncated = true; return result; } page_options.max_blocks = max - result.blocks_decoded; } const page = try runPage(io, allocator, page_options, handler); if (pinned_tip == null) pinned_tip = page.sealed_tip_seq; result.sealed_tip_seq = pinned_tip.?; result.planned_through_seq = page.planned_through_seq; result.events_delivered += page.events_delivered; result.blocks_decoded += page.blocks_decoded; result.bytes_fetched += page.bytes_fetched; if (page.position) |position| result.position = position; if (page.last_time_us) |time| result.last_time_us = time; result.stopped = page.stopped; result.truncated = page.truncated; if (page.stopped or page.truncated or page.planned_through_seq >= pinned_tip.?) return result; if (page.planned_through_seq <= (page_options.after_seq orelse 0)) return error.PlanStalled; page_options.after_seq = page.planned_through_seq; // Snapshot membership is fixed by page one, even if the server seals // more data while downloading. The live tail covers that later data. page_options.before_seq = pinned_tip; } } fn runPage(io: Io, allocator: Allocator, options: Options, handler: anytype) !Result { var transport = HttpTransport.init(io, allocator); defer transport.deinit(); const authorization = try bearerHeader(allocator, options.api_key); defer if (authorization) |a| allocator.free(a); var plan_arena = std.heap.ArenaAllocator.init(allocator); defer plan_arena.deinit(); const plan = try fetchPlan(plan_arena.allocator(), &transport, options, authorization); var result: Result = .{ .planned_through_seq = plan.planned_through_seq, .sealed_tip_seq = plan.sealed_tip_seq, }; const concurrency = @min(@max(options.concurrency, 1), max_concurrency); var window = try FetchWindow.init(io, allocator, concurrency); defer window.deinit(); window.authorization = authorization; for (plan.segments) |segment| { if (segment.blocks) |ranges| { 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)) { window.discardAll(); return result; } if (window.isFull()) window.deliverOldest(allocator, options, handler, &result) catch |err| { window.discardAll(); if (err == error.Stopped) { result.stopped = true; return result; } return err; }; window.start(options, segment.name, .{ .segment_index = segment.index, .block_index = bi, }) catch |err| { window.discardAll(); return err; }; } } } else { if (options.start_after) |sa| if (segment.index < sa.segment_index) continue; // whole-segment fetches are serial; flush pending blocks first so // delivery stays in plan order window.deliverAll(allocator, options, handler, &result) catch |err| { window.discardAll(); if (err == error.Stopped) { result.stopped = true; return result; } return err; }; if (limitReached(options, &result)) return result; const url = try std.fmt.allocPrint( allocator, "{s}/xrpc/network.bsky.jetstream.getSegment?name={s}", .{ options.host, segment.name }, ); defer allocator.free(url); var fetched = fetchWithRetry(io, &transport, url, authorization, .segment) catch |err| { reportArchiveError(handler, err) catch |reported| { if (reported != error.Stopped) return reported; result.stopped = true; return result; }; continue; }; defer fetched.deinit(allocator); result.bytes_fetched += fetched.body.len; deliverSegment(allocator, fetched.body, options, segment, handler, &result) catch |err| { reportArchiveError(handler, err) catch |reported| { if (reported != error.Stopped) return reported; result.stopped = true; return result; }; }; } } window.deliverAll(allocator, options, handler, &result) catch |err| { window.discardAll(); if (err == error.Stopped) { result.stopped = true; return result; } return err; }; return result; } const max_concurrency = 16; const FetchError = anyerror; /// No callback means no permission to skip data. Cancellation and allocator /// failure are not corrupt archive entries and must never become recoverable. fn reportArchiveError(handler: anytype, err: anyerror) !void { switch (err) { error.Canceled, error.OutOfMemory, error.Stopped => return err, else => {}, } if (comptime @hasDecl(@TypeOf(handler.*), "onError")) { const decision = handler.onError(err); // Legacy archive handlers only log errors. Keep them source-compatible // without treating a notification as permission to skip an entry. if (@TypeOf(decision) == void) return err; const accepted = if (@TypeOf(decision) == bool) decision else try decision; if (!accepted) return error.Stopped; } else return err; } fn fetchJob(io: Io, transport: *HttpTransport, url: []const u8, authorization: ?[]const u8) FetchError!HttpTransport.FetchResult { return fetchWithRetry(io, transport, url, authorization, .xrpc); } /// a FIFO window of in-flight block fetches. starting is out-of-order I/O; /// delivery always awaits the oldest first, so handler order equals plan /// order regardless of concurrency. each slot owns a transport (one /// connection) reused as the window slides. const FetchWindow = struct { io: Io, allocator: Allocator, transports: []HttpTransport, slots: []Inflight, /// borrowed from run(); outlives every in-flight fetch authorization: ?[]const u8 = null, head: usize = 0, count: usize = 0, failed_segment: ?u32 = null, const Inflight = struct { future: std.Io.Future(FetchError!HttpTransport.FetchResult), url: []u8, position: Position, }; fn init(io: Io, allocator: Allocator, concurrency: usize) !FetchWindow { const transports = try allocator.alloc(HttpTransport, concurrency); errdefer allocator.free(transports); const slots = try allocator.alloc(Inflight, concurrency); for (transports) |*t| t.* = HttpTransport.init(io, allocator); return .{ .io = io, .allocator = allocator, .transports = transports, .slots = slots }; } fn deinit(self: *FetchWindow) void { self.discardAll(); for (self.transports) |*t| t.deinit(); self.allocator.free(self.transports); self.allocator.free(self.slots); } fn isFull(self: *const FetchWindow) bool { return self.count == self.slots.len; } fn start(self: *FetchWindow, options: Options, segment_name: []const u8, position: Position) !void { if (self.failed_segment == position.segment_index) return; std.debug.assert(!self.isFull()); const url = try std.fmt.allocPrint( self.allocator, "{s}/xrpc/network.bsky.jetstream.getBlock?segment={s}&blockIndex={d}", .{ options.host, segment_name, position.block_index }, ); errdefer self.allocator.free(url); const slot = (self.head + self.count) % self.slots.len; self.slots[slot] = .{ .future = self.io.async(fetchJob, .{ self.io, &self.transports[slot], url, self.authorization }), .url = url, .position = position, }; self.count += 1; } fn deliverOldest(self: *FetchWindow, allocator: Allocator, options: Options, handler: anytype, result: *Result) !void { std.debug.assert(self.count > 0); const slot = &self.slots[self.head]; self.head = (self.head + 1) % self.slots.len; self.count -= 1; defer self.allocator.free(slot.url); if (self.failed_segment == slot.position.segment_index) { if (slot.future.cancel(self.io)) |fetched| { var discarded = fetched; discarded.deinit(allocator); } else |_| {} return; } var fetched = slot.future.await(self.io) catch |err| { self.failed_segment = slot.position.segment_index; return reportArchiveError(handler, err); }; defer fetched.deinit(allocator); result.bytes_fetched += fetched.body.len; deliverFrame(allocator, fetched.body, options, slot.position, handler, result) catch |err| { try reportArchiveError(handler, err); }; } /// drain in order, stopping (and discarding the rest) if max_blocks hits fn deliverAll(self: *FetchWindow, allocator: Allocator, options: Options, handler: anytype, result: *Result) !void { while (self.count > 0) { if (limitReached(options, result)) { self.discardAll(); return; } try self.deliverOldest(allocator, options, handler, result); } } /// Cancel and release pending requests when the caller stops or a budget /// is reached; cleanup must not wait through unused retry schedules. fn discardAll(self: *FetchWindow) void { while (self.count > 0) { const slot = &self.slots[self.head]; self.head = (self.head + 1) % self.slots.len; self.count -= 1; if (slot.future.cancel(self.io)) |fetched| { var discarded = fetched; discarded.deinit(self.allocator); } else |_| {} self.allocator.free(slot.url); } } }; /// decode every block of a sealed in-memory jss segment through `handler` — /// the offline counterpart of `run` for a segment already on disk (local /// archive copies, getSegment bodies fetched out-of-band). same delivery /// contract; Options.host is unused. honors collections/dids/max_blocks/ /// verify_checksums. pub fn replaySegment(allocator: Allocator, body: []const u8, options: Options, handler: anytype) !Result { var result: Result = .{ .planned_through_seq = 0, .sealed_tip_seq = 0 }; try deliverSegment(allocator, body, options, .{ .name = "", .index = 0, .checksum = null, .blocks = null, }, 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 RetryKind = enum { xrpc, segment }; fn retryDelay(kind: RetryKind, attempt: u32, hint_ns: u64, jitter: f64) u64 { const cap = 30 * std.time.ns_per_s; const base = @as(u64, 500 * std.time.ns_per_ms) << @as(u6, @intCast(@min(attempt - 1, 6))); const scheduled: u64 = @min(cap, if (kind == .xrpc) @as(u64, @intFromFloat(@as(f64, @floatFromInt(base)) * (1 + 0.2 * jitter))) else base); if (kind == .xrpc and hint_ns > 0) return @min(cap, hint_ns); return @min(cap, @max(scheduled, hint_ns)); } fn fillRandom(io: *Io, bytes: []u8) void { io.random(bytes); } const fetch_attempts = 3; /// Match the Go bulk-segment and XRPC retry policies separately. fn fetchWithRetry(io: Io, transport: *HttpTransport, url: []const u8, authorization: ?[]const u8, kind: RetryKind) !HttpTransport.FetchResult { const previous_attempts = transport.dead_connection_attempts; transport.dead_connection_attempts = 1; defer transport.dead_connection_attempts = previous_attempts; var random_io = io; const random = std.Random.init(&random_io, fillRandom); var attempt: u32 = 0; while (true) { attempt += 1; var hint_ns: u64 = 0; if (transport.fetch(.{ .url = url, .authorization = authorization })) |result| { if (result.status == .ok) return result; var discarded = result; const code = @backingInt(result.status); hint_ns = retryHint(result.rate_limit, Io.Clock.real.now(io).toNanoseconds()); discarded.deinit(transport.allocator); const retryable = if (kind == .segment) code >= 500 or code == 429 else code == 429 or code == 500 or code == 502 or code == 503 or code == 504; if (!retryable) return error.FetchFailed; if (attempt >= fetch_attempts) return error.FetchRetriesExhausted; } else |err| { if (err == error.Canceled or err == error.OutOfMemory) return err; if (attempt >= fetch_attempts) return err; } const jitter = if (kind == .xrpc) random.float(f64) else 0; const delay = retryDelay(kind, attempt, hint_ns, jitter); try io.sleep(Io.Duration.fromNanoseconds(@intCast(delay)), .awake); } } fn retryHint(headers: HttpTransport.RateLimitHeaders, now_ns: i96) u64 { const nanoseconds: i128 = if (headers.reset) |reset| @as(i128, reset) * std.time.ns_per_s - now_ns else if (headers.retry_after) |delta| @as(i128, delta) * std.time.ns_per_s else if (headers.retry_after_at) |date| @as(i128, date) * std.time.ns_per_s - now_ns else 0; return @intCast(@min(@max(nanoseconds, 0), 30 * std.time.ns_per_s)); } fn limitReached(options: Options, result: *Result) bool { const max = options.max_blocks orelse return false; if (result.blocks_decoded < max) return false; result.truncated = true; return true; } // === plan === 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, }; const Plan = struct { planned_through_seq: u64, sealed_tip_seq: u64, segments: []const PlannedSegment, }; 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(); try stringify.objectField("collections"); try stringify.beginArray(); for (options.collections) |c| try stringify.write(c); try stringify.endArray(); try stringify.objectField("dids"); 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); } if (options.kinds.len > 0) { try stringify.objectField("kinds"); try stringify.beginArray(); for (options.kinds) |k| try stringify.write(@tagName(k)); try stringify.endArray(); } try stringify.endObject(); return body.written(); } fn bearerHeader(allocator: Allocator, api_key: ?[]const u8) !?[]u8 { const key = api_key orelse return null; return try std.fmt.allocPrint(allocator, "Bearer {s}", .{key}); } fn fetchPlan(arena: Allocator, transport: *HttpTransport, options: Options, authorization: ?[]const u8) !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, .content_type = "application/json", .authorization = authorization, }); defer fetched.deinit(transport.allocator); if (fetched.status != .ok) { log.warn("planSnapshot failed: {d} {s}", .{ @backingInt(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 }; } /// re-anchor point for entering a DIFFERENT instance's seq space from a /// witnessed-time floor (cross-instance failover: seqs don't transfer, /// witnessed time approximately does). `after_seq` is where the sweep /// starts on the new host; rows older than the floor are the caller's to /// trim (segment granularity is coarse — hours — so the seq bound decides /// what to FETCH and the time floor decides what to DELIVER). pub const Anchor = struct { after_seq: ?u64, /// false when no sealed segment reaches the floor: the window lives in /// the unsealed tail (after_seq = sealed tip, live cutover covers the /// rest) or the host has no sealed archive at all (after_seq = null, /// start live from the tip). covered: bool, }; /// reduce sealed spans to a failover anchor for `floor_us`: the sweep must /// start at or before the first row witnessed >= floor. covered → one seq /// below the earliest intersecting segment; not covered but spans exist → /// the sealed tip (the floor lies in the unsealed tail; the live cutover's /// server-side cold replay covers it); no spans → null (live from tip). pub fn anchorForFloor(spans: []const SegmentSpan, floor_us: i64) Anchor { const bounds = seqBoundsForWindow(spans, floor_us, null); if (bounds.covered) return .{ .after_seq = bounds.after_seq, .covered = true }; var tip: ?u64 = null; for (spans) |span| { if (tip == null or span.max_seq > tip.?) tip = span.max_seq; } return .{ .after_seq = tip, .covered = false }; } /// fetch every sealed segment's metadata from `host` into `arena` pub fn fetchSegmentSpans(io: Io, allocator: Allocator, arena: Allocator, host: []const u8, api_key: ?[]const u8) ![]SegmentSpan { var transport = HttpTransport.init(io, allocator); defer transport.deinit(); const authorization = try bearerHeader(allocator, api_key); defer if (authorization) |a| allocator.free(a); 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, authorization, .xrpc); 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 spans.items; } /// 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, api_key: ?[]const u8) !SeqBounds { var arena_state = std.heap.ArenaAllocator.init(allocator); defer arena_state.deinit(); const spans = try fetchSegmentSpans(io, allocator, arena_state.allocator(), host, api_key); return seqBoundsForWindow(spans, start_us, end_us); } /// fetch segment metadata from `host` and derive the failover anchor for /// `floor_us` (see anchorForFloor) pub fn fetchAnchor(io: Io, allocator: Allocator, host: []const u8, floor_us: i64, api_key: ?[]const u8) !Anchor { var arena_state = std.heap.ArenaAllocator.init(allocator); defer arena_state.deinit(); const spans = try fetchSegmentSpans(io, allocator, arena_state.allocator(), host, api_key); return anchorForFloor(spans, floor_us); } fn parsePlan(arena: Allocator, body: []const u8) !Plan { const parsed = try json.parseFromSliceLeaky(json.Value, arena, body, .{}); const root = switch (parsed) { .object => |o| o, else => return error.MalformedPlan, }; const segments_val = root.get("segments") orelse return error.MalformedPlan; const segments_json = switch (segments_val) { .array => |a| a.items, else => return error.MalformedPlan, }; const segments = try arena.alloc(PlannedSegment, segments_json.len); for (segments_json, segments) |seg_val, *out| { const seg = switch (seg_val) { .object => |o| o, else => return error.MalformedPlan, }; const name = switch (seg.get("name") orelse return error.MalformedPlan) { .string => |s| s, else => return error.MalformedPlan, }; const mode = switch (seg.get("mode") orelse return error.MalformedPlan) { .string => |s| s, else => return error.MalformedPlan, }; var blocks: ?[]const BlockRange = null; if (mem.eql(u8, mode, "blocks")) { const ranges_json = switch (seg.get("blocks") orelse return error.MalformedPlan) { .array => |a| a.items, else => return error.MalformedPlan, }; const ranges = try arena.alloc(BlockRange, ranges_json.len); for (ranges_json, ranges) |range_val, *r| { r.* = .{ .first = try getJsonU32(range_val, "first"), .last = try getJsonU32(range_val, "last"), }; } blocks = ranges; } 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 .{ .planned_through_seq = try getJsonU64(.{ .object = root }, "plannedThroughSeq"), .sealed_tip_seq = try getJsonU64(.{ .object = root }, "sealedTipSeq"), .segments = segments, }; } fn getJsonU64(val: json.Value, key: []const u8) !u64 { const obj = switch (val) { .object => |o| o, else => return error.MalformedPlan, }; return switch (obj.get(key) orelse return error.MalformedPlan) { .integer => |i| if (i < 0) error.MalformedPlan else @intCast(i), else => error.MalformedPlan, }; } fn getJsonU32(val: json.Value, key: []const u8) !u32 { return std.math.cast(u32, try getJsonU64(val, key)) orelse error.MalformedPlan; } // === jss v1 segment walk === const header_size = 256; 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, 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; const offset = mem.readInt(u64, body[entry_off..][0..8], .little); 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; deliverFrame(allocator, body[frame_start..][0..compressed_size], options, .{ .segment_index = segment.index, .block_index = @intCast(bi), }, handler, result) catch |err| { try reportArchiveError(handler, err); }; } } /// 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, position: Position, handler: anytype, result: *Result) !void { var arena_state = std.heap.ArenaAllocator.init(allocator); defer arena_state.deinit(); const arena = arena_state.allocator(); const raw = zstd.decompressAlloc(arena, compressed) catch |err| switch (err) { error.OutOfMemory => return error.OutOfMemory, else => return error.BlockDecompressFailed, }; 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) const kind_create = 1; const kind_update = 2; const kind_delete = 3; const kind_create_resync = 7; const kind_identity = 4; const kind_account = 5; const kind_sync = 6; /// decode one decompressed columnar block and deliver matching rows. /// `arena` backs per-row record decoding; the caller resets it per block. fn deliverBlock(arena: Allocator, raw: []const u8, options: Options, handler: anytype, result: *Result) !void { if (raw.len < 4) return error.MalformedBlock; const n: usize = mem.readInt(u32, raw[0..4], .little); if (n == 0) { if (raw.len != 4) return error.MalformedBlock; return; } if (n > max_block_event_count) return error.MalformedBlock; var o: usize = 4; const seq_base = o; o += 8 * n; // seq const wit_base = o; o += 8 * n; // witnessed_at o += 8 * n; // indexed_at (unused) const kind_base = o; o += n; const col_len_base = o; o += n; const did_len_base = o; o += 2 * n; const rkey_len_base = o; o += n; const rev_len_base = o; o += n; const pay_len_base = o; o += 4 * n; if (o > raw.len) return error.MalformedBlock; // blob section bases from length totals var col_total: usize = 0; var did_total: usize = 0; var rkey_total: usize = 0; var rev_total: usize = 0; var pay_total: usize = 0; for (0..n) |i| { col_total += raw[col_len_base + i]; did_total += mem.readInt(u16, raw[did_len_base + 2 * i ..][0..2], .little); rkey_total += raw[rkey_len_base + i]; rev_total += raw[rev_len_base + i]; pay_total += mem.readInt(u32, raw[pay_len_base + 4 * i ..][0..4], .little); } if (o + col_total + did_total + rkey_total + rev_total + pay_total != raw.len) return error.MalformedBlock; var col_off = o; var did_off = col_off + col_total; var rkey_off = did_off + did_total; var rev_off = rkey_off + rkey_total; var pay_off = rev_off + rev_total; var record_error: ?anyerror = null; for (0..n) |i| { const col_len = raw[col_len_base + i]; const did_len = mem.readInt(u16, raw[did_len_base + 2 * i ..][0..2], .little); const rkey_len = raw[rkey_len_base + i]; const rev_len = raw[rev_len_base + i]; const pay_len = mem.readInt(u32, raw[pay_len_base + 4 * i ..][0..4], .little); const collection = raw[col_off..][0..col_len]; const did = raw[did_off..][0..did_len]; const rkey = raw[rkey_off..][0..rkey_len]; const rev = raw[rev_off..][0..rev_len]; const payload = raw[pay_off..][0..pay_len]; col_off += col_len; did_off += did_len; rkey_off += rkey_len; rev_off += rev_len; pay_off += pay_len; const kind_byte = raw[kind_base + i]; 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 // upstream's backfill (a folding consumer needs the #sync tombstone // to know to refold). the classic onEvent handler below keeps its // original commits-only, seq-less contract. if (comptime @hasDecl(@TypeOf(handler.*), "onRow")) { if (options.dids.len > 0 and !containsString(options.dids, did)) continue; const row_event: ?livedecode.Event = switch (kind_byte) { kind_create, kind_create_resync, kind_update, kind_delete => blk: { if (!filter.collectionMatches(options.collections, collection)) break :blk null; const op: livedecode.Operation = switch (kind_byte) { kind_update => .update, kind_delete => .delete, else => .create, }; const rec: ?json.Value = if (options.decode_records and op != .delete and payload.len > 0) rec: { const decoded = cbor.decode(arena, payload) catch |err| { if (err == error.OutOfMemory) return err; record_error = error.MalformedRecord; continue; }; break :rec try livedecode.cborToJson(arena, decoded.value); } else null; break :blk .{ .seq = seq_row, .time_us = witnessed_at_row, .did = did, .payload = .{ .commit = .{ .operation = op, .collection = collection, .rkey = rkey, .rev = rev, .record = rec, // upstream Event.RecordCBOR on the archive path: // the stored payload verbatim, aliasing the block // buffer (same lifetime as the string fields) .record_cbor = if (op != .delete) payload else "", }, }, }; }, kind_identity, kind_account, kind_sync => blk: { // marker payloads are the archived upstream envelope as a // CBOR map; fields legitimately absent on synthetic rows // (e.g. an async-resync #sync) default to zero values, // matching upstream's generated decoder const env: ?cbor.Value = if (payload.len > 0) (cbor.decode(arena, payload) catch |err| { if (err == error.OutOfMemory) return err; record_error = error.MalformedRecord; continue; }).value else null; break :blk switch (kind_byte) { kind_identity => .{ .seq = seq_row, .time_us = witnessed_at_row, .did = did, .payload = .{ .identity = .{ .did = did, .handle = if (env) |e| e.getString("handle") else null, .seq = if (env) |e| e.getInt("seq") else null, .time = if (env) |e| e.getString("time") else null, } } }, kind_account => .{ .seq = seq_row, .time_us = witnessed_at_row, .did = did, .payload = .{ .account = .{ .did = did, .active = if (env) |e| e.getBool("active") orelse false else false, .status = if (env) |e| e.getString("status") else null, .seq = if (env) |e| e.getInt("seq") else null, .time = if (env) |e| e.getString("time") else null, } } }, else => .{ .seq = seq_row, .time_us = witnessed_at_row, .did = did, .payload = .{ .sync = .{ .did = did, .rev = if (env) |e| e.getString("rev") orelse rev else rev, .seq = if (env) |e| e.getInt("seq") else null, .time = if (env) |e| e.getString("time") else null, } } }, }; }, else => null, // unknown jss kind from a newer writer: skip }; if (row_event) |ev| { if (!handler.onRow(ev)) return error.Stopped; result.events_delivered += 1; result.last_time_us = witnessed_at_row; } continue; } const operation: sync.CommitAction = switch (kind_byte) { kind_create, kind_create_resync => .create, kind_update => .update, kind_delete => .delete, // identity/account/sync marker rows carry no record; the live-tail // pass re-covers them, so the archive pass skips them. else => continue, }; // planned blocks can contain other collections' rows — always filter client-side if (!filter.collectionMatches(options.collections, collection)) continue; if (options.dids.len > 0 and !containsString(options.dids, did)) continue; const record: ?json.Value = if (options.decode_records and operation != .delete and payload.len > 0) blk: { const decoded = cbor.decode(arena, payload) catch |err| { if (err == error.OutOfMemory) return err; record_error = error.MalformedRecord; continue; }; break :blk try livedecode.cborToJson(arena, decoded.value); } else null; const witnessed_at = witnessed_at_row; handler.onEvent(.{ .commit = .{ .did = did, .time_us = witnessed_at, .rev = if (rev_len > 0) rev else null, .operation = operation, .collection = collection, .rkey = rkey, .record = record, } }); result.events_delivered += 1; result.last_time_us = witnessed_at; } if (record_error) |err| try reportArchiveError(handler, err); } fn containsString(haystack: []const []const u8, needle: []const u8) bool { for (haystack) |s| if (mem.eql(u8, s, needle)) return true; return false; } // === tests === const testing = std.testing; const CollectingHandler = struct { arena: Allocator, events: std.array_list.Managed(jetstream.CommitEvent), fn init(arena: Allocator) CollectingHandler { return .{ .arena = arena, .events = .init(arena) }; } pub fn onEvent(self: *CollectingHandler, event: Event) void { // copy string fields — block memory dies after deliverBlock returns var commit = event.commit; commit.did = self.arena.dupe(u8, commit.did) catch unreachable; commit.collection = self.arena.dupe(u8, commit.collection) catch unreachable; commit.rkey = self.arena.dupe(u8, commit.rkey) catch unreachable; commit.record = null; self.events.append(commit) catch unreachable; } }; /// build a columnar jss block from rows for tests (shared with the /// loopback e2e fixtures) pub fn buildTestBlock(arena: Allocator, rows: []const struct { seq: u64, witnessed_at: i64, kind: u8, collection: []const u8, did: []const u8, rkey: []const u8, rev: []const u8, payload: []const u8, }) ![]u8 { var out: std.Io.Writer.Allocating = .init(arena); const w = &out.writer; try w.writeInt(u32, @intCast(rows.len), .little); for (rows) |r| try w.writeInt(u64, r.seq, .little); for (rows) |r| try w.writeInt(i64, r.witnessed_at, .little); for (rows) |_| try w.writeInt(i64, 0, .little); // indexed_at for (rows) |r| try w.writeByte(r.kind); for (rows) |r| try w.writeByte(@intCast(r.collection.len)); for (rows) |r| try w.writeInt(u16, @intCast(r.did.len), .little); for (rows) |r| try w.writeByte(@intCast(r.rkey.len)); for (rows) |r| try w.writeByte(@intCast(r.rev.len)); for (rows) |r| try w.writeInt(u32, @intCast(r.payload.len), .little); for (rows) |r| try w.writeAll(r.collection); for (rows) |r| try w.writeAll(r.did); for (rows) |r| try w.writeAll(r.rkey); for (rows) |r| try w.writeAll(r.rev); for (rows) |r| try w.writeAll(r.payload); return out.written(); } // {"$type": "tech.waow.pollz.vote", "answer": 2} pub const test_record_cbor = [_]u8{ 0xa2, // map(2) 0x65, '$', 't', 'y', 'p', 'e', // "$type" 0x74, 't', 'e', 'c', 'h', '.', 'w', 'a', 'o', 'w', '.', 'p', 'o', 'l', 'l', 'z', '.', 'v', 'o', 't', 'e', 0x66, 'a', 'n', 's', 'w', 'e', 'r', // "answer" 0x02, }; test "deliverBlock decodes rows, filters collections, maps kinds" { 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 = "tech.waow.pollz.vote", .did = "did:plc:alice", .rkey = "r1", .rev = "rev1", .payload = &test_record_cbor }, .{ .seq = 11, .witnessed_at = 1001, .kind = kind_create, .collection = "app.bsky.feed.post", .did = "did:plc:bob", .rkey = "r2", .rev = "rev2", .payload = &test_record_cbor }, .{ .seq = 12, .witnessed_at = 1002, .kind = kind_delete, .collection = "tech.waow.pollz.vote", .did = "did:plc:carol", .rkey = "r3", .rev = "", .payload = "" }, .{ .seq = 13, .witnessed_at = 1003, .kind = kind_create_resync, .collection = "tech.waow.pollz.vote", .did = "did:plc:dave", .rkey = "r4", .rev = "rev4", .payload = &test_record_cbor }, .{ .seq = 14, .witnessed_at = 1004, .kind = 5, .collection = "$account", .did = "did:plc:eve", .rkey = "", .rev = "", .payload = "" }, }); var handler = CollectingHandler.init(arena); var result: Result = .{ .planned_through_seq = 0, .sealed_tip_seq = 0 }; try deliverBlock(arena, raw, .{ .host = "https://example.test", .collections = &.{"tech.waow.pollz.vote"}, }, &handler, &result); try testing.expectEqual(@as(u64, 3), result.events_delivered); try testing.expectEqual(@as(i64, 1003), result.last_time_us.?); try testing.expectEqualStrings("did:plc:alice", handler.events.items[0].did); try testing.expectEqual(sync.CommitAction.create, handler.events.items[0].operation); try testing.expectEqual(@as(i64, 1000), handler.events.items[0].time_us); try testing.expectEqual(sync.CommitAction.delete, handler.events.items[1].operation); try testing.expectEqualStrings("did:plc:carol", handler.events.items[1].did); // create_resync maps to create try testing.expectEqual(sync.CommitAction.create, handler.events.items[2].operation); try testing.expectEqualStrings("did:plc:dave", handler.events.items[2].did); } test "deliverBlock decodes DAG-CBOR record into json.Value" { var arena_state = std.heap.ArenaAllocator.init(testing.allocator); defer arena_state.deinit(); const arena = arena_state.allocator(); const raw = try buildTestBlock(arena, &.{ .{ .seq = 1, .witnessed_at = 42, .kind = kind_create, .collection = "tech.waow.pollz.vote", .did = "did:plc:alice", .rkey = "r1", .rev = "rev1", .payload = &test_record_cbor }, }); const RecordHandler = struct { answer: ?i64 = null, type_ok: bool = false, pub fn onEvent(self: *@This(), event: Event) void { const record = event.commit.record orelse return; const obj = record.object; self.type_ok = mem.eql(u8, obj.get("$type").?.string, "tech.waow.pollz.vote"); self.answer = obj.get("answer").?.integer; } }; var handler = RecordHandler{}; var result: Result = .{ .planned_through_seq = 0, .sealed_tip_seq = 0 }; try deliverBlock(arena, raw, .{ .host = "https://example.test" }, &handler, &result); try testing.expect(handler.type_ok); try testing.expectEqual(@as(i64, 2), handler.answer.?); } test "deliverBlock rejects trailing bytes and truncation" { var arena_state = std.heap.ArenaAllocator.init(testing.allocator); defer arena_state.deinit(); const arena = arena_state.allocator(); const raw = try buildTestBlock(arena, &.{ .{ .seq = 1, .witnessed_at = 42, .kind = kind_create, .collection = "c", .did = "d", .rkey = "r", .rev = "", .payload = "" }, }); const NoopHandler = struct { pub fn onEvent(_: *@This(), _: Event) void {} }; var handler = NoopHandler{}; var result: Result = .{ .planned_through_seq = 0, .sealed_tip_seq = 0 }; const opts = Options{ .host = "https://example.test" }; const with_trailing = try mem.concat(arena, u8, &.{ raw, "x" }); try testing.expectError(error.MalformedBlock, deliverBlock(arena, with_trailing, opts, &handler, &result)); try testing.expectError(error.MalformedBlock, deliverBlock(arena, raw[0 .. raw.len - 1], opts, &handler, &result)); } test "deliverBlock accepts empty compacted block" { const NoopHandler = struct { pub fn onEvent(_: *@This(), _: Event) void {} }; var handler = NoopHandler{}; var result: Result = .{ .planned_through_seq = 0, .sealed_tip_seq = 0 }; try deliverBlock(testing.allocator, &.{ 0, 0, 0, 0 }, .{ .host = "https://example.test" }, &handler, &result); try testing.expectEqual(@as(u64, 0), result.events_delivered); } test "deliverFrame decompresses a plain zstd frame" { var arena_state = std.heap.ArenaAllocator.init(testing.allocator); defer arena_state.deinit(); const arena = arena_state.allocator(); // empty block (event_count 0) wrapped in a hand-built zstd frame: // magic, FHD (single-segment, 1-byte FCS), content size 4, // block header (last, raw, size 4), payload const frame = [_]u8{ 0x28, 0xb5, 0x2f, 0xfd, 0x20, 0x04, 0x21, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00 }; const NoopHandler = struct { pub fn onEvent(_: *@This(), _: Event) void {} }; var handler = NoopHandler{}; var result: Result = .{ .planned_through_seq = 0, .sealed_tip_seq = 0 }; 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" { var arena_state = std.heap.ArenaAllocator.init(testing.allocator); defer arena_state.deinit(); const arena = arena_state.allocator(); const bytes_json = try livedecode.cborToJson(arena, .{ .bytes = "hi" }); try testing.expectEqualStrings("aGk", bytes_json.object.get("$bytes").?.string); const cid_raw = [_]u8{ 0x01, 0x71, 0x12, 0x04, 0xde, 0xad, 0xbe, 0xef }; const cid_json = try livedecode.cborToJson(arena, .{ .cid = .{ .raw = &cid_raw } }); const link = cid_json.object.get("$link").?.string; try testing.expect(link.len > 1 and link[0] == 'b'); const nested = try livedecode.cborToJson(arena, .{ .array = &.{ .{ .unsigned = 7 }, .null, .{ .boolean = true } } }); try testing.expectEqual(@as(i64, 7), nested.array.items[0].integer); try testing.expectEqual(json.Value.null, nested.array.items[1]); try testing.expectEqual(true, nested.array.items[2].bool); } test "plan request carries kinds only when set" { var arena_state = std.heap.ArenaAllocator.init(testing.allocator); defer arena_state.deinit(); const arena = arena_state.allocator(); const with_kinds = try planRequestBody(arena, .{ .host = "h", .kinds = &.{ .commit, .sync } }); try testing.expect(std.mem.indexOf(u8, with_kinds, "\"kinds\":[\"commit\",\"sync\"]") != null); const without = try planRequestBody(arena, .{ .host = "h" }); try testing.expect(std.mem.indexOf(u8, without, "kinds") == null); } 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 "anchorForFloor: covered floor, unsealed-tail floor, empty archive" { 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 }, }; // floor inside the sealed range: sweep from one below the earliest // intersecting segment (the time floor trims the over-fetch) const covered = anchorForFloor(&spans, 2500); try std.testing.expect(covered.covered); try std.testing.expectEqual(@as(?u64, 100), covered.after_seq); // floor past every sealed segment: the window lives in the unsealed // tail — anchor at the sealed tip and let the live cutover's cold // replay cover the rest const tail = anchorForFloor(&spans, 9000); try std.testing.expect(!tail.covered); try std.testing.expectEqual(@as(?u64, 200), tail.after_seq); // no sealed archive at all: nothing to anchor; live from the tip const empty = anchorForFloor(&.{}, 2500); try std.testing.expect(!empty.covered); try std.testing.expectEqual(@as(?u64, null), empty.after_seq); } 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(); const arena = arena_state.allocator(); const plan = try parsePlan(arena, \\{"plannedThroughSeq": 12345, "sealedTipSeq": 12000, "segments": [ \\ {"name": "seg_0000000001", "index": 0, "checksum": "abc", "minSeq": 1, "maxSeq": 100, \\ "mode": "blocks", "blocks": [{"first": 0, "last": 2}, {"first": 5, "last": 5}]}, \\ {"name": "seg_0000000002", "index": 1, "checksum": "def", "minSeq": 101, "maxSeq": 200, \\ "mode": "segment"} \\]} ); try testing.expectEqual(@as(u64, 12345), plan.planned_through_seq); 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); try testing.expectEqual(@as(u32, 5), ranges[1].first); try testing.expect(plan.segments[1].blocks == null); } test "deliverSegment walks header and block index" { var arena_state = std.heap.ArenaAllocator.init(testing.allocator); defer arena_state.deinit(); const arena = arena_state.allocator(); // one empty block, hand-built zstd frame as in the frame test const frame = [_]u8{ 0x28, 0xb5, 0x2f, 0xfd, 0x20, 0x04, 0x21, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00 }; var out: std.Io.Writer.Allocating = .init(arena); const w = &out.writer; try w.writeAll("jss0"); try w.splatByteAll(0, 10); // checksum u64 + version u16 try w.writeInt(u32, 1, .little); // block_count @14 try w.splatByteAll(0, 72); // through offset 90 const frame_offset: u64 = header_size; const block_index_offset: u64 = header_size + 8 + frame.len; try w.writeInt(u64, block_index_offset, .little); // block_index_offset @90 try w.splatByteAll(0, 158); // reserved → header complete @256 try w.writeInt(u64, frame.len, .little); // block length prefix try w.writeAll(&frame); // block index entry try w.writeInt(u64, frame_offset, .little); try w.writeInt(u32, frame.len, .little); try w.writeInt(u32, 4, .little); // uncompressed_size try w.writeInt(u32, 0, .little); // event_count try w.splatByteAll(0, 32); // seq/witnessed bounds const NoopHandler = struct { pub fn onEvent(_: *@This(), _: Event) void {} }; var handler = NoopHandler{}; var result: Result = .{ .planned_through_seq = 0, .sealed_tip_seq = 0 }; 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" { const NoopHandler = struct { pub fn onEvent(_: *@This(), _: Event) void {} }; var handler = NoopHandler{}; var result: Result = .{ .planned_through_seq = 0, .sealed_tip_seq = 0 }; const bogus: [header_size]u8 = @splat(0); try testing.expectError(error.MalformedSegment, deliverSegment( testing.allocator, &bogus, .{ .host = "https://example.test" }, .{ .name = "seg_test", .index = 0, .checksum = null, .blocks = null }, &handler, &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); } // Upstream filter_test.go: TestMatcherCollectionExactAndWildcard, // TestMatcherWildcardBoundary, and marker/DID filter contracts. test "conformance: archive wildcards preserve namespace boundaries and markers" { var arena_state = std.heap.ArenaAllocator.init(testing.allocator); defer arena_state.deinit(); const arena = arena_state.allocator(); const raw = try buildTestBlock(arena, &.{ .{ .seq = 1, .witnessed_at = 1, .kind = kind_create, .collection = "app.bsky.feed.post", .did = "did:plc:a", .rkey = "r", .rev = "r", .payload = &test_record_cbor }, .{ .seq = 2, .witnessed_at = 2, .kind = kind_create, .collection = "app.bsky.graph.follow", .did = "did:plc:a", .rkey = "r", .rev = "r", .payload = &test_record_cbor }, .{ .seq = 3, .witnessed_at = 3, .kind = kind_create, .collection = "app.bsky.graphient.thing", .did = "did:plc:a", .rkey = "r", .rev = "r", .payload = &test_record_cbor }, .{ .seq = 4, .witnessed_at = 4, .kind = kind_account, .collection = "$account", .did = "did:plc:a", .rkey = "", .rev = "", .payload = "" }, .{ .seq = 5, .witnessed_at = 5, .kind = kind_identity, .collection = "$identity", .did = "did:plc:other", .rkey = "", .rev = "", .payload = "" }, }); const Handler = struct { seqs: std.ArrayList(u64) = .empty, allocator: Allocator, pub fn onRow(self: *@This(), event: livedecode.Event) bool { self.seqs.append(self.allocator, event.seq) catch @panic("out of memory"); return true; } }; var handler = Handler{ .allocator = arena }; var result: Result = .{ .planned_through_seq = 0, .sealed_tip_seq = 0 }; try deliverBlock(arena, raw, .{ .host = "", .collections = &.{ "app.bsky.feed.post", "app.bsky.graph.*" }, .dids = &.{"did:plc:a"}, }, &handler, &result); try testing.expectEqualSlices(u64, &.{ 1, 2, 4 }, handler.seqs.items); } // Upstream WithRawRecords/decodeCommitInto: skip map materialization, // preserve raw bytes for creates, and provide no record for deletes. test "conformance: raw archive records skip decoding and preserve bytes" { var arena_state = std.heap.ArenaAllocator.init(testing.allocator); defer arena_state.deinit(); const arena = arena_state.allocator(); const invalid_cbor = "\xff"; const raw = try buildTestBlock(arena, &.{ .{ .seq = 1, .witnessed_at = 1, .kind = kind_create, .collection = "app.bsky.feed.post", .did = "did:plc:a", .rkey = "r", .rev = "r", .payload = invalid_cbor }, .{ .seq = 2, .witnessed_at = 2, .kind = kind_delete, .collection = "app.bsky.feed.post", .did = "did:plc:a", .rkey = "r", .rev = "r", .payload = "" }, }); const Handler = struct { count: usize = 0, valid: bool = true, pub fn onRow(self: *@This(), event: livedecode.Event) bool { const commit = event.payload.commit; self.valid = self.valid and commit.record == null and mem.eql(u8, if (commit.operation == .delete) "" else invalid_cbor, commit.record_cbor); self.count += 1; return true; } }; var handler = Handler{}; var result: Result = .{ .planned_through_seq = 0, .sealed_tip_seq = 0 }; try deliverBlock(arena, raw, .{ .host = "", .decode_records = false }, &handler, &result); try testing.expectEqual(@as(usize, 2), handler.count); try testing.expect(handler.valid); try testing.expectError(error.MalformedRecord, deliverBlock(arena, raw, .{ .host = "" }, &handler, &result)); } // Upstream TestDownloadTransformKeepsPayloadWithMalformedRecordError: // valid rows from a mixed block survive and the decode error is reported. test "archive malformed records preserve valid rows before reporting the error" { var arena_state = std.heap.ArenaAllocator.init(testing.allocator); defer arena_state.deinit(); const arena = arena_state.allocator(); const raw = try buildTestBlock(arena, &.{ .{ .seq = 1, .witnessed_at = 1, .kind = kind_create, .collection = "app.bsky.feed.post", .did = "did:plc:a", .rkey = "r", .rev = "r", .payload = &test_record_cbor }, .{ .seq = 2, .witnessed_at = 2, .kind = kind_create, .collection = "app.bsky.feed.post", .did = "did:plc:a", .rkey = "r", .rev = "r", .payload = "\xff" }, .{ .seq = 3, .witnessed_at = 3, .kind = kind_create, .collection = "app.bsky.feed.post", .did = "did:plc:a", .rkey = "r", .rev = "r", .payload = &test_record_cbor }, }); const Handler = struct { items: [4]u64 = undefined, len: usize = 0, err: ?anyerror = null, pub fn onRow(self: *@This(), event: livedecode.Event) bool { self.items[self.len] = event.seq; self.len += 1; return true; } pub fn onError(self: *@This(), err: anyerror) bool { self.err = err; self.items[self.len] = 0; self.len += 1; return true; } }; var handler = Handler{}; var result: Result = .{ .planned_through_seq = 0, .sealed_tip_seq = 0 }; try deliverBlock(arena, raw, .{ .host = "" }, &handler, &result); try testing.expectEqualSlices(u64, &.{ 1, 3, 0 }, handler.items[0..handler.len]); try testing.expectEqual(@as(?anyerror, error.MalformedRecord), handler.err); } test "archive recovery never swallows cancellation allocation failure or absent callbacks" { const Handler = struct { errors: usize = 0, pub fn onError(self: *@This(), _: anyerror) bool { self.errors += 1; return true; } }; var handler = Handler{}; try testing.expectError(error.Canceled, reportArchiveError(&handler, error.Canceled)); try testing.expectError(error.OutOfMemory, reportArchiveError(&handler, error.OutOfMemory)); try testing.expectEqual(@as(usize, 0), handler.errors); var no_callback: struct {} = .{}; try testing.expectError(error.FetchFailed, reportArchiveError(&no_callback, error.FetchFailed)); const NotifyOnly = struct { errors: usize = 0, pub fn onError(self: *@This(), _: anyerror) void { self.errors += 1; } }; var notify_only = NotifyOnly{}; try testing.expectError(error.FetchFailed, reportArchiveError(¬ify_only, error.FetchFailed)); try testing.expectEqual(@as(usize, 1), notify_only.errors); } test "download retries use upstream's three attempt default" { try std.testing.expectEqual(@as(u32, 3), fetch_attempts); } test "retry timing matches upstream segment and XRPC policies" { const ms = std.time.ns_per_ms; const sec = std.time.ns_per_s; try std.testing.expectEqual(@as(u64, 500 * ms), retryDelay(.segment, 1, 0, 0)); try std.testing.expectEqual(@as(u64, sec), retryDelay(.segment, 2, 0, 0)); try std.testing.expectEqual(@as(u64, 30 * sec), retryDelay(.segment, 100, 0, 0)); try std.testing.expectEqual(@as(u64, 5 * sec), retryDelay(.segment, 1, 5 * sec, 0)); try std.testing.expectEqual(@as(u64, sec), retryDelay(.segment, 2, 100 * ms, 0)); try std.testing.expectEqual(@as(u64, 500 * ms), retryDelay(.xrpc, 1, 0, 0)); try std.testing.expect(retryDelay(.xrpc, 1, 0, 0.999999) < 600 * ms); try std.testing.expectEqual(@as(u64, 550 * ms), retryDelay(.xrpc, 1, 0, 0.5)); try std.testing.expectEqual(@as(u64, 100 * ms), retryDelay(.xrpc, 2, 100 * ms, 0)); try std.testing.expectEqual(@as(u64, 30 * sec), retryDelay(.xrpc, 1, 60 * sec, 0)); } test "retry hints prefer absolute reset and bound missing expired and large values" { const sec = std.time.ns_per_s; try std.testing.expectEqual(@as(u64, 0), retryHint(.{}, 100 * sec)); try std.testing.expectEqual(@as(u64, 250 * std.time.ns_per_ms), retryHint(.{ .reset = 101 }, 100 * sec + 750 * std.time.ns_per_ms)); try std.testing.expectEqual(@as(u64, 5 * sec), retryHint(.{ .retry_after = 5 }, 100 * sec)); try std.testing.expectEqual(@as(u64, 3 * sec), retryHint(.{ .reset = 103, .retry_after = 9 }, 100 * sec)); try std.testing.expectEqual(@as(u64, 0), retryHint(.{ .reset = 99, .retry_after = 9 }, 100 * sec)); try std.testing.expectEqual(@as(u64, 2 * sec), retryHint(.{ .retry_after_at = 102 }, 100 * sec)); try std.testing.expectEqual(@as(u64, 0), retryHint(.{ .retry_after_at = 90 }, 100 * sec)); try std.testing.expectEqual(@as(u64, 30 * sec), retryHint(.{ .retry_after = std.math.maxInt(u64) }, 100 * sec)); } test "negative reset takes precedence over retry after" { try testing.expectEqual(@as(u64, 0), retryHint(.{ .reset = -1, .retry_after = 9 }, 100 * std.time.ns_per_s)); }