#!/usr/bin/env -S PYTHONUNBUFFERED=1 uv run --script --quiet # /// script # requires-python = ">=3.12" # dependencies = [] # /// """ indigo relay (relay.waow.tech) operator probe. queries Prometheus *through the Grafana datasource-proxy over the public ingress* (no port-forward to drop), reads pod state via kubectl, and for the deep subcommands SSHes to the node (the relay is hostNetwork, so pprof on :2471 is reachable on the node) + uses `go tool pprof`. subcommands: snapshot full-stack health + mem + rates + verdict (the "is it up right now" answer). ~5s, no SSH. heap heap decomposition: inuse / alloc / next_gc / idle / goroutines (live vs garbage at a glance). gc forced-GC live-vs-garbage test (pprof ?gc=1 on :2471): if heap_alloc barely drops, it's real. pprof grab a heap profile + `go tool pprof -top -inuse_space` (what structure holds the heap). trend [hours=24] working_set / heap / ingest range series (when did it change). usage: just indigo probe snapshot just indigo probe heap just indigo probe gc just indigo probe pprof just indigo probe trend 12 env: RELAY_DOMAIN (default: relay.waow.tech) GRAFANA_HOST (default: relay-metrics.waow.tech) KUBECONFIG (default: /indigo/kubeconfig.yaml; set by the indigo justfile automatically) RELAY_NAMESPACE (default: relay) note: WS firehose delivery is checked with `scripts/jetstream` (python websockets). do NOT use websocat — it is broken on macOS arm64. """ from __future__ import annotations import base64 import json import os import subprocess import sys import time import urllib.error import urllib.parse import urllib.request from pathlib import Path DOMAIN = os.environ.get("RELAY_DOMAIN", "relay.waow.tech") GRAFANA = os.environ.get("GRAFANA_HOST", "relay-metrics.waow.tech") NS = os.environ.get("RELAY_NAMESPACE", "relay") REPO = Path(__file__).resolve().parent.parent PPROF_PORT = 2471 # indigo relay pprof + metrics live here, NOT 2470 SEL = 'pod=~"relay-.*",container="main"' # cAdvisor container_* metrics GSEL = 'pod=~"relay-.*"' # go_memstats_*/app metrics (no container label) # ---------------- plumbing ---------------- def sh(*args: str, timeout: int = 20) -> str: return subprocess.run( args, capture_output=True, text=True, timeout=timeout ).stdout.strip() def grafana_pw() -> str: raw = sh( "kubectl", "get", "secret", "-n", "monitoring", "kube-prometheus-stack-grafana", "-o", "jsonpath={.data.admin-password}", ) return base64.b64decode(raw).decode() if raw else "" _PW: str | None = None def _auth() -> str: global _PW if _PW is None: _PW = grafana_pw() return "Basic " + base64.b64encode(f"admin:{_PW}".encode()).decode() def promql(query: str, path: str = "query", **params) -> dict: base = f"https://{GRAFANA}/api/datasources/proxy/uid/prometheus/api/v1/{path}" params["query"] = query url = base + "?" + urllib.parse.urlencode(params) req = urllib.request.Request(url, headers={"Authorization": _auth()}) with urllib.request.urlopen(req, timeout=15) as r: return json.load(r) def scalar(query: str, div: float = 1.0) -> float | None: res = promql(query).get("data", {}).get("result", []) return float(res[0]["value"][1]) / div if res else None def fmt(v: float | None, suffix: str = "", nd: int = 2) -> str: return f"{v:,.{nd}f}{suffix}" if v is not None else "na" def node_ip() -> str: """relay is hostNetwork → node IP = kubeconfig apiserver host.""" url = sh("kubectl", "config", "view", "--minify", "-o", "jsonpath={.clusters[0].cluster.server}") return urllib.parse.urlparse(url).hostname or "" def pod_line(name: str) -> str: out = sh("kubectl", "get", "pod", "-n", NS, "-l", f"app.kubernetes.io/name={name}", "--no-headers") if not out: return "MISSING" f = out.split() return f"{f[1]} {f[2]} r={f[3]}" def http_code(url: str) -> str: try: req = urllib.request.Request(url, method="GET") with urllib.request.urlopen(req, timeout=5) as r: return str(r.status) except urllib.error.HTTPError as e: # type: ignore[attr-defined] return str(e.code) except Exception: return "DOWN" GIB = 2 ** 30 # ---------------- subcommands ---------------- def cmd_snapshot() -> None: print(f"════ relay stack snapshot {time.strftime('%H:%M:%SZ', time.gmtime())} ════") ws = scalar(f"max(container_memory_working_set_bytes{{namespace=\"{NS}\",{SEL}}})", GIB) ingest = scalar("sum(rate(indigo_repo_stream_events_received_total[2m]))") bcast = scalar("sum(rate(indigo_events_broadcast_total[2m]))") backlog = scalar( "max(indigo_events_enqueued_for_broadcast_total)-max(indigo_events_broadcast_total)" ) print(f" relay {http_code(f'https://{DOMAIN}/xrpc/_health')} {pod_line('relay')}" f" mem={fmt(ws,' GiB',1)}/12 in={fmt(ingest,'/s',0)} out={fmt(bcast,'/s',0)}") print(f" jetstream {http_code(f'https://jetstream.{DOMAIN.split('.',1)[1]}/')} {pod_line('jetstream')}") print(f" lightrail - {pod_line('lightrail')}") # verdict flags = [] if ws and ws > 10: flags.append(f"mem {ws:.1f} GiB near 12 cap") if backlog and backlog > 1000: flags.append(f"broadcast backlog {backlog:,.0f}") print(" verdict: " + ("⚠ " + "; ".join(flags) if flags else "✓ healthy")) print(" firehose delivery → run: scripts/jetstream --url " f"wss://jetstream.{DOMAIN.split('.',1)[1]} --duration 3") def cmd_heap() -> None: print("── heap decomposition (live vs garbage vs GC target) ──") for label, q in [ ("heap_inuse", f"max(go_memstats_heap_inuse_bytes{{namespace=\"{NS}\",{GSEL}}})"), ("heap_alloc", f"max(go_memstats_heap_alloc_bytes{{namespace=\"{NS}\",{GSEL}}})"), ("next_gc", f"max(go_memstats_next_gc_bytes{{namespace=\"{NS}\",{GSEL}}})"), ("heap_idle", f"max(go_memstats_heap_idle_bytes{{namespace=\"{NS}\",{GSEL}}})"), ]: print(f" {label:12} {fmt(scalar(q, GIB), ' GiB')}") print(f" {'goroutines':12} {fmt(scalar(f'max(go_goroutines{{namespace=\"{NS}\",{GSEL}}})'), '', 0)}") print(" read: heap_alloc≈next_gc with idle≈0 → dense-live; next_gc pinned at " "GOMEMLIMIT → live > limit/2 (real growth, not pacing).") def cmd_gc() -> None: print("── forced-GC live-vs-garbage (pprof ?gc=1 on :2471) ──") ip = node_ip() script = ( f"b=$(curl -s --max-time 6 http://localhost:{PPROF_PORT}/metrics | " "awk '/^go_memstats_heap_alloc_bytes /{print $2}'); " f"curl -s --max-time 15 'http://localhost:{PPROF_PORT}/debug/pprof/heap?gc=1' -o /dev/null; " f"a=$(curl -s --max-time 6 http://localhost:{PPROF_PORT}/metrics | " "awk '/^go_memstats_heap_alloc_bytes /{print $2}'); echo \"$b $a\"" ) out = sh("ssh", "-o", "ConnectTimeout=8", f"root@{ip}", script, timeout=40) try: b, a = (float(x) for x in out.split()) print(f" before={b/GIB:.2f} GiB after_forced_GC={a/GIB:.2f} GiB freed={(b-a)/GIB:.2f} GiB") print(" read: freed≈0 → the heap is genuinely live (leak/retention), not GC-pacing garbage.") except ValueError: print(f" could not read memstats (relay up?): {out!r}") def cmd_pprof() -> None: print("── pprof heap: top inuse_space (what holds the heap) ──") ip = node_ip() sh("ssh", "-o", "ConnectTimeout=8", f"root@{ip}", f"curl -s --max-time 20 http://localhost:{PPROF_PORT}/debug/pprof/heap -o /tmp/relay-heap.pb.gz", timeout=35) sh("scp", "-o", "ConnectTimeout=8", f"root@{ip}:/tmp/relay-heap.pb.gz", "/tmp/relay-heap.pb.gz") out = sh("go", "tool", "pprof", "-top", "-inuse_space", "-nodecount=8", "/tmp/relay-heap.pb.gz", timeout=40) for line in out.splitlines(): if line.startswith(("File:", "Build", "Type:", "Time:", "Showing", "Dropped")): continue print(" " + line) Path("/tmp/relay-heap.pb.gz").unlink(missing_ok=True) def cmd_trend(hours: int = 24) -> None: print(f"── trend: working_set / ingest, last {hours}h ──") end = int(time.time()); start = end - hours * 3600; step = max(1800, hours * 3600 // 48) series = promql( f"max(container_memory_working_set_bytes{{namespace=\"{NS}\",{SEL}}})/2^30", path="query_range", start=start, end=end, step=step, ).get("data", {}).get("result", []) if not series: print(" no data"); return for ts, v in series[0]["values"]: t = time.strftime("%m-%d %H:%M", time.gmtime(float(ts))) print(f" {t}Z {float(v):5.1f} GiB") CMDS = {"snapshot": cmd_snapshot, "heap": cmd_heap, "gc": cmd_gc, "pprof": cmd_pprof, "trend": cmd_trend} def main() -> int: args = sys.argv[1:] cmd = args[0] if args else "snapshot" if cmd not in CMDS: print(f"unknown subcommand {cmd!r}; one of: {', '.join(CMDS)}", file=sys.stderr) return 2 if cmd == "trend" and len(args) > 1: cmd_trend(int(args[1])) else: CMDS[cmd]() return 0 if __name__ == "__main__": import urllib.error # noqa: E402 (used in http_code) sys.exit(main())