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>
180 lines
8.0 KiB
Python
180 lines
8.0 KiB
Python
"""Talking to the runners (host sidecar / paired workers) and the run-exit watcher."""
|
|
|
|
import datetime
|
|
import json
|
|
import time
|
|
import urllib.error
|
|
import urllib.request
|
|
import uuid as uuidlib
|
|
|
|
from fastapi import HTTPException
|
|
|
|
import workers as workers_mod
|
|
from core.config import SIDECAR_TOKEN, SIDECAR_URL
|
|
from core.state import AppState
|
|
|
|
|
|
# ── spawn / resume / interrupt `claude -p` sessions via the host sidecar ──────
|
|
def _sidecar_call(deps: AppState, path: str, payload: dict | None = None, *,
|
|
method: str = "POST", timeout: float = 15,
|
|
worker: str | None = None) -> dict:
|
|
"""Call a runner (bearer auth) and return its JSON reply: the host sidecar,
|
|
or — with ``worker`` — that paired remote worker (workers.py).
|
|
|
|
Raises 503 when the sidecar isn't configured, 502 when it's unreachable,
|
|
and passes a sidecar error through with its own status. The detail is
|
|
always an object ``{message, kind, hint?, account?}`` (see the sidecar's
|
|
``claude_cli.error_detail``), so the viewer can tell an expired login from
|
|
a rate limit from a broken CLI without parsing prose.
|
|
|
|
A worker runs on exactly one Claude login — its own — so ``account`` is
|
|
dropped from the payload (the worker's sidecar takes its default)."""
|
|
if worker:
|
|
w = deps.worker_store.get(worker)
|
|
if not w:
|
|
raise HTTPException(409, {"message": f"unknown worker: {worker!r} "
|
|
"(unpaired?)",
|
|
"kind": "worker_unknown"})
|
|
if payload is not None:
|
|
payload = {k: v for k, v in payload.items() if k != "account"}
|
|
try:
|
|
return workers_mod.sidecar_request(
|
|
w["url"], w["token"], path, payload, method=method,
|
|
timeout=timeout, label=f"worker {w.get('name') or worker}")
|
|
except HTTPException as e:
|
|
if isinstance(e.detail, dict) and \
|
|
e.detail.get("kind") == "worker_offline":
|
|
deps.worker_store.set_status(worker, status="offline",
|
|
lastError=e.detail.get("message"))
|
|
deps.hub.publish({"type": "workers"})
|
|
raise
|
|
if not SIDECAR_TOKEN:
|
|
raise HTTPException(503, {"message": "sidecar not configured (no "
|
|
"SIDECAR_TOKEN)", "kind": "sidecar"})
|
|
return workers_mod.sidecar_request(SIDECAR_URL, SIDECAR_TOKEN, path,
|
|
payload, method=method, timeout=timeout)
|
|
|
|
|
|
def _detail_text(e: HTTPException) -> str:
|
|
"""An HTTPException's detail as one line of text (structured or not)."""
|
|
d = e.detail
|
|
return d.get("message", str(d)) if isinstance(d, dict) else str(d)
|
|
|
|
|
|
_sidecar_detail = workers_mod.error_detail
|
|
|
|
|
|
def _valid_sid(raw: str) -> str:
|
|
sid = (raw or "").strip()
|
|
try:
|
|
uuidlib.UUID(sid)
|
|
except ValueError:
|
|
raise HTTPException(400, "sessionId must be a valid UUID")
|
|
return sid
|
|
|
|
|
|
# ── run-exit detection ────────────────────────────────────────────────────────
|
|
# The sidecar knows exactly when a run's process dies (its pidfiles), so the
|
|
# viewer doesn't have to guess "finished" from the DONE marker or wait out the
|
|
# recency window. A background thread polls the sidecar's live-session list and
|
|
# stamps a conversation finished the moment its run vanishes from it.
|
|
|
|
# Sessions with an in-flight /resume or /interrupt: those deliberately stop or
|
|
# replace the run, so the exit watcher must not race the vanish and stamp
|
|
# "finished" over the state the endpoint is about to write.
|
|
_sidecar_busy: set[str] = set()
|
|
|
|
|
|
def _sidecar_live_sids(deps: AppState, worker: str | None = None) -> set[str] | None:
|
|
"""Session ids with a live process behind them, per the runner's
|
|
``/sessions`` (which also reaps finished runs' pidfiles as it answers) —
|
|
the host sidecar, or the paired ``worker``. None = unconfigured/unreachable
|
|
— the caller must treat that as "unknown", never as "nothing is running"."""
|
|
if worker:
|
|
w = deps.worker_store.get(worker)
|
|
if not w or not deps.worker_store.online(worker):
|
|
return None
|
|
base, token = w["url"], w["token"]
|
|
elif not SIDECAR_TOKEN:
|
|
return None
|
|
else:
|
|
base, token = SIDECAR_URL, SIDECAR_TOKEN
|
|
req = urllib.request.Request(
|
|
f"{base}/sessions", headers={"Authorization": f"Bearer {token}"})
|
|
try:
|
|
with urllib.request.urlopen(req, timeout=3) as r:
|
|
out = json.loads(r.read().decode())
|
|
return {s["sessionId"] for s in out.get("sessions", [])
|
|
if s.get("sessionId")}
|
|
except (urllib.error.URLError, OSError, ValueError, KeyError):
|
|
return None
|
|
|
|
|
|
def _mark_run_exited(deps: AppState, sid: str) -> None:
|
|
"""A sidecar-launched run's process is gone: stamp the conversation
|
|
``finished`` (with the exit time) unless an endpoint/skill already moved it
|
|
to another state. The transcript's own trailing markers still outrank the
|
|
stamp in ``_resolve_state`` (interrupted, DONE+notified, …), and
|
|
``exitedAt`` lets it be distrusted if the session later continues outside
|
|
the sidecar (see ``_outlived_exit``)."""
|
|
if deps.meta_store.get(sid).get("state") not in (None, "running"):
|
|
return
|
|
now = datetime.datetime.now(datetime.timezone.utc) \
|
|
.strftime("%Y-%m-%dT%H:%M:%SZ")
|
|
deps.meta_store.update(sid, {"state": "finished", "exitedAt": now})
|
|
deps.hub.publish({"type": "meta", "id": sid})
|
|
|
|
|
|
def _watch_sidecar_runs(deps: AppState, interval: float = 2.0) -> None:
|
|
"""Poll the sidecar's live-session list; a session id vanishing from it
|
|
means its process exited — flip the conversation to ``finished`` within a
|
|
tick instead of after the recency window.
|
|
|
|
The first successful poll only sets the baseline (a fresh backend can't
|
|
tell "exited while we were down" from "never sidecar-launched" — those
|
|
stay on the recency fallback). Ids with an in-flight resume/interrupt are
|
|
deferred a tick so the endpoint's own stamp wins; if the resume then
|
|
failed, the id is still absent next tick and gets flipped then.
|
|
|
|
One baseline **per runner** (the host sidecar = ``None``, then each paired
|
|
worker): an unreachable runner answers None and keeps *its* baseline, so a
|
|
Mac dropping off the tailnet never reads as "all its runs exited" — and
|
|
never disturbs the lab sidecar's."""
|
|
prev: dict[str | None, set[str]] = {}
|
|
while True:
|
|
time.sleep(interval)
|
|
runners: list[str | None] = [None] + [w["id"]
|
|
for w in deps.worker_store.list()]
|
|
for runner in runners:
|
|
try:
|
|
cur = _sidecar_live_sids(deps, runner)
|
|
if cur is None:
|
|
continue # unknown — keep the baseline, flip nothing
|
|
if runner in prev:
|
|
for sid in prev[runner] - cur:
|
|
if sid in _sidecar_busy:
|
|
cur.add(sid) # decide next tick
|
|
continue
|
|
_mark_run_exited(deps, sid)
|
|
prev[runner] = cur
|
|
if runner:
|
|
deps.worker_store.set_status(runner, running=len(cur))
|
|
except Exception:
|
|
pass
|
|
for gone in [r for r in prev if r not in runners]:
|
|
prev.pop(gone, None) # unpaired
|
|
|
|
|
|
def _sidecar_get(path: str, timeout: float = 25) -> dict | None:
|
|
"""GET from the host sidecar; None when it's unconfigured/unreachable."""
|
|
if not SIDECAR_TOKEN:
|
|
return None
|
|
req = urllib.request.Request(
|
|
f"{SIDECAR_URL}{path}",
|
|
headers={"Authorization": f"Bearer {SIDECAR_TOKEN}"})
|
|
try:
|
|
with urllib.request.urlopen(req, timeout=timeout) as r:
|
|
return json.loads(r.read().decode())
|
|
except (urllib.error.URLError, OSError, ValueError):
|
|
return None
|