CI and images / lint (push) Successful in 3s
CI and images / extension-version (push) Successful in 3s
CI and images / frontend-build (push) Successful in 27s
CI and images / backend-lint-and-test (push) Successful in 34s
CI and images / integration (push) Successful in 2m22s
CI and images / sign-extension (push) Successful in 4s
CI and images / build-web (push) Successful in 2m46s
CI and images / smoke-web (push) Successful in 50s
CI and images / build-agent (push) Successful in 6m35s
CI and images / promote (push) Successful in 2s
torch and onnxruntime both fall back to the CPU without raising, so the agent that ran CPU-bound for weeks after a driver update leased and checked in like a healthy one. - The agent sends its startup accel report on every lease and heartbeat. - The server keeps a bounded copy on the roster row. A running agent with a runtime off the GPU becomes `degraded`, with a sentence naming the runtime and the reason. - The top nav shows it amber. - The agent page carries a banner, and its pill reads "CPU only". Also: the bandwidth field gets the page's − / + stepper, and both number fields drop the browser's spin arrows. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LVjrnpQjRgHdvq95rASoiR
244 lines
9.8 KiB
Python
244 lines
9.8 KiB
Python
"""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,
|
|
},
|
|
})
|