"""Aggregated dashboard data: stack-health funnel + summary widgets. Everything here is read-only and cheap by construction: one container *summary* list (no per-container inspect) feeds the whole funnel, image freshness comes from the cache the update-service background loop already maintains, and the daily uptime sample is appended lazily on read. """ from __future__ import annotations import asyncio import json import os import time from datetime import datetime, timedelta, timezone from typing import Optional from sqlmodel import Session, select from config import settings from docker_client import DockerError, get_client, safe_call from models.audit import AuditLog from models.auto_update import AutoUpdate from services import compose_service, update_service COMPOSE_LABEL = compose_service.COMPOSE_LABEL DOCKER_TIMEOUT = 5.0 # seconds — a slow daemon must not stall the dashboard FUNNEL_TTL = 30.0 UPTIME_FILE = os.path.join(settings.DATA_DIR, "uptime.jsonl") UPTIME_DAYS = 30 UPTIME_SAMPLE_INTERVAL = 300.0 # seconds between uptime samples _funnel_cache: dict = {"data": None, "ts": 0.0} # --------------------------------------------------------------------------- # # Container summary (single Docker round-trip) # --------------------------------------------------------------------------- # def _list_containers() -> list[dict]: client = get_client() return safe_call(client.api.containers, all=True) async def _containers_with_timeout() -> list[dict]: return await asyncio.wait_for(asyncio.to_thread(_list_containers), timeout=DOCKER_TIMEOUT) 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 def _auto_managed_ids(session: Session) -> set[str]: """Local stacks with an enabled auto-update policy.""" rows = session.exec( select(AutoUpdate.stack_id).where( AutoUpdate.enabled == True, # noqa: E712 — SQL expression, not identity AutoUpdate.agent_id == None, # noqa: E711 ) ).all() return set(rows) async def compute_funnel(session: Session, refresh: bool = False) -> dict: now = time.time() if not refresh and _funnel_cache["data"] and now - _funnel_cache["ts"] < FUNNEL_TTL: return _funnel_cache["data"] discovered_ids = compose_service.discover_stacks() try: raw = await _containers_with_timeout() except (DockerError, asyncio.TimeoutError): raw = [] by_project = _group_by_project(raw) update_cache = update_service.get_cache() auto_managed = _auto_managed_ids(session) running = healthy = updated = monitored = 0 for stack_id in discovered_ids: containers = by_project.get(stack_id, []) states = [c.get("State", "") for c in containers] if not states or any(s != "running" for s in states): continue running += 1 if not _is_healthy(containers): continue healthy += 1 if not _is_updated(containers, update_cache): continue updated += 1 # Final stage: the stack also keeps itself fresh (auto-update enabled). if stack_id in auto_managed: monitored += 1 data = { "discovered": len(discovered_ids), "running": running, "healthy": healthy, "updated": updated, "monitored": monitored, "as_of": datetime.now(timezone.utc).isoformat(), } _funnel_cache["data"] = data _funnel_cache["ts"] = now return data # --------------------------------------------------------------------------- # # Uptime series (sampled every few minutes, aggregated to a daily mean) # # Each sample is one JSONL line {"ts": iso, "value": pct}; a background loop # samples every UPTIME_SAMPLE_INTERVAL and summary reads top up opportunistically, # so the daily value approximates real uptime instead of a once-a-day snapshot. # Legacy pre-0.31.1 lines {"date": d, "value": pct} still count as one sample. # --------------------------------------------------------------------------- # def _read_uptime() -> list[dict]: if not os.path.isfile(UPTIME_FILE): return [] entries = [] with open(UPTIME_FILE, "r", encoding="utf-8") as fh: for line in fh: line = line.strip() if not line: continue try: entries.append(json.loads(line)) except json.JSONDecodeError: continue return entries def _append_uptime(entry: dict) -> None: os.makedirs(settings.DATA_DIR, exist_ok=True) with open(UPTIME_FILE, "a", encoding="utf-8") as fh: fh.write(json.dumps(entry) + "\n") def _entry_date(e: dict) -> Optional[str]: if "date" in e: return e["date"] ts = e.get("ts") return ts[:10] if isinstance(ts, str) and len(ts) >= 10 else None def _latest_sample_ts(entries: list[dict]) -> float: latest = 0.0 for e in entries: ts = e.get("ts") if not isinstance(ts, str): continue try: latest = max(latest, datetime.fromisoformat(ts).timestamp()) except ValueError: continue return latest def _sample_uptime(raw: list[dict]) -> Optional[dict]: """Record a sample if the last one is older than the sample interval. Uptime% = share of compose containers currently running. With no compose containers at all there is nothing to be up — skip rather than fake 100%. """ entries = _read_uptime() now = datetime.now(timezone.utc) if now.timestamp() - _latest_sample_ts(entries) < UPTIME_SAMPLE_INTERVAL: return None labelled = [c for c in raw if (c.get("Labels") or {}).get(COMPOSE_LABEL)] if not labelled: return None running = sum(1 for c in labelled if c.get("State") == "running") value = round(running / len(labelled) * 100, 1) entry = {"ts": now.isoformat(timespec="seconds"), "value": value} _append_uptime(entry) return entry def prune_uptime_file() -> None: """Drop samples older than twice the chart window (called at startup).""" entries = _read_uptime() if not entries: return cutoff = (datetime.now(timezone.utc).date() - timedelta(days=UPTIME_DAYS * 2)).isoformat() kept = [e for e in entries if (_entry_date(e) or cutoff) >= cutoff] if len(kept) == len(entries): return os.makedirs(settings.DATA_DIR, exist_ok=True) tmp = UPTIME_FILE + ".tmp" with open(tmp, "w", encoding="utf-8") as fh: for e in kept: fh.write(json.dumps(e) + "\n") os.replace(tmp, UPTIME_FILE) def uptime_series(raw: list[dict]) -> list[dict]: _sample_uptime(raw) by_date: dict[str, list[float]] = {} for e in _read_uptime(): day = _entry_date(e) value = e.get("value") if day and isinstance(value, (int, float)): by_date.setdefault(day, []).append(float(value)) series = [] today = datetime.now(timezone.utc).date() last_value: Optional[float] = None for i in range(UPTIME_DAYS - 1, -1, -1): day = (today - timedelta(days=i)).isoformat() samples = by_date.get(day) if samples: last_value = round(sum(samples) / len(samples), 1) # Days before monitoring started (or gaps) reuse the last known value # so the chart doesn't show artificial dips. series.append({"date": day, "value": last_value}) return series async def uptime_sampler_loop() -> None: """Background task: keep the uptime series fed even when nobody is looking at the dashboard.""" prune_uptime_file() while True: try: raw = await _containers_with_timeout() await asyncio.to_thread(_sample_uptime, raw) except (DockerError, asyncio.TimeoutError, OSError): pass await asyncio.sleep(UPTIME_SAMPLE_INTERVAL) # --------------------------------------------------------------------------- # # Ops (audit-log) activity # --------------------------------------------------------------------------- # _WEEKDAYS = ["Monday", "Tuesday", "Wednesday", "Thursday", "Friday", "Saturday", "Sunday"] def ops_activity(session: Session) -> tuple[list[dict], Optional[str]]: today = datetime.now(timezone.utc).date() cutoff = datetime.combine(today - timedelta(days=UPTIME_DAYS - 1), datetime.min.time(), timezone.utc) timestamps = session.exec( select(AuditLog.timestamp).where(AuditLog.timestamp >= cutoff) ).all() per_day: dict[str, int] = {} per_weekday = [0] * 7 for ts in timestamps: per_day[ts.date().isoformat()] = per_day.get(ts.date().isoformat(), 0) + 1 per_weekday[ts.weekday()] += 1 series = [] for i in range(UPTIME_DAYS - 1, -1, -1): day = (today - timedelta(days=i)).isoformat() series.append({"date": day, "count": per_day.get(day, 0)}) peak = _WEEKDAYS[per_weekday.index(max(per_weekday))] if any(per_weekday) else None return series, peak async def compute_summary(session: Session) -> dict: try: raw = await _containers_with_timeout() except (DockerError, asyncio.TimeoutError): raw = [] labelled = [c for c in raw if (c.get("Labels") or {}).get(COMPOSE_LABEL)] ops_series, peak = ops_activity(session) return { "total_containers": sum(1 for c in labelled if c.get("State") == "running"), "containers_total": len(labelled), "uptime_series": uptime_series(raw), "ops_last_30d": ops_series, "ops_peak_day": peak, "as_of": datetime.now(timezone.utc).isoformat(), }