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>
149 lines
6.0 KiB
Python
149 lines
6.0 KiB
Python
"""Remote workers: pair, list, patch, probe, unpair."""
|
|
|
|
import time
|
|
|
|
from fastapi import APIRouter, HTTPException
|
|
from pydantic import BaseModel
|
|
|
|
import accounts as accounts_mod
|
|
import schemas
|
|
import workers as workers_mod
|
|
from core.config import HUB_NAME
|
|
from core.http import _r
|
|
from core.state import AppState, State
|
|
from runs.sidecar_client import _sidecar_get
|
|
|
|
router = APIRouter()
|
|
|
|
|
|
# ── remote workers (workers.py): pair, list, probe, unpair ────────────────────
|
|
class PairWorkerBody(BaseModel):
|
|
pairing: str # host:port/code[/login] from the installer
|
|
account: str | None = None # override the login → account match
|
|
name: str | None = None
|
|
# Send the account's new runs here whenever it is online (unset ⇒ only
|
|
# runs the composer pins to it — the lab stays the default runner).
|
|
autoRoute: bool = False
|
|
|
|
|
|
class WorkerPatchBody(BaseModel):
|
|
name: str | None = None
|
|
account: str | None = None
|
|
autoRoute: bool | None = None
|
|
|
|
|
|
def _account_for_login(login: str) -> str | None:
|
|
"""Which of the hub's accounts is signed in as ``login`` (an email, or an
|
|
org name), per the host sidecar's ``claude auth status`` of each."""
|
|
want = (login or "").strip().lower()
|
|
if not want:
|
|
return None
|
|
live = _sidecar_get("/accounts", timeout=10) or {}
|
|
for a in live.get("accounts") or []:
|
|
for k in ("email", "orgName"):
|
|
if (a.get(k) or "").strip().lower() == want:
|
|
return accounts_mod.valid(a.get("id"))
|
|
return None
|
|
|
|
|
|
def _workers_payload(deps: AppState) -> dict:
|
|
return {"workers": [deps.worker_store.public(w) for w in deps.worker_store.list()]}
|
|
|
|
|
|
def _worker_or_404(deps: AppState, worker_id: str) -> dict:
|
|
w = deps.worker_store.get(worker_id)
|
|
if not w:
|
|
raise HTTPException(404, {"message": f"unknown worker: {worker_id!r}",
|
|
"kind": "worker_unknown"})
|
|
return w
|
|
|
|
|
|
@router.get("/api/workers", responses=_r(schemas.WorkersResponse))
|
|
def list_workers(deps: State):
|
|
"""Paired remote runners with their live status (online / offline /
|
|
unauthorized), last seen, running-run count and mirror progress."""
|
|
return _workers_payload(deps, )
|
|
|
|
|
|
@router.post("/api/workers", responses=_r(schemas.Worker))
|
|
def pair_worker(deps: State, body: PairWorkerBody):
|
|
"""Pair with a worker from the string its installer printed. The account
|
|
comes from the string's Claude login (matched against this hub's accounts)
|
|
unless ``account`` overrides it; the one-time code is only spent once the
|
|
account is settled. Re-pairing a known machine replaces its record."""
|
|
w = workers_mod.pair(deps.worker_store, body.pairing, account=body.account,
|
|
name=body.name, hub_name=HUB_NAME,
|
|
resolve_login=_account_for_login,
|
|
valid_account=accounts_mod.valid,
|
|
auto_route=body.autoRoute)
|
|
deps.hub.publish({"type": "workers"})
|
|
return deps.worker_store.public(w)
|
|
|
|
|
|
@router.patch("/api/workers/{worker_id}", responses=_r(schemas.Worker))
|
|
def patch_worker(deps: State, worker_id: str, body: WorkerPatchBody):
|
|
_worker_or_404(deps, worker_id)
|
|
changes: dict = {}
|
|
if body.name is not None and body.name.strip():
|
|
changes["name"] = body.name.strip()[:60]
|
|
if body.account is not None:
|
|
a = accounts_mod.valid(body.account)
|
|
if not a:
|
|
raise HTTPException(400, {"message": f"unknown account: "
|
|
f"{body.account!r}",
|
|
"kind": "account"})
|
|
changes["account"] = a
|
|
if body.autoRoute is not None:
|
|
changes["autoRoute"] = bool(body.autoRoute)
|
|
w = deps.worker_store.patch(worker_id, changes)
|
|
deps.hub.publish({"type": "workers"})
|
|
return deps.worker_store.public(w)
|
|
|
|
|
|
_worker_skills_cache: dict[str, tuple[float, list[dict]]] = {}
|
|
_WORKER_SKILLS_TTL_S = 60.0
|
|
|
|
|
|
@router.get("/api/workers/{worker_id}/skills")
|
|
def worker_skills(deps: State, worker_id: str, fresh: bool = False):
|
|
"""The skills a run on that worker can use — its sidecar's ``/skills``
|
|
(user-level + default-cwd skills on *that* machine), cached a minute. The
|
|
composer lists these instead of the lab catalog when "Runs on" resolves
|
|
to the worker. An offline or pre-``/skills`` worker answers an empty list
|
|
rather than an error: the chips simply don't show."""
|
|
w = deps.worker_store.get(worker_id)
|
|
if not w:
|
|
raise HTTPException(404, "unknown worker")
|
|
hit = _worker_skills_cache.get(worker_id)
|
|
if hit and not fresh and time.monotonic() - hit[0] < _WORKER_SKILLS_TTL_S:
|
|
return {"skills": hit[1], "count": len(hit[1]), "worker": worker_id}
|
|
skills: list[dict] = []
|
|
if deps.worker_store.online(worker_id):
|
|
try:
|
|
out = workers_mod.sidecar_request(
|
|
w["url"], w["token"], "/skills", method="GET", timeout=6,
|
|
label=f"worker {w.get('name') or worker_id}")
|
|
skills = [s for s in (out.get("skills") or []) if isinstance(s, dict)]
|
|
except HTTPException:
|
|
skills = []
|
|
_worker_skills_cache[worker_id] = (time.monotonic(), skills)
|
|
return {"skills": skills, "count": len(skills), "worker": worker_id}
|
|
|
|
|
|
@router.post("/api/workers/{worker_id}/probe", responses=_r(schemas.Worker))
|
|
def probe_worker(deps: State, worker_id: str):
|
|
"""Health-check one worker now instead of waiting for the next poll."""
|
|
_worker_or_404(deps, worker_id)
|
|
deps.worker_monitor.check(worker_id)
|
|
return deps.worker_store.public(deps.worker_store.get(worker_id))
|
|
|
|
|
|
@router.delete("/api/workers/{worker_id}")
|
|
def unpair_worker(deps: State, worker_id: str):
|
|
"""Forget a worker (and tell it to drop the bearer, if it's reachable).
|
|
Its mirrored conversations stay in the archive, read-only."""
|
|
_worker_or_404(deps, worker_id)
|
|
workers_mod.unpair(deps.worker_store, worker_id)
|
|
deps.hub.publish({"type": "workers"})
|
|
return {"ok": True}
|