jetstream v2 in zig stream.waow.tech
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104#!/usr/bin/env python3"""Write /srv/dash/data.json from live experiment state. Run every 60s."""import json, subprocess, time, re, os, urllib.parseOUT = "/srv/dash/data.json"PROM = "http://127.0.0.1:9090/api/v1/query"START = 1785214907 # bootstrap start, 2026-07-28T05:01:47ZDEADLINE_FILE_VALUE = 1786035707 # mirrored; watchdog host is authoritativeCAPTURE_STOPPED = 1785740698 # live capture torn down at "phase -> merging"RELAY_RETENTION_H = 72 # bsky.network firehose retention (operator-confirmed)def sh(cmd): try: return subprocess.run(cmd, shell=True, capture_output=True, text=True, timeout=25).stdout.strip() except Exception: return ""def prom(q): u = PROM + "?query=" + urllib.parse.quote(q) raw = sh(f"cd /opt/stream-experiment && docker compose exec -T prometheus wget -qO- '{u}'") try: r = json.loads(raw)["data"]["result"] return r except Exception: return []def scalar(q): r = prom(q) try: return float(r[0]["value"][1]) except Exception: return Nonenow = int(time.time())# --- lifecycle phase -------------------------------------------------------phase_n = scalar("jetstream_orchestrator_phase")PHASES = {1: "bootstrap", 2: "merging", 3: "steady_state"}phase = PHASES.get(int(phase_n)) if phase_n is not None else None# --- repo buckets ----------------------------------------------------------buckets = {}for x in prom("stream_backfill_repos_durable"): buckets[x["metric"].get("status", "?")] = int(float(x["value"][1]))# --- merge counters --------------------------------------------------------merge = { "kept": scalar("jetstream_orchestrator_merge_events_kept_total"), "dropped": scalar("jetstream_orchestrator_merge_events_dropped_total"), "segments": scalar("jetstream_orchestrator_merge_segments_consumed_total"), "revs_updated": scalar("jetstream_orchestrator_merge_repo_revs_updated_total"),}# --- host facts ------------------------------------------------------------df = sh("df -B1 --output=size,used,avail /data | tail -1").split()disk = {"size": int(df[0]), "used": int(df[1]), "avail": int(df[2])} if len(df) == 3 else Nonerss = scalar("process_resident_memory_bytes")restarts = sh('docker inspect -f "{{.RestartCount}}" stream-experiment-stream-1')segs = sh("find /data/stream/segments -name '*.jss' | wc -l")tmps = sh("find /data/stream -name '*.tmp' | wc -l")backfill_tree = sh("du -sb /data/stream/backfill 2>/dev/null | cut -f1")# --- serving gate ----------------------------------------------------------code = sh("cd /opt/stream-experiment && docker compose exec -T prometheus " "wget -S -qO- http://172.18.0.2:8080/xrpc/jetstream.listSegments 2>&1 " "| grep -m1 -oE 'HTTP/1.1 [0-9]+' | grep -oE '[0-9]+$'")serving = int(code) if code.isdigit() else None# --- cost, from Hetzner list prices (gross EUR) ----------------------------VOL_GB_H = 0.0767 / 730GROW = 1785438000 # 2026-07-30T19:00Z, 2.5TB -> 3.0TBcost = ( (now - 1785214981) / 3600 * 0.528 # ccx43 workload + (now - 1785214980) / 3600 * 0.0288 # cpx12 observer + (now - 1785222411) / 3600 * 0.0288 # cpx12 watchdog + (GROW - 1785214979) / 3600 * 2500 * VOL_GB_H + (now - GROW) / 3600 * 3000 * VOL_GB_H)rate_h = 0.528 + 2 * 0.0288 + 3000 * VOL_GB_Hjson.dump({ "generated": now, "start": START, "deadline": DEADLINE_FILE_VALUE, "elapsed_h": (now - START) / 3600, "phase": phase, "phase_n": phase_n, "buckets": buckets, "merge": merge, "disk": disk, "rss": rss, "restarts": int(restarts) if restarts.isdigit() else None, "segments": int(segs) if segs.isdigit() else None, "stale_tmp": int(tmps) if tmps.isdigit() else None, "backfill_tree": int(backfill_tree) if backfill_tree.isdigit() else None, "serving_http": serving, "capture_stopped": CAPTURE_STOPPED, "cursor_expiry": CAPTURE_STOPPED + RELAY_RETENTION_H * 3600, "cost_eur": cost, "rate_eur_h": rate_h,}, open(OUT + ".tmp", "w"))os.replace(OUT + ".tmp", OUT)