#!/usr/bin/env python3 """Observe queue drain, API responsiveness, and inference idle reaping live.""" from __future__ import annotations import argparse import json import time from pathlib import Path from typing import Any import httpx def call(client: httpx.Client, path: str) -> tuple[dict[str, Any], float]: started = time.perf_counter() response = client.get(path) elapsed_ms = round((time.perf_counter() - started) * 1000, 1) if response.is_error: raise RuntimeError(f"GET {path} -> {response.status_code}: {response.text[:1000]}") return response.json(), elapsed_ms def main() -> None: parser = argparse.ArgumentParser() parser.add_argument("--base-url", required=True) parser.add_argument("--password", required=True) parser.add_argument("--output", type=Path, required=True) parser.add_argument("--timeout", type=int, default=900) parser.add_argument("--idle-seconds", type=int, default=135) args = parser.parse_args() client = httpx.Client( base_url=args.base_url.rstrip("/"), timeout=httpx.Timeout(30, connect=10), ) login = client.post( "/api/v1/auth/login", json={"password": args.password, "remember_device": False}, ) login.raise_for_status() samples: list[dict[str, Any]] = [] idle_started: float | None = None deadline = time.monotonic() + args.timeout last_report = 0.0 while time.monotonic() < deadline: resources, resources_ms = call(client, "/api/v1/system/resources") diagnostics, diagnostics_ms = call(client, "/api/v1/system/diagnostics") status, status_ms = call(client, "/api/v1/status") lanes = resources.get("lanes") or {} active = sum( int(lane.get("running") or 0) + int(lane.get("queued") or 0) for lane in lanes.values() ) now = time.monotonic() if active: idle_started = None elif idle_started is None: idle_started = now idle_elapsed = 0 if idle_started is None else round(now - idle_started, 1) sample = { "at": time.time(), "version": status.get("version"), "active_jobs": active, "lanes": lanes, "cpu_percent": resources.get("cpu_percent"), "memory_available_bytes": resources.get("memory_available_bytes"), "process_rss_bytes": diagnostics.get("process_rss_bytes"), "event_loop_lag_ms": diagnostics.get("event_loop_lag_ms"), "event_loop_max_lag_ms": diagnostics.get("event_loop_max_lag_ms"), "request_p95_ms": diagnostics.get("request_p95_ms"), "inference": diagnostics.get("inference"), "database": diagnostics.get("database"), "latency_ms": { "resources": resources_ms, "diagnostics": diagnostics_ms, "status": status_ms, }, "idle_elapsed_seconds": idle_elapsed, } samples.append(sample) if now - last_report >= 10: print( json.dumps( { "active_jobs": active, "idle_seconds": idle_elapsed, "cpu": sample["cpu_percent"], "rss_mb": round(int(sample["process_rss_bytes"] or 0) / 1024**2, 1), "inference_running": bool((sample["inference"] or {}).get("running")), "max_request_ms": max(sample["latency_ms"].values()), }, ensure_ascii=False, ), flush=True, ) last_report = now if idle_elapsed >= args.idle_seconds and not (sample["inference"] or {}).get("running"): break time.sleep(2) else: raise TimeoutError("background queues did not drain and release inference before timeout") final = samples[-1] all_latencies = [value for sample in samples for value in sample["latency_ms"].values()] report = { "samples": len(samples), "duration_seconds": round(samples[-1]["at"] - samples[0]["at"], 1), "active_jobs_initial": samples[0]["active_jobs"], "active_jobs_final": final["active_jobs"], "cpu_max_percent": max(float(sample["cpu_percent"] or 0) for sample in samples), "memory_available_min_gb": round( min(int(sample["memory_available_bytes"] or 0) for sample in samples) / 1024**3, 2, ), "process_rss_initial_mb": round(int(samples[0]["process_rss_bytes"] or 0) / 1024**2, 1), "process_rss_final_mb": round(int(final["process_rss_bytes"] or 0) / 1024**2, 1), "event_loop_lag_final_ms": final["event_loop_lag_ms"], "event_loop_max_lag_ms": final["event_loop_max_lag_ms"], "request_latency_max_ms": max(all_latencies), "inference_final": final["inference"], "database_final": final["database"], } if final["active_jobs"] != 0 or (final["inference"] or {}).get("running"): raise AssertionError(f"resources were not released: {report}") if max(all_latencies) > 5000: raise AssertionError(f"API latency exceeded 5 seconds during soak: {report}") args.output.parent.mkdir(parents=True, exist_ok=True) args.output.write_text(json.dumps(report, ensure_ascii=False, indent=2) + "\n", encoding="utf-8") print(json.dumps(report, ensure_ascii=False, indent=2)) if __name__ == "__main__": main()