"""Fleet aggregate for the dashboard cockpit. One read-only call rolls up every host (local + agents) into a "needs attention" list, headline KPIs and a per-host resource view. It is cheap by construction: a single container *summary* list (no per-container inspect) drives the local figures, image freshness comes from the cache the update-service background loop already maintains, and per agent it makes a small bounded set of HTTP fetches that degrade gracefully on failure. """ from __future__ import annotations import asyncio import time from datetime import datetime, timedelta, timezone from sqlmodel import Session, select from database import engine from docker_client import DockerError, get_client, safe_call from models.agent import Agent from models.backup_schedule import BackupSchedule from services import agent_service, compose_service, update_service COMPOSE_LABEL = compose_service.COMPOSE_LABEL # Fleet aggregate: cached briefly because each call may fan out to every agent. FLEET_TTL = 25.0 DISK_PRESSURE = 0.85 # disk used fraction above which a host needs attention MEM_PRESSURE = 0.90 # memory used fraction above which a host needs attention _fleet_cache: dict = {"data": None, "ts": 0.0} # --------------------------------------------------------------------------- # # Container summary (single Docker round-trip) + shared classifiers # --------------------------------------------------------------------------- # def _list_containers() -> list[dict]: client = get_client() return safe_call(client.api.containers, all=True) def _group_by_project(raw: list[dict]) -> dict[str, list[dict]]: by_project: dict[str, list[dict]] = {} for c in raw: project = (c.get("Labels") or {}).get(COMPOSE_LABEL) if project: by_project.setdefault(project, []).append(c) return by_project def _is_healthy(containers: list[dict]) -> bool: """All containers that *have* a healthcheck report healthy. The summary ``Status`` string carries the health suffix — "(healthy)", "(unhealthy)" or "(health: starting)" — only for containers with a healthcheck configured, so its absence simply means "no healthcheck". """ for c in containers: status = c.get("Status", "") or "" if "(unhealthy)" in status or "(health:" in status: return False return True def _is_updated(containers: list[dict], cache: dict[str, dict]) -> bool: for c in containers: st = cache.get(c.get("Image", "")) if st and st.get("update_available"): return False return True # --------------------------------------------------------------------------- # # Fleet aggregate (one call → "needs attention" + KPIs across every host) # # The dashboard used to poll each agent individually and recombine the numbers # client-side. ``compute_fleet`` does the fan-out server-side instead: one local # Docker pass plus, per online agent, a small set of system/stacks/updates # fetches — all wrapped so a slow or broken agent degrades to "offline" rather # than stalling the whole view. Cached for FLEET_TTL because the fan-out is not # free. # --------------------------------------------------------------------------- # # Stack-status buckets the status bar / KPIs are built from. Anything reporting # "error" or "dead" containers counts as a problem stack. _PROBLEM_STATUSES = {"error", "dead"} def _bucket_statuses(statuses: list[str]) -> dict[str, int]: return { "running": sum(1 for s in statuses if s == "running"), "partial": sum(1 for s in statuses if s == "partial"), "stopped": sum(1 for s in statuses if s in ("stopped", "exited")), "error": sum(1 for s in statuses if s in _PROBLEM_STATUSES), "total": len(statuses), } def _attn(severity: str, kind: str, host: str, title: str, detail: str, link: str) -> dict: return { "severity": severity, "kind": kind, "host": host, "title": title, "detail": detail, "link": link, } def _resource_attention(name: str, link: str, mem_used: int, mem_total: int, disk_used: int, disk_total: int) -> list[dict]: items: list[dict] = [] if mem_total and mem_used / mem_total >= MEM_PRESSURE: pct = round(mem_used / mem_total * 100) items.append(_attn("warn", "mem_pressure", name, f"{name}: memory at {pct}%", "Free memory or move stacks.", link)) if disk_total and disk_used / disk_total >= DISK_PRESSURE: pct = round(disk_used / disk_total * 100) items.append(_attn("warn", "disk_pressure", name, f"{name}: disk at {pct}%", "Prune images/volumes or add capacity.", link)) return items def _local_host() -> tuple[dict, list[dict]]: """Local host card + attention items from a single Docker container pass.""" discovered = compose_service.discover_stacks() try: raw = _list_containers() except DockerError: raw = [] by_project = _group_by_project(raw) update_cache = update_service.get_cache() statuses: list[str] = [] unhealthy: list[str] = [] updates = 0 for stack_id in discovered: containers = by_project.get(stack_id, []) status = compose_service._status_from_states([c.get("State", "") for c in containers]) statuses.append(status) if status == "running" and not _is_healthy(containers): unhealthy.append(stack_id) if not _is_updated(containers, update_cache): updates += 1 buckets = _bucket_statuses(statuses) labelled = [c for c in raw if (c.get("Labels") or {}).get(COMPOSE_LABEL)] # Local resource figures (lazy import keeps dashboard_service free of a # router dependency at module load time). from routers.system import _cpu_count, _disk_usage, _mem_info mem = _mem_info() disk = _disk_usage() host = { "id": "local", "name": "local", "online": True, "status": "online", "cpu_cores": _cpu_count(), "mem_used": mem["used"], "mem_total": mem["total"], "disk_used": disk["used"], "disk_total": disk["total"], "stacks": buckets, "containers_running": sum(1 for c in labelled if c.get("State") == "running"), "containers_total": len(labelled), "unhealthy": len(unhealthy), "updates_available": updates, } attention: list[dict] = [] for sid in unhealthy: attention.append(_attn("error", "unhealthy", "local", f"{sid} is unhealthy", "A container is failing its healthcheck.", f"/stacks/{sid}")) if buckets["error"]: attention.append(_attn("error", "stack_error", "local", f"{buckets['error']} stack(s) in error", "Containers are dead.", "/stacks")) if buckets["partial"]: attention.append(_attn("warn", "stack_partial", "local", f"{buckets['partial']} stack(s) partially running", "Some services are down.", "/stacks")) if updates: attention.append(_attn("warn", "updates", "local", f"{updates} stack(s) have image updates", "Pull the newer images.", "/images")) attention += _resource_attention("local", "/", mem["used"], mem["total"], disk["used"], disk["total"]) return host, attention def _offline_host(agent_id: int, name: str, status: str) -> dict: return {"id": agent_id, "name": name, "online": False, "status": status, "cpu_cores": 0, "mem_used": 0, "mem_total": 0, "disk_used": 0, "disk_total": 0, "stacks": _bucket_statuses([]), "containers_running": 0, "containers_total": 0, "unhealthy": 0, "updates_available": 0} async def _agent_host(agent_id: int) -> tuple[dict, list[dict]]: """One agent's host card + attention items, tolerant of partial failure. Opens its own Session so the fan-out across agents stays concurrency-safe (a shared SQLModel session is not), and fetches sequentially within the agent because :func:`agent_service.call` commits a status update each time. """ with Session(engine) as session: agent = session.get(Agent, agent_id) if agent is None: return _offline_host(agent_id, str(agent_id), "unknown"), [] name = agent.name # Agent stacks live in the host section of the dashboard, not a route # of their own, so aggregate agent items deep-link back to it. link = "/" if agent.status != "online": return _offline_host(agent_id, name, agent.status), [ _attn("error", "agent_offline", name, f"{name} is {agent.status}", "Check it under Settings → Remote hosts.", "/settings")] async def _fetch(path: str): try: return await agent_service.call(session, agent, "GET", path) except Exception: # AgentError or transport — degrade gracefully return None sys_data = await _fetch("/agent/system") stacks = await _fetch("/agent/stacks") updates = await _fetch("/agent/stacks/updates") if sys_data is None and stacks is None: # Couldn't reach it at all — agent_service.call already marked it offline. return _offline_host(agent_id, name, "offline"), [ _attn("error", "agent_offline", name, f"{name} is unreachable", "Check it under Settings → Remote hosts.", "/settings")] sys_data = sys_data or {} statuses = [s.get("status", "") for s in (stacks or [])] buckets = _bucket_statuses(statuses) update_count = sum(1 for v in (updates or {}).values() if isinstance(v, dict) and v.get("update_available")) host = { "id": agent_id, "name": name, "online": True, "status": "online", "cpu_cores": sys_data.get("cpu_cores", 0), "mem_used": sys_data.get("mem_used", 0), "mem_total": sys_data.get("mem_total", 0), "disk_used": sys_data.get("disk_used", 0), "disk_total": sys_data.get("disk_total", 0), "stacks": buckets, "containers_running": sys_data.get("compose_running", sys_data.get("containers_running", 0)), "containers_total": sys_data.get("containers_total", 0), "unhealthy": buckets["error"], # agents expose no healthcheck rollup; error is the proxy "updates_available": update_count, } attention: list[dict] = [] if buckets["error"]: attention.append(_attn("error", "stack_error", name, f"{name}: {buckets['error']} stack(s) in error", "Containers are dead.", link)) if buckets["partial"]: attention.append(_attn("warn", "stack_partial", name, f"{name}: {buckets['partial']} stack(s) partially running", "Some services are down.", link)) if update_count: attention.append(_attn("warn", "updates", name, f"{name}: {update_count} stack(s) have image updates", "Pull the newer images.", link)) attention += _resource_attention(name, link, host["mem_used"], host["mem_total"], host["disk_used"], host["disk_total"]) return host, attention def _backup_attention(session: Session, agent_names: dict[int, str]) -> list[dict]: """Flag enabled backup schedules whose last run failed or is overdue.""" now = datetime.now(timezone.utc) items: list[dict] = [] for sch in session.exec(select(BackupSchedule).where(BackupSchedule.enabled == True)).all(): # noqa: E712 host = agent_names.get(sch.agent_id, "local") if sch.agent_id else "local" status = (sch.last_status or "").lower() if status and not status.startswith("ok"): items.append(_attn("error", "backup_failed", host, f"Backup of {sch.stack_id} failed", sch.last_status or "", "/settings")) elif sch.next_run and _aware(sch.next_run) < now - timedelta(hours=1): items.append(_attn("warn", "backup_overdue", host, f"Backup of {sch.stack_id} is overdue", "Scheduled run did not happen.", "/settings")) return items def _aware(dt: datetime) -> datetime: """Treat naive DB timestamps as UTC (they're stored that way).""" return dt if dt.tzinfo else dt.replace(tzinfo=timezone.utc) _SEVERITY_ORDER = {"error": 0, "warn": 1} async def compute_fleet(session: Session, refresh: bool = False) -> dict: now = time.time() if not refresh and _fleet_cache["data"] and now - _fleet_cache["ts"] < FLEET_TTL: return _fleet_cache["data"] local_host, attention = await asyncio.to_thread(_local_host) agents = session.exec(select(Agent)).all() agent_names = {a.id: a.name for a in agents} agent_ids = [a.id for a in agents] agent_results = await asyncio.gather(*[_agent_host(aid) for aid in agent_ids]) hosts = [local_host] for host, items in agent_results: hosts.append(host) attention += items attention += _backup_attention(session, agent_names) attention.sort(key=lambda a: _SEVERITY_ORDER.get(a["severity"], 9)) kpis = { "hosts_online": sum(1 for h in hosts if h["online"]), "hosts_total": len(hosts), "stacks_running": sum(h["stacks"]["running"] for h in hosts), "stacks_partial": sum(h["stacks"]["partial"] for h in hosts), "stacks_total": sum(h["stacks"]["total"] for h in hosts), "containers_running": sum(h["containers_running"] for h in hosts), "containers_total": sum(h["containers_total"] for h in hosts), "unhealthy": sum(h["unhealthy"] for h in hosts), "updates_available": sum(h["updates_available"] for h in hosts), "backups_failing": sum(1 for a in attention if a["kind"] in ("backup_failed", "backup_overdue")), } status_totals = { "running": kpis["stacks_running"], "partial": kpis["stacks_partial"], "stopped": sum(h["stacks"]["stopped"] for h in hosts), "error": sum(h["stacks"]["error"] for h in hosts), } data = { "as_of": datetime.now(timezone.utc).isoformat(), "hosts": hosts, "kpis": kpis, "status_totals": status_totals, "attention": attention, } _fleet_cache["data"] = data _fleet_cache["ts"] = now return data