"""Is every part of FabledCurator running? One verdict, one endpoint. Milestone 365. The nav indicator and the System page both read this and nothing else — composing a verdict is this module's job, not the UI's. ## Two kinds of part, answered two different ways **Learned** — celery roles and the GPU agent, from `service_seen`. The question is "how long since it checked in", and these are the parts that can be ABSENT, which is the whole point: `celery inspect` alone reports presence, so a dead worker is a shorter list rather than a red light. **Probed live** — Postgres and Redis. Always expected, never learned, and a last-seen for them would be actively misleading: that Redis answered thirty seconds ago says nothing about now. ## This endpoint must never fail because something it checks has failed The inversion is easy to write by accident and it destroys the feature exactly when it is needed — a 500 when Redis is down, instead of `redis: down`. Every probe is wrapped, every wait has a deadline (rule 156), and the roster refresh swallows its own errors. The worst case is a part reported `unknown`, which is a true statement. """ from __future__ import annotations import asyncio import time from datetime import UTC, datetime from quart import Blueprint, jsonify from sqlalchemy import select, text from ..config import get_config from ..extensions import get_session from ..models import ServiceSeen from ..services.worker_lanes import SWEEP_PERIOD_SECONDS system_health_bp = Blueprint("system_health", __name__, url_prefix="/api/system") # How long a learned part may go quiet before it is doubted, then disbelieved. # # These are deliberately generous, and the reason is a deploy rather than a # worker: `docker compose up -d` rolls start-first, so a role is briefly served # by two containers and then by neither while the old one drains. Thresholds # tight enough to catch a crash in seconds would paint the page red every time # the stack is updated, and an alarm that cries wolf on every deploy is one # nobody reads. Tune down only after watching a real deploy pass through. STALE_AFTER_SECONDS = 90 DOWN_AFTER_SECONDS = 300 # The celery roster is written by `size_worker_lanes` and by nothing else, so # these thresholds are only meaningful against ITS cadence. Asserted at import # rather than left to a reader, because this is precisely the comparison that # was never made for the GPU agent: its lease poll backed off to 900s while # the roster called it stopped at 300s, and both numbers were individually # correct, in different directions, in different files (lesson #4355). # # Two clear sweeps before a part is even called STALE. One missed tick is # routine — the sweep rides the maintenance queue and does an inspect that can # take eleven seconds — and must not turn the page yellow. _SWEEPS_BEFORE_STALE = 2 assert STALE_AFTER_SECONDS >= SWEEP_PERIOD_SECONDS * _SWEEPS_BEFORE_STALE, ( f"a {SWEEP_PERIOD_SECONDS}s sweep cannot keep a roster fresh against a " f"{STALE_AFTER_SECONDS}s stale threshold: raise the threshold or shorten " f"the sweep" ) # Probes cross a process boundary, so they carry deadlines. A hung Postgres # must make this endpoint say "postgres: down", not hang alongside it. PROBE_TIMEOUT_SECONDS = 2.0 _OK, _STALE, _DOWN, _UNKNOWN = "ok", "stale", "down", "unknown" # Checking in, but working at a fraction of its speed: a GPU agent whose # runtimes fell back to the CPU (#4410). Below stale — a part that may have # stopped is the more urgent question — and above unknown, because this one # IS known to be wrong. _DEGRADED = "degraded" # Worst-first, so an overall verdict is just the max. _SEVERITY = {_OK: 0, _UNKNOWN: 1, _DEGRADED: 2, _STALE: 3, _DOWN: 4} def _age_state(age_seconds: float) -> str: if age_seconds >= DOWN_AFTER_SECONDS: return _DOWN if age_seconds >= STALE_AFTER_SECONDS: return _STALE return _OK def _describe_learned(name: str, state: str, age: float, details: dict) -> str: """Say what the state MEANS. A red chip tells an operator less than a sentence does at the moment they are deciding whether to go and look.""" if state == _OK: replicas = details.get("replicas") if replicas and replicas > 1: return f"{name} is running ({replicas} replicas)" return f"{name} is running" mins = int(age // 60) ago = f"{mins} min" if mins else f"{int(age)}s" if state == _STALE: return f"{name} has not checked in for {ago}" return f"{name} has not checked in for {ago} — treat it as stopped" def _cpu_runtimes(details: dict) -> list[str]: """The runtimes an agent reported as NOT on the GPU, with why. Both torch and onnxruntime fall back to the CPU without raising, so an agent in that state leases, works and checks in exactly like a healthy one. On 2026-09-24 one had been doing so since a driver update left a stale CDI spec; the only sign was a line in the agent's own log. """ accel = details.get("accel") if not isinstance(accel, dict): return [] out = [] for name, entry in sorted(accel.items()): if not isinstance(entry, dict) or entry.get("device") == "cuda": continue why = entry.get("error") or entry.get("device") or "unknown" out.append(f"{name} ({why})") return out def _learned_state(name: str, state: str, age: float, details: dict) -> tuple[str, str]: """A roster row's state and its sentence, degraded included.""" if state == _OK: cpu = _cpu_runtimes(details) if cpu: return _DEGRADED, ( f"{name} is running on the CPU — not on the GPU: {'; '.join(cpu)}. " "After a driver update, regenerate the agent host's CDI spec " "(agent README)." ) return state, _describe_learned(name, state, age, details) async def _probe_postgres(session) -> dict: started = time.monotonic() try: await asyncio.wait_for( session.execute(text("SELECT 1")), timeout=PROBE_TIMEOUT_SECONDS ) except Exception as exc: # noqa: BLE001 — a probe reports, it never raises return { "key": "postgres", "kind": "datastore", "name": "PostgreSQL", "state": _DOWN, "detail": f"not answering: {type(exc).__name__}", } return { "key": "postgres", "kind": "datastore", "name": "PostgreSQL", "state": _OK, "detail": "answering", "latency_ms": round((time.monotonic() - started) * 1000, 1), } def _ping_redis_sync() -> None: import redis # local import; mirrors system_activity's pattern client = redis.Redis.from_url( get_config().celery_broker_url, socket_connect_timeout=PROBE_TIMEOUT_SECONDS, socket_timeout=PROBE_TIMEOUT_SECONDS, ) client.ping() async def _probe_redis() -> dict: started = time.monotonic() try: await asyncio.wait_for( asyncio.to_thread(_ping_redis_sync), timeout=PROBE_TIMEOUT_SECONDS * 2 ) except Exception as exc: # noqa: BLE001 return { "key": "redis", "kind": "datastore", "name": "Redis", "state": _DOWN, "detail": f"not answering: {type(exc).__name__} — queues and workers " f"cannot be reached either", } return { "key": "redis", "kind": "datastore", "name": "Redis", "state": _OK, "detail": "answering", "latency_ms": round((time.monotonic() - started) * 1000, 1), } @system_health_bp.route("/health", methods=["GET"]) async def system_health(): """Every part, its state, and one overall verdict. Response: {overall, parts: [{key, kind, name, state, detail, last_seen_at, …}], checked_at} """ parts: list[dict] = [] now = datetime.now(UTC) async with get_session() as session: # Postgres first, and if it is unreachable nothing else can be read — # say so rather than failing, because "the database is down" is the # single most useful thing this endpoint can ever report. pg = await _probe_postgres(session) parts.append(pg) if pg["state"] == _OK: # A PURE READ since 2026-09-23. This used to refresh the celery # roster here, rate-limited to once per 20s — so the roster only # advanced while somebody had a browser open, and a broadcast rode # on a request. `size_worker_lanes` writes it now, on a timer, and # the assertion below is what keeps that cadence honest. rows = ( await session.execute(select(ServiceSeen).order_by(ServiceSeen.display_name)) ).scalars().all() for row in rows: age = (now - row.last_seen_at).total_seconds() state, detail = _learned_state( row.display_name, _age_state(age), age, row.details or {}, ) parts.append({ "key": row.key, "kind": row.kind, "name": row.display_name, "state": state, "detail": detail, "last_seen_at": row.last_seen_at.isoformat(), "first_seen_at": row.first_seen_at.isoformat(), **{k: v for k, v in (row.details or {}).items() if k != "agent_id"}, }) parts.append(await _probe_redis()) overall = max((p["state"] for p in parts), key=lambda s: _SEVERITY[s], default=_UNKNOWN) return jsonify({ "overall": overall, "parts": sorted(parts, key=lambda p: (-_SEVERITY[p["state"]], p["name"])), "checked_at": now.isoformat(), # So the UI can explain a `stale` without hard-coding the same numbers # in a second place. "thresholds": { "stale_after_seconds": STALE_AFTER_SECONDS, "down_after_seconds": DOWN_AFTER_SECONDS, }, })