Sealed segments
{d} sealed segments, {d} events, {d} bytes on disk.
", + .{ rows.len, events, bytes }, + ); + try w.writeAll("| Index | " ++ + "Events | Unique DIDs | Blocks | " ++ + "Min seq | Max seq | Bytes |
|---|---|---|---|---|---|---|
| {d} | {d} | {d} | " ++ + "{d} | {d} | {d} | " ++ + "{d} |
diff --git a/docs/semantic-parity.md b/docs/semantic-parity.md index 360cf3b..2f23c06 100644 --- a/docs/semantic-parity.md +++ b/docs/semantic-parity.md @@ -66,7 +66,7 @@ blockers remain. The detailed bootstrap audit is in | Delete/update compaction | **verified** | Rewrite selection, survivor correctness, manifest-before-watermark ordering, fsync/rename cutpoints, bounded workers, and reason counters have physical-file tests. The versioned watermark now distinguishes absence from Store failure, rejects wrong width/version, and refuses to initialize over corrupt bytes; focused tests inject the production Store read fault. | This row does not close merge orchestration or cold-serving gaps. | | Timestamp import | **partial** | CSV parsing, rule precedence, patch topology preservation, job persistence, restart, HTTP authentication, and write/fsync/rename power-loss cuts have concrete tests. | Public visibility after import still depends on the blocked cold-replay path. No audit has yet proved behavior when timestamp metadata reads are corrupt rather than absent. | | Status: hosts and accounts | **partial** | Durable host aggregates and account lookup/verification have focused production HTTP tests, including rate limiting and restart. | This is not the complete upstream status surface. | -| Status: summary, collections, segments | **blocked** | Stream has a custom summary. | Upstream exposes phase entry time, backfill timing/cursor/counts, live cursors, storage/manifest data, imports, and collections/segments views. Stream implements no collections or segments tabs and lacks several underlying persisted fields. | +| Status: summary, collections, segments | **partial** | Stream has a custom summary. Phase entry time and backfill timing are now persisted (`phase/entered_at`, `backfill/timing/*`) in the same synced batch as the phase. Segments and collections views exist at `?tab=segments` and `?tab=collections`, rendering from resident manifest metadata only — no segment file is opened, so cost is bounded by segment count. The archive contract cross-checks every row against `listSegments` and asserts DID-marker sentinels are excluded; the status contract asserts the empty state says so explicitly rather than rendering a blank table. | The summary itself is still Stream's own shape rather than upstream's, and the persisted timing fields are not yet surfaced on it. Live cursors, storage/manifest totals and import history remain absent from the page. | | Prometheus metric exposition | **partial** | Metrics are emitted in Prometheus syntax, and focused tests show that a subset increments at their intended boundaries. | The dashboard contract checks metric-family presence, not producer semantics, label equivalence, monotonicity, or non-placeholder behavior. Several metrics described as durable progress were process-local or tied to the wrong batch boundary. | | Grafana dashboard reuse | **partial** | The upstream dashboard is checksum-pinned; deliberate runtime-only panels are separated from upstream panels; the public experiment dashboard can be provisioned read-only. | A query finding a named family does not prove semantic equivalence. Progress panels must be derived only from genuinely durable, restart-stable counters. Dashboard admission remains blocked by producer semantics. | | Public/debug listeners | **verified** | Production-process tests establish route isolation, disabled debug binding, public WebSocket admission, readiness, and the historical combined listener mode. | This row does not establish lifecycle readiness or data completeness. | diff --git a/src/internal/server.zig b/src/internal/server.zig index 3465483..c0ab838 100644 --- a/src/internal/server.zig +++ b/src/internal/server.zig @@ -996,6 +996,20 @@ pub const Handler = struct { observation.code = 200; return respondStatus(conn, if (wants_html) "text/html; charset=utf-8" else "text/plain; charset=utf-8", page, status_head, now_us); } + if (status_page.segmentTabFromQuery(query)) |which| { + const archive = hub.archive orelse { + observation.code = 500; + return respond(conn, "500 Internal Server Error", "text/plain", "archive diagnostics unavailable\n"); + }; + const now_us = Io.Timestamp.now(hub.io, .real).toMicroseconds(); + const page = renderArchiveTab(hub, archive, which, wants_html) catch { + observation.code = 500; + return respond(conn, "500 Internal Server Error", "text/plain", "failed to read archive diagnostics\n"); + }; + defer hub.allocator.free(page); + observation.code = 200; + return respondStatus(conn, if (wants_html) "text/html; charset=utf-8" else "text/plain; charset=utf-8", page, status_head, now_us); + } if (status_page.hostSortFromQuery(query)) |sort| { const store = hub.repo_status orelse { observation.code = 500; @@ -1353,6 +1367,69 @@ fn statusLine(status: u16) []const u8 { }; } +/// Render the segments or collections view from resident manifest metadata. +/// Both walk sealed segments only; neither opens a segment file, so the cost +/// is bounded by segment count rather than by archive size. +fn renderArchiveTab( + hub: *Hub, + archive: *archive_mod.Archive, + which: []const u8, + html: bool, +) ![]u8 { + var arena = std.heap.ArenaAllocator.init(hub.allocator); + defer arena.deinit(); + const alloc = arena.allocator(); + + if (std.mem.eql(u8, which, "segments")) { + const summaries = try archive.manifest.all(alloc); + const rows = try alloc.alloc(status_page.SegmentRow, summaries.len); + for (summaries, rows) |meta, *row| { + row.* = .{ + .index = meta.idx, + .size_bytes = meta.size, + .event_count = meta.header.event_count, + .unique_did_count = meta.header.unique_did_count, + .block_count = meta.header.block_count, + .min_seq = meta.header.min_seq, + .max_seq = meta.header.max_seq, + .min_witnessed_at = meta.header.min_witnessed_at, + .max_witnessed_at = meta.header.max_witnessed_at, + }; + } + return status_page.renderSegmentsAlloc(hub.allocator, rows, html); + } + + // Collections: fold each sealed segment's resident collection index into + // one archive-wide table. `$`-prefixed entries are DID-level marker + // sentinels, not real collections, so they stay out of the operator view. + const resident = try archive.manifest.planningSegments(alloc); + defer for (resident) |*item| item.deinit(); + var totals: std.StringHashMapUnmanaged(status_page.CollectionRow) = .empty; + for (resident) |*sealed| { + for (sealed.collection_index.entries) |entry| { + if (entry.nsid.len == 0 or entry.nsid[0] == '$') continue; + const gop = try totals.getOrPut(alloc, entry.nsid); + if (gop.found_existing) { + gop.value_ptr.event_count += entry.count; + gop.value_ptr.segment_count += 1; + } else { + gop.key_ptr.* = try alloc.dupe(u8, entry.nsid); + gop.value_ptr.* = .{ + .nsid = gop.key_ptr.*, + .event_count = entry.count, + .segment_count = 1, + }; + } + } + } + const rows = try alloc.alloc(status_page.CollectionRow, totals.count()); + var it = totals.valueIterator(); + var at: usize = 0; + while (it.next()) |value| : (at += 1) rows[at] = value.*; + status_page.sortCollections(rows); + return status_page.renderCollectionsAlloc(hub.allocator, rows, html); +} + fn respond(conn: *websocket.Conn, status: []const u8, content_type: []const u8, resp_body: []const u8) void { respondMaybeHead(conn, status, content_type, resp_body, false); } diff --git a/src/internal/status_page.zig b/src/internal/status_page.zig index 6af67c0..0ddb08b 100644 --- a/src/internal/status_page.zig +++ b/src/internal/status_page.zig @@ -662,3 +662,194 @@ test "account lookup validates DID and reports a resolved DID missing locally" { try testing.expect(!missing.found); try testing.expectEqualStrings("did:plc:missing", missing.did); } + +// === segments and collections === +// +// Upstream's status surface exposes a segments view and a collections view. +// Both render from metadata the manifest already holds resident, so a request +// costs a walk over sealed segments rather than any disk read — the same +// bounded-cost rule the host view follows. + +pub const SegmentRow = struct { + index: u64, + size_bytes: u64, + event_count: u32, + unique_did_count: u32, + block_count: u32, + min_seq: u64, + max_seq: u64, + min_witnessed_at: i64, + max_witnessed_at: i64, +}; + +pub const CollectionRow = struct { + nsid: []const u8, + event_count: u64, + segment_count: u32, +}; + +pub fn segmentTabFromQuery(query: []const u8) ?[]const u8 { + var it = std.mem.splitScalar(u8, query, '&'); + while (it.next()) |pair| { + const eq = std.mem.indexOfScalar(u8, pair, '=') orelse continue; + if (!std.mem.eql(u8, pair[0..eq], "tab")) continue; + const value = pair[eq + 1 ..]; + if (std.mem.eql(u8, value, "segments") or std.mem.eql(u8, value, "collections")) return value; + } + return null; +} + +pub fn renderSegmentsAlloc(allocator: Allocator, rows: []const SegmentRow, html: bool) ![]u8 { + var output: Writer.Allocating = .init(allocator); + errdefer output.deinit(); + if (html) try renderSegmentsHtml(&output.writer, rows) else try renderSegmentsText(&output.writer, rows); + return output.toOwnedSlice(); +} + +pub fn renderCollectionsAlloc(allocator: Allocator, rows: []const CollectionRow, html: bool) ![]u8 { + var output: Writer.Allocating = .init(allocator); + errdefer output.deinit(); + if (html) try renderCollectionsHtml(&output.writer, rows) else try renderCollectionsText(&output.writer, rows); + return output.toOwnedSlice(); +} + +/// Descending event count, then nsid, so the busiest collections lead and the +/// order is stable for an operator refreshing the page. +pub fn sortCollections(rows: []CollectionRow) void { + std.mem.sort(CollectionRow, rows, {}, struct { + fn lessThan(_: void, a: CollectionRow, b: CollectionRow) bool { + if (a.event_count != b.event_count) return a.event_count > b.event_count; + return std.mem.order(u8, a.nsid, b.nsid) == .lt; + } + }.lessThan); +} + +fn renderSegmentsText(w: *Writer, rows: []const SegmentRow) !void { + try w.writeAll("stream sealed segments\n\n"); + if (rows.len == 0) { + try w.writeAll("No sealed segments yet.\n"); + return; + } + try w.writeAll("index\tevents\tunique dids\tblocks\tmin seq\tmax seq\tbytes\n"); + var events: u64 = 0; + var bytes: u64 = 0; + for (rows) |row| { + try w.print("{d}\t{d}\t{d}\t{d}\t{d}\t{d}\t{d}\n", .{ + row.index, row.event_count, row.unique_did_count, + row.block_count, row.min_seq, row.max_seq, + row.size_bytes, + }); + events += row.event_count; + bytes += row.size_bytes; + } + try w.print("\n{d} sealed segments, {d} events, {d} bytes\n", .{ rows.len, events, bytes }); +} + +fn renderCollectionsText(w: *Writer, rows: []const CollectionRow) !void { + try w.writeAll("stream archived collections\n\n"); + if (rows.len == 0) { + try w.writeAll("No collections have been archived yet.\n"); + return; + } + try w.writeAll("collection\tevents\tsegments\n"); + for (rows) |row| { + try w.print("{s}\t{d}\t{d}\n", .{ row.nsid, row.event_count, row.segment_count }); + } + try w.print("\n{d} collections\n", .{rows.len}); +} + +const page_head = + \\
+ \\ + \\ +; + +/// The hosts and accounts pages already ship this stylesheet; segments and +/// collections reuse it verbatim so the operator surface stays one design. +const page_style = + \\{d} sealed segments, {d} events, {d} bytes on disk.
", + .{ rows.len, events, bytes }, + ); + try w.writeAll("| Index | " ++ + "Events | Unique DIDs | Blocks | " ++ + "Min seq | Max seq | Bytes |
|---|---|---|---|---|---|---|
| {d} | {d} | {d} | " ++ + "{d} | {d} | {d} | " ++ + "{d} |