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,
    }