Spaces:
Running
Running
| """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)], | |
| ) | |
| 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} | |
| 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} | |