""" system.py -- the hub. Core logic that isn't specific to any one /LIBRARY module: the Supabase client, the storage bucket bridge, logging, error capture, and writing a handler's response back to the `inbound` table. `dispatch` is the trigger table Claude/Snickers write into to reach this Space (a Postgres trigger fires on insert and calls /command) -- it is never a write target from inside the Space. Writing to it here would re-fire the trigger and create a loop. BUCKET ------ Reads and writes the attached Hugging Face bucket (1990two/PROJECT92-v2-storage by default, override with the BUCKET_ID env var). Needs an HF_TOKEN Space secret with write access to the bucket. This is where any module's durable state lives now -- LIBRARY/AGENTS/task_lock and LIBRARY/SYS/validate both persist here (as small JSON files) instead of local sqlite, because the Space's own container disk is wiped on every restart/redeploy and the bucket doesn't. bucket_load()/bucket_load_json() return None for a path that doesn't exist (not an exception) -- callers treat "no file yet" as a normal state, not an error. System-level loops, heartbeats, and polling belong here as they get added -- this is the one file that's allowed to know about Supabase or the bucket. """ import json import os import logging import tempfile from datetime import datetime, timezone from typing import Optional from supabase import create_client, Client from huggingface_hub import batch_bucket_files, download_bucket_files, list_bucket_tree log = logging.getLogger("system") _client: Optional[Client] = None BUCKET_ID = os.environ.get("BUCKET_ID", "1990two/PROJECT92-v2-storage") def init(): global _client url = os.environ["SUPABASE_URL"] key = os.environ["SUPABASE_SERVICE_KEY"] _client = create_client(url, key) log.info("Supabase client initialized against %s", url) def client() -> Client: if _client is None: raise RuntimeError("system.init() has not been called") return _client def route_response(target: str, command: str, response: dict) -> dict: row = { "target": target, "command": command, "payload": response.get("payload", {}), "agent_system": response.get("agent_system", "claude"), "status": "not_read", } client().table("inbound").insert(row).execute() return {"routed_to": "inbound", **row} def log_error(target: str, command: str, error: str): try: client().table("inbound").insert({ "target": target, "command": command, "payload": {"error": error}, "agent_system": "system", "status": "not_read", }).execute() except Exception: log.exception("failed to write error to inbound table") # ------------------------------------------------------------------ jobs def claim_next_job() -> Optional[dict]: """Atomically claim the oldest pending job via the claim_next_job() SQL function.""" result = client().rpc("claim_next_job").execute() return result.data[0] if result.data else None def mark_job_running(job_id: str) -> None: client().table("jobs").update({"status": "running"}).eq("id", job_id).execute() def complete_job(job_id: str, result: dict) -> None: client().table("jobs").update({ "status": "done", "result": result, "completed_at": datetime.now(timezone.utc).isoformat(), }).eq("id", job_id).execute() def fail_job(job_id: str, error: str) -> None: client().table("jobs").update({ "status": "failed", "error": error, "completed_at": datetime.now(timezone.utc).isoformat(), }).eq("id", job_id).execute() async def poll_loop(interval: float = 30.0) -> None: """Background coroutine. Every tick: drains one pending job, then advances all running orchestrator tasks.""" import asyncio from LIBRARY.SYS import orchestrator # lazy import from LIBRARY.SYS import scheduler # lazy import -- avoids circular at module load time log.info("poll_loop started (interval=%.0fs)", interval) while True: try: result = scheduler.poll_once() if result.get("status") != "idle": log.info("poll_loop scheduler: %s", result) except Exception: log.exception("poll_loop: unhandled error in scheduler.poll_once()") try: orch_result = orchestrator.poll_once() if orch_result.get("tasks_processed", 0) > 0: log.info("poll_loop orchestrator: %s", orch_result) except Exception: log.exception("poll_loop: unhandled error in orchestrator.poll_once()") await asyncio.sleep(interval) # ------------------------------------------------------------------ bucket def bucket_save(path: str, content) -> None: """content may be str or bytes.""" data = content.encode() if isinstance(content, str) else content batch_bucket_files(BUCKET_ID, add=[(data, path)]) def bucket_load(path: str): """Returns raw bytes, or None if path doesn't exist in the bucket.""" with tempfile.TemporaryDirectory() as tmp: local_path = os.path.join(tmp, os.path.basename(path) or "file") try: download_bucket_files(BUCKET_ID, files=[(path, local_path)]) except Exception: return None if not os.path.exists(local_path): return None with open(local_path, "rb") as f: return f.read() def bucket_list(prefix: str = ""): items = list_bucket_tree(BUCKET_ID, prefix=prefix, recursive=True) return [{"path": i.path, "size": i.size} for i in items if i.type == "file"] def bucket_delete(paths: list) -> None: batch_bucket_files(BUCKET_ID, delete=paths) def bucket_save_json(path: str, obj) -> None: bucket_save(path, json.dumps(obj)) def bucket_load_json(path: str): """Returns the parsed object, or None if path doesn't exist.""" raw = bucket_load(path) if raw is None: return None text = raw.decode() if isinstance(raw, (bytes, bytearray)) else raw return json.loads(text)