File size: 3,326 Bytes
5176170
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
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
"""Stato operativo amministrativo protetto da JWT Supabase admin."""
from __future__ import annotations

from datetime import datetime, timedelta, timezone
from typing import Any

from fastapi import APIRouter, Depends, Query

from .auth_guard import require_admin_user
from .private_state import _MAX_TASK_PAGE, _as_epoch_ms, _call, _json_object

router = APIRouter(
    prefix="/api/admin/state",
    tags=["admin"],
    dependencies=[Depends(require_admin_user)],
)


@router.get("/sessions")
async def admin_sessions(
    max_age_ms: int = Query(default=300_000, ge=10_000, le=3_600_000),
    limit: int = Query(default=100, ge=1, le=200),
) -> dict[str, object]:
    cutoff = (datetime.now(timezone.utc) - timedelta(milliseconds=max_age_ms)).isoformat()

    def operation(client: Any):
        return client.table("agent_tasks").select("task_id,context,updated_at").eq("status", "__session__").gte("updated_at", cutoff).order("updated_at", desc=True).limit(limit).execute()

    result = await _call(operation)
    sessions = []
    for row in result.data or []:
        context = _json_object(row.get("context"))
        session_id = str(context.get("sessionId") or row.get("task_id") or "").strip()
        if not session_id:
            continue
        claimed = context.get("claimedFiles")
        sessions.append({
            "session_id": session_id,
            "session_name": str(context.get("sessionName") or session_id)[:160],
            "sprint": str(context["sprint"])[:120] if context.get("sprint") else None,
            "claimed_files": [str(item)[:300] for item in claimed[:100]] if isinstance(claimed, list) else [],
            "last_heartbeat": _as_epoch_ms(context.get("lastHeartbeat")) or _as_epoch_ms(row.get("updated_at")),
            "current_task": str(context["currentTask"])[:500] if context.get("currentTask") else None,
        })
    return {"sessions": sessions}


@router.get("/tasks")
async def admin_tasks(
    limit: int = Query(default=20, ge=1, le=_MAX_TASK_PAGE),
    offset: int = Query(default=0, ge=0, le=10_000),
    status: str | None = Query(default=None, max_length=64),
) -> dict[str, object]:
    normalized_status = status.strip().upper() if status else ""

    def operation(client: Any):
        query = client.table("agent_tasks").select("task_id,goal,status,updated_at").neq("status", "__session__").neq("status", "__config__")
        if normalized_status:
            query = query.eq("status", normalized_status)
        page = query.order("updated_at", desc=True).range(offset, offset + limit - 1).execute()
        all_statuses = client.table("agent_tasks").select("status").neq("status", "__session__").neq("status", "__config__").limit(2_000).execute()
        return page, all_statuses

    page, all_statuses = await _call(operation)
    counts: dict[str, int] = {}
    for row in all_statuses.data or []:
        key = str(row.get("status") or "UNKNOWN").upper()
        counts[key] = counts.get(key, 0) + 1
    tasks = [{
        "task_id": str(row.get("task_id") or ""),
        "goal": str(row.get("goal") or "")[:1_000],
        "status": str(row.get("status") or "UNKNOWN"),
        "updated_at": _as_epoch_ms(row.get("updated_at")),
    } for row in page.data or []]
    return {"tasks": tasks, "counts": counts, "offset": offset, "limit": limit}