From 824bd5aeffe42e6d43e3ad989a260128b7294db8 Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Sat, 25 Jul 2026 00:55:35 -0500 Subject: [PATCH] status: add the segments and collections operator views Upstream's status surface has segments and collections views; Stream had neither, so an operator watching a backfill had no way to see what the archive actually contained. Both render from metadata the manifest already holds resident. Segments comes straight off the sealed summaries; collections folds each segment's resident collection index into one table. Neither opens a segment file, so a request costs a walk over sealed segments rather than anything proportional to archive size -- the same bounded-cost rule the host view follows. DID-level marker sentinels are excluded: they carry no real collection and would read as one. The archive contract cross-checks every segment row against listSegments so the view cannot drift into a second source of truth, and asserts no sentinel leaks. The status contract covers the empty archive, where the views must say so plainly rather than render a blank table that looks like a healthy zero. zig build test (304), ReleaseSafe, differential-oracle, archive-contract, status-contract. Audit: status summary/collections/segments blocked -> partial; the summary's own shape and the unsurfaced timing fields are named as what remains. Co-Authored-By: Claude Opus 5 (1M context) --- docs/semantic-parity.md | 2 +- src/internal/server.zig | 77 ++++++++++++++ src/internal/status_page.zig | 191 +++++++++++++++++++++++++++++++++++ tests/backfill_api.py | 27 +++++ tests/status_contract.py | 8 ++ 5 files changed, 304 insertions(+), 1 deletion(-) 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 = + \\
+; + +fn writeNav(w: *Writer, active: []const u8) !void { + try w.writeAll(""); +} + +fn renderSegmentsHtml(w: *Writer, rows: []const SegmentRow) !void { + try w.writeAll(page_head); + try w.writeAll("stream — segments" ++ page_style); + try w.writeAll("

Sealed segments

"); + try writeNav(w, "segments"); + try w.writeAll("
"); + if (rows.len == 0) { + try w.writeAll("
No sealed segments yet.
"); + return; + } + var events: u64 = 0; + var bytes: u64 = 0; + for (rows) |row| { + events += row.event_count; + bytes += row.size_bytes; + } + try w.print( + "

{d} sealed segments, {d} events, {d} bytes on disk.

", + .{ rows.len, events, bytes }, + ); + try w.writeAll("
" ++ + "" ++ + ""); + for (rows) |row| { + try w.print( + "" ++ + "" ++ + "", + .{ row.index, row.event_count, row.unique_did_count, row.block_count, row.min_seq, row.max_seq, row.size_bytes }, + ); + } + try w.writeAll("
IndexEventsUnique DIDsBlocksMin seqMax seqBytes
{d}{d}{d}{d}{d}{d}{d}
"); +} + +fn renderCollectionsHtml(w: *Writer, rows: []const CollectionRow) !void { + try w.writeAll(page_head); + try w.writeAll("stream — collections" ++ page_style); + try w.writeAll("

Archived collections

"); + try writeNav(w, "collections"); + try w.writeAll("
"); + if (rows.len == 0) { + try w.writeAll("
No collections have been archived yet.
"); + return; + } + try w.print("

{d} collections across the sealed archive.

", .{rows.len}); + try w.writeAll(""); + try w.writeAll("
" ++ + ""); + for (rows) |row| { + try w.writeAll("", .{ row.event_count, row.segment_count }); + } + try w.writeAll("
CollectionEventsSegments
"); + try writeHtmlEscaped(w, row.nsid); + try w.print("{d}{d}
"); +} diff --git a/tests/backfill_api.py b/tests/backfill_api.py index 20f9166..c30004b 100644 --- a/tests/backfill_api.py +++ b/tests/backfill_api.py @@ -292,3 +292,30 @@ assert malformed.headers.get("Content-Range") is None assert malformed.read() == b"invalid range\n" print("servecontent conditional/range semantics: PASS") + +# Operator archive views. These render from resident manifest metadata only, +# so they must agree exactly with what listSegments reports rather than being +# a second, drifting source of truth. +seg_text = urllib.request.urlopen(SERVER + "/status?tab=segments").read().decode() +assert "stream sealed segments" in seg_text, seg_text[:200] +for seg in segs: + assert f"{seg['index']}\t{seg['eventCount']}" in seg_text, (seg, seg_text[:400]) +total_events = sum(s["eventCount"] for s in segs) +assert f"{len(segs)} sealed segments, {total_events} events" in seg_text, seg_text[-200:] + +col_text = urllib.request.urlopen(SERVER + "/status?tab=collections").read().decode() +assert "stream archived collections" in col_text, col_text[:200] +assert "app.bsky.feed.post" in col_text, col_text[:400] +# DID-level marker sentinels are not collections and must not be listed. +assert "$" not in col_text, col_text[:400] + +req = urllib.request.Request( + SERVER + "/status?tab=segments", headers={"Accept": "text/html"} +) +seg_html = urllib.request.urlopen(req).read().decode() +assert "Sealed segments" in seg_html +# The new views must join the existing nav, not replace it. +for label in ("Summary", "Hosts", "Accounts", "Segments", "Collections"): + assert f">{label}" in seg_html, label + +print(f"archive status views: PASS ({len(segs)} segments, {total_events} events)") diff --git a/tests/status_contract.py b/tests/status_contract.py index 700a07c..4f72146 100644 --- a/tests/status_contract.py +++ b/tests/status_contract.py @@ -75,3 +75,11 @@ print( "STATUS HTTP CONTRACT PASS " "(durable host/account rows, restart, sorts, filtering, aliases, HEAD, escaping)" ) + +# The archive views exist even with no archive, and say so plainly rather than +# rendering an empty table that reads like a healthy zero. +_, seg_empty = get("/status?tab=segments", "text/plain") +assert "No sealed segments yet." in seg_empty, seg_empty[:200] +_, col_empty = get("/status?tab=collections", "text/plain") +assert "No collections have been archived yet." in col_empty, col_empty[:200] +print("STATUS ARCHIVE VIEWS PASS (empty-state is explicit, not a blank table)") -- 2.51.2