#!/usr/bin/env python3 """Gathers real miner/tunnel/chain status plus the full per-task processing pipeline (validator request -> download -> encode -> upload -> response), appends a snapshot to history, and renders the public status dashboard. No fabricated data -- anything that can't be reliably read is shown as unknown rather than guessed. Run standalone (one refresh) or via scripts/tracker_loop.sh (repeated). """ from __future__ import annotations import html import json import re import subprocess import sys from datetime import datetime, timezone from pathlib import Path REPO = Path("/root/vidaio-subnet") HISTORY_PATH = Path("/tmp/vidaio-status-history.jsonl") RESET_MARKER_PATH = Path("/tmp/vidaio-dashboard-reset-at") WATCHDOG_EVENTS_PATH = Path("/tmp/vidaio-watchdog-events.jsonl") MINER_LOG_PATH = Path("/tmp/vidaio-miner-process.log") OUTPUT_HTML = Path("/tmp/vidaio-dashboard-public/index.html") HISTORY_MAX = 1000 BATCH_LIMIT = 6 TUNNEL_HOST = "159.223.110.159" TUNNEL_PORT = 39518 AXON_LOCAL_PORT = 8091 def run(cmd: str) -> str: try: return subprocess.run( cmd, shell=True, capture_output=True, text=True, timeout=15 ).stdout.strip() except Exception: return "" def run_raw(cmd: str, timeout: int = 20) -> str: try: return subprocess.run( cmd, shell=True, capture_output=True, text=True, timeout=timeout ).stdout except Exception: return "" def now_iso() -> str: return datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") def strip_ansi(text: str) -> str: return re.sub(r"\x1b\[[0-9;]*m", "", text) # --------------------------------------------------------------------------- # Task pipeline parsing # --------------------------------------------------------------------------- _RECEIVING_RE = re.compile( r"(\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2})\S* \| INFO\s*\| " r"__main__:forward_compression_requests:\d+ - .*Receiving CompressionRequest " r"from validator: (\S+) with uid: (\d+) \| queries=(\d+) \| VMAF: ([\d.]+) \| " r"Codec: (\w+) \| Mode: (\w+) \| Bitrate: ([\d.]+) Mbps" ) _RESPONSE_RE = re.compile( r"(\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2})\S* \| INFO\s*\| " r"__main__:forward_compression_requests:\d+ - .*Returning Response, " r"Processed in ([\d.]+) seconds" ) _DOWNLOAD_RE = re.compile( r"(\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2})\S* \| INFO\s*\| " r"__main__:_download_to_shared_volume:\d+ - Downloading validator payload " r"to shared volume: \S*/([0-9a-f]{8,16})_input\.mp4" ) _UPLOAD_RE = re.compile( r"(\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2})\S* \| INFO\s*\| " r"__main__:_upload_processed_video:\d+ - Uploading processed compression " r"output: processing/compression/([0-9a-f]{8,16})/" ) _QUEUED_RE = re.compile( r"(\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2}) \| INFO\s*\| \[([0-9a-f]{8,16})\] " r"Queued compression \(codec=(\S+), cq=(\d+)" ) _PROBE_RE = re.compile( r"(\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2}) \| INFO\s*\| \[([0-9a-f]{8,16})\] " r"cq search probe cq=(\d+) vmaf=([\d.]+) score_est=([\d.]+)" ) _SELECTED_RE = re.compile( r"(\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2}) \| INFO\s*\| \[([0-9a-f]{8,16})\] " r"adaptive cq search selected cq=(\d+)" ) _SINGLEPASS_RE = re.compile( r"\[([0-9a-f]{8,16})\] single-pass compression: ffmpeg .*-c:v (\S+) " r"(?:-cq (\d+)|-crf (\d+)|-b:v (\d+))" ) _COMPLETE_RE = re.compile( r"(\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2}) \| INFO\s*\| \[([0-9a-f]{8,16})\] " r"Compression complete \((\w+)\)" ) _TS_FMT = "%Y-%m-%d %H:%M:%S" def _parse_ts(s: str) -> datetime: return datetime.strptime(s, _TS_FMT).replace(tzinfo=timezone.utc) def reset_marker() -> datetime | None: """User-requested cutoff for the batch table/counters -- entries from before this point are real (the persistent miner log keeps them), but are excluded from display so the dashboard counts fresh from the reset point instead of replaying pre-reset history on every refresh.""" if not RESET_MARKER_PATH.exists(): return None try: return _parse_ts(RESET_MARKER_PATH.read_text().strip()) except (ValueError, OSError): return None def parse_task_pipeline(miner_log: str, compression_log: str) -> list[dict]: """Correlates the miner process log (validator request/download/upload/ response) with the compression container log (per-item adaptive search and encode) by task ID, grouped into recent batches.""" miner_log = strip_ansi(miner_log) receivings = [ { "ts": m.group(1), "validator_hotkey": m.group(2), "validator_uid": m.group(3), "n_queries": int(m.group(4)), "vmaf_target": float(m.group(5)), "codec": m.group(6), "mode": m.group(7), "bitrate_mbps": float(m.group(8)), } for m in _RECEIVING_RE.finditer(miner_log) ] marker = reset_marker() if marker is not None: receivings = [r for r in receivings if _parse_ts(r["ts"]) >= marker] responses = [ {"ts": m.group(1), "seconds": float(m.group(2))} for m in _RESPONSE_RE.finditer(miner_log) ] downloads = [(m.group(1), m.group(2)) for m in _DOWNLOAD_RE.finditer(miner_log)] uploads = [(m.group(1), m.group(2)) for m in _UPLOAD_RE.finditer(miner_log)] # per-task compression-container detail tasks: dict[str, dict] = {} def task(tid: str) -> dict: return tasks.setdefault( tid, {"probes": [], "queued_ts": None, "selected_cq": None, "final_codec": None, "final_param": None, "complete_ts": None} ) for m in _QUEUED_RE.finditer(compression_log): t = task(m.group(2)) t["queued_ts"] = m.group(1) t["queued_codec"] = m.group(3) for m in _PROBE_RE.finditer(compression_log): task(m.group(2))["probes"].append( {"ts": m.group(1), "cq": int(m.group(3)), "vmaf": float(m.group(4)), "score": float(m.group(5))} ) for m in _SELECTED_RE.finditer(compression_log): task(m.group(2))["selected_cq"] = int(m.group(3)) for m in _SINGLEPASS_RE.finditer(compression_log): t = task(m.group(1)) t["final_codec"] = m.group(2) cq, crf, brate = m.group(3), m.group(4), m.group(5) if brate: t["final_param"] = f"{int(brate)/1_000_000:.1f} Mbps (VBR)" elif crf: t["final_param"] = f"CRF {crf}" elif cq: t["final_param"] = f"CQ {cq}" for m in _COMPLETE_RE.finditer(compression_log): t = task(m.group(2)) t["complete_ts"] = m.group(1) t["complete_mode"] = m.group(3) # group into batches: each "Receiving" starts a batch, closed by the next # "Returning Response" that follows it chronologically batches = [] resp_idx = 0 for i, rec in enumerate(receivings): rec_dt = _parse_ts(rec["ts"]) resp = None while resp_idx < len(responses) and _parse_ts(responses[resp_idx]["ts"]) < rec_dt: resp_idx += 1 if resp_idx < len(responses): resp = responses[resp_idx] resp_idx += 1 window_end = _parse_ts(resp["ts"]) if resp else rec_dt batch_downloads = [ (ts, tid) for ts, tid in downloads if rec_dt <= _parse_ts(ts) <= window_end ] batch_uploads = [ (ts, tid) for ts, tid in uploads if rec_dt <= _parse_ts(ts) <= window_end ] task_ids = sorted({tid for _, tid in batch_downloads} | {tid for _, tid in batch_uploads}) items = [] for tid in task_ids: dl_ts = next((ts for ts, t in batch_downloads if t == tid), None) up_ts = next((ts for ts, t in batch_uploads if t == tid), None) detail = tasks.get(tid, {}) items.append({ "task_id": tid, "download_ts": dl_ts, "upload_ts": up_ts, "queued_ts": detail.get("queued_ts"), "complete_ts": detail.get("complete_ts"), "probes": detail.get("probes", []), "selected_cq": detail.get("selected_cq"), "final_codec": detail.get("final_codec"), "final_param": detail.get("final_param"), }) batches.append({ "received_ts": rec["ts"], "validator_hotkey": rec["validator_hotkey"], "validator_uid": rec["validator_uid"], "n_queries": rec["n_queries"], "vmaf_target": rec["vmaf_target"], "codec": rec["codec"], "mode": rec["mode"], "bitrate_mbps": rec["bitrate_mbps"], "response_ts": resp["ts"] if resp else None, "response_seconds": resp["seconds"] if resp else None, "complete": resp is not None, "items": items, }) batches.reverse() # newest first return batches[:BATCH_LIMIT] # --------------------------------------------------------------------------- # Snapshot / history # --------------------------------------------------------------------------- def gather_snapshot() -> dict: snap = {"ts": now_iso()} # --- compression backend --- health_raw = run("curl -sf http://localhost:8004/health") try: health = json.loads(health_raw) if health_raw else None except json.JSONDecodeError: health = None snap["compression_healthy"] = health is not None if health: snap["compression_active_tasks"] = health.get("active_tasks") snap["compression_queued_tasks"] = health.get("queued_tasks") snap["compression_max_concurrent"] = health.get("max_concurrent") storage = health.get("storage", {}) snap["storage_provider"] = storage.get("provider") snap["storage_configured"] = all( storage.get(k) for k in ( "bucket_configured", "access_key_configured", "secret_key_configured", "endpoint_configured", ) ) container_running = run( "docker inspect -f '{{.State.Running}}' miner-compression-1 2>/dev/null" ) snap["compression_container_running"] = container_running.strip() == "true" restarts = run("docker inspect -f '{{.RestartCount}}' miner-compression-1 2>/dev/null") snap["compression_restarts"] = int(restarts) if restarts.isdigit() else None stats = run( "docker stats miner-compression-1 --no-stream --format '{{.CPUPerc}}|{{.MemUsage}}' 2>/dev/null" ) if "|" in stats: cpu, mem = stats.split("|", 1) snap["compression_cpu"] = cpu.strip() snap["compression_mem"] = mem.strip() # --- GPU --- gpu = run( "nvidia-smi --query-gpu=utilization.gpu,memory.used,memory.total," "temperature.gpu,power.draw --format=csv,noheader,nounits" ) if gpu: parts = [p.strip() for p in gpu.split(",")] if len(parts) == 5: snap["gpu_util"] = parts[0] snap["gpu_mem_used"] = parts[1] snap["gpu_mem_total"] = parts[2] snap["gpu_temp"] = parts[3] snap["gpu_power"] = parts[4] # --- host --- disk = run("df -h / | tail -1") disk_parts = disk.split() if len(disk_parts) >= 5: snap["disk_used"] = disk_parts[2] snap["disk_total"] = disk_parts[1] snap["disk_pct"] = disk_parts[4] load = run("uptime") m = re.search(r"load average:\s*([\d.]+),\s*([\d.]+),\s*([\d.]+)", load) if m: snap["load1"], snap["load5"], snap["load15"] = m.group(1), m.group(2), m.group(3) # --- axon / miner process --- snap["axon_process_running"] = bool(run("pgrep -f 'neurons/miner.py'")) snap["axon_port_listening"] = bool(run(f"ss -tln 2>/dev/null | grep ':{AXON_LOCAL_PORT} '")) snap["bore_running"] = bool(run(f"pgrep -f 'bore local {AXON_LOCAL_PORT}'")) snap["tunnel_host"] = TUNNEL_HOST snap["tunnel_port"] = TUNNEL_PORT miner_log_raw = MINER_LOG_PATH.read_text(errors="replace") if MINER_LOG_PATH.exists() else "" text = strip_ansi(miner_log_raw) events = re.findall( r"(\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2})\S* \| \s*ERROR\s* \| bittensor:axon\.py:\d+ \| " r"UnknownSynapseError#[0-9a-f-]+: Synapse name '([^']*)' not found", text, ) snap["validator_query_count"] = len(events) if events: snap["last_validator_query_ts"] = events[-1][0] snap["last_validator_query_synapse"] = events[-1][1] or "(empty)" served = re.findall(r"Axon served with: AxonInfo\([^,]+, ([\d.]+:\d+)\)", text) if served: snap["last_served_address"] = served[-1] # --- real task pipeline --- compression_log_raw = run_raw("docker logs miner-compression-1 2>&1") batches = parse_task_pipeline(miner_log_raw, compression_log_raw) marker = reset_marker() receiving_matches = list(_RECEIVING_RE.finditer(text)) download_matches = list(_DOWNLOAD_RE.finditer(text)) if marker is not None: receiving_matches = [m for m in receiving_matches if _parse_ts(m.group(1)) >= marker] download_matches = [m for m in download_matches if _parse_ts(m.group(1)) >= marker] snap["batches_total"] = len(receiving_matches) snap["items_total"] = len(download_matches) if marker is not None: snap["reset_at"] = marker.strftime(_TS_FMT) completed = [b for b in batches if b["complete"]] if completed: snap["avg_batch_seconds"] = sum(b["response_seconds"] for b in completed) / len(completed) return snap, batches def load_history() -> list[dict]: if not HISTORY_PATH.exists(): return [] out = [] for line in HISTORY_PATH.read_text().splitlines(): line = line.strip() if not line: continue try: out.append(json.loads(line)) except json.JSONDecodeError: continue return out def append_history(snap: dict) -> list[dict]: history = load_history() history.append(snap) history = history[-HISTORY_MAX:] with HISTORY_PATH.open("w") as f: for row in history: f.write(json.dumps(row) + "\n") return history def load_watchdog_events(limit: int = 10) -> list[dict]: if not WATCHDOG_EVENTS_PATH.exists(): return [] out = [] for line in WATCHDOG_EVENTS_PATH.read_text().splitlines(): line = line.strip() if not line: continue try: out.append(json.loads(line)) except json.JSONDecodeError: continue return out[-limit:] def uptime_bar(history: list[dict], key: str, width: int = 60) -> str: recent = history[-width:] cells = [] for row in recent: ok = row.get(key) cells.append("█" if ok else ("·" if ok is False else " ")) return "".join(cells).rjust(width) def pct_true(history: list[dict], key: str, window: int = 200) -> float | None: recent = [row for row in history[-window:] if key in row] if not recent: return None return 100.0 * sum(1 for row in recent if row.get(key)) / len(recent) def esc(v) -> str: return html.escape(str(v)) if v is not None else "—" def clock(ts: str | None) -> str: return ts[11:19] if ts else "—" # --------------------------------------------------------------------------- # Rendering # --------------------------------------------------------------------------- STAGES = ["received", "downloaded", "encoded", "uploaded", "responded"] STAGE_LABELS = { "received": "Received", "downloaded": "Downloaded", "encoded": "Encoded", "uploaded": "Uploaded", "responded": "Responded", } def render_batch_card(batch: dict, idx: int) -> str: n = batch["n_queries"] items = batch["items"] all_downloaded = all(it["download_ts"] for it in items) and len(items) >= n all_uploaded = all(it["upload_ts"] for it in items) and len(items) >= n all_encoded = all(it["complete_ts"] for it in items) and len(items) >= n stage_done = { "received": True, "downloaded": all_downloaded, "encoded": all_encoded, "uploaded": all_uploaded, "responded": batch["complete"], } stage_ts = { "received": batch["received_ts"], "downloaded": max((it["download_ts"] for it in items if it["download_ts"]), default=None), "encoded": max((it["complete_ts"] for it in items if it["complete_ts"]), default=None), "uploaded": max((it["upload_ts"] for it in items if it["upload_ts"]), default=None), "responded": batch["response_ts"], } active_stage = None for s in STAGES: if not stage_done[s]: active_stage = s break stepper = "" for i, s in enumerate(STAGES): state = "done" if stage_done[s] else ("active" if s == active_stage else "pending") stepper += ( f'
VideoCompressionProtocol request arrives — showing every '
"stage from receipt through response, with per-item CQ search detail.