diff --git a/CHANGELOG.md b/CHANGELOG.md --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,8 @@ ## 0.3.5 +- **feat**: firehose `#commit` events now expose the raw `blocks` CAR bytes, `prevData`, operation `prev` CIDs, and a `toMstOperations()` helper so consumers can call `verifyCommitDiff`/`verifyCommitCar` without re-decoding raw frames. Added `#sync` event decoding and a named `LoadedCommitCar` return type for `loadCommitFromCAR`. +- **bench**: added `zig build firehose-decode-bench -- ` for measuring `FirehoseClient.decodeFrame` directly against the atproto-bench fixture corpus. - **perf**: align MST hot paths with Atmos while preserving AT Protocol root parity. The tree now uses nullable child pointers with CID-only lazy stubs, direct DAG-CBOR node serialization, cached clean-node CIDs, chunked key comparison, borrowed-key insertion for benchmark/import paths, and Atmos-style ordered-tree lookup that loads only selected lazy child nodes. The atproto-bench apples-to-apples MST run now has Zat ahead of Atmos on both insert+root and lookup while producing the same root bytes. - **feat**: `Mst.collectBlocks` and `Mst.walk` add the missing middle layer for consumers that need commit-CAR MST blocks or ordered repo traversal without reaching through `Node` internals. This is additive public API; existing callers should not break. - **test**: ZDS main now consumes the new MST middle layer in record writes and repo import, pinned to Zat commit `485f1d485a9b8e7b703e8627a6b6a8c3e3c36a0e`, and is deployed as `atcr.io/zat.dev/zds:latest` from merged ZDS commit `8d1a8182dace`. diff --git a/build.zig b/build.zig --- a/build.zig +++ b/build.zig @@ -94,6 +94,24 @@ 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); + // firehose decodeFrame benchmark over atproto-bench fixtures + const firehose_decode_bench = b.addExecutable(.{ + .name = "firehose-decode-bench", + .root_module = b.createModule(.{ + .root_source_file = b.path("scripts/firehose_decode_bench.zig"), + .target = target, + .optimize = optimize, + .link_libc = true, + .imports = &.{.{ .name = "zat", .module = mod }}, + }), + }); + b.installArtifact(firehose_decode_bench); + + const run_firehose_decode_bench = b.addRunArtifact(firehose_decode_bench); + if (b.args) |args| run_firehose_decode_bench.addArgs(args); + const firehose_decode_bench_step = b.step("firehose-decode-bench", "benchmark FirehoseClient.decodeFrame over atproto-bench fixtures"); + firehose_decode_bench_step.dependOn(&run_firehose_decode_bench.step); + // CBOR codec benchmarks const cbor_bench = b.addExecutable(.{ .name = "cbor-bench", diff --git a/scripts/firehose_decode_bench.zig b/scripts/firehose_decode_bench.zig new file mode 100644 --- /dev/null +++ b/scripts/firehose_decode_bench.zig @@ -0,0 +1,179 @@ +//! benchmark FirehoseClient.decodeFrame over an atproto-bench firehose corpus. + +const std = @import("std"); +const zat = @import("zat"); + +const Allocator = std.mem.Allocator; + +const warmup_passes: usize = 2; +const measured_passes: usize = 5; +const default_fixture = "../zzstoatzz.io/atproto-bench/fixtures/firehose-frames.bin"; + +const Corpus = struct { + frames: []const []const u8, + total_bytes: usize, + min_frame: usize, + max_frame: usize, +}; + +const PassResult = struct { + frames: usize, + commits: usize, + records: usize, + errors: usize, + elapsed_ns: u64, +}; + +pub fn main(init: std.process.Init) !void { + const allocator = init.gpa; + var args = std.process.Args.Iterator.init(init.minimal.args); + _ = args.skip(); + const fixture = args.next() orelse default_fixture; + + const corpus = try loadCorpus(allocator, fixture); + + std.debug.print("\n=== zat firehose decodeFrame benchmark ===\n\n", .{}); + std.debug.print("corpus: {d} frames, {d} bytes total\n", .{ corpus.frames.len, corpus.total_bytes }); + std.debug.print(" frame sizes: {d}..{d} bytes\n", .{ corpus.min_frame, corpus.max_frame }); + std.debug.print(" passes: {d} warmup, {d} measured\n\n", .{ warmup_passes, measured_passes }); + + { + var arena = std.heap.ArenaAllocator.init(allocator); + defer arena.deinit(); + const result = decodeOne(arena.allocator(), corpus.frames[0]); + std.debug.print("first frame: commits={d} records={d} errors={d}\n\n", .{ + result.commits, result.records, result.errors, + }); + } + + var arena = std.heap.ArenaAllocator.init(allocator); + defer arena.deinit(); + + for (0..warmup_passes) |_| { + for (corpus.frames) |frame| { + _ = arena.reset(.retain_capacity); + _ = decodeOne(arena.allocator(), frame); + } + } + + var results: [measured_passes]PassResult = undefined; + var total_commits: usize = 0; + var total_records: usize = 0; + var total_errors: usize = 0; + + for (0..measured_passes) |pass| { + var commits: usize = 0; + var records: usize = 0; + var errors: usize = 0; + const start_ns = nowNs(); + for (corpus.frames) |frame| { + _ = arena.reset(.retain_capacity); + const result = decodeOne(arena.allocator(), frame); + commits += result.commits; + records += result.records; + errors += result.errors; + } + results[pass] = .{ + .frames = corpus.frames.len, + .commits = commits, + .records = records, + .errors = errors, + .elapsed_ns = nowNs() - start_ns, + }; + total_commits += commits; + total_records += records; + total_errors += errors; + } + + report(corpus, &results, total_commits, total_records, total_errors); +} + +fn decodeOne(allocator: Allocator, frame: []const u8) struct { commits: usize, records: usize, errors: usize } { + const event = zat.firehose.decodeFrame(allocator, frame) catch { + return .{ .commits = 0, .records = 0, .errors = 1 }; + }; + return switch (event) { + .commit => |commit| blk: { + var records: usize = 0; + for (commit.ops) |op| { + if (op.record != null) records += 1; + } + break :blk .{ .commits = 1, .records = records, .errors = 0 }; + }, + else => .{ .commits = 0, .records = 0, .errors = 0 }, + }; +} + +fn report( + corpus: Corpus, + results: []const PassResult, + total_commits: usize, + total_records: usize, + total_errors: usize, +) void { + var fps_values: [measured_passes]f64 = undefined; + var total_ns: u64 = 0; + for (results, 0..) |r, i| { + const elapsed_s = @as(f64, @floatFromInt(r.elapsed_ns)) / 1_000_000_000.0; + fps_values[i] = @as(f64, @floatFromInt(r.frames)) / elapsed_s; + total_ns += r.elapsed_ns; + } + std.mem.sort(f64, &fps_values, {}, std.sort.asc(f64)); + + const elapsed_s = @as(f64, @floatFromInt(total_ns)) / 1_000_000_000.0; + const total_bytes = @as(f64, @floatFromInt(corpus.total_bytes)) * @as(f64, @floatFromInt(measured_passes)); + const mb_s = total_bytes / (1024.0 * 1024.0) / elapsed_s; + + std.debug.print("decodeFrame {d:>10.0} frames/sec {d:>8.1} MB/s commits={d} records={d} errors={d}\n", .{ + fps_values[measured_passes / 2], + mb_s, + total_commits, + total_records, + total_errors, + }); + std.debug.print(" variance: min={d:.0} median={d:.0} max={d:.0} frames/sec\n", .{ + fps_values[0], + fps_values[measured_passes / 2], + fps_values[measured_passes - 1], + }); +} + +fn nowNs() u64 { + return @intCast(std.Io.Clock.awake.now(std.Options.debug_io).toNanoseconds()); +} + +fn loadCorpus(allocator: Allocator, path: []const u8) !Corpus { + const io = std.Options.debug_io; + const data = std.Io.Dir.cwd().readFileAlloc(io, path, allocator, .limited(50 * 1024 * 1024)) catch |err| { + std.debug.print("cannot open {s}: {s}\n", .{ path, @errorName(err) }); + return err; + }; + if (data.len < 4) return error.InvalidFormat; + + const frame_count = std.mem.readInt(u32, data[0..4], .big); + var frames: std.ArrayListUnmanaged([]const u8) = .empty; + var pos: usize = 4; + var total_bytes: usize = 0; + var min_frame: usize = std.math.maxInt(usize); + var max_frame: usize = 0; + + for (0..frame_count) |_| { + if (pos + 4 > data.len) return error.InvalidFormat; + const frame_len = std.mem.readInt(u32, data[pos..][0..4], .big); + pos += 4; + if (pos + frame_len > data.len) return error.InvalidFormat; + const frame = data[pos..][0..frame_len]; + try frames.append(allocator, frame); + pos += frame_len; + total_bytes += frame_len; + min_frame = @min(min_frame, frame_len); + max_frame = @max(max_frame, frame_len); + } + + return .{ + .frames = try frames.toOwnedSlice(allocator), + .total_bytes = total_bytes, + .min_frame = min_frame, + .max_frame = max_frame, + }; +} diff --git a/src/root.zig b/src/root.zig --- a/src/root.zig +++ b/src/root.zig @@ -48,6 +48,7 @@ // sync 1.1: commit diff verification pub const MstOperation = mst.Operation; pub const Commit = repo_verifier.Commit; +pub const LoadedCommitCar = repo_verifier.LoadedCommitCar; pub const loadCommitFromCAR = repo_verifier.loadCommitFromCAR; pub const verifyCommitDiff = repo_verifier.verifyCommitDiff; pub const CommitDiffResult = repo_verifier.CommitDiffResult; @@ -83,5 +84,6 @@ _ = @import("internal/repo/car_test.zig"); _ = @import("internal/repo/cbor_rfc8949_test.zig"); _ = @import("internal/repo/mst_test.zig"); + _ = @import("internal/streaming/firehose.zig"); } } diff --git a/src/internal/repo/repo_verifier.zig b/src/internal/repo/repo_verifier.zig --- a/src/internal/repo/repo_verifier.zig +++ b/src/internal/repo/repo_verifier.zig @@ -299,15 +299,17 @@ prev: ?[]const u8, // raw CID bytes — previous commit CID (null for first commit) }; -/// lightweight: parse CAR, find root block, decode commit CBOR. -/// no MST loading. reusable for both #commit and #sync frames. -/// pre-computes unsigned commit bytes for signature verification (avoids re-decode). -pub fn loadCommitFromCAR(allocator: Allocator, car_bytes: []const u8) !struct { +pub const LoadedCommitCar = struct { commit: Commit, commit_cid: []const u8, unsigned_commit_bytes: []const u8, repo_car: car.Car, -} { +}; + +/// lightweight: parse CAR, find root block, decode commit CBOR. +/// no MST loading. reusable for both #commit and #sync frames. +/// pre-computes unsigned commit bytes for signature verification (avoids re-decode). +pub fn loadCommitFromCAR(allocator: Allocator, car_bytes: []const u8) !LoadedCommitCar { const repo_car = car.readWithOptions(allocator, car_bytes, .{}) catch return error.InvalidCommit; if (repo_car.roots.len == 0) return error.NoRootsInCar; diff --git a/src/internal/streaming/firehose.zig b/src/internal/streaming/firehose.zig --- a/src/internal/streaming/firehose.zig +++ b/src/internal/streaming/firehose.zig @@ -13,6 +13,7 @@ const websocket = @import("websocket"); const cbor = @import("../repo/cbor.zig"); const car = @import("../repo/car.zig"); +const mst = @import("../repo/mst.zig"); const sync = @import("sync.zig"); const mem = std.mem; @@ -40,6 +41,7 @@ /// decoded firehose event pub const Event = union(enum) { commit: CommitEvent, + sync: SyncEvent, identity: IdentityEvent, account: AccountEvent, info: InfoEvent, @@ -47,6 +49,7 @@ pub fn seq(self: Event) ?i64 { return switch (self) { .commit => |c| c.seq, + .sync => |s| s.seq, .identity => |i| i.seq, .account => |a| a.seq, .info => null, @@ -61,17 +64,46 @@ time: []const u8, // datetime — when event was received since: ?[]const u8 = null, // TID — rev of preceding commit (null = full repo export) commit: ?cbor.Cid = null, // CID of the commit object + blocks: []const u8 = &.{}, // raw CAR diff bytes ops: []const RepoOp, + prev_data: ?cbor.Cid = null, // MST root CID of the previous revision blobs: []const cbor.Cid = &.{}, // new blobs referenced by records in this commit too_big: bool = false, + + pub fn toMstOperations(self: CommitEvent, allocator: Allocator) Allocator.Error![]mst.Operation { + var ops: std.ArrayList(mst.Operation) = .empty; + errdefer ops.deinit(allocator); + for (self.ops) |op| { + const path = if (op.path.len != 0) + try allocator.dupe(u8, op.path) + else + try std.fmt.allocPrint(allocator, "{s}/{s}", .{ op.collection, op.rkey }); + try ops.append(allocator, .{ + .path = path, + .value = if (op.cid) |cid| cid.raw else null, + .prev = if (op.prev) |prev| prev.raw else null, + }); + } + return ops.toOwnedSlice(allocator); + } }; pub const RepoOp = struct { action: CommitAction, + path: []const u8 = "", collection: []const u8, rkey: []const u8, cid: ?cbor.Cid = null, // CID of the record (null for deletes) + prev: ?cbor.Cid = null, // CID of the previous record for updates/deletes record: ?cbor.Value = null, // decoded DAG-CBOR record from CAR block +}; + +pub const SyncEvent = struct { + seq: i64, + did: []const u8, + rev: []const u8, + time: []const u8, + blocks: []const u8, // raw CAR bytes containing the current commit object }; pub const IdentityEvent = struct { @@ -132,6 +164,8 @@ if (mem.eql(u8, t, "#commit")) { return try decodeCommit(allocator, payload); + } else if (mem.eql(u8, t, "#sync")) { + return decodeSync(payload); } else if (mem.eql(u8, t, "#identity")) { return decodeIdentity(payload); } else if (mem.eql(u8, t, "#account")) { @@ -152,14 +186,9 @@ const rev = payload.getString("rev") orelse return error.MissingField; const time = payload.getString("time") orelse return error.MissingField; - // parse commit CID - var commit_cid: ?cbor.Cid = null; - if (payload.get("commit")) |commit_val| { - switch (commit_val) { - .cid => |c| commit_cid = c, - else => {}, - } - } + const commit_cid = payload.getCid("commit") orelse return error.MissingField; + const blocks_bytes = payload.getBytes("blocks") orelse return error.MissingField; + const prev_data = payload.getCid("prevData") orelse return error.MissingField; // parse blobs array (array of CID links) var blobs: std.ArrayList(cbor.Cid) = .empty; @@ -172,12 +201,9 @@ } } - // parse CAR blocks - const blocks_bytes = payload.getBytes("blocks"); - var parsed_car: ?car.Car = null; - if (blocks_bytes) |b| { - parsed_car = car.read(allocator, b) catch null; - } + // parse CAR blocks for record hydration. CAR verification failure should not + // hide the wire event; consumers can reject by verifying `blocks` explicitly. + const parsed_car: ?car.Car = car.read(allocator, blocks_bytes) catch null; // parse ops const ops_array = payload.getArray("ops"); @@ -196,6 +222,7 @@ // extract CID from op and look up record from CAR blocks var op_cid: ?cbor.Cid = null; + var op_prev: ?cbor.Cid = null; var record: ?cbor.Value = null; if (op_val.get("cid")) |cid_val| { switch (cid_val) { @@ -210,12 +237,20 @@ else => {}, } } + if (op_val.get("prev")) |prev_val| { + switch (prev_val) { + .cid => |cid| op_prev = cid, + else => {}, + } + } try ops.append(allocator, .{ .action = action, + .path = path, .collection = collection, .rkey = rkey, .cid = op_cid, + .prev = op_prev, .record = record, }); } @@ -228,9 +263,21 @@ .time = time, .since = payload.getString("since"), .commit = commit_cid, + .blocks = blocks_bytes, .ops = try ops.toOwnedSlice(allocator), + .prev_data = prev_data, .blobs = try blobs.toOwnedSlice(allocator), .too_big = payload.getBool("tooBig") orelse false, + } }; +} + +fn decodeSync(payload: cbor.Value) DecodeError!Event { + return .{ .sync = .{ + .seq = payload.getInt("seq") orelse return error.MissingField, + .did = payload.getString("did") orelse return error.MissingField, + .rev = payload.getString("rev") orelse return error.MissingField, + .time = payload.getString("time") orelse return error.MissingField, + .blocks = payload.getBytes("blocks") orelse return error.MissingField, } }; } @@ -263,6 +310,7 @@ const tag = switch (event) { .commit => "#commit", + .sync => "#sync", .identity => "#identity", .account => "#account", .info => "#info", @@ -278,6 +326,7 @@ // encode payload based on event type switch (event) { .commit => |commit| try encodeCommitPayload(allocator, &aw.writer, commit), + .sync => |sync_event| try encodeSyncPayload(allocator, &aw.writer, sync_event), .identity => |id| try encodeIdentityPayload(allocator, &aw.writer, id), .account => |acct| try encodeAccountPayload(allocator, &aw.writer, acct), .info => |inf| try encodeInfoPayload(allocator, &aw.writer, inf), @@ -297,7 +346,10 @@ for (commit.ops) |op| { const action_str: []const u8 = @tagName(op.action); - const path = try std.fmt.allocPrint(allocator, "{s}/{s}", .{ op.collection, op.rkey }); + const path = if (op.path.len != 0) + op.path + else + try std.fmt.allocPrint(allocator, "{s}/{s}", .{ op.collection, op.rkey }); if (op.record) |record| { // encode record, create CID, add to CAR blocks @@ -313,25 +365,35 @@ try root_cids.append(allocator, cid); } - try op_values.append(allocator, .{ .map = @constCast(&[_]cbor.Value.MapEntry{ - .{ .key = "action", .value = .{ .text = action_str } }, - .{ .key = "cid", .value = .{ .cid = cid } }, - .{ .key = "path", .value = .{ .text = path } }, - }) }); + var op_entries: std.ArrayList(cbor.Value.MapEntry) = .empty; + defer op_entries.deinit(allocator); + try op_entries.append(allocator, .{ .key = "action", .value = .{ .text = action_str } }); + try op_entries.append(allocator, .{ .key = "cid", .value = .{ .cid = cid } }); + try op_entries.append(allocator, .{ .key = "path", .value = .{ .text = path } }); + if (op.prev) |prev| { + try op_entries.append(allocator, .{ .key = "prev", .value = .{ .cid = prev } }); + } + try op_values.append(allocator, .{ .map = try op_entries.toOwnedSlice(allocator) }); } else { - try op_values.append(allocator, .{ .map = @constCast(&[_]cbor.Value.MapEntry{ - .{ .key = "action", .value = .{ .text = action_str } }, - .{ .key = "path", .value = .{ .text = path } }, - }) }); + var op_entries: std.ArrayList(cbor.Value.MapEntry) = .empty; + defer op_entries.deinit(allocator); + try op_entries.append(allocator, .{ .key = "action", .value = .{ .text = action_str } }); + try op_entries.append(allocator, .{ .key = "cid", .value = .null }); + try op_entries.append(allocator, .{ .key = "path", .value = .{ .text = path } }); + if (op.prev) |prev| { + try op_entries.append(allocator, .{ .key = "prev", .value = .{ .cid = prev } }); + } + try op_values.append(allocator, .{ .map = try op_entries.toOwnedSlice(allocator) }); } } - // build CAR file from blocks - const car_data = car.Car{ - .roots = root_cids.items, - .blocks = car_blocks.items, + const blocks_bytes = if (commit.blocks.len != 0) commit.blocks else blk: { + const car_data = car.Car{ + .roots = root_cids.items, + .blocks = car_blocks.items, + }; + break :blk try car.writeAlloc(allocator, car_data); }; - const blocks_bytes = try car.writeAlloc(allocator, car_data); // build blobs array var blob_values: std.ArrayList(cbor.Value) = .empty; @@ -345,11 +407,14 @@ defer entries.deinit(allocator); try entries.append(allocator, .{ .key = "blocks", .value = .{ .bytes = blocks_bytes } }); - if (commit.commit) |c| { - try entries.append(allocator, .{ .key = "commit", .value = .{ .cid = c } }); + if (commit.commit) |commit_cid| { + try entries.append(allocator, .{ .key = "commit", .value = .{ .cid = commit_cid } }); } try entries.append(allocator, .{ .key = "blobs", .value = .{ .array = blob_values.items } }); try entries.append(allocator, .{ .key = "ops", .value = .{ .array = op_values.items } }); + if (commit.prev_data) |prev_data| { + try entries.append(allocator, .{ .key = "prevData", .value = .{ .cid = prev_data } }); + } try entries.append(allocator, .{ .key = "repo", .value = .{ .text = commit.repo } }); try entries.append(allocator, .{ .key = "rev", .value = .{ .text = commit.rev } }); try entries.append(allocator, .{ .key = "seq", .value = .{ .unsigned = @intCast(commit.seq) } }); @@ -360,6 +425,19 @@ if (commit.too_big) { try entries.append(allocator, .{ .key = "tooBig", .value = .{ .boolean = true } }); } + + try cbor.encode(allocator, writer, .{ .map = entries.items }); +} + +fn encodeSyncPayload(allocator: Allocator, writer: anytype, sync_event: SyncEvent) !void { + var entries: std.ArrayList(cbor.Value.MapEntry) = .empty; + defer entries.deinit(allocator); + + try entries.append(allocator, .{ .key = "blocks", .value = .{ .bytes = sync_event.blocks } }); + try entries.append(allocator, .{ .key = "did", .value = .{ .text = sync_event.did } }); + try entries.append(allocator, .{ .key = "rev", .value = .{ .text = sync_event.rev } }); + try entries.append(allocator, .{ .key = "seq", .value = .{ .unsigned = @intCast(sync_event.seq) } }); + try entries.append(allocator, .{ .key = "time", .value = .{ .text = sync_event.time } }); try cbor.encode(allocator, writer, .{ .map = entries.items }); } @@ -557,8 +635,8 @@ // simulate a frame: header {op: 1, t: "#info"} + payload {name: "OutdatedCursor"} const header_bytes = [_]u8{ 0xa2, // map(2) - 0x62, 'o', 'p', 0x01, // "op": 1 0x61, 't', 0x65, '#', 'i', 'n', 'f', 'o', // "t": "#info" + 0x62, 'o', 'p', 0x01, // "op": 1 }; const payload_bytes = [_]u8{ 0xa1, // map(1) @@ -681,6 +759,8 @@ .{ .key = "$type", .value = .{ .text = "app.bsky.feed.post" } }, .{ .key = "text", .value = .{ .text = "hello firehose" } }, } }; + const commit_cid = try cbor.Cid.forDagCbor(alloc, "commit"); + const prev_data = try cbor.Cid.forDagCbor(alloc, "prev-data"); const original = Event{ .commit = .{ .seq = 999, @@ -688,6 +768,8 @@ .rev = "3k2abc000000", .time = "2024-01-15T10:30:00Z", .since = "3k2abd000000", + .commit = commit_cid, + .prev_data = prev_data, .ops = &.{.{ .action = .create, .collection = "app.bsky.feed.post", @@ -705,11 +787,15 @@ try std.testing.expectEqualStrings("3k2abc000000", commit.rev); try std.testing.expectEqualStrings("2024-01-15T10:30:00Z", commit.time); try std.testing.expectEqualStrings("3k2abd000000", commit.since.?); + try std.testing.expectEqualSlices(u8, commit_cid.raw, commit.commit.?.raw); + try std.testing.expect(commit.blocks.len > 0); + try std.testing.expectEqualSlices(u8, prev_data.raw, commit.prev_data.?.raw); try std.testing.expectEqual(@as(usize, 0), commit.blobs.len); try std.testing.expectEqual(@as(usize, 1), commit.ops.len); const op = commit.ops[0]; try std.testing.expectEqual(CommitAction.create, op.action); + try std.testing.expectEqualStrings("app.bsky.feed.post/3k2abc", op.path); try std.testing.expectEqualStrings("app.bsky.feed.post", op.collection); try std.testing.expectEqualStrings("3k2abc", op.rkey); try std.testing.expect(op.cid != null); @@ -718,22 +804,34 @@ const rec = op.record.?; try std.testing.expectEqualStrings("hello firehose", rec.getString("text").?); try std.testing.expectEqualStrings("app.bsky.feed.post", rec.getString("$type").?); + + const mst_ops = try commit.toMstOperations(alloc); + try std.testing.expectEqual(@as(usize, 1), mst_ops.len); + try std.testing.expectEqualStrings("app.bsky.feed.post/3k2abc", mst_ops[0].path); + try std.testing.expectEqualSlices(u8, op.cid.?.raw, mst_ops[0].value.?); + try std.testing.expect(mst_ops[0].prev == null); } test "encode → decode commit with delete (no record)" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const alloc = arena.allocator(); + const commit_cid = try cbor.Cid.forDagCbor(alloc, "commit"); + const prev_data = try cbor.Cid.forDagCbor(alloc, "prev-data"); + const prev_record = try cbor.Cid.forDagCbor(alloc, "prev-record"); const original = Event{ .commit = .{ .seq = 500, .repo = "did:plc:deleter", .rev = "3k2xyz000000", .time = "2024-01-15T10:30:00Z", + .commit = commit_cid, + .prev_data = prev_data, .ops = &.{.{ .action = .delete, .collection = "app.bsky.feed.post", .rkey = "abc123", + .prev = prev_record, .record = null, }}, } }; @@ -747,5 +845,28 @@ try std.testing.expectEqual(@as(usize, 1), decoded.commit.ops.len); try std.testing.expectEqual(CommitAction.delete, decoded.commit.ops[0].action); try std.testing.expect(decoded.commit.ops[0].cid == null); + try std.testing.expectEqualSlices(u8, prev_record.raw, decoded.commit.ops[0].prev.?.raw); try std.testing.expect(decoded.commit.ops[0].record == null); +} + +test "encode → decode sync frame" { + var arena = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena.deinit(); + const alloc = arena.allocator(); + + const original = Event{ .sync = .{ + .seq = 777, + .did = "did:plc:sync", + .rev = "3k2sync00000", + .time = "2024-01-15T10:30:00Z", + .blocks = "car bytes", + } }; + + const frame = try encodeFrame(alloc, original); + const decoded = try decodeFrame(alloc, frame); + + try std.testing.expectEqual(@as(i64, 777), decoded.sync.seq); + try std.testing.expectEqualStrings("did:plc:sync", decoded.sync.did); + try std.testing.expectEqualStrings("3k2sync00000", decoded.sync.rev); + try std.testing.expectEqualStrings("car bytes", decoded.sync.blocks); }