The worker installer no longer assumes the Orus work seat. It asks for the default working dir — the directory it is run from, confirmed with `y`, or a typed path (taken as is without a terminal) — and for the hub-side account, proposed from `claude auth status` (a gmail login → personal, else work). `--update` keeps the installed agent's cwd and account. Remote Control stays off on every worker. Pairing a personal Mac on `personal` would have taken every personal spawn while it was awake (for_account), so workers now carry an `autoRoute` flag: off at pairing (the lab stays the runner; the worker only gets what the composer's Runs-on chip pins to it), toggled from the Settings → Workers card or the pairing dialog. Records from before the flag keep auto-routing. Docs genericized (worker README, CLAUDE.md Workers, notify SKILL, GOAL). Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
391 lines
17 KiB
Python
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)
|