Spaces:
Running
Running
File size: 5,019 Bytes
28a08e7 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 | """backend/api/daemon_status.py — Live status of the Telegram session daemon.
Espone lo stato del session-daemon (scripts/session-daemon.mjs) leggendo
le sessioni attive da Supabase agent_tasks (status = "__session__").
Endpoint:
GET /api/daemon/status — sessioni attive, ultimo heartbeat, uptime, task corrente
"""
from __future__ import annotations
import asyncio
import json
import logging
import time
from typing import Any
from fastapi import APIRouter, Depends
from .auth_guard import require_role, AuthRole
_logger = logging.getLogger("api.daemon_status")
router = APIRouter(prefix="/api/daemon", tags=["daemon-status"], dependencies=[Depends(require_role(AuthRole.MACHINE))]) # GAP-1-fix: router-level auth
# ACTIVE_TTL: allineato a session-daemon.mjs (5 * 60_000 ms)
_ACTIVE_TTL_MS = 5 * 60 * 1000 # 5 minuti
def _parse_context(raw: Any) -> dict:
"""Decodifica il campo context (può essere str JSON o dict)."""
if isinstance(raw, dict):
return raw
if isinstance(raw, str):
try:
return json.loads(raw)
except Exception:
return {}
return {}
def _fmt_uptime(started_at_ms: int | None) -> str:
if not started_at_ms:
return "unknown"
elapsed_s = int((time.time() * 1000 - started_at_ms) / 1000)
h, rem = divmod(elapsed_s, 3600)
m, s = divmod(rem, 60)
if h:
return f"{h}h {m}m"
if m:
return f"{m}m {s}s"
return f"{s}s"
@router.get("/status")
async def daemon_status() -> dict:
"""
Ritorna lo stato live del session-daemon leggendo Supabase agent_tasks.
Una sessione è ATTIVA se updated_at < 5 minuti fa (ACTIVE_TTL del daemon).
Il daemon fa heartbeat ogni 30s — se non si vede da >5min è stale/crashato.
Response:
ok — False se Supabase non raggiungibile
supabase_ok — True se query OK
active_sessions — conteggio sessioni con heartbeat < 5min
sessions[] — lista sessioni con: session_name, active, last_heartbeat,
uptime, current_task, pid, head_sha, claimed_files
checked_at — timestamp UTC della verifica
"""
from .state import _sb
checked_at = time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime())
now_ms = int(time.time() * 1000)
if not _sb:
return {
"ok": False,
"supabase_ok": False,
"error": "Supabase non configurato (SUPABASE_URL / SUPABASE_KEY mancanti)",
"sessions": [],
"active_sessions": 0,
"checked_at": checked_at,
}
try:
result = await asyncio.to_thread(
lambda: _sb.table("agent_tasks")
.select("task_id, goal, context, updated_at, created_at")
.eq("status", "__session__")
.order("updated_at", desc=True)
.limit(10)
.execute()
)
rows = result.data or []
except Exception as exc:
_logger.warning("daemon_status: Supabase query failed: %s", exc)
return {
"ok": False,
"supabase_ok": False,
"error": str(exc),
"sessions": [],
"active_sessions": 0,
"checked_at": checked_at,
}
sessions = []
for row in rows:
ctx = _parse_context(row.get("context"))
updated_ms: int = row.get("updated_at") or 0
age_ms = now_ms - updated_ms
# CLOCK-SKEW-FIX: Railway clock può essere leggermente avanti di HF Space.
# age_ms negativo = heartbeat *appena* avvenuto → trattare come attivo.
# Tolleriamo fino a 60s di skew (ben oltre il tipico 1-2s osservato).
is_active = -60_000 < age_ms < _ACTIVE_TTL_MS
ago_s = int(age_ms / 1000)
if ago_s < 0:
last_seen = "appena ora"
elif ago_s < 60:
last_seen = f"{ago_s}s fa"
elif ago_s < 3600:
last_seen = f"{ago_s // 60}m {ago_s % 60}s fa"
else:
last_seen = f"{ago_s // 3600}h {(ago_s % 3600) // 60}m fa"
head = (ctx.get("headSha") or "")
sessions.append({
"session_id": row.get("task_id", ""),
"session_name": row.get("goal") or ctx.get("sessionName", ""),
"active": is_active,
"last_heartbeat": last_seen,
"uptime": _fmt_uptime(ctx.get("startedAt")),
"current_task": ctx.get("sprint") or "idle",
"pid": ctx.get("pid"),
"head_sha": head[:10] if head else None,
"claimed_files": ctx.get("claimedFiles", []),
"updated_at_ms": updated_ms,
})
active_count = sum(1 for s in sessions if s["active"])
return {
"ok": True,
"supabase_ok": True,
"active_sessions": active_count,
"total_sessions": len(sessions),
"sessions": sessions,
"active_ttl_s": _ACTIVE_TTL_MS // 1000,
"checked_at": checked_at,
}
|