Luca448's picture
Savestate: Fix Dockerfile and add Browser_Agent for HF build, bump version to 1.8.334
508316a
Raw History Blame Contribute Delete
4.54 kB
"""Durable task storage backed by Supabase with an in-process safety cache."""
from __future__ import annotations
import asyncio
import logging
from typing import Any
import httpx
from config.settings import settings
logger = logging.getLogger(__name__)
class TaskStore:
"""Small persistence boundary used by the queue and API."""
def __init__(self) -> None:
self._memory: dict[str, dict[str, Any]] = {}
self._lock = asyncio.Lock()
self._remote_ready = False
self._last_error: str | None = None
self._client: httpx.AsyncClient | None = None
if settings.supabase_url and settings.supabase_service_role_key:
base_url = settings.supabase_url.rstrip("/") + "/rest/v1"
self._client = httpx.AsyncClient(
base_url=base_url,
timeout=10.0,
headers={
"apikey": settings.supabase_service_role_key,
"Authorization": f"Bearer {settings.supabase_service_role_key}",
"Content-Type": "application/json",
},
)
@property
def configured(self) -> bool:
return self._client is not None
@property
def ready(self) -> bool:
return self._remote_ready
@property
def last_error(self) -> str | None:
return self._last_error
async def initialize(self) -> None:
if not self._client:
return
try:
response = await self._client.get("/agent_tasks", params={"select": "id", "limit": "1"})
response.raise_for_status()
self._remote_ready = True
self._last_error = None
except Exception as exc:
self._remote_ready = False
self._last_error = type(exc).__name__
logger.error("Supabase task store unavailable: %s", type(exc).__name__)
async def save(self, payload: dict[str, Any]) -> None:
async with self._lock:
self._memory[str(payload["id"])] = dict(payload)
if not self._client:
return
try:
response = await self._client.post(
"/agent_tasks",
params={"on_conflict": "id"},
headers={"Prefer": "resolution=merge-duplicates,return=minimal"},
json=payload,
)
response.raise_for_status()
self._remote_ready = True
self._last_error = None
except Exception as exc:
self._remote_ready = False
self._last_error = type(exc).__name__
logger.error("Task %s could not be persisted: %s", payload.get("id"), type(exc).__name__)
async def append_event(self, payload: dict[str, Any]) -> None:
if not self._client:
return
try:
response = await self._client.post(
"/agent_events",
headers={"Prefer": "resolution=ignore-duplicates,return=minimal"},
json=payload,
)
response.raise_for_status()
self._remote_ready = True
except Exception as exc:
self._remote_ready = False
self._last_error = type(exc).__name__
logger.error("Task event could not be persisted: %s", type(exc).__name__)
async def load_active(self) -> list[dict[str, Any]]:
if not self._client:
async with self._lock:
return list(self._memory.values())
active = (
"pending,planning,running,retrying,waiting_for_confirmation,"
"waiting_for_handoff"
)
try:
response = await self._client.get(
"/agent_tasks",
params={
"select": "*",
"status": f"in.({active})",
"order": "created_at.asc",
},
)
response.raise_for_status()
rows = response.json()
self._remote_ready = True
self._last_error = None
return rows if isinstance(rows, list) else []
except Exception as exc:
self._remote_ready = False
self._last_error = type(exc).__name__
logger.error("Active tasks could not be restored: %s", type(exc).__name__)
async with self._lock:
return list(self._memory.values())
async def close(self) -> None:
if self._client:
await self._client.aclose()
task_store = TaskStore()