diff --git a/README.md b/README.md index 5dfdc3d..4392e0f 100644 --- a/README.md +++ b/README.md @@ -67,6 +67,7 @@ start the server or open RocksDB: ``` stream version stream inspect-segment --blocks=table --blocks-truncate=100 +stream inspect-all --data-dir=./data --collections-truncate=100 ``` `inspect-segment` accepts `summary`, `table`, and `full` detail, handles both @@ -74,6 +75,10 @@ sealed and active files, and reports a readable checksum-corrupt sealed file as invalid without weakening the checksum enforcement used by archive reads. The committed upstream-produced JSS fixture pins its rendered output byte for byte (apart from the file path). +`inspect-all` walks both the steady archive and Stream's bootstrap live-capture +tree, aggregates network/tree/collection totals, tolerates a racing tail file, +and supports upstream's `--skip-unsealed` and collection-table truncation +controls. runtime flags are listed in `src/main.zig`; deploy notes in `deploy/README.md`. The Zig build pins and statically compiles its native dependencies, including diff --git a/docs/configuration-parity.md b/docs/configuration-parity.md index 617e505..18e835d 100644 --- a/docs/configuration-parity.md +++ b/docs/configuration-parity.md @@ -28,7 +28,7 @@ the same runtime mechanism and an offline receipt observes that behavior. | public/debug bind addresses, shutdown/client drain timeouts | open | Stream has one port and fixed shutdown behavior | | log level/format and OpenTelemetry | open | canonical metrics exist; upstream process configuration and tracing surface do not | | `JETSTREAM_*` environment sources and unknown-variable rejection | open | Stream currently accepts CLI arguments only | -| serve/version/inspect command surface | partial | `serve` is accepted explicitly while retaining Stream's historical flag-only invocation. `version` emits the upstream build-information shape, and `inspect-segment` matches the pinned renderer byte-for-byte on an upstream-produced sealed fixture while also walking active files, surfacing partial tails, and labelling readable checksum corruption. Database-wide `inspect-all` remains open. | +| serve/version/inspect command surface | closed | `serve` is accepted explicitly while retaining Stream's historical flag-only invocation. `version` emits the upstream build-information shape. `inspect-segment` matches the pinned renderer byte-for-byte on an upstream-produced sealed fixture while also walking active files, surfacing partial tails, and labelling readable checksum corruption. `inspect-all` folds the real steady and bootstrap trees with upstream's missing-root, racing-tail, active-skip, aggregation, sorting, truncation, and text-rendering semantics; a copied upstream golden report is byte-exact. Listener and root logging controls remain tracked by their separate rows. | This table tracks configuration semantics only. It does not override the behavioral gates in `semantic-parity.md` or `bootstrap-semantic-parity.md`. diff --git a/docs/jss-format-v1.md b/docs/jss-format-v1.md index 1838d0c..bc509a8 100644 --- a/docs/jss-format-v1.md +++ b/docs/jss-format-v1.md @@ -86,5 +86,5 @@ JSS metadata checksum still passes, and requires block decode to fail. ## ops notes - no CLI flag forces small segments; sealed files come from the bootstrap→merge cutover (`--max-backfill-repos N` gives a small sealed set). steady-state shutdown does NOT seal. -- paths: `/segments/seg_<10-char base36>.jss` (+ bootstrap `backfill/live_segments/`) +- paths: `/segments/seg_<10-char base36>.jss` (+ Stream's bootstrap capture at `backfill/live/segments/`; upstream uses `backfill/live_segments/`) - ground-truth tooling: `jetstream inspect-segment --blocks=full`, `inspect-all` diff --git a/src/internal/inspect_all.zig b/src/internal/inspect_all.zig new file mode 100644 index 0000000..44b9efb --- /dev/null +++ b/src/internal/inspect_all.zig @@ -0,0 +1,654 @@ +//! Database-wide offline JSS inspection. +//! +//! This is the Stream storage-layout equivalent of upstream `inspect-all`: +//! aggregate every segment in the steady tree and the bootstrap live-capture +//! tree without opening RocksDB or starting any network component. + +const std = @import("std"); +const operational = @import("operational.zig"); +const segment_writer = @import("segment_writer.zig"); + +const Allocator = std.mem.Allocator; +const Io = std.Io; + +pub const Options = struct { + skip_unsealed: bool = false, +}; + +pub const SegmentSummary = struct { + index: u64, + sealed: bool, + event_count: u64, + block_count: usize, + size_bytes: usize, +}; + +pub const TreeAggregate = struct { + dir: []const u8, + sealed_count: usize = 0, + active_count: usize = 0, + compressed_bytes: u64 = 0, + uncompressed_bytes: u64 = 0, + disk_bytes: u64 = 0, + event_count: u64 = 0, + block_count: u64 = 0, + oldest_mtime_us: i64 = 0, + newest_mtime_us: i64 = 0, + min_seq: u64 = 0, + max_seq: u64 = 0, + min_witnessed_at: i64 = 0, + max_witnessed_at: i64 = 0, + latest_segment: ?SegmentSummary = null, +}; + +pub const CollectionAggregate = struct { + nsid: []const u8, + event_count: u64 = 0, + segment_count: usize = 0, + block_count: u64 = 0, +}; + +pub const NetworkTotals = struct { + segments: usize = 0, + sealed_segments: usize = 0, + active_segments: usize = 0, + blocks: u64 = 0, + events: u64 = 0, + collections: usize = 0, + compressed_bytes: u64 = 0, + uncompressed_bytes: u64 = 0, + disk_bytes: u64 = 0, + min_seq: u64 = 0, + max_seq: u64 = 0, + min_witnessed_at: i64 = 0, + max_witnessed_at: i64 = 0, +}; + +pub const Aggregate = struct { + arena: std.heap.ArenaAllocator, + trees: []const TreeAggregate, + collections: []const CollectionAggregate, + network: NetworkTotals, + warnings: []const []const u8, + + pub fn deinit(self: *Aggregate) void { + self.arena.deinit(); + } +}; + +const FileRecord = struct { + index: u64, + name: []const u8, + mtime_us: i64, +}; + +pub fn inspectAll(allocator: Allocator, io: Io, data_dir: []const u8, options: Options) !Aggregate { + var arena = std.heap.ArenaAllocator.init(allocator); + errdefer arena.deinit(); + const alloc = arena.allocator(); + + const steady = try std.fs.path.join(alloc, &.{ data_dir, "segments" }); + const bootstrap = try std.fs.path.join(alloc, &.{ data_dir, "backfill", "live", "segments" }); + const roots = [_][]const u8{ steady, bootstrap }; + var trees: std.ArrayList(TreeAggregate) = .empty; + var collections: std.ArrayList(CollectionAggregate) = .empty; + var collection_ids: std.StringHashMapUnmanaged(usize) = .empty; + var warnings: std.ArrayList([]const u8) = .empty; + for (roots) |root| { + try trees.append(alloc, try scanTree( + alloc, + io, + root, + options, + &collections, + &collection_ids, + &warnings, + )); + } + + std.mem.sort(CollectionAggregate, collections.items, {}, struct { + fn lessThan(_: void, a: CollectionAggregate, b: CollectionAggregate) bool { + return a.event_count > b.event_count or + (a.event_count == b.event_count and std.mem.order(u8, a.nsid, b.nsid) == .lt); + } + }.lessThan); + + const collection_slice = try collections.toOwnedSlice(alloc); + const tree_slice = try trees.toOwnedSlice(alloc); + return .{ + .arena = arena, + .trees = tree_slice, + .collections = collection_slice, + .network = computeNetwork(tree_slice, collection_slice.len), + .warnings = try warnings.toOwnedSlice(alloc), + }; +} + +fn scanTree( + allocator: Allocator, + io: Io, + root: []const u8, + options: Options, + collections: *std.ArrayList(CollectionAggregate), + collection_ids: *std.StringHashMapUnmanaged(usize), + warnings: *std.ArrayList([]const u8), +) !TreeAggregate { + var tree: TreeAggregate = .{ .dir = root }; + var dir = Io.Dir.cwd().openDir(io, root, .{ .iterate = true }) catch |err| switch (err) { + error.FileNotFound => return tree, + else => return err, + }; + defer dir.close(io); + + var files: std.ArrayList(FileRecord) = .empty; + var iterator = dir.iterate(); + while (try iterator.next(io)) |entry| { + if (entry.kind == .directory) continue; + const index = parseSegmentIndex(entry.name) orelse continue; + const stat = try dir.statFile(io, entry.name, .{}); + try files.append(allocator, .{ + .index = index, + .name = try allocator.dupe(u8, entry.name), + .mtime_us = stat.mtime.toMicroseconds(), + }); + } + if (files.items.len == 0) return tree; + std.mem.sort(FileRecord, files.items, {}, struct { + fn lessThan(_: void, a: FileRecord, b: FileRecord) bool { + return a.index < b.index; + } + }.lessThan); + tree.oldest_mtime_us = files.items[0].mtime_us; + tree.newest_mtime_us = files.items[0].mtime_us; + + for (files.items, 0..) |file, at| { + tree.oldest_mtime_us = @min(tree.oldest_mtime_us, file.mtime_us); + tree.newest_mtime_us = @max(tree.newest_mtime_us, file.mtime_us); + const path = try std.fs.path.join(allocator, &.{ root, file.name }); + var inspection = operational.inspectSegment(allocator, io, path) catch |err| { + // The highest-index file may rotate between readdir and inspect. + if (at + 1 != files.items.len) + try warnings.append(allocator, try std.fmt.allocPrint(allocator, "{s}: {s}", .{ path, @errorName(err) })); + continue; + }; + defer inspection.deinit(); + + tree.disk_bytes += inspection.file_size; + if (inspection.sealed) + tree.sealed_count += 1 + else + tree.active_count += 1; + tree.latest_segment = .{ + .index = file.index, + .sealed = inspection.sealed, + .event_count = inspection.total_events, + .block_count = inspection.blocks.len, + .size_bytes = inspection.file_size, + }; + if (!inspection.sealed and options.skip_unsealed) continue; + try foldInspection(&tree, &inspection, allocator, collections, collection_ids); + } + return tree; +} + +fn foldInspection( + tree: *TreeAggregate, + inspection: *const operational.Inspection, + allocator: Allocator, + collections: *std.ArrayList(CollectionAggregate), + collection_ids: *std.StringHashMapUnmanaged(usize), +) !void { + tree.event_count += inspection.total_events; + tree.block_count += inspection.blocks.len; + for (inspection.blocks) |block| { + tree.compressed_bytes += block.compressed_size; + tree.uncompressed_bytes += block.uncompressed_size; + } + if (inspection.total_events > 0) { + if (tree.min_seq == 0 or inspection.min_seq < tree.min_seq) tree.min_seq = inspection.min_seq; + tree.max_seq = @max(tree.max_seq, inspection.max_seq); + if (tree.min_witnessed_at == 0 or inspection.min_witnessed_at < tree.min_witnessed_at) + tree.min_witnessed_at = inspection.min_witnessed_at; + tree.max_witnessed_at = @max(tree.max_witnessed_at, inspection.max_witnessed_at); + } + + const block_counts = try allocator.alloc(u64, inspection.collections.len); + @memset(block_counts, 0); + for (inspection.block_collections) |ids| for (ids) |id| { + if (id < block_counts.len) block_counts[id] += 1; + }; + for (inspection.collections, 0..) |collection, collection_id| { + if (isSentinel(collection.nsid)) continue; + const aggregate_id = collection_ids.get(collection.nsid) orelse new: { + const owned = try allocator.dupe(u8, collection.nsid); + const id = collections.items.len; + try collections.append(allocator, .{ .nsid = owned }); + try collection_ids.put(allocator, owned, id); + break :new id; + }; + collections.items[aggregate_id].event_count += collection.event_count; + collections.items[aggregate_id].segment_count += 1; + collections.items[aggregate_id].block_count += block_counts[collection_id]; + } +} + +fn computeNetwork(trees: []const TreeAggregate, collection_count: usize) NetworkTotals { + var network: NetworkTotals = .{ .collections = collection_count }; + for (trees) |tree| { + network.segments += tree.sealed_count + tree.active_count; + network.sealed_segments += tree.sealed_count; + network.active_segments += tree.active_count; + network.blocks += tree.block_count; + network.events += tree.event_count; + network.compressed_bytes += tree.compressed_bytes; + network.uncompressed_bytes += tree.uncompressed_bytes; + network.disk_bytes += tree.disk_bytes; + if (tree.event_count == 0) continue; + if (network.min_seq == 0 or tree.min_seq < network.min_seq) network.min_seq = tree.min_seq; + network.max_seq = @max(network.max_seq, tree.max_seq); + if (network.min_witnessed_at == 0 or tree.min_witnessed_at < network.min_witnessed_at) + network.min_witnessed_at = tree.min_witnessed_at; + network.max_witnessed_at = @max(network.max_witnessed_at, tree.max_witnessed_at); + } + return network; +} + +pub fn renderInspectAll( + writer: *Io.Writer, + data_dir: []const u8, + aggregate: *const Aggregate, + generated_us: i64, + truncate: usize, +) !void { + var time_buffer: [40]u8 = undefined; + try writer.writeAll("inspect-all\n"); + try writer.print("data-dir: {s}\n", .{data_dir}); + try writer.print("generated: {s}\n", .{formatMicros(&time_buffer, generated_us)}); + try renderNetwork(writer, aggregate.network); + try renderTrees(writer, aggregate.trees); + try renderCollections(writer, aggregate.collections, truncate); + if (aggregate.warnings.len > 0) { + try writer.print("\nwarnings ({d}):\n", .{aggregate.warnings.len}); + for (aggregate.warnings) |warning| try writer.print(" {s}\n", .{warning}); + } +} + +fn renderNetwork(writer: *Io.Writer, network: NetworkTotals) !void { + try writer.writeAll("\nnetwork totals:\n"); + try writer.print(" segments: {d} ({d} sealed, {d} active)\n", .{ network.segments, network.sealed_segments, network.active_segments }); + try writer.print(" blocks: {f}\n", .{commaInt(network.blocks)}); + try writer.print(" events: {f}\n", .{commaInt(network.events)}); + try writer.print(" collections: {d} distinct NSIDs\n", .{network.collections}); + if (network.events > 0) { + try writer.print(" seq range: [{d}, {d}]\n", .{ network.min_seq, network.max_seq }); + var min_buffer: [40]u8 = undefined; + var max_buffer: [40]u8 = undefined; + try writer.print(" witnessed_at range: {s} → {s}\n", .{ + formatMicros(&min_buffer, network.min_witnessed_at), + formatMicros(&max_buffer, network.max_witnessed_at), + }); + } + try writer.print(" payload (uncompressed): {f}\n", .{humanBytes(network.uncompressed_bytes)}); + try writer.print(" payload (compressed): {f}\n", .{humanBytes(network.compressed_bytes)}); + try writer.print(" disk usage: {f}\n", .{humanBytes(network.disk_bytes)}); + if (network.compressed_bytes > 0) { + const ratio: f64 = @as(f64, @floatFromInt(network.uncompressed_bytes)) / + @as(f64, @floatFromInt(network.compressed_bytes)); + try writer.print(" compression ratio: {d:.2}x\n", .{ratio}); + } +} + +fn renderTrees(writer: *Io.Writer, trees: []const TreeAggregate) !void { + try writer.writeAll("\ntrees:\n"); + if (trees.len == 0) return writer.writeAll(" (none)\n"); + for (trees, 0..) |tree, index| { + try writer.print(" [{d}] {s}\n", .{ index, tree.dir }); + if (tree.sealed_count + tree.active_count == 0) { + try writer.writeAll(" (empty)\n"); + continue; + } + try writer.print(" files: {d} sealed + {d} active\n", .{ tree.sealed_count, tree.active_count }); + try writer.print(" events: {f}\n", .{commaInt(tree.event_count)}); + try writer.print(" blocks: {f}\n", .{commaInt(tree.block_count)}); + if (tree.event_count > 0) { + try writer.print(" seq range: [{d}, {d}]\n", .{ tree.min_seq, tree.max_seq }); + var min_buffer: [40]u8 = undefined; + var max_buffer: [40]u8 = undefined; + try writer.print(" witnessed_at: {s} → {s}\n", .{ + formatMicros(&min_buffer, tree.min_witnessed_at), + formatMicros(&max_buffer, tree.max_witnessed_at), + }); + } + if (tree.oldest_mtime_us != 0) { + var oldest_buffer: [40]u8 = undefined; + var newest_buffer: [40]u8 = undefined; + try writer.print(" oldest mtime: {s}\n", .{formatMicros(&oldest_buffer, tree.oldest_mtime_us)}); + try writer.print(" newest mtime: {s}\n", .{formatMicros(&newest_buffer, tree.newest_mtime_us)}); + } + try writer.print(" compressed: {f}\n", .{humanBytes(tree.compressed_bytes)}); + try writer.print(" uncompressed: {f}\n", .{humanBytes(tree.uncompressed_bytes)}); + try writer.print(" disk: {f}\n", .{humanBytes(tree.disk_bytes)}); + if (tree.latest_segment) |latest| { + try writer.print(" latest: idx={d} {s} events={f} blocks={d} size={f}\n", .{ + latest.index, + if (latest.sealed) "sealed" else "active", + commaInt(latest.event_count), + latest.block_count, + humanBytes(latest.size_bytes), + }); + } + } +} + +fn renderCollections(writer: *Io.Writer, collections: []const CollectionAggregate, truncate: usize) !void { + try writer.print("\ncollections ({d} distinct NSIDs):\n", .{collections.len}); + if (collections.len == 0) return writer.writeAll(" (none)\n"); + var nsid_width: usize = "NSID".len; + for (collections) |collection| nsid_width = @max(nsid_width, collection.nsid.len); + + const half = truncate / 2; + for (collections, 0..) |collection, index| { + if (truncate != 0 and collections.len > truncate and index == half) + try writer.print(" ... ({d} rows omitted) ...\n", .{collections.len - 2 * half}); + if (truncate != 0 and collections.len > truncate and index >= half and index < collections.len - half) continue; + var event_buffer: [32]u8 = undefined; + var block_buffer: [32]u8 = undefined; + try writer.print(" [{d:>3}] {s:<[5]} events: {s:>12} segments: {d:>5} blocks: {s:>12}\n", .{ + index, + collection.nsid, + formatComma(&event_buffer, collection.event_count), + collection.segment_count, + formatComma(&block_buffer, collection.block_count), + nsid_width, + }); + } +} + +const CommaInt = struct { + value: u64, + + pub fn format(self: CommaInt, writer: *Io.Writer) Io.Writer.Error!void { + var buffer: [32]u8 = undefined; + try writer.writeAll(formatComma(&buffer, self.value)); + } +}; + +fn commaInt(value: u64) CommaInt { + return .{ .value = value }; +} + +const HumanBytes = struct { + value: u64, + + pub fn format(self: HumanBytes, writer: *Io.Writer) Io.Writer.Error!void { + if (self.value < 1024) return writer.print("{d} B", .{self.value}); + const suffixes = [_][]const u8{ "KiB", "MiB", "GiB", "TiB", "PiB" }; + var divisor: u64 = 1024; + var exponent: usize = 0; + var reduced = self.value / 1024; + while (reduced >= 1024 and exponent + 1 < suffixes.len) : (reduced /= 1024) { + divisor *= 1024; + exponent += 1; + } + const scaled: f64 = @as(f64, @floatFromInt(self.value)) / @as(f64, @floatFromInt(divisor)); + try writer.print("{d:.2} {s}", .{ scaled, suffixes[exponent] }); + } +}; + +fn humanBytes(value: u64) HumanBytes { + return .{ .value = value }; +} + +fn formatComma(buffer: []u8, value: u64) []const u8 { + var raw_buffer: [32]u8 = undefined; + const raw = std.fmt.bufPrint(&raw_buffer, "{d}", .{value}) catch unreachable; + const commas = if (raw.len == 0) 0 else (raw.len - 1) / 3; + const out = buffer[0 .. raw.len + commas]; + var source = raw.len; + var target = out.len; + var digits: usize = 0; + while (source > 0) { + if (digits == 3) { + target -= 1; + out[target] = ','; + digits = 0; + } + source -= 1; + target -= 1; + out[target] = raw[source]; + digits += 1; + } + return out; +} + +fn formatMicros(buffer: []u8, micros: i64) []const u8 { + if (micros == 0) return "0"; + const seconds: u64 = @intCast(@divFloor(micros, std.time.us_per_s)); + const fraction: u64 = @intCast(@mod(micros, std.time.us_per_s)); + const epoch_seconds: std.time.epoch.EpochSeconds = .{ .secs = seconds }; + const year_day = epoch_seconds.getEpochDay().calculateYearDay(); + const month_day = year_day.calculateMonthDay(); + const day_seconds = epoch_seconds.getDaySeconds(); + return std.fmt.bufPrint(buffer, "{d:0>4}-{d:0>2}-{d:0>2}T{d:0>2}:{d:0>2}:{d:0>2}.{d:0>6}Z", .{ + year_day.year, + month_day.month.numeric(), + month_day.day_index + 1, + day_seconds.getHoursIntoDay(), + day_seconds.getMinutesIntoHour(), + day_seconds.getSecondsIntoMinute(), + fraction, + }) catch unreachable; +} + +fn isSentinel(nsid: []const u8) bool { + return std.mem.eql(u8, nsid, "$account") or + std.mem.eql(u8, nsid, "$identity") or + std.mem.eql(u8, nsid, "$sync"); +} + +fn parseSegmentIndex(name: []const u8) ?u64 { + if (!std.mem.startsWith(u8, name, "seg_") or !std.mem.endsWith(u8, name, ".jss")) return null; + const digits = name[4 .. name.len - 4]; + if (digits.len != 10) return null; + return std.fmt.parseInt(u64, digits, 36) catch null; +} + +test "inspect-all scans steady and real bootstrap trees" { + var tmp = std.testing.tmpDir(.{}); + defer tmp.cleanup(); + var steady = try tmp.dir.createDirPathOpen(std.testing.io, "data/segments", .{}); + steady.close(std.testing.io); + var bootstrap = try tmp.dir.createDirPathOpen(std.testing.io, "data/backfill/live/segments", .{}); + bootstrap.close(std.testing.io); + const fixture = try Io.Dir.cwd().readFileAlloc(std.testing.io, "tests/fixtures/seg_0000000000.jss", std.testing.allocator, .limited(1 << 20)); + defer std.testing.allocator.free(fixture); + try tmp.dir.writeFile(std.testing.io, .{ .sub_path = "data/segments/seg_0000000000.jss", .data = fixture }); + try tmp.dir.writeFile(std.testing.io, .{ .sub_path = "data/backfill/live/segments/seg_0000000000.jss", .data = fixture }); + const data_dir = try tmp.dir.realPathFileAlloc(std.testing.io, "data", std.testing.allocator); + defer std.testing.allocator.free(data_dir); + + var aggregate = try inspectAll(std.testing.allocator, std.testing.io, data_dir, .{}); + defer aggregate.deinit(); + try std.testing.expectEqual(@as(usize, 2), aggregate.network.segments); + try std.testing.expectEqual(@as(u64, 14_036), aggregate.network.events); + try std.testing.expectEqual(@as(usize, 5), aggregate.network.collections); + try std.testing.expectEqual(@as(u64, 8_210), aggregate.collections[0].event_count); + try std.testing.expectEqualStrings("app.bsky.feed.post", aggregate.collections[0].nsid); +} + +test "inspect-all skip-unsealed keeps file bookkeeping but excludes active payload" { + var active = try segment_writer.ActiveWriter.init(std.testing.allocator); + defer active.deinit(); + try active.append(.{ + .seq = 1, + .witnessed_at = 1_700_000_000_000_001, + .indexed_at = 0, + .kind = .create, + .did = "did:plc:active", + .collection = "app.bsky.feed.post", + .rkey = "r", + .rev = "v", + .payload = "p", + }); + try active.flush(); + + var tmp = std.testing.tmpDir(.{}); + defer tmp.cleanup(); + var steady = try tmp.dir.createDirPathOpen(std.testing.io, "data/segments", .{}); + steady.close(std.testing.io); + try tmp.dir.writeFile(std.testing.io, .{ .sub_path = "data/segments/seg_0000000000.jss", .data = active.bytes() }); + const data_dir = try tmp.dir.realPathFileAlloc(std.testing.io, "data", std.testing.allocator); + defer std.testing.allocator.free(data_dir); + + var complete = try inspectAll(std.testing.allocator, std.testing.io, data_dir, .{}); + defer complete.deinit(); + try std.testing.expectEqual(@as(u64, 1), complete.network.events); + try std.testing.expectEqual(@as(u64, 1), complete.network.blocks); + + var skipped = try inspectAll(std.testing.allocator, std.testing.io, data_dir, .{ .skip_unsealed = true }); + defer skipped.deinit(); + try std.testing.expectEqual(@as(usize, 1), skipped.network.active_segments); + try std.testing.expect(skipped.network.disk_bytes > 0); + try std.testing.expectEqual(@as(u64, 0), skipped.network.events); + try std.testing.expectEqual(@as(u64, 0), skipped.network.blocks); + try std.testing.expectEqual(@as(u64, 1), skipped.trees[0].latest_segment.?.event_count); +} + +test "inspect-all warns on corrupt non-tail and silently tolerates corrupt tail" { + const fixture = try Io.Dir.cwd().readFileAlloc(std.testing.io, "tests/fixtures/seg_0000000000.jss", std.testing.allocator, .limited(1 << 20)); + defer std.testing.allocator.free(fixture); + var tmp = std.testing.tmpDir(.{}); + defer tmp.cleanup(); + var steady = try tmp.dir.createDirPathOpen(std.testing.io, "data/segments", .{}); + steady.close(std.testing.io); + try tmp.dir.writeFile(std.testing.io, .{ .sub_path = "data/segments/seg_0000000000.jss", .data = fixture }); + try tmp.dir.writeFile(std.testing.io, .{ .sub_path = "data/segments/seg_0000000001.jss", .data = "bad" }); + try tmp.dir.writeFile(std.testing.io, .{ .sub_path = "data/segments/seg_0000000002.jss", .data = fixture }); + try tmp.dir.writeFile(std.testing.io, .{ .sub_path = "data/segments/seg_0000000003.jss", .data = "bad tail" }); + const data_dir = try tmp.dir.realPathFileAlloc(std.testing.io, "data", std.testing.allocator); + defer std.testing.allocator.free(data_dir); + + var aggregate = try inspectAll(std.testing.allocator, std.testing.io, data_dir, .{}); + defer aggregate.deinit(); + try std.testing.expectEqual(@as(usize, 2), aggregate.network.sealed_segments); + try std.testing.expectEqual(@as(usize, 1), aggregate.warnings.len); + try std.testing.expect(std.mem.indexOf(u8, aggregate.warnings[0], "seg_0000000001.jss") != null); + try std.testing.expect(std.mem.indexOf(u8, aggregate.warnings[0], "seg_0000000003.jss") == null); +} + +test "inspect-all collection truncation keeps both ends" { + const collections = [_]CollectionAggregate{ + .{ .nsid = "a", .event_count = 5 }, + .{ .nsid = "b", .event_count = 4 }, + .{ .nsid = "c", .event_count = 3 }, + .{ .nsid = "d", .event_count = 2 }, + .{ .nsid = "e", .event_count = 1 }, + }; + var out: Io.Writer.Allocating = .init(std.testing.allocator); + defer out.deinit(); + try renderCollections(&out.writer, &collections, 2); + try std.testing.expect(std.mem.indexOf(u8, out.written(), "[ 0] a") != null); + try std.testing.expect(std.mem.indexOf(u8, out.written(), "... (3 rows omitted) ...") != null); + try std.testing.expect(std.mem.indexOf(u8, out.written(), "[ 4] e") != null); + try std.testing.expect(std.mem.indexOf(u8, out.written(), "[ 2] c") == null); +} + +test "inspect-all renderer matches upstream golden text" { + var arena = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena.deinit(); + const trees = [_]TreeAggregate{ + .{ + .dir = "/data/segments", + .sealed_count = 3, + .active_count = 1, + .compressed_bytes = 2 * 1024 * 1024, + .uncompressed_bytes = 6 * 1024 * 1024, + .disk_bytes = 3 * 1024 * 1024, + .event_count = 1234, + .block_count = 12, + .min_seq = 10, + .max_seq = 1243, + .min_witnessed_at = 1_779_840_000_000_000, + .max_witnessed_at = 1_779_926_400_000_000, + .oldest_mtime_us = 1_779_840_060_000_000, + .newest_mtime_us = 1_779_926_460_000_000, + .latest_segment = .{ .index = 4, .sealed = false, .event_count = 300, .block_count = 3, .size_bytes = 768 * 1024 }, + }, + .{ .dir = "/data/backfill/live_segments" }, + }; + const collections = [_]CollectionAggregate{ + .{ .nsid = "app.bsky.feed.post", .event_count = 900, .segment_count = 4, .block_count = 9 }, + .{ .nsid = "app.bsky.feed.like", .event_count = 300, .segment_count = 2, .block_count = 2 }, + .{ .nsid = "app.bsky.graph.follow", .event_count = 34, .segment_count = 1, .block_count = 1 }, + }; + const warnings = [_][]const u8{"/data/segments/seg_0000000002.jss: corrupt segment: bad magic \"XXXX\""}; + var aggregate: Aggregate = .{ + .arena = arena, + .trees = &trees, + .collections = &collections, + .network = .{ + .segments = 4, + .sealed_segments = 3, + .active_segments = 1, + .blocks = 12, + .events = 1234, + .collections = 3, + .compressed_bytes = 2 * 1024 * 1024, + .uncompressed_bytes = 6 * 1024 * 1024, + .disk_bytes = 3 * 1024 * 1024, + .min_seq = 10, + .max_seq = 1243, + .min_witnessed_at = 1_779_840_000_000_000, + .max_witnessed_at = 1_779_926_400_000_000, + }, + .warnings = &warnings, + }; + // The local defer owns this arena; prevent Aggregate.deinit here. + _ = &aggregate; + var out: Io.Writer.Allocating = .init(std.testing.allocator); + defer out.deinit(); + try renderInspectAll(&out.writer, "/data", &aggregate, 1_779_971_696_000_000, 100); + const expected = + \\inspect-all + \\data-dir: /data + \\generated: 2026-05-28T12:34:56.000000Z + \\ + \\network totals: + \\ segments: 4 (3 sealed, 1 active) + \\ blocks: 12 + \\ events: 1,234 + \\ collections: 3 distinct NSIDs + \\ seq range: [10, 1243] + \\ witnessed_at range: 2026-05-27T00:00:00.000000Z → 2026-05-28T00:00:00.000000Z + \\ payload (uncompressed): 6.00 MiB + \\ payload (compressed): 2.00 MiB + \\ disk usage: 3.00 MiB + \\ compression ratio: 3.00x + \\ + \\trees: + \\ [0] /data/segments + \\ files: 3 sealed + 1 active + \\ events: 1,234 + \\ blocks: 12 + \\ seq range: [10, 1243] + \\ witnessed_at: 2026-05-27T00:00:00.000000Z → 2026-05-28T00:00:00.000000Z + \\ oldest mtime: 2026-05-27T00:01:00.000000Z + \\ newest mtime: 2026-05-28T00:01:00.000000Z + \\ compressed: 2.00 MiB + \\ uncompressed: 6.00 MiB + \\ disk: 3.00 MiB + \\ latest: idx=4 active events=300 blocks=3 size=768.00 KiB + \\ [1] /data/backfill/live_segments + \\ (empty) + \\ + \\collections (3 distinct NSIDs): + \\ [ 0] app.bsky.feed.post events: 900 segments: 4 blocks: 9 + \\ [ 1] app.bsky.feed.like events: 300 segments: 2 blocks: 2 + \\ [ 2] app.bsky.graph.follow events: 34 segments: 1 blocks: 1 + \\ + \\warnings (1): + \\ /data/segments/seg_0000000002.jss: corrupt segment: bad magic "XXXX" + \\ + ; + try std.testing.expectEqualStrings(expected, out.written()); +} diff --git a/src/main.zig b/src/main.zig index d23d489..c86d119 100644 --- a/src/main.zig +++ b/src/main.zig @@ -28,6 +28,7 @@ const timestamp_rules = @import("internal/timestamp/rules.zig"); const timestamp_jobs = @import("internal/timestamp/jobs.zig"); const timestamp_manager = @import("internal/timestamp/manager.zig"); const operational = @import("internal/operational.zig"); +const inspect_all = @import("internal/inspect_all.zig"); const Io = std.Io; const log = std.log.scoped(.stream); @@ -306,6 +307,45 @@ fn runInspectSegment(allocator: std.mem.Allocator, io: Io, args: []const []const try stdout.interface.flush(); } +fn runInspectAll(allocator: std.mem.Allocator, io: Io, args: []const []const u8) !void { + var data_dir: []const u8 = "./data"; + var skip_unsealed = false; + var truncate: usize = 100; + var i: usize = 0; + while (i < args.len) : (i += 1) { + const arg = args[i]; + if (std.mem.startsWith(u8, arg, "--data-dir=")) { + data_dir = arg["--data-dir=".len..]; + } else if (std.mem.eql(u8, arg, "--data-dir")) { + i += 1; + if (i >= args.len) return error.BadArgs; + data_dir = args[i]; + } else if (std.mem.eql(u8, arg, "--skip-unsealed")) { + skip_unsealed = true; + } else if (std.mem.startsWith(u8, arg, "--collections-truncate=")) { + truncate = try parseNonNegativeUsize(arg["--collections-truncate=".len..]); + } else if (std.mem.eql(u8, arg, "--collections-truncate")) { + i += 1; + if (i >= args.len) return error.BadArgs; + truncate = try parseNonNegativeUsize(args[i]); + } else { + return error.BadArgs; + } + } + var aggregate = try inspect_all.inspectAll(allocator, io, data_dir, .{ .skip_unsealed = skip_unsealed }); + defer aggregate.deinit(); + var buffer: [16 * 1024]u8 = undefined; + var stdout = Io.File.stdout().writer(io, &buffer); + try inspect_all.renderInspectAll( + &stdout.interface, + data_dir, + &aggregate, + Io.Timestamp.now(io, .real).toMicroseconds(), + truncate, + ); + try stdout.interface.flush(); +} + fn parseBlocksMode(raw: []const u8) ?operational.BlocksMode { if (std.mem.eql(u8, raw, "summary")) return .summary; if (std.mem.eql(u8, raw, "table")) return .table; @@ -329,14 +369,20 @@ pub fn main(init: std.process.Init.Minimal) !void { _ = arg_it.next(); // program name var first_arg = arg_it.next(); if (first_arg) |command| { - if (std.mem.eql(u8, command, "version") or std.mem.eql(u8, command, "inspect-segment")) { + if (std.mem.eql(u8, command, "version") or + std.mem.eql(u8, command, "inspect-segment") or + std.mem.eql(u8, command, "inspect-all")) + { var command_args: std.ArrayList([]const u8) = .empty; defer command_args.deinit(allocator); while (arg_it.next()) |arg| try command_args.append(allocator, arg); - if (std.mem.eql(u8, command, "version")) - try writeVersion(io, command_args.items) - else + if (std.mem.eql(u8, command, "version")) { + try writeVersion(io, command_args.items); + } else if (std.mem.eql(u8, command, "inspect-segment")) { try runInspectSegment(allocator, io, command_args.items); + } else { + try runInspectAll(allocator, io, command_args.items); + } return; } // Like upstream's explicit `serve` subcommand, while retaining the diff --git a/src/tests.zig b/src/tests.zig index 9c0af31..339ebff 100644 --- a/src/tests.zig +++ b/src/tests.zig @@ -36,6 +36,7 @@ comptime { _ = @import("internal/compact/steady.zig"); _ = @import("internal/segment_fixture_test.zig"); _ = @import("internal/operational.zig"); + _ = @import("internal/inspect_all.zig"); _ = @import("internal/verify.zig"); _ = @import("internal/wire.zig"); _ = @import("internal/convert.zig");