//! repo download + conversion — listRepos pages and getRepo CARs to rows //! //! pure-ish core of the backfill (docs/upstream-bootstrap-spec.md §2): //! a full-repo CAR becomes one KindCreate row per record, every row stamped //! with the HEAD commit rev — that whole-repo watermark is what the merge //! filter later compares live-captured revs against, so the coupling is //! deliberate and load-bearing. one witnessed_at per repo (per-record times //! would imply false ordering). downloads go via the relay (which redirects //! to the PDS); like upstream, this path does not verify signatures. //! //! validation gates (counter-only, never fatal to the run): //! - bad head rev → the whole repo fails (retryable) //! - bad NSID/rkey in an MST key → that record drops //! - field too long for the column format → that record drops //! - missing record block in the CAR → the repo fails (retryable) const std = @import("std"); const zat = @import("zat"); const segment = @import("../../storage/segment.zig"); const segment_writer = @import("../../storage/segment_writer.zig"); const Allocator = std.mem.Allocator; const log = std.log.scoped(.stream); pub const RepoRef = struct { did: []const u8, rev: []const u8, active: bool = true, }; pub const ListReposPage = struct { repos: []RepoRef, cursor: ?[]const u8, }; /// parse a com.atproto.sync.listRepos JSON response. all strings borrow /// from `arena`. pub fn parseListRepos(arena: Allocator, body: []const u8) !ListReposPage { const parsed = try std.json.parseFromSliceLeaky(std.json.Value, arena, body, .{}); if (parsed != .object) return error.BadResponse; const repos_val = parsed.object.get("repos") orelse return error.BadResponse; if (repos_val != .array) return error.BadResponse; var repos = try arena.alloc(RepoRef, repos_val.array.items.len); var n: usize = 0; for (repos_val.array.items) |item| { if (item != .object) continue; const did = item.object.get("did") orelse continue; if (did != .string) continue; const rev: []const u8 = if (item.object.get("rev")) |r| (if (r == .string) r.string else "") else ""; const active = if (item.object.get("active")) |a| (a == .bool and a.bool) else true; repos[n] = .{ .did = did.string, .rev = rev, .active = active }; n += 1; } const cursor: ?[]const u8 = if (parsed.object.get("cursor")) |c| (if (c == .string and c.string.len > 0) c.string else null) else null; return .{ .repos = repos[0..n], .cursor = cursor }; } pub const DropCounts = struct { invalid_collection: u64 = 0, invalid_rkey: u64 = 0, field_too_long: u64 = 0, }; pub const RepoError = error{ InvalidCommitRev, MissingBlock, IncompleteCar, InvalidCar, ArchiveAppendFailed, OutOfMemory, }; pub const AuthenticationError = error{ InvalidCommit, SignatureVerificationFailed, OutOfMemory, }; /// walk a full-repo CAR and hand one row per record to `sink` /// (fn(ctx, segment.Event) !void), stamped with `kind` (.create for /// bootstrap, .create_resync for repo resync). rows borrow from `arena`; /// the caller keeps it alive until the sink is done. returns the head rev /// (arena-owned) for the completion watermark. pub fn carToRows( arena: Allocator, car_bytes: []const u8, expected_did: []const u8, witnessed_at: i64, kind: segment.Kind, drops: *DropCounts, sink_ctx: anytype, comptime sink: fn (@TypeOf(sink_ctx), segment.Event) anyerror!void, ) RepoError![]const u8 { var prepared = try prepareRepo(arena, car_bytes); return emitPrepared(arena, &prepared, expected_did, witnessed_at, kind, drops, sink_ctx, sink); } /// Parse/index a full CAR and prove every MST node and referenced record block /// is present before any caller is allowed to append rows. pub fn prepareRepo(arena: Allocator, car_bytes: []const u8) RepoError!PreparedRepo { const loaded = loadRepoCommit(arena, car_bytes) catch |err| return switch (err) { error.OutOfMemory => error.OutOfMemory, error.UnexpectedEof => error.IncompleteCar, else => error.InvalidCar, }; if (zat.Tid.parse(loaded.rev) == null) return error.InvalidCommitRev; const reader_ctx = try arena.create(CarBlockReader); reader_ctx.* = .{ .blocks = loaded.blocks }; var mst = zat.mst.Mst.loadLazy(arena, loaded.data_cid, .{ .ctx = reader_ctx, .getFn = CarBlockReader.get, // loadLazy fails only on allocation. A missing block still reports as // one, but everything else here is exhaustion, and calling it // `InvalidCar` would retire a good repository permanently and blame // its PDS -- see the classification test below. }) catch |err| switch (err) { error.OutOfMemory => return if (reader_ctx.missing) error.MissingBlock else error.OutOfMemory, }; try validateComplete(&mst, loaded.blocks, reader_ctx); return .{ .embedded_did = loaded.did, .rev = loaded.rev, .data_cid = loaded.data_cid, .commit_version = loaded.commit_version, .signature = loaded.signature, .unsigned_commit = loaded.unsigned_commit, .blocks = loaded.blocks, .mst = mst, }; } pub const PreparedRepo = struct { embedded_did: []const u8, rev: []const u8, /// Signed post-state root from the repository commit. Bootstrap does not /// consume it, but verifier-triggered repair must promote this exact CID /// only after the fetched commit's signature and staleness gates pass. data_cid: []const u8, commit_version: i64, signature: []const u8, /// Canonical commit bytes with `sig` removed. Retaining this from the /// initial CAR parse lets verifier repair authenticate the exact prepared /// head without reparsing the CAR or walking the MST a second time. unsigned_commit: []const u8, blocks: BlockMap, mst: zat.mst.Mst, pub fn authenticate( self: *const PreparedRepo, expected_did: []const u8, public_key: zat.multicodec.PublicKey, ) AuthenticationError!void { if (self.commit_version != 3 or zat.Did.parse(self.embedded_did) == null) return error.InvalidCommit; if (!std.mem.eql(u8, self.embedded_did, expected_did)) return error.InvalidCommit; switch (public_key.key_type) { .p256 => zat.jwt.verifyP256(self.unsigned_commit, self.signature, public_key.raw) catch return error.SignatureVerificationFailed, .secp256k1 => zat.jwt.verifySecp256k1(self.unsigned_commit, self.signature, public_key.raw) catch return error.SignatureVerificationFailed, } } }; /// Emit an already-completeness-checked repository. Reuses the MST populated /// by prepareRepo, so resync can safely append its #sync marker immediately /// before this walk without a second parse/index/allocation phase. pub fn emitPrepared( arena: Allocator, prepared: *PreparedRepo, expected_did: []const u8, witnessed_at: i64, kind: segment.Kind, drops: *DropCounts, sink_ctx: anytype, comptime sink: fn (@TypeOf(sink_ctx), segment.Event) anyerror!void, ) RepoError![]const u8 { var collector: RowCollector(@TypeOf(sink_ctx), sink) = .{ .arena = arena, .blocks = prepared.blocks, // The relay's listRepos DID is authoritative for row attribution, // matching upstream's handler. Never let an embedded commit DID route // rows into another account's stream. .did = if (expected_did.len > 0) expected_did else prepared.embedded_did, .rev = prepared.rev, .witnessed_at = witnessed_at, .kind = kind, .drops = drops, .sink_ctx = sink_ctx, }; prepared.mst.walk(.{ .ctx = &collector, .entryFn = @TypeOf(collector).entry }) catch |err| { return collector.failure orelse switch (err) { error.OutOfMemory => error.OutOfMemory, else => error.InvalidCar, }; }; return prepared.rev; } fn validateComplete(mst: *zat.mst.Mst, blocks: BlockMap, reader: *CarBlockReader) RepoError!void { var validator: CompletenessValidator = .{ .blocks = blocks }; mst.walk(.{ .ctx = &validator, .entryFn = CompletenessValidator.entry }) catch |err| { if (reader.missing) return error.MissingBlock; return validator.failure orelse switch (err) { error.OutOfMemory => error.OutOfMemory, else => error.InvalidCar, }; }; } const CompletenessValidator = struct { blocks: BlockMap, failure: ?RepoError = null, fn entry(ctx: *anyopaque, e: zat.mst.WalkEntry) anyerror!void { const self: *@This() = @ptrCast(@alignCast(ctx)); if (self.blocks.get(e.value.raw) == null) { self.failure = error.MissingBlock; return error.MissingBlock; } } }; const RepoCommit = struct { did: []const u8, rev: []const u8, data_cid: []const u8, commit_version: i64, signature: []const u8, unsigned_commit: []const u8, blocks: BlockMap, }; const BlockMap = std.StringHashMapUnmanaged([]const u8); /// parse a full-repo CAR's root commit via zat's streaming reader /// (zat.loadCommitFromCAR is firehose-frame-scoped: 2 MB / 10k blocks; /// repo exports run to hundreds of MB). one pass builds the CID -> bytes /// index the MST walk needs; slices borrow from `car_bytes`. Signature inputs /// are retained for verifier-triggered repair, while ordinary bootstrap still /// intentionally does not authenticate repositories (matching upstream). fn loadRepoCommit(arena: Allocator, car_bytes: []const u8) !RepoCommit { var streamed = try zat.car.streamBlocks(arena, car_bytes, .{ .max_size = car_bytes.len }); if (streamed.roots.len == 0) return error.InvalidCar; var blocks: BlockMap = .empty; while (try streamed.iter.next()) |block| { try blocks.put(arena, block.cid_raw, block.data); } const commit_data = blocks.get(streamed.roots[0].raw) orelse return error.InvalidCar; const commit = try zat.cbor.decodeAll(arena, commit_data); const did = commit.getString("did") orelse return error.InvalidCar; const rev = commit.getString("rev") orelse return error.InvalidCar; const commit_version = commit.getInt("version") orelse return error.InvalidCar; const signature = commit.getBytes("sig") orelse return error.InvalidCar; const unsigned_commit = try encodeUnsignedCommit(arena, commit); const data_cid = switch (commit.get("data") orelse return error.InvalidCar) { .cid => |c| c.raw, else => return error.InvalidCar, }; return .{ .did = did, .rev = rev, .data_cid = data_cid, .commit_version = commit_version, .signature = signature, .unsigned_commit = unsigned_commit, .blocks = blocks, }; } fn encodeUnsignedCommit(allocator: Allocator, commit: zat.cbor.Value) ![]u8 { const entries = switch (commit) { .map => |map| map, else => return error.InvalidCar, }; var unsigned: std.ArrayList(zat.cbor.Value.MapEntry) = .empty; for (entries) |entry| { if (!std.mem.eql(u8, entry.key, "sig")) try unsigned.append(allocator, entry); } return zat.cbor.encodeAlloc(allocator, .{ .map = unsigned.items }); } const CarBlockReader = struct { blocks: BlockMap, missing: bool = false, fn get(ctx: *anyopaque, cid_raw: []const u8) anyerror!?[]const u8 { const self: *CarBlockReader = @ptrCast(@alignCast(ctx)); const block = self.blocks.get(cid_raw); if (block == null) self.missing = true; return block; } }; fn RowCollector(comptime Ctx: type, comptime sink: fn (Ctx, segment.Event) anyerror!void) type { return struct { arena: Allocator, blocks: BlockMap, did: []const u8, rev: []const u8, witnessed_at: i64, kind: segment.Kind, drops: *DropCounts, sink_ctx: Ctx, failure: ?RepoError = null, fn entry(ctx: *anyopaque, e: zat.mst.WalkEntry) anyerror!void { const self: *@This() = @ptrCast(@alignCast(ctx)); const slash = std.mem.indexOfScalar(u8, e.key, '/') orelse { self.drops.invalid_collection += 1; return; }; const collection = e.key[0..slash]; const rkey = e.key[slash + 1 ..]; if (zat.Nsid.parse(collection) == null) { self.drops.invalid_collection += 1; return; } if (zat.Rkey.parse(rkey) == null) { self.drops.invalid_rkey += 1; return; } const payload = self.blocks.get(e.value.raw) orelse { // a repo export missing a record block is broken: fail the // repo (retryable) rather than silently archive a hole self.failure = error.MissingBlock; return error.MissingBlock; }; const event: segment.Event = .{ .seq = 0, // archive assigns .witnessed_at = self.witnessed_at, .indexed_at = 0, .kind = self.kind, .did = self.did, .collection = collection, .rkey = rkey, .rev = self.rev, .payload = payload, }; segment_writer.validateFields(event) catch { self.drops.field_too_long += 1; return; }; sink(self.sink_ctx, event) catch { // Archive/storage failures are lifecycle-fatal. Calling them // a malformed repo would leave a partially appended repo in // the low bootstrap sequence range and continue crawling. self.failure = error.ArchiveAppendFailed; return error.ArchiveAppendFailed; }; } }; } // === tests === const testing = std.testing; test "parseListRepos shapes" { var arena = std.heap.ArenaAllocator.init(testing.allocator); defer arena.deinit(); const page = try parseListRepos(arena.allocator(), \\{"cursor":"abc","repos":[ \\ {"did":"did:plc:one","head":"bafy...","rev":"3l3qo2vutsw2b","active":true}, \\ {"did":"did:plc:two","rev":"3l3qo2vutsw2c"} \\]} ); try testing.expectEqual(@as(usize, 2), page.repos.len); try testing.expectEqualStrings("did:plc:one", page.repos[0].did); try testing.expect(page.repos[1].active); try testing.expectEqualStrings("abc", page.cursor.?); const done = try parseListRepos(arena.allocator(), "{\"repos\":[]}"); try testing.expectEqual(@as(usize, 0), done.repos.len); try testing.expect(done.cursor == null); } test "truncated CAR is retryable and emits nothing" { var arena = std.heap.ArenaAllocator.init(testing.allocator); defer arena.deinit(); var drops: DropCounts = .{}; const Sink = struct { rows: usize = 0, fn emit(self: *@This(), _: segment.Event) anyerror!void { self.rows += 1; } }; var sink: Sink = .{}; // Header declares two bytes but only one follows: an unambiguous // transport truncation, not a permanently malformed repository. try testing.expectError( error.IncompleteCar, carToRows(arena.allocator(), &.{ 0x02, 0xa1 }, "did:plc:expected", 0, .create, &drops, &sink, Sink.emit), ); try testing.expectEqual(@as(usize, 0), sink.rows); } // Allocation failure anywhere in prepare/emit must stay `OutOfMemory`. The // engine retries that, but classifies `InvalidCar` as a permanent per-repo // failure attributed to the remote's CAR -- so an exhaustion mistaken for a // bad record silently drops a good repository from a whole-network archive // and blames its PDS. `Mst.loadLazy`'s catch collapsed every error, including // this one, into `InvalidCar`. test "allocation failure never masquerades as a malformed repository" { var fixture_arena = std.heap.ArenaAllocator.init(testing.allocator); defer fixture_arena.deinit(); const fixture = fixture_arena.allocator(); const record = "oom classification fixture"; const record_cid = try zat.cbor.Cid.forDagCbor(fixture, record); var tree = zat.mst.Mst.init(fixture); try tree.put("app.bsky.feed.post/3k2abcdefghij", record_cid); const data_cid = try tree.rootCid(); const keypair = try zat.Keypair.fromSecretKey(.p256, @splat(17)); const did = try keypair.did(fixture); const signed = try zat.signCommit(fixture, .{ .did = did, .rev = "3k2abcdefghij", .data = data_cid, }, &keypair); var blocks: std.ArrayList(zat.car.Block) = .empty; try blocks.append(fixture, .{ .cid_raw = signed.cid.raw, .data = signed.bytes }); try tree.collectBlocks(&blocks); try blocks.append(fixture, .{ .cid_raw = record_cid.raw, .data = record }); const car_bytes = try zat.car.writeAlloc(fixture, .{ .roots = &.{signed.cid}, .blocks = blocks.items, }); const Sink = struct { rows: usize = 0, fn emit(self: *@This(), _: segment.Event) anyerror!void { self.rows += 1; } }; var completed: usize = 0; var exhausted: usize = 0; for (0..512) |fail_index| { var failing = testing.FailingAllocator.init(testing.allocator, .{ .fail_index = fail_index }); // The arena is the cleanup contract for a prepared repo: it owns every // byte the lazy MST and record rows reference, so unwinding at any // ownership transfer frees through it rather than per-object deinit. var arena = std.heap.ArenaAllocator.init(failing.allocator()); defer arena.deinit(); var drops: DropCounts = .{}; var sink: Sink = .{}; if (carToRows(arena.allocator(), car_bytes, did, 0, .create, &drops, &sink, Sink.emit)) |_| { completed += 1; } else |err| switch (err) { error.OutOfMemory => exhausted += 1, else => { std.debug.print( "fail_index {d}: exhaustion surfaced as {s}\n", .{ fail_index, @errorName(err) }, ); return error.ExhaustionMisclassified; }, } } // Both arms must be exercised, or the sweep proved nothing: too narrow a // range would only ever exhaust, and a broken fixture would only ever pass. try testing.expect(exhausted > 0); try testing.expect(completed > 0); } test "listRepos DID remains authoritative over a valid CAR's embedded DID" { var arena = std.heap.ArenaAllocator.init(testing.allocator); defer arena.deinit(); const allocator = arena.allocator(); const record = "hostile embedded DID fixture"; const record_cid = try zat.cbor.Cid.forDagCbor(allocator, record); var tree = zat.mst.Mst.init(allocator); try tree.put("app.bsky.feed.post/3k2abcdefghij", record_cid); const data_cid = try tree.rootCid(); const keypair = try zat.Keypair.fromSecretKey(.p256, @splat(13)); const embedded_did = try keypair.did(allocator); const signed = try zat.signCommit(allocator, .{ .did = embedded_did, .rev = "3k2abcdefghij", .data = data_cid, }, &keypair); var blocks: std.ArrayList(zat.car.Block) = .empty; try blocks.append(allocator, .{ .cid_raw = signed.cid.raw, .data = signed.bytes }); try tree.collectBlocks(&blocks); try blocks.append(allocator, .{ .cid_raw = record_cid.raw, .data = record }); const car_bytes = try zat.car.writeAlloc(allocator, .{ .roots = &.{signed.cid}, .blocks = blocks.items, }); const expected_did = "did:plc:authoritative-listrepos-did"; try testing.expect(!std.mem.eql(u8, embedded_did, expected_did)); var drops: DropCounts = .{}; const Sink = struct { count: usize = 0, did: []const u8 = "", fn emit(self: *@This(), event: segment.Event) anyerror!void { self.count += 1; self.did = event.did; } }; var sink: Sink = .{}; const rev = try carToRows( allocator, car_bytes, expected_did, 0, .create, &drops, &sink, Sink.emit, ); try testing.expectEqualStrings("3k2abcdefghij", rev); try testing.expectEqual(@as(usize, 1), sink.count); try testing.expectEqualStrings(expected_did, sink.did); try testing.expectEqual(@as(u64, 0), drops.invalid_collection); try testing.expectEqual(@as(u64, 0), drops.invalid_rkey); try testing.expectEqual(@as(u64, 0), drops.field_too_long); }