Files
ai-agent/backend/workers/store.py
Gabriel Vidal 0ff9e40242 refactor(backend): split main.py into domain packages with app.state injection
main.py (5146 lines, 106 routes) becomes an assembly only: one package per
domain — core, files, settings, dashboards, conversations, diff,
notifications, forms, runs, accounts, workers, models, cron, agents,
projects, services, goals, memories, plans, templates — each exposing an
APIRouter; the flat domain modules move into their package behind a barrel
that keeps the old `import conversations` / `import projects` spellings.

The shared singletons (store, meta_store, hub, indexer, …) are built once by
core.state.build_state() and attached to app.state.ai; routes take them as
the `deps: State` dependency and helpers as an explicit `deps: AppState`.
conversations/pricing.py carries the per-model rates out of the parser.

Verified: route table and OpenAPI byte-identical; 90 read endpoints
golden-diffed against the monolith on a copy of the live data (identical);
write routes smoke-tested; 66 backend tests pass.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
2026-10-06 23:55:47 +02:00

391 lines
17 KiB
Python

"""
Workers — runners on other machines the hub drives like its own sidecar.
The homelab's host sidecar is the default runner: the backend POSTs to
``SIDECAR_URL`` and reads the transcripts through a mount. A **worker** is the
very same ``sidecar/sidecar.py`` running somewhere else (the Orus MacBook, as a
launchd agent — see ``worker/``), reached over the tailnet with its own bearer.
Every run is still "a POST to a sidecar URL"; a worker is just a second
``(url, token)`` pair chosen per run.
The direction never flips: **the hub pulls, the worker never calls back.**
Spawn/resume/fork/interrupt are POSTs to the worker, liveness is a poll of its
``/health`` + ``/sessions``, and its transcripts come over its ``/feed``
(``remotefeed.py`` mirrors them into a dir the indexer already watches).
Pairing (``pair``): the worker's installer prints ``host:port/code/login``.
The hub resolves which of *its* accounts that login is (``login`` is the
email the Mac's ``claude`` is signed in with — no pick needed when it matches
one), trades the one-time code for the worker's bearer at ``POST /pair``, and
stores the worker in ``/data/workers.json``. Only persisted fields go to disk;
the live status (online / offline / unauthorized, last seen, running count)
lives in memory and is re-learnt by the poller within a tick of a restart.
"""
from __future__ import annotations
import datetime
import json
import os
import pathlib
import re
import threading
import time
import urllib.error
import urllib.parse
import urllib.request
from fastapi import HTTPException
DEFAULT_PORT = 8790
# The fields a worker record keeps on disk. Everything else is runtime status.
_PERSISTED = ("id", "name", "host", "port", "url", "token", "account", "email",
"orgName", "hostname", "platform", "machine", "version", "cwd",
"permissionMode", "pairedAt", "autoRoute")
_CODE_RE = re.compile(r"^[A-Za-z0-9]{6,32}$")
_HOST_RE = re.compile(r"^[A-Za-z0-9.-]{1,253}$")
def _now_iso() -> str:
return datetime.datetime.now(datetime.timezone.utc) \
.strftime("%Y-%m-%dT%H:%M:%SZ")
# ── the sidecar HTTP call (shared with main's local-sidecar path) ─────────────
def error_detail(body: str) -> dict:
"""A sidecar error body → the structured detail. An older sidecar (or a
FastAPI validation error) answers a plain string / list — wrap it."""
try:
detail = json.loads(body).get("detail")
except (ValueError, AttributeError):
detail = body[:500]
if isinstance(detail, dict) and detail.get("message"):
return detail
return {"message": str(detail)[:500] or "sidecar error", "kind": "unknown"}
def sidecar_request(base: str, token: str, path: str,
payload: dict | None = None, *, method: str = "POST",
timeout: float = 15, label: str = "sidecar") -> dict:
"""Call a runner (bearer auth) and return its JSON reply.
502 when it's unreachable, and a runner error passes through with its own
status — except 401/403 (its bearer check), which become 502: the viewer
reads any 401 as an expired Authelia session and reloads into the login
page. The detail is always ``{message, kind, hint?, account?}``."""
data = json.dumps(payload).encode() if payload is not None else None
req = urllib.request.Request(
f"{base}{path}", data=data, method=method,
headers={"Content-Type": "application/json",
"Authorization": f"Bearer {token}"})
try:
with urllib.request.urlopen(req, timeout=timeout) as r:
return json.loads(r.read().decode())
except urllib.error.HTTPError as e:
body = e.read().decode(errors="replace")
if e.code in (401, 403):
if label == "sidecar":
raise HTTPException(502, error_detail(body))
raise HTTPException(502, {
"message": f"{label} refused the hub's token — re-pair it",
"kind": "worker_auth"})
status = e.code if 400 <= e.code < 600 else 502
raise HTTPException(status, error_detail(body))
except (urllib.error.URLError, OSError, ValueError) as e:
kind = "worker_offline" if label != "sidecar" else "sidecar"
raise HTTPException(502, {"message": f"{label} unreachable: {e}",
"kind": kind})
def _get_json(url: str, token: str | None = None, timeout: float = 4) -> dict:
headers = {"Authorization": f"Bearer {token}"} if token else {}
with urllib.request.urlopen(urllib.request.Request(url, headers=headers),
timeout=timeout) as r:
return json.loads(r.read().decode())
# ── pairing string ────────────────────────────────────────────────────────────
def parse_pairing(raw: str) -> dict:
"""``100.80.162.92:8790/K7QMX4PJ2R/gabriel@orus.insure`` → its parts.
The login segment is optional (an older installer, a Mac with no login
yet); a leading ``http://`` is tolerated because people paste URLs. 400
with ``kind: "pairing"`` on anything else."""
s = (raw or "").strip()
s = re.sub(r"^https?://", "", s)
head, _, rest = s.partition("/")
code, _, login = rest.partition("/")
host, _, port = head.rpartition(":") if ":" in head else (head, "", "")
try:
port_n = int(port) if port else DEFAULT_PORT
except ValueError:
port_n = -1
if not (host and _HOST_RE.match(host) and 0 < port_n < 65536
and _CODE_RE.match(code or "")):
raise HTTPException(400, {
"message": "not a pairing string — expected host:port/code "
"(the line install-macos.sh prints)",
"kind": "pairing"})
return {"host": host, "port": port_n, "code": code.upper(),
"login": urllib.parse.unquote(login).strip() or None}
# ── the store ─────────────────────────────────────────────────────────────────
class WorkerStore:
"""``/data/workers.json`` — MetaStore/CronStore discipline: one lock,
atomic saves, mtime reload (the deploy-cutover twin shares the file)."""
def __init__(self, path: str):
self.path = pathlib.Path(path)
self._lock = threading.Lock()
self._mtime: float | None = None
self._data: dict[str, dict] = {}
self._status: dict[str, dict] = {}
self.version = 0
self._load()
def _load(self) -> None:
try:
self._data = json.loads(self.path.read_text(encoding="utf-8")) or {}
self._mtime = self.path.stat().st_mtime
except (OSError, ValueError):
self._data, self._mtime = {}, None
# Records paired before the flag existed auto-routed by construction
# (that was the only behaviour) — keep them doing so.
for rec in self._data.values():
rec.setdefault("autoRoute", True)
def _save(self) -> None:
self.path.parent.mkdir(parents=True, exist_ok=True)
tmp = self.path.with_suffix(".json.tmp")
fd = os.open(tmp, os.O_WRONLY | os.O_CREAT | os.O_TRUNC, 0o600)
with os.fdopen(fd, "w", encoding="utf-8") as f:
json.dump(self._data, f, indent=2, sort_keys=True)
os.replace(tmp, self.path)
try:
self._mtime = self.path.stat().st_mtime
except OSError:
pass
self.version += 1
def _reload_if_changed(self) -> None:
try:
m = self.path.stat().st_mtime
except OSError:
if self._mtime is not None:
self._data, self._mtime = {}, None
return
if m != self._mtime:
self._load()
self.version += 1
def list(self) -> list[dict]:
with self._lock:
self._reload_if_changed()
return [dict(w) for w in self._data.values()]
def get(self, wid: str | None) -> dict | None:
if not wid:
return None
with self._lock:
self._reload_if_changed()
w = self._data.get(wid)
return dict(w) if w else None
def put(self, w: dict) -> dict:
rec = {k: w.get(k) for k in _PERSISTED}
with self._lock:
self._reload_if_changed()
self._data[rec["id"]] = rec
self._save()
return dict(rec)
def patch(self, wid: str, changes: dict) -> dict | None:
with self._lock:
self._reload_if_changed()
cur = self._data.get(wid)
if not cur:
return None
cur.update({k: v for k, v in changes.items() if k in _PERSISTED})
self._save()
return dict(cur)
def delete(self, wid: str) -> bool:
with self._lock:
self._reload_if_changed()
gone = self._data.pop(wid, None) is not None
if gone:
self._save()
self._status.pop(wid, None)
return gone
# runtime status (memory only)
def status(self, wid: str) -> dict:
with self._lock:
return dict(self._status.get(wid) or {"status": "unknown"})
def set_status(self, wid: str, **st) -> bool:
"""Merge status fields; True when ``status`` itself changed."""
with self._lock:
cur = self._status.setdefault(wid, {"status": "unknown"})
changed = "status" in st and st["status"] != cur.get("status")
cur.update(st)
return changed
def online(self, wid: str | None) -> bool:
return bool(wid) and self.status(wid).get("status") == "online"
def public(self, w: dict) -> dict:
"""What the API serves: no token, plus the live status."""
out = {k: w.get(k) for k in _PERSISTED if k != "token"}
st = self.status(w["id"])
out.update({"status": st.get("status", "unknown"),
"lastSeen": st.get("lastSeen"),
"lastError": st.get("lastError"),
"running": st.get("running", 0),
"mirrored": st.get("mirrored", 0),
"syncing": bool(st.get("syncing"))})
return out
def for_account(self, account: str | None) -> dict | None:
"""The first *online*, auto-routing worker carrying ``account`` —
where a new run on that account goes when the composer didn't pick a
runner. A worker without ``autoRoute`` (the default at pairing) only
runs what the composer sends it explicitly: the lab stays the runner,
a personal laptop that happens to be awake doesn't take over."""
for w in self.list():
if (w.get("account") == account and w.get("autoRoute")
and self.online(w["id"])):
return w
return None
# ── pairing ───────────────────────────────────────────────────────────────────
def pair(store: WorkerStore, raw: str, *, account: str | None,
name: str | None, hub_name: str,
resolve_login, valid_account, auto_route: bool = False) -> dict:
"""Pair with a worker from its pairing string. The account is decided
**before** the code is spent (explicit pick → the string's login matched
against the hub's accounts → the login segment naming an account id), so a
string whose login matches nothing asks for a pick instead of burning the
one-time code."""
p = parse_pairing(raw)
acct = valid_account(account) if account else None
if account and not acct:
raise HTTPException(400, {"message": f"unknown account: {account!r}",
"kind": "account"})
if not acct and p["login"]:
acct = resolve_login(p["login"]) or valid_account(p["login"])
if not acct:
raise HTTPException(400, {
"message": (f"no account on this hub is signed in as "
f"{p['login']} — pick the account" if p["login"] else
"the pairing string names no Claude login — pick the "
"account this worker runs on"),
"kind": "account", "login": p["login"]})
url = f"http://{p['host']}:{p['port']}"
try:
req = urllib.request.Request(
f"{url}/pair", method="POST",
data=json.dumps({"code": p["code"], "hubName": hub_name}).encode(),
headers={"Content-Type": "application/json"})
with urllib.request.urlopen(req, timeout=10) as r:
out = json.loads(r.read().decode())
except urllib.error.HTTPError as e:
d = error_detail(e.read().decode(errors="replace"))
hint = {"bad_code": "check the code, or run `install-macos.sh --pair` "
"on the worker for a new one",
"no_code": "run `install-macos.sh --pair` on the worker"}
raise HTTPException(409 if e.code in (403, 409) else 502,
{**d, "kind": d.get("kind") or "pairing",
"hint": hint.get(d.get("kind"))})
except (urllib.error.URLError, OSError, ValueError) as e:
raise HTTPException(502, {
"message": f"worker unreachable at {url}: {e}",
"kind": "worker_offline",
"hint": "is the worker running (install-macos.sh --status) and "
"is this box on the same tailnet?"})
if not out.get("token") or not out.get("workerId"):
raise HTTPException(502, {"message": "worker answered /pair without a "
"token", "kind": "pairing"})
login = out.get("account") or {}
w = store.put({
"id": out["workerId"],
"name": (name or "").strip()[:60] or out.get("name")
or out.get("hostname") or p["host"],
"host": p["host"], "port": p["port"], "url": url,
"token": out["token"], "account": acct,
"email": login.get("email") or p["login"],
"orgName": login.get("orgName"),
"hostname": out.get("hostname"), "platform": out.get("platform"),
"machine": out.get("machine"), "version": out.get("version"),
"cwd": out.get("cwd"), "permissionMode": out.get("permissionMode"),
"pairedAt": _now_iso(), "autoRoute": bool(auto_route),
})
store.set_status(w["id"], status="online", lastSeen=_now_iso(),
lastError=None)
return w
def unpair(store: WorkerStore, wid: str) -> bool:
"""Forget a worker; tell it to drop the bearer too (best-effort — an
offline worker keeps a token nobody holds any more)."""
w = store.get(wid)
if not w:
return False
try:
sidecar_request(w["url"], w["token"], "/unpair", {}, timeout=5,
label="worker")
except HTTPException:
pass
return store.delete(wid)
# ── liveness ──────────────────────────────────────────────────────────────────
def probe(store: WorkerStore, w: dict) -> bool:
"""One health check: ``/health`` (is it up) + ``/sessions`` (does our token
still work — a re-pair elsewhere rotates it). Updates the status; True
when the status changed."""
try:
h = _get_json(f"{w['url']}/health")
s = _get_json(f"{w['url']}/sessions", w["token"])
changes = {k: h.get(k) for k in ("version", "permissionMode", "cwd")
if h.get(k) is not None and h.get(k) != w.get(k)}
if changes:
store.patch(w["id"], changes)
return store.set_status(w["id"], status="online", lastSeen=_now_iso(),
lastError=None, running=s.get("count", 0))
except urllib.error.HTTPError as e:
st = "unauthorized" if e.code in (401, 403) else "offline"
return store.set_status(w["id"], status=st, lastError=f"HTTP {e.code}")
except (urllib.error.URLError, OSError, ValueError) as e:
return store.set_status(w["id"], status="offline",
lastError=str(getattr(e, "reason", e))[:200])
class WorkerMonitor(threading.Thread):
"""Polls every paired worker's health; publishes a ``workers`` SSE event
when one goes online/offline (the Settings cards, the composer chip and
the home banner all key off it)."""
def __init__(self, store: WorkerStore, publish, interval: float = 15.0):
super().__init__(daemon=True, name="worker-monitor")
self.store, self.publish, self.interval = store, publish, interval
def check(self, wid: str | None = None) -> None:
changed = False
for w in self.store.list():
if wid and w["id"] != wid:
continue
changed |= probe(self.store, w)
if changed or wid:
self.publish({"type": "workers"})
def run(self) -> None:
while True:
try:
self.check()
except Exception:
pass
time.sleep(self.interval)