diff --git a/build.zig b/build.zig index e5ff85c..3ad3512 100644 --- a/build.zig +++ b/build.zig @@ -108,6 +108,23 @@ pub fn build(b: *std.Build) void { const firehose_smoke_step = b.step("firehose-smoke", "run firehose smoke test (CBOR/CAR/CID on live data)"); firehose_smoke_step.dependOn(&run_firehose_smoke.step); + // archive backfill smoke test (plan/fetch/decode on live archive data) + const archive_smoke = b.addExecutable(.{ + .name = "archive-backfill-smoke", + .root_module = b.createModule(.{ + .root_source_file = b.path("scripts/archive_backfill_smoke.zig"), + .target = target, + .optimize = optimize, + .link_libc = true, + .imports = &.{.{ .name = "zat", .module = mod }}, + }), + }); + b.installArtifact(archive_smoke); + + const run_archive_smoke = b.addRunArtifact(archive_smoke); + const archive_smoke_step = b.step("archive-smoke", "run archive backfill smoke test (live archive data)"); + archive_smoke_step.dependOn(&run_archive_smoke.step); + // runnable examples (compile-checked usage recipes) const example_parse_identifiers = b.addExecutable(.{ .name = "example-parse-identifiers", diff --git a/scripts/archive_backfill_smoke.zig b/scripts/archive_backfill_smoke.zig new file mode 100644 index 0000000..fb4c4fe --- /dev/null +++ b/scripts/archive_backfill_smoke.zig @@ -0,0 +1,52 @@ +const std = @import("std"); +const zat = @import("zat"); + +pub fn main() !void { + var da: std.heap.DebugAllocator(.{}) = .init; + defer _ = da.deinit(); + const allocator = da.allocator(); + + std.debug.print("archive backfill smoke: pollz collections from stream.waow.tech\n", .{}); + + var handler = Handler{}; + const result = try zat.ArchiveBackfill.run(std.Options.debug_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, + }, &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 }, + ); + if (result.events_delivered == 0) return error.NoEvents; + if (handler.missing_record != 0) return error.MissingRecords; +} + +const Handler = struct { + polls: u64 = 0, + votes: u64 = 0, + deletes: u64 = 0, + missing_record: u64 = 0, + printed: u64 = 0, + + pub fn onEvent(self: *Handler, event: zat.JetstreamEvent) void { + const commit = event.commit; + if (commit.operation == .delete) { + self.deletes += 1; + return; + } + if (commit.record == null) { + self.missing_record += 1; + return; + } + if (std.mem.endsWith(u8, commit.collection, ".poll")) self.polls += 1 else self.votes += 1; + if (self.printed < 3) { + self.printed += 1; + const question = zat.json.getString(commit.record.?, "question") orelse + zat.json.getString(commit.record.?, "answer") orelse "?"; + std.debug.print(" sample: {s} {s}/{s} time_us={d} field={s}\n", .{ commit.did, commit.collection, commit.rkey, commit.time_us, question }); + } + } +}; diff --git a/src/internal/streaming/archive_backfill.zig b/src/internal/streaming/archive_backfill.zig new file mode 100644 index 0000000..07d9c1b --- /dev/null +++ b/src/internal/streaming/archive_backfill.zig @@ -0,0 +1,700 @@ +//! 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 +//! 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. +//! +//! jss v1 format: https://tangled.org/zat.dev/stream → docs/jss-format-v1.md + +const std = @import("std"); +const jetstream = @import("jetstream.zig"); +const sync = @import("sync.zig"); +const cbor = @import("../repo/cbor.zig"); +const multibase = @import("../crypto/multibase.zig"); +const HttpTransport = @import("../xrpc/transport.zig").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 = &.{}, + /// dids to replay; empty = all + dids: []const []const u8 = &.{}, + /// 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, +}; + +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, + events_delivered: u64 = 0, + blocks_decoded: u64 = 0, + /// true when Options.max_blocks stopped the run before the plan was exhausted + truncated: bool = false, + /// 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. blocks until the plan is exhausted. +pub fn run(io: Io, allocator: Allocator, options: Options, handler: anytype) !Result { + var transport = HttpTransport.init(io, allocator); + defer transport.deinit(); + + var plan_arena = std.heap.ArenaAllocator.init(allocator); + defer plan_arena.deinit(); + + const plan = try fetchPlan(plan_arena.allocator(), &transport, options); + + var result: Result = .{ + .planned_through_seq = plan.planned_through_seq, + .sealed_tip_seq = plan.sealed_tip_seq, + }; + + for (plan.segments) |segment| { + if (segment.blocks) |ranges| { + for (ranges) |range| { + var bi = range.first; + while (bi <= range.last) : (bi += 1) { + if (limitReached(options, &result)) return result; + const url = try std.fmt.allocPrint( + allocator, + "{s}/xrpc/network.bsky.jetstream.getBlock?segment={s}&blockIndex={d}", + .{ options.host, segment.name, bi }, + ); + defer allocator.free(url); + var fetched = try transport.fetch(.{ .url = url }); + defer fetched.deinit(allocator); + if (fetched.status != .ok) return error.BlockFetchFailed; + try deliverFrame(allocator, fetched.body, options, handler, &result); + } + } + } else { + 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 = try transport.fetch(.{ .url = url }); + defer fetched.deinit(allocator); + if (fetched.status != .ok) return error.SegmentFetchFailed; + try deliverSegment(allocator, fetched.body, options, handler, &result); + } + } + + return result; +} + +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, + /// null = whole-segment mode + blocks: ?[]const BlockRange, +}; + +const Plan = struct { + planned_through_seq: u64, + sealed_tip_seq: u64, + segments: []const PlannedSegment, +}; + +fn fetchPlan(arena: Allocator, transport: *HttpTransport, options: Options) !Plan { + 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(); + try stringify.endObject(); + + const url = try std.fmt.allocPrint(arena, "{s}/xrpc/network.bsky.jetstream.planBackfill", .{options.host}); + var fetched = try transport.fetch(.{ + .url = url, + .method = .POST, + .payload = body.written(), + }); + defer fetched.deinit(transport.allocator); + if (fetched.status != .ok) { + log.err("planBackfill failed: {d} {s}", .{ @intFromEnum(fetched.status), fetched.body }); + return error.PlanFailed; + } + return parsePlan(arena, fetched.body); +} + +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; + } + out.* = .{ .name = name, .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, 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; + + for (0..block_count) |bi| { + 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; + try deliverFrame(allocator, body[frame_start..][0..compressed_size], options, handler, result); + } +} + +fn deliverFrame(allocator: Allocator, compressed: []const u8, options: Options, handler: anytype, result: *Result) !void { + var arena_state = std.heap.ArenaAllocator.init(allocator); + defer arena_state.deinit(); + const arena = arena_state.allocator(); + + var out: std.Io.Writer.Allocating = .init(arena); + var in: std.Io.Reader = .fixed(compressed); + var decompress: std.compress.zstd.Decompress = .init(&in, &.{}, .{}); + _ = decompress.reader.streamRemaining(&out.writer) catch return error.BlockDecompressFailed; + + try deliverBlock(arena, out.written(), options, handler, result); + result.blocks_decoded += 1; +} + +// jss event kinds (1..7) +const kind_create = 1; +const kind_update = 2; +const kind_delete = 3; +const kind_create_resync = 7; + +/// 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; + + 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 operation: sync.CommitAction = switch (raw[kind_base + i]) { + 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 (options.collections.len > 0 and !containsString(options.collections, collection)) continue; + if (options.dids.len > 0 and !containsString(options.dids, did)) continue; + + const record: ?json.Value = if (operation != .delete and payload.len > 0) blk: { + const decoded = cbor.decode(arena, payload) catch return error.MalformedRecord; + break :blk try cborToJson(arena, decoded.value); + } else null; + + const witnessed_at = mem.readInt(i64, raw[wit_base + 8 * i ..][0..8], .little); + _ = mem.readInt(u64, raw[seq_base + 8 * i ..][0..8], .little); + + 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; + } +} + +fn containsString(haystack: []const []const u8, needle: []const u8) bool { + for (haystack) |s| if (mem.eql(u8, s, needle)) return true; + return false; +} + +// === DAG-CBOR record → the json.Value shape the live client delivers === + +/// bytes and links use the atproto lex-JSON conventions the live jetstream +/// emits: {"$bytes": base64-no-pad} and {"$link": cid-base32-string}. +fn cborToJson(arena: Allocator, value: cbor.Value) Allocator.Error!json.Value { + return switch (value) { + .null => .null, + .boolean => |b| .{ .bool = b }, + .unsigned => |u| if (std.math.cast(i64, u)) |i| + .{ .integer = i } + else + .{ .float = @floatFromInt(u) }, + .negative => |i| .{ .integer = i }, + .text => |s| .{ .string = s }, + .bytes => |b| blk: { + const b64 = std.base64.standard_no_pad.Encoder; + const encoded = try arena.alloc(u8, b64.calcSize(b.len)); + _ = b64.encode(encoded, b); + var obj: json.ObjectMap = .empty; + try obj.put(arena, "$bytes", .{ .string = encoded }); + break :blk .{ .object = obj }; + }, + .cid => |c| blk: { + const encoded = multibase.base32lower.encode(arena, c.raw) catch return error.OutOfMemory; + var obj: json.ObjectMap = .empty; + try obj.put(arena, "$link", .{ .string = encoded }); + break :blk .{ .object = obj }; + }, + .array => |items| blk: { + var arr = json.Array.init(arena); + try arr.ensureTotalCapacityPrecise(items.len); + for (items) |item| arr.appendAssumeCapacity(try cborToJson(arena, item)); + break :blk .{ .array = arr }; + }, + .map => |entries| blk: { + var obj: json.ObjectMap = .empty; + try obj.ensureTotalCapacity(arena, entries.len); + for (entries) |entry| obj.putAssumeCapacity(entry.key, try cborToJson(arena, entry.value)); + break :blk .{ .object = obj }; + }, + }; +} + +// === 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 +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} +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" }, &handler, &result); + try testing.expectEqual(@as(u64, 0), result.events_delivered); +} + +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 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 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 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 "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); + 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 }; + try deliverSegment(arena, out.written(), .{ .host = "https://example.test" }, &handler, &result); + try testing.expectEqual(@as(u64, 0), result.events_delivered); +} + +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 = [_]u8{0} ** header_size; + try testing.expectError(error.MalformedSegment, deliverSegment( + testing.allocator, + &bogus, + .{ .host = "https://example.test" }, + &handler, + &result, + )); +} diff --git a/src/root.zig b/src/root.zig index c60f863..b4f65d3 100644 --- a/src/root.zig +++ b/src/root.zig @@ -66,6 +66,8 @@ pub const JetstreamClient = jetstream.JetstreamClient; pub const JetstreamEvent = jetstream.Event; // firehose (raw CBOR event stream) +pub const ArchiveBackfill = @import("internal/streaming/archive_backfill.zig"); + pub const firehose = @import("internal/streaming/firehose.zig"); pub const FirehoseClient = firehose.FirehoseClient; pub const FirehoseEvent = firehose.Event; @@ -103,6 +105,7 @@ comptime { _ = @import("internal/repo/mst.zig"); _ = @import("internal/repo/mst_test.zig"); _ = @import("internal/repo/repo.zig"); + _ = @import("internal/streaming/archive_backfill.zig"); _ = @import("internal/streaming/firehose.zig"); _ = @import("internal/streaming/jetstream.zig"); _ = @import("internal/streaming/sync.zig");