"""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()