#!/usr/bin/env python3 """Deterministically adapt Jetstream's dashboard for the Stream runtime.""" from __future__ import annotations import copy import hashlib import json import pathlib import sys ROOT = pathlib.Path(__file__).resolve().parents[1] SOURCE = ROOT / "deploy/grafana/jetstream-upstream.json" OUTPUT = ROOT / "deploy/grafana/stream.json" SOURCE_SHA256 = "ef6033c182d8c3f9d88af7d91e61126724efbe72e6b940984c9de12e9530cd91" def target(template: dict, expr: str, legend: str, ref_id: str) -> dict: value = copy.deepcopy(template) value.update(expr=expr, legendFormat=legend, refId=ref_id) return value def pipeline_panel(template: dict, panel_id: int, title: str, description: str, x: int) -> dict: panel = copy.deepcopy(template) panel.update( id=panel_id, title=title, description=description, gridPos={"h": 8, "w": 8, "x": x, "y": 176}, ) panel["fieldConfig"]["defaults"]["unit"] = "short" panel["fieldConfig"]["defaults"]["min"] = 0 panel["fieldConfig"]["overrides"] = [] return panel def progress_panel(template: dict) -> dict: panel = copy.deepcopy(template) panel.update( id=81, type="timeseries", title="Backfill repositories over time", description=( "Absolute repository counts across the current backfill run. listRepos is fetched in pages " "of 1,000; Stream follows upstream Jetstream V2 by accumulating about 100,000 " "entries before shuffling and dispatching a batch. Durably committed advances only " "after that batch checkpoint. These lines are not a percentage of the full network." ), gridPos={"h": 8, "w": 24, "x": 0, "y": 171}, ) panel["fieldConfig"]["defaults"]["unit"] = "short" panel["fieldConfig"]["defaults"]["min"] = 0 panel["options"] = { "legend": { "calcs": ["lastNotNull"], "displayMode": "table", "placement": "bottom", "showLegend": True, }, "tooltip": {"mode": "multi", "sort": "desc"}, } discovered = panel["targets"][0] discovered.update( expr=( 'sum(jetstream_backfill_discovered_total' '{job=~"$job",instance=~"$instance"})' ), legendFormat="discovered", refId="A", range=True, instant=False, ) processed = copy.deepcopy(discovered) processed.update( expr=( 'sum(jetstream_backfill_completion_queued_total' '{job=~"$job",instance=~"$instance"})' ), legendFormat="processed", refId="B", ) # jetstream_backfill_progress_completed is process-local: it resets to zero # on restart and climbs again, so a series labelled "durably committed" was # showing this process's work rather than the archive's. That mislabel is # what made a restart read as progress in July. The durable series is read # from the metadata store at scrape time and continues across a restart. durable = copy.deepcopy(discovered) durable.update( expr=( 'sum(stream_backfill_repos_durable{status="complete",' 'job=~"$job",instance=~"$instance"})' ), legendFormat="durably committed", refId="C", range=True, instant=False, ) panel["targets"] = [discovered, processed, durable] return panel def adapt(raw: bytes) -> bytes: digest = hashlib.sha256(raw).hexdigest() if digest != SOURCE_SHA256: raise ValueError(f"upstream dashboard checksum mismatch: expected {SOURCE_SHA256}, got {digest}") dashboard = json.loads(raw) panels = {panel.get("id"): panel for panel in dashboard["panels"]} expected = { 3: "Uptime", 69: "CPU", 70: "Memory", 71: "Goroutines & threads", 72: "File descriptors", 73: "GC pauses", 74: "Process network I/O", 75: "Process & Go runtime", } for panel_id, title in expected.items(): if panels.get(panel_id, {}).get("title") != title: raise ValueError(f"upstream runtime panel {panel_id} is no longer {title!r}") selector = '{job=~"$job",instance=~"$instance"}' rate = "[$__rate_interval]" panels[75]["title"] = "Process runtime" # Standard process collector families retain their upstream meaning. panels[69]["targets"] = [copy.deepcopy(panels[69]["targets"][0])] memory_template = panels[70]["targets"][0] panels[70]["targets"] = [ target(memory_template, f"max(process_resident_memory_bytes{selector})", "RSS", "A"), target(memory_template, f"max(process_virtual_memory_bytes{selector})", "virtual", "B"), ] thread_template = panels[71]["targets"][0] panels[71]["title"] = "Threads" panels[71]["targets"] = [ target(thread_template, f"max(stream_process_threads{selector})", "OS threads", "A") ] fault_template = panels[73]["targets"][0] panels[73]["title"] = "Page faults" panels[73]["targets"] = [ target(fault_template, f"sum(rate(process_minor_page_faults_total{selector}{rate}))", "minor/s", "A"), target(fault_template, f"sum(rate(process_major_page_faults_total{selector}{rate}))", "major/s", "B"), ] panels[73]["fieldConfig"]["defaults"]["unit"] = "ops" switch_template = panels[74]["targets"][0] panels[74]["title"] = "Context switches" panels[74]["targets"] = [ target( switch_template, f"sum(rate(process_voluntary_context_switches_total{selector}{rate}))", "voluntary/s", "A", ), target( switch_template, f"sum(rate(process_involuntary_context_switches_total{selector}{rate}))", "involuntary/s", "B", ), ] panels[74]["fieldConfig"]["defaults"]["unit"] = "ops" # Show absolute repository counts over time. The upstream completion ratio # remains in its original panel, but it is approximate while listRepos is # still expanding the denominator. progress_row = copy.deepcopy(panels[75]) progress_row.update( id=80, title="Backfill progress over time", description="The final network denominator is unknown until listRepos discovery finishes.", gridPos={"h": 1, "w": 24, "x": 0, "y": 170}, ) progress_template = next(panel for panel in dashboard["panels"] if panel.get("type") == "timeseries") progress = progress_panel(progress_template) # Upstream has no scheduler between its websocket reader and archive # writer. Stream does, so retain every upstream protocol/lifecycle panel # byte-for-byte and append a clearly Stream-specific diagnostics row. diagnostics_row = copy.deepcopy(panels[75]) diagnostics_row.update( id=76, title="Stream pipeline diagnostics", description="Live during bootstrap and steady state; gaps and terminal failures are never inferred from archive throughput.", gridPos={"h": 1, "w": 24, "x": 0, "y": 175}, ) tickets = pipeline_panel( panels[69], 77, "Pipeline tickets", "Cumulative scheduler boundaries. Emitted is the historical name for every settled ticket, including explicit drops and verifier rejections.", 0, ) tickets["targets"] = [ target( panels[69]["targets"][0], f'max(stream_pipeline_ticket{selector}) by (stage)', "{{stage}}", "A", ) ] gaps = pipeline_panel( panels[69], 78, "Pipeline stage gaps", "Current unsettled work split between the unclaimed verifier queue and claimed work awaiting settlement. Their sum is submitted-to-emitted.", 8, ) gaps["targets"] = [ target( panels[69]["targets"][0], f'max(stream_pipeline_stage_gap{selector}) by (boundary)', "{{boundary}}", "A", ) ] failures = pipeline_panel( panels[69], 79, "Pipeline terminal failures", "Any non-zero value is deployment-significant. The gauge marks the current pipeline fatal; the counters retain its first failing stage.", 16, ) failures["targets"] = [ target( panels[69]["targets"][0], f'max(stream_pipeline_fatal{selector})', "fatal", "A", ), target( panels[69]["targets"][0], f'sum(stream_pipeline_failures_total{selector}) by (stage)', "{{stage}} failures", "B", ), ] failures["fieldConfig"]["defaults"]["thresholds"] = { "mode": "absolute", "steps": [ {"color": "green", "value": None}, {"color": "red", "value": 1}, ], } failures["fieldConfig"]["defaults"]["custom"]["thresholdsStyle"] = {"mode": "line"} dashboard["panels"].extend([progress_row, progress, diagnostics_row, tickets, gaps, failures]) dashboard["title"] = "Stream / Jetstream V2" dashboard["uid"] = "stream-jetstream-v2" dashboard["tags"] = ["stream", "jetstream", "atproto"] dashboard["description"] = ( "Stream operational dashboard adapted from the checksum-pinned Jetstream V2 dashboard. " "Protocol and lifecycle rows are unchanged; the runtime row uses Stream process metrics, " "and final Stream-only rows expose backfill work over time and scheduler health." ) return (json.dumps(dashboard, indent=2, ensure_ascii=False) + "\n").encode() def main() -> int: generated = adapt(SOURCE.read_bytes()) if len(sys.argv) == 2 and sys.argv[1] == "--check": if not OUTPUT.exists() or OUTPUT.read_bytes() != generated: print(f"{OUTPUT} is stale; run tools/adapt_dashboard.py", file=sys.stderr) return 1 print(f"stream dashboard sha256={hashlib.sha256(generated).hexdigest()}") return 0 if len(sys.argv) != 1: print("usage: tools/adapt_dashboard.py [--check]", file=sys.stderr) return 2 OUTPUT.write_bytes(generated) print(f"wrote {OUTPUT} sha256={hashlib.sha256(generated).hexdigest()}") return 0 if __name__ == "__main__": raise SystemExit(main())