Spaces:
Running
Running
Download Browser_Agent/job_queue/task_store.py from Luca448/APP-Backend: direct link, hf CLI and curl.
- Browser
- Download file 4.54 kB
-
https://huggingface.co/spaces/Luca448/APP-Backend/resolve/main/Browser_Agent/job_queue/task_store.py
- Command line
-
hf download hf://spaces/Luca448/APP-Backend/Browser_Agent/job_queue/task_store.py
-
curl -L -o task_store.py https://huggingface.co/spaces/Luca448/APP-Backend/resolve/main/Browser_Agent/job_queue/task_store.py
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", | |
| }, | |
| ) | |
| def configured(self) -> bool: | |
| return self._client is not None | |
| def ready(self) -> bool: | |
| return self._remote_ready | |
| 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() | |