From 36e5c64d50febb7072df12452e0ef4ca87d13237 Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Fri, 24 Jul 2026 15:59:55 -0500 Subject: [PATCH] planBackfill: one-sided planning and no silent short plans The planner could return a successful plan that omitted matching data three different ways. Each is invisible to a client, which reads a short plan as "no such data" and never revisits the range. - A block whose collection bitmask row was absent (or a segment carrying no collection index at all) was pruned. Upstream's blockHasAnyCollection returns true past the slice; degraded metadata must fail open. The fail-open now lives in the planner, not in the accessor, so inspect-segment keeps reporting membership as fact. - matchBlocks failure was `catch continue`, silently dropping the segment. It now answers 500, matching the sibling allocation path. - Allocation failure while parsing the request answered 400, telling a caller its valid request was permanently malformed. Response encoding also used `catch return` throughout, so an allocation failure mid-encode returned without ever responding. Rendering moved into renderPlan, which propagates; the handler answers 500. Tests: fail-open on absent row and absent index, exact reporting stays fail-closed, matchBlocks propagates OOM, and a fail-index sweep over the production handler asserting every response is the exact baseline plan or 5xx -- never a short 200, a 400, or a dropped request. zig build test (294), ReleaseSafe, differential-oracle, plan-config-contract. Audit: planBackfill blocked -> verified. listSegments verified -> partial; it still has the drop-on-encode pattern this commit removed. Co-Authored-By: Claude Opus 5 (1M context) --- docs/semantic-parity.md | 4 +- src/internal/segment.zig | 14 +- src/internal/xrpcapi.zig | 299 ++++++++++++++++++++++++++++++++------- 3 files changed, 261 insertions(+), 56 deletions(-) diff --git a/docs/semantic-parity.md b/docs/semantic-parity.md index 8e10c14..123fd1a 100644 --- a/docs/semantic-parity.md +++ b/docs/semantic-parity.md @@ -56,8 +56,8 @@ blockers remain. The detailed bootstrap audit is in | Failed-repository healing | **blocked** | A real global/per-host worker gate exists, final redirect hosts are recorded, and retry state is stored. | Pass-level RocksDB/archive failures are logged and hidden from service health. Every failed DID and host gate is materialized in RAM before work. Persisted host parking is not upstream behavior, and park-read errors mean “not parked.” | | JSS sealed format interoperability | **verified** | Upstream-produced sealed fixtures are parsed; Stream-produced sealed files are consumed by the pinned Go reader; header/footer, block index, blooms, collections, compression, and checksums have reciprocal fixtures. | This verifies sealed-format compatibility only. It does not verify startup recovery or replay completeness. | | Archive startup recovery | **verified** | Startup validates the active header, walks complete frames, propagates read/decode/allocation failures, truncates and fsyncs only a framing-torn suffix, reconstructs block/event/sequence state, and resumes the same active file and segment index. Sealed high-water floors are accepted only through the checksum-verifying parser, including the empty-active/lower-sealed case. Physical tests prove repeated same-file restart, exact torn-tail truncation, byte-preserving failure on a complete corrupt frame, and rejection of a forged sealed `max_seq`. | This row establishes startup recovery. Cold replay and post-startup rewrite paths remain separate surfaces. | -| `listSegments`, `getSegment`, `getBlock` | **verified** | The pinned Go client and direct HTTP tests cover normal responses, byte identity, checksums/ETags, conditional requests, ranges, cache headers, and storage-error responses through the production server. | The claim is limited to these three methods and the tested file states. | -| `planBackfill` | **blocked** | Normal DID/collection/sequence filtering, whole-segment versus block mode, pagination, and configured limits match the pinned client fixtures. | `matchBlocks` errors are handled with `catch continue`, creating a successful plan with a false negative. Upstream's planner returns the error. Add an injected allocation/internal-error contract. | +| `listSegments`, `getSegment`, `getBlock` | **partial** | The pinned Go client and direct HTTP tests cover normal responses, byte identity, checksums/ETags, conditional requests, ranges, cache headers, and storage-error responses through the production server. | Downgraded from verified on 2026-07-24: `listSegments` builds its JSON with `catch return`, so an allocation failure mid-encode returns without ever calling `respond` and drops the request instead of answering `500`. The same pattern was just removed from `planBackfill`. Needs the identical treatment plus a fail-index sweep. | +| `planBackfill` | **verified** | Normal DID/collection/sequence filtering, whole-segment versus block mode, pagination, and configured limits match the pinned client fixtures. The planner is one-sided like upstream: a block whose collection bitmask row is absent, or a segment carrying no collection index, fails open instead of being pruned (upstream `blockHasAnyCollection`). Internal failure can no longer become a truthful-looking short plan — a fail-index sweep over the whole production handler proves every response is either the exact baseline plan or `5xx`, never a 200 with fewer segments, never a 400, and never a dropped request. | Exact collection reporting (`inspect-segment`) deliberately stays fail-closed; that split matches upstream and is asserted separately. | | Subscribe v1/v2 handshake and wire formats | **partial** | Cursor parsing, v1/v2 filter shapes, WebSocket framing, v1 deflate, v2 dictionary zstd, dictionary validators, size caps, pings, write deadlines, and shutdown close frames have production-path tests. | These tests do not establish gap-free cold replay or remote-data encoding isolation. | | Cold archive replay | **blocked** | Normal sealed and active-prefix replay reaches the hot tail in order in existing fixtures. Timestamp cursor corruption is surfaced for the selected block. | Both cold walkers skip a manifest-listed segment when `openFile` fails. They may then advance beyond durable events never delivered. Upstream returns an error and enforces progress/hole invariants below the readable-log floor. | | Live Sync 1.1 verification | **partial** | DID/signature checks, MST inversion, op-CID checks, replay/future-rev gates, durable chain/hosting updates, and several repair-trigger cases use real RocksDB, CARs, and sockets. Relay cursor corruption and read failures now fail closed. | The surface cannot be called closed while an accepted live row can persist and then kill the process during eager wire encoding. | diff --git a/src/internal/segment.zig b/src/internal/segment.zig index d28f0e4..84a0d6f 100644 --- a/src/internal/segment.zig +++ b/src/internal/segment.zig @@ -170,11 +170,19 @@ pub const CollectionIndex = struct { allocator.free(self.buffer); } + /// null when this block has no resident bitmask row. A planner must treat + /// null as "may contain anything" and fail open; an exact reporter must + /// not invent membership from it. + pub fn blockRow(self: *const CollectionIndex, block: usize) ?[]const u8 { + if (self.bitmask_len == 0) return null; + const at = block * self.bitmask_len; + if (at + self.bitmask_len > self.bitmasks.len) return null; + return self.bitmasks[at..][0..self.bitmask_len]; + } + pub fn blockHasCollection(self: *const CollectionIndex, block: usize, collection_id: usize) bool { if (collection_id >= self.entries.len) return false; - const at = block * self.bitmask_len; - if (at + self.bitmask_len > self.bitmasks.len) return false; - const row = self.bitmasks[at..][0..self.bitmask_len]; + const row = self.blockRow(block) orelse return false; return (row[collection_id / 8] >> @intCast(collection_id % 8)) & 1 == 1; } }; diff --git a/src/internal/xrpcapi.zig b/src/internal/xrpcapi.zig index 99a66af..ddbf2b0 100644 --- a/src/internal/xrpcapi.zig +++ b/src/internal/xrpcapi.zig @@ -934,8 +934,12 @@ pub fn planBackfill(self: *Api, respond: anytype, body: []const u8) void { defer arena.deinit(); const alloc = arena.allocator(); - const req = parsePlanRequest(alloc, body, self.plan) catch - return respond.err(400, "InvalidRequest", "malformed plan request"); + // Exhaustion is ours, not the caller's: a 400 tells a client its valid + // request is permanently bad and it will never retry the range. + const req = parsePlanRequest(alloc, body, self.plan) catch |err| switch (err) { + error.OutOfMemory => return respond.err(500, "InternalError", "oom"), + else => return respond.err(400, "InvalidRequest", "malformed plan request"), + }; self.plan_mutex.lockUncancelable(self.io); var plan_locked = true; @@ -946,10 +950,28 @@ pub fn planBackfill(self: *Api, respond: anytype, body: []const u8) void { defer for (resident) |*item| item.deinit(); var out: Io.Writer.Allocating = .init(alloc); + renderPlan(self.plan, alloc, req, resident, &out) catch + return respond.err(500, "InternalError", "oom"); + // Do not let a slow HTTP reader hold the planner lock while the + // already-materialized JSON is written. + plan_locked = false; + self.plan_mutex.unlock(self.io); + respond.json(200, out.written()); +} + +/// Render the planBackfill body. Every failure is propagated so the caller +/// answers 500; returning without responding would leave the request dropped. +fn renderPlan( + plan: PlanConfig, + alloc: Allocator, + req: PlanRequest, + resident: []@import("manifest.zig").PlanningSegment, + out: *Io.Writer.Allocating, +) !void { var s: std.json.Stringify = .{ .writer = &out.writer }; var segments_json: Io.Writer.Allocating = .init(alloc); var sj: std.json.Stringify = .{ .writer = &segments_json.writer }; - sj.beginArray() catch return; + try sj.beginArray(); var sealed_tip: u64 = 0; var planned_through: u64 = 0; @@ -973,7 +995,7 @@ pub fn planBackfill(self: *Api, respond: anytype, body: []const u8) void { if (header.min_seq > req.before_seq) continue; if (truncated) continue; // still walk for sealed_tip, plan nothing more - const matched = matchBlocks(alloc, sealed, req) catch continue; + const matched = try matchBlocks(alloc, sealed, req); var match_count: usize = 0; for (matched) |m| { if (m) match_count += 1; @@ -981,18 +1003,17 @@ pub fn planBackfill(self: *Api, respond: anytype, body: []const u8) void { if (match_count == 0) continue; const density = @as(f64, @floatFromInt(match_count)) / @as(f64, @floatFromInt(matched.len)); - const whole = density >= self.plan.whole_segment_threshold; + const whole = density >= plan.whole_segment_threshold; var admitted = matched; if (whole) { - if (self.plan.max_entries > 0 and entries >= self.plan.max_entries and entries > 0) { + if (plan.max_entries > 0 and entries >= plan.max_entries and entries > 0) { truncated = true; continue; } entries += 1; last_unit_max_seq = header.max_seq; } else { - admitted = alloc.alloc(bool, matched.len) catch - return respond.err(500, "InternalError", "oom"); + admitted = try alloc.alloc(bool, matched.len); @memset(admitted, false); var at: usize = 0; while (at < matched.len) { @@ -1002,7 +1023,7 @@ pub fn planBackfill(self: *Api, respond: anytype, body: []const u8) void { } var last = at; while (last + 1 < matched.len and matched[last + 1]) last += 1; - if (self.plan.max_entries > 0 and entries >= self.plan.max_entries and entries > 0) { + if (plan.max_entries > 0 and entries >= plan.max_entries and entries > 0) { truncated = true; break; } @@ -1015,57 +1036,52 @@ pub fn planBackfill(self: *Api, respond: anytype, body: []const u8) void { if (countRanges(admitted) == 0) continue; } - sj.beginObject() catch return; - sj.objectField("name") catch return; - sj.write(name) catch return; - sj.objectField("index") catch return; - sj.write(idx) catch return; - sj.objectField("checksum") catch return; + try sj.beginObject(); + try sj.objectField("name"); + try sj.write(name); + try sj.objectField("index"); + try sj.write(idx); + try sj.objectField("checksum"); var hex_buf: [16]u8 = undefined; - sj.write(checksumHex(&hex_buf, header.checksum)) catch return; - sj.objectField("minSeq") catch return; - sj.write(header.min_seq) catch return; - sj.objectField("maxSeq") catch return; - sj.write(header.max_seq) catch return; - sj.objectField("mode") catch return; - sj.write(if (whole) "segment" else "blocks") catch return; + try sj.write(checksumHex(&hex_buf, header.checksum)); + try sj.objectField("minSeq"); + try sj.write(header.min_seq); + try sj.objectField("maxSeq"); + try sj.write(header.max_seq); + try sj.objectField("mode"); + try sj.write(if (whole) "segment" else "blocks"); if (!whole) { - sj.objectField("blocks") catch return; - writeBlockRanges(&sj, admitted) catch return; + try sj.objectField("blocks"); + try writeBlockRanges(&sj, admitted); } - sj.endObject() catch return; + try sj.endObject(); segs_matched += 1; if (whole) blocks_matched += match_count; } - sj.endArray() catch return; + try sj.endArray(); planned_through = if (truncated) last_unit_max_seq else sealed_tip; - s.beginObject() catch return; - s.objectField("plannedThroughSeq") catch return; - s.write(planned_through) catch return; - s.objectField("sealedTipSeq") catch return; - s.write(sealed_tip) catch return; - s.objectField("segments") catch return; - s.beginWriteRaw() catch return; - s.writer.writeAll(segments_json.written()) catch return; + try s.beginObject(); + try s.objectField("plannedThroughSeq"); + try s.write(planned_through); + try s.objectField("sealedTipSeq"); + try s.write(sealed_tip); + try s.objectField("segments"); + try s.beginWriteRaw(); + try s.writer.writeAll(segments_json.written()); s.endWriteRaw(); - s.objectField("stats") catch return; - s.beginObject() catch return; - s.objectField("segmentsExamined") catch return; - s.write(segs_examined) catch return; - s.objectField("segmentsMatched") catch return; - s.write(segs_matched) catch return; - s.objectField("blocksMatched") catch return; - s.write(blocks_matched) catch return; - s.objectField("entries") catch return; - s.write(entries) catch return; - s.endObject() catch return; - s.endObject() catch return; - // Do not let a slow HTTP reader hold the planner lock while the - // already-materialized JSON is written. - plan_locked = false; - self.plan_mutex.unlock(self.io); - respond.json(200, out.written()); + try s.objectField("stats"); + try s.beginObject(); + try s.objectField("segmentsExamined"); + try s.write(segs_examined); + try s.objectField("segmentsMatched"); + try s.write(segs_matched); + try s.objectField("blocksMatched"); + try s.write(blocks_matched); + try s.objectField("entries"); + try s.write(entries); + try s.endObject(); + try s.endObject(); } /// which blocks of `sealed` may contain matching rows: seq-window pruned, @@ -1096,6 +1112,7 @@ fn matchBlocks(alloc: Allocator, sealed: *const @import("manifest.zig").Planning } for (matched, 0..) |*m, block| { if (!m.*) continue; + if (ci.entries.len == 0 or ci.blockRow(block) == null) continue; // missing metadata fails open m.* = false; for (wanted, 0..) |w, id| { if (w and ci.blockHasCollection(block, id)) { @@ -1242,6 +1259,186 @@ test "plan block ranges use lexicon first and last fields" { ); } +// Upstream's planner is one-sided: degraded resident metadata must include a +// block (`blockHasAnyCollection` returns true past the slice), never prune it. +// Pruning here is a silent false negative in a collection-filtered backfill. +// `.resident` is untouched because these requests carry no DIDs. +fn planFixture( + blocks: []const segment.BlockIndexEntry, + ci: *const segment.CollectionIndex, +) @import("manifest.zig").PlanningSegment { + return .{ + .resident = undefined, + .summary = .{ .idx = 0, .size = 0, .mtime = undefined, .header = undefined }, + .blocks = blocks, + .segment_bloom = &.{}, + .segment_bloom_valid = false, + .collection_index = ci, + }; +} + +fn oneBlock(seq: u64) segment.BlockIndexEntry { + return .{ + .offset = 0, + .compressed_size = 1, + .uncompressed_size = 1, + .event_count = 1, + .min_seq = seq, + .max_seq = seq, + .min_witnessed_at = 0, + .max_witnessed_at = 0, + }; +} + +test "collection filter fails open when a block has no resident bitmask row" { + var arena = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena.deinit(); + const alloc = arena.allocator(); + + const blocks = [_]segment.BlockIndexEntry{ oneBlock(1), oneBlock(2) }; + var entries = [_]segment.CollectionEntry{.{ .nsid = "app.bsky.feed.post", .count = 1 }}; + // one row of bitmasks for two blocks: block 1 has no resident row. + var bitmasks = [_]u8{0}; + const ci: segment.CollectionIndex = .{ + .entries = &entries, + .bitmasks = &bitmasks, + .bitmask_len = 1, + .buffer = &.{}, + }; + const sealed = planFixture(&blocks, &ci); + + const matched = try matchBlocks(alloc, &sealed, .{ .collections = &.{"app.bsky.feed.post"} }); + try std.testing.expect(!matched[0]); // resident row says the collection is absent + try std.testing.expect(matched[1]); // no resident row must not prune +} + +test "collection filter fails open when a segment carries no collection index" { + var arena = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena.deinit(); + const alloc = arena.allocator(); + + const blocks = [_]segment.BlockIndexEntry{ oneBlock(1), oneBlock(2) }; + const ci: segment.CollectionIndex = .{ + .entries = &.{}, + .bitmasks = &.{}, + .bitmask_len = 0, + .buffer = &.{}, + }; + const sealed = planFixture(&blocks, &ci); + + const matched = try matchBlocks(alloc, &sealed, .{ .collections = &.{"app.bsky.feed.post"} }); + for (matched) |m| try std.testing.expect(m); +} + +test "exact collection reporting does not fail open" { + var entries = [_]segment.CollectionEntry{.{ .nsid = "app.bsky.feed.post", .count = 1 }}; + var bitmasks = [_]u8{0}; + const ci: segment.CollectionIndex = .{ + .entries = &entries, + .bitmasks = &bitmasks, + .bitmask_len = 1, + .buffer = &.{}, + }; + // inspect-segment reports membership as fact; a missing row is not a hit. + try std.testing.expect(!ci.blockHasCollection(1, 0)); + try std.testing.expect(ci.blockRow(1) == null); +} + +test "matchBlocks propagates allocation failure instead of planning a false negative" { + const blocks = [_]segment.BlockIndexEntry{oneBlock(1)}; + const ci: segment.CollectionIndex = .{ + .entries = &.{}, + .bitmasks = &.{}, + .bitmask_len = 0, + .buffer = &.{}, + }; + const sealed = planFixture(&blocks, &ci); + + var failing = std.testing.FailingAllocator.init(std.testing.allocator, .{ .fail_index = 0 }); + try std.testing.expectError( + error.OutOfMemory, + matchBlocks(failing.allocator(), &sealed, .{ .dids = &.{"did:plc:abc"} }), + ); +} + +const PlanResponder = struct { + code: u16 = 0, + body: [8192]u8 = undefined, + body_len: usize = 0, + + fn err(self: *PlanResponder, code: u16, _: []const u8, _: []const u8) void { + self.code = code; + } + fn json(self: *PlanResponder, code: u16, body: []const u8) void { + self.code = code; + self.body_len = @min(body.len, self.body.len); + @memcpy(self.body[0..self.body_len], body[0..self.body_len]); + } + fn written(self: *const PlanResponder) []const u8 { + return self.body[0..self.body_len]; + } +}; + +// An allocation failure anywhere inside planning must never surface as a 200 +// carrying a smaller plan. A client cannot distinguish that from "no matching +// data" and would skip the missing range forever. Sweeping the fail index +// keeps this honest as the allocation sequence changes. +test "planBackfill never answers 200 with a short plan under allocation failure" { + const testing = std.testing; + var threaded: Io.Threaded = .init(testing.allocator, .{}); + defer threaded.deinit(); + const io = threaded.io(); + + var tmp = testing.tmpDir(.{}); + defer tmp.cleanup(); + var path_buffer: [Io.Dir.max_path_bytes]u8 = undefined; + const path_len = try tmp.dir.realPath(io, &path_buffer); + var archive = try archive_mod.Archive.init(testing.allocator, io, path_buffer[0..path_len]); + defer archive.deinit(); + + for (0..6) |i| { + _ = try archive.append(.{ + .seq = i, + .witnessed_at = @intCast(i + 1), + .indexed_at = 0, + .kind = .create, + .did = "did:plc:planfault", + .collection = "app.bsky.feed.post", + .rkey = "3l3qo2vuowo2b", + .rev = "3l3qo2vutsw2b", + .payload = "\xa1\x64\x74\x65\x78\x74\x62\x68\x69", + }, 1); + if (i == 2) try archive.rotate(); + } + try archive.rotate(); + + const request = "{\"collections\":[\"app.bsky.feed.post\"]}"; + + var baseline_api: Api = .{ .allocator = testing.allocator, .io = io, .archive = &archive }; + var baseline: PlanResponder = .{}; + planBackfill(&baseline_api, &baseline, request); + try testing.expectEqual(@as(u16, 200), baseline.code); + try testing.expect(std.mem.indexOf(u8, baseline.written(), "\"segmentsMatched\":2") != null); + + var unanswered: usize = 0; + for (0..256) |fail_index| { + var failing = testing.FailingAllocator.init(testing.allocator, .{ .fail_index = fail_index }); + var api: Api = .{ .allocator = failing.allocator(), .io = io, .archive = &archive }; + var response: PlanResponder = .{}; + planBackfill(&api, &response, request); + if (response.code == 200) { + try testing.expectEqualStrings(baseline.written(), response.written()); + } else if (response.code == 0) { + unanswered += 1; + } else { + try testing.expect(response.code >= 500); + } + } + // Encoding paths that return without responding are a separate defect; + // this asserts the planner does not invent a truthful-looking short plan. + try testing.expect(unanswered == 0); +} + test "plan request rejects invalid filters and windows" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); -- 2.51.2