From 9055eb651ea9183d02841864e648f8eea7737803 Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Thu, 23 Jul 2026 13:23:55 -0500 Subject: [PATCH] ingest: isolate repair tail encode failures --- src/internal/ingest.zig | 113 ++++++++++++++++++++++++++++++++-------- 1 file changed, 92 insertions(+), 21 deletions(-) diff --git a/src/internal/ingest.zig b/src/internal/ingest.zig index bffe28d..5043849 100644 --- a/src/internal/ingest.zig +++ b/src/internal/ingest.zig @@ -259,27 +259,7 @@ pub const Consumer = struct { for (sink.batch.items, 0..) |row, i| { var stamped = row; stamped.seq = first + i; - // Tail.append deep-copies both wire encodings. Allocate - // their complete CBOR/CID scratch graph in a per-row - // arena; the repair-lifetime arena must not retain two - // whole-repo JSON copies on top of the CAR and MST. - var wire_arena = std.heap.ArenaAllocator.init(sink.consumer.allocator); - defer wire_arena.deinit(); - const wire_alloc = wire_arena.allocator(); - const json = (try cold.encodeEvent(wire_alloc, stamped)) orelse ""; - const json_v2 = (try cold.encodeEventV2(wire_alloc, stamped)) orelse ""; - const kind: @import("wire.zig").Kind = switch (row.kind) { - .create, .update, .delete, .create_resync => .commit, - .identity => .identity, - .account => .account, - .sync => .sync, - }; - const skip_v1: tail_mod.SkipV1 = switch (row.kind) { - .sync => .sync, - .create_resync => .resync_replacement, - else => .none, - }; - try sink.consumer.tail.append(kind, row.did, row.collection, stamped.displayTime(), stamped.seq, json, json_v2, skip_v1); + try appendArchivedTailRow(sink.consumer, stamped); } sink.last_seq = first + sink.batch.items.len - 1; sink.batch.clearRetainingCapacity(); @@ -770,6 +750,52 @@ pub const Consumer = struct { } }; +/// Publish one already-archived row into the hot tail. Archive payloads are +/// opaque bytes; a repository can therefore contain a record that a +/// subscriber encoder cannot render as AT Protocol JSON. Upstream memoizes +/// that encoding error on the row and each subscriber skips it. Preserve the +/// durable row and publish an empty body for the affected wire instead of +/// turning remote record bytes into a process-fatal repair failure. +fn appendArchivedTailRow(consumer: *Consumer, row: segment.Event) !void { + // Tail.append deep-copies both wire encodings. Allocate their complete + // CBOR/CID scratch graph in a per-row arena; a repair-lifetime arena must + // not retain two whole-repo JSON copies on top of the CAR and MST. + var wire_arena = std.heap.ArenaAllocator.init(consumer.allocator); + defer wire_arena.deinit(); + const alloc = wire_arena.allocator(); + const json = cold.encodeEvent(alloc, row) catch |err| blk: { + log.warn("hot-tail v1 encode skipped row seq={d}: {s}", .{ row.seq, @errorName(err) }); + if (consumer.stats) |stats| _ = stats.subscribe_encode_errors_total.fetchAdd(1, .monotonic); + break :blk ""; + } orelse ""; + const json_v2 = cold.encodeEventV2(alloc, row) catch |err| blk: { + log.warn("hot-tail v2 encode skipped row seq={d}: {s}", .{ row.seq, @errorName(err) }); + if (consumer.stats) |stats| _ = stats.subscribe_encode_errors_total.fetchAdd(1, .monotonic); + break :blk ""; + } orelse ""; + const kind: @import("wire.zig").Kind = switch (row.kind) { + .create, .update, .delete, .create_resync => .commit, + .identity => .identity, + .account => .account, + .sync => .sync, + }; + const skip_v1: tail_mod.SkipV1 = switch (row.kind) { + .sync => .sync, + .create_resync => .resync_replacement, + else => .none, + }; + try consumer.tail.append( + kind, + row.did, + row.collection, + row.displayTime(), + row.seq, + json, + json_v2, + skip_v1, + ); +} + const CursorPoint = struct { local_seq: u64, upstream_seq: i64 }; const ChainPoint = struct { local_seq: u64, did: []u8, state: verify.ChainUpdate }; const Persisted = struct { @@ -846,6 +872,51 @@ test "identity durable envelope preserves optional handle" { try testing.expect(decoded_absent.getString("handle") == null); } +test "repair hot tail skips an undecodable archived record without failing ingest" { + const testing = std.testing; + var threaded: Io.Threaded = .init(testing.allocator, .{}); + defer threaded.deinit(); + const io = threaded.io(); + var tail = tail_mod.Tail.init(testing.allocator, io, 1 << 20); + defer tail.deinit(); + var stats: metrics.Stats = .{}; + var consumer: Consumer = .{ + .allocator = testing.allocator, + .io = io, + .tail = &tail, + .upstream = "ws://offline.invalid", + .stats = &stats, + }; + defer consumer.deinit(); + + // A complete repository can contain opaque record bytes outside + // DAG-CBOR. This exact float decoder error killed the experiment while a + // verifier repair was filling the hot tail from its fetched repository. + const float_record = "\xa1\x61x\xfb\x3f\xf0\x00\x00\x00\x00\x00\x00"; + try appendArchivedTailRow(&consumer, .{ + .seq = 42, + .witnessed_at = 100, + .indexed_at = 0, + .kind = .create_resync, + .did = "did:plc:float-record", + .collection = "app.bsky.feed.post", + .rkey = "3l3qo2vuowo2b", + .rev = "3l3qo2vutsw2b", + .payload = float_record, + }); + + try testing.expectEqual(@as(u64, 1), stats.subscribe_encode_errors_total.load(.monotonic)); + var arena = std.heap.ArenaAllocator.init(testing.allocator); + defer arena.deinit(); + var index: u64 = 0; + var stop: std.atomic.Value(bool) = .init(false); + const result = try tail.readBatch(arena.allocator(), &index, &stop, true, 1); + try testing.expectEqual(@as(usize, 1), result.batch.entries.len); + try testing.expectEqual(@as(u64, 42), result.batch.entries[0].seq); + try testing.expectEqual(@as(usize, 0), result.batch.entries[0].json.len); + try testing.expectEqual(@as(u64, 1), index); +} + fn logNewDrops(before: *const convert.Drops, after: *const convert.Drops) void { inline for (comptime std.enums.values(convert.DropReason)) |reason| { const delta = after.get(reason) - before.get(reason); -- 2.51.2