PROJECT92-v2 / system.py
1990two's picture
Update system.py
dc787e5 verified
Raw History Blame Contribute Delete
6.14 kB
"""
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)