Files
ai-agent/backend/workers/routes.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

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}