Files
ai-agent/backend/runs/service.py
Gabriel Vidal 68a3206641 Merge branch 'main' into split-big-files
# Conflicts:
#	backend/main.py
#	backend/schemas.py
#	sidecar/sidecar.py
#	sidecar/test_claude_args.py
2026-10-06 23:57:03 +02:00

279 lines
14 KiB
Python

"""The one send path behind both composers: spawn, resume, runner and account routing."""
import datetime
import os
import pathlib
import time
import uuid as uuidlib
from fastapi import HTTPException
import accounts as accounts_mod
from conversations.catalog import _session_conv_map
from core.config import TRANSCRIPTS_DIR, WORKERS_SOURCE_DIR
from core.state import AppState
from runs.sidecar_client import (
_sidecar_busy,
_sidecar_call,
_sidecar_live_sids,
_valid_sid,
)
# A worker session whose transcript grew this recently with no worker-launched
# run behind it is probably still open in a terminal on that machine.
TERMINAL_LIVE_S = float(os.environ.get("TERMINAL_LIVE_S", "60"))
def _session_worker(deps: AppState, sid: str) -> str | None:
"""The worker a session lives on: the one the hub launched it on
(``meta.worker``), else the one the mirror imported it from
(``meta.workerSource``). None = the host sidecar's."""
m = deps.meta_store.peek(sid)
return m.get("worker") or m.get("workerSource")
def _mirror_transcript(sid: str) -> pathlib.Path | None:
if not WORKERS_SOURCE_DIR or not WORKERS_SOURCE_DIR.is_dir():
return None
return next(WORKERS_SOURCE_DIR.glob(f"*/{sid}.jsonl"), None)
def _terminal_live_guard(deps: AppState, sid: str, worker: str) -> None:
"""Refuse (409 ``terminal_live``) to resume a worker session a terminal on
that machine may still hold: two writers on one transcript interleave its
parent chain. The worker can't see the terminal's process, so it is a
heuristic — the mirror copy grew in the last TERMINAL_LIVE_S and the worker
has no run of its own for the session. The UI turns it into a confirm."""
ap = _mirror_transcript(sid)
if ap is None:
return
try:
age = time.time() - ap.stat().st_mtime
except OSError:
return
if age > TERMINAL_LIVE_S:
return
live = _sidecar_live_sids(deps, worker) or set()
if sid in live:
return # our own run — a normal force-resume
# Growth our own last run explains: it wrote the file, then the exit
# watcher stamped exitedAt (or hasn't ticked yet and it still reads
# running). Only writes *after* that are a terminal's.
m = deps.meta_store.peek(sid)
if m.get("worker") and m.get("state") == "running":
return
exited = m.get("exitedAt")
if exited:
try:
ex = datetime.datetime.strptime(exited, "%Y-%m-%dT%H:%M:%SZ") \
.replace(tzinfo=datetime.timezone.utc).timestamp()
if ap.stat().st_mtime <= ex + 5:
return
except (ValueError, OSError):
pass
w = deps.worker_store.get(worker) or {}
raise HTTPException(409, {
"message": f"this session was active in a terminal on "
f"{w.get('name') or 'the worker'} {int(age)}s ago",
"kind": "terminal_live", "ageSeconds": int(age),
"hint": "close it there first, or continue anyway — two writers on "
"one session garble its history"})
# The explicit "run it on the lab's own sidecar" runner pick.
LOCAL_RUNNER = "local"
def _pick_worker(deps: AppState, choice: str | None, account: str, harness: str) -> str | None:
"""Which runner a *new* run goes to: an explicit worker id (must be
online — a run silently landing elsewhere is the worse surprise), "local"
for the host sidecar, unset ⇒ the account's online worker, else local.
Workers run claude only (no pi runner there)."""
c = (choice or "").strip()
if c == LOCAL_RUNNER:
return None
if c:
w = deps.worker_store.get(c)
if not w:
raise HTTPException(400, {"message": f"unknown worker: {c!r}",
"kind": "worker_unknown"})
if harness != "claude":
raise HTTPException(400, {"message": "workers run claude only",
"kind": "worker_harness"})
if not deps.worker_store.online(c):
raise HTTPException(409, {
"message": f"worker {w.get('name') or c} is offline",
"kind": "worker_offline",
"hint": "wake the machine, or send it to the lab instead"})
return c
if harness != "claude":
return None
w = deps.worker_store.for_account(account)
return w["id"] if w else None
def _session_account(deps: AppState, sid: str, cwd: str | None = None) -> str:
"""The account an existing session runs on — the one a resume/fork must
launch on, since its transcript lives only in that account's ``projects/``.
Same order as ``accounts.resolve``, fed from the metadata stamp and the
stored summary's cwd (falling back to the cwd the caller sent, for a
session whose transcript hasn't been parsed yet)."""
conv = _session_conv_map(deps, ).get(sid)
summary = deps.store.get_summary(str(TRANSCRIPTS_DIR / conv)) if conv else None
return accounts_mod.resolve(summary or {"cwd": cwd}, deps.meta_store.peek(sid))
def _send_message(deps: AppState, prompt: str, *, session_id: str | None = None,
model: str | None = None, cwd: str | None = None,
harness: str | None = None,
thinking: bool | None = None,
effort: str | None = None,
account: str | None = None,
projects: list[str] | None = None,
worker: str | None = None,
takeover: bool = False) -> dict:
"""The one send path behind both composers.
With no ``session_id`` this starts a new conversation: we mint the UUID here
so the frontend can start watching for the transcript to sync (~2s) and
redirect to it without waiting for the run. With one, it continues that
session — and a continue *always* takes the prompt: if a run is still live
(or parked on a ScheduleWakeup/Monitor), the sidecar stops it first (SIGINT
the group, wait for it to die) and then launches the ``--resume`` with the
new turn, so sending into a "running"/"paused" conversation restarts it on
the new message. Either way the actual agent process (`claude -p` or the pi
runner, per ``harness``) is launched by the host sidecar (the container
can't reach the host's authenticated CLIs), and we stamp ``running``
immediately so the UI flips before the first turn has synced.
**Account.** A new claude run bills ``account`` when given (the composer's
chip), else the account owning one of ``projects`` (orus-monorepo → work),
else the default one; a non-default account also starts in its own repo
unless ``cwd`` says otherwise. A continue always keeps the account the
session started on — whatever the caller asks. The choice is stamped as
``meta.account`` next to ``harness``. pi runs aren't a Claude login: they
bill the default account.
**Worker.** A new run goes to ``worker`` when the composer picked one
("local" = the host sidecar), else to the account's online worker, else
the host (``_pick_worker``); a worker run bills the worker's own login.
A continue goes where the session lives (``_session_worker``) — never
elsewhere, its transcript isn't anywhere else — and a session a terminal
may still hold needs ``takeover`` (``_terminal_live_guard``). Stamped as
``meta.worker``."""
prompt = (prompt or "").strip()
if not prompt:
raise HTTPException(400, "empty prompt")
if session_id:
sid = _valid_sid(session_id)
# A conversation resumes on the harness it started on — and with the
# thinking setting it last ran with — unless the caller overrides. Both
# were stamped into the metadata sidecar by the spawn.
prev = deps.meta_store.get(sid)
harness = (harness or prev.get("harness") or "claude").lower()
if thinking is None:
thinking = prev.get("thinking")
if effort is None:
effort = prev.get("effort")
account = (_session_account(deps, sid, cwd) if harness == "claude"
else accounts_mod.DEFAULT)
run_on = _session_worker(deps, sid) if harness == "claude" else None
if run_on:
w = deps.worker_store.get(run_on)
if w and not deps.worker_store.online(run_on):
raise HTTPException(409, {
"message": f"this session lives on {w.get('name') or run_on}"
", which is offline",
"kind": "worker_offline",
"hint": "wake the machine — a session can't move runners"})
if not takeover:
_terminal_live_guard(deps, sid, run_on)
if not cwd: # --resume is cwd-scoped: the session's own, on the Mac
conv = _session_conv_map(deps, ).get(sid)
summ = deps.store.get_summary(str(TRANSCRIPTS_DIR / conv)) \
if conv else None
cwd = (summ or {}).get("cwd")
# A run that is still alive takes the follow-up through its inbox
# socket (sidecar /message): its background work — a build, a
# subagent, a Monitor on a form — survives, and the model reads the
# message between tool calls, or at once when it is idle-waiting.
# Only when nothing is live (404: a plain --resume) or the live run
# can't take it (409 inbox_unavailable: a pi run, an older CLI) does
# the message restart the session the old way, below. A model /
# thinking / effort switch doesn't apply through the inbox — the run
# keeps what it was launched with; Stop it first to switch.
if harness == "claude":
try:
out = _sidecar_call(deps, "/message", {"sessionId": sid,
"prompt": prompt},
worker=run_on)
except HTTPException as e:
if e.status_code not in (404, 409):
raise
else:
deps.meta_store.update(sid, {"state": "running"})
deps.hub.publish({"type": "meta", "id": sid})
return {"sessionId": sid, "pid": out.get("pid"),
"delivered": out.get("delivered") or "inbox"}
# A force-resume swaps the session's process; shield the swap from the
# exit watcher so the old run's vanish isn't stamped "finished" over
# the "running" written below.
_sidecar_busy.add(sid)
try:
out = _sidecar_call(deps, "/resume", {"sessionId": sid, "prompt": prompt,
"model": model, "cwd": cwd,
"harness": harness,
"thinking": thinking,
"effort": effort,
"account": account,
"force": True}, worker=run_on)
if out.get("reused"):
# force=True makes the sidecar stop a live run instead of
# answering `reused`; only a sidecar predating the flag still
# can. The turn was dropped either way — a 409, not a success,
# so the composer says so instead of waiting on an answer that
# isn't coming.
deps.meta_store.update(sid, {"state": "running"})
deps.hub.publish({"type": "meta", "id": sid})
raise HTTPException(409, "this session is still running and "
"the sidecar is too old to stop it — "
"restart the claude-sidecar service")
# Stamp running (and re-stamp harness + thinking) before lifting the
# shield — the UI flips before the transcript syncs. A None thinking
# is dropped by the merge, so the session keeps what it had.
deps.meta_store.update(sid, {"state": "running", "harness": harness,
"thinking": thinking, "effort": effort,
"account": account, "worker": run_on})
finally:
_sidecar_busy.discard(sid)
else:
sid = str(uuidlib.uuid4())
harness = (harness or "claude").lower()
if account and not accounts_mod.valid(account):
raise HTTPException(400, f"unknown account: {account!r}")
if harness == "claude":
account = (accounts_mod.valid(account)
or accounts_mod.account_for_projects(projects)
or accounts_mod.DEFAULT)
else:
account = accounts_mod.DEFAULT
run_on = _pick_worker(deps, worker, account, harness)
if run_on:
# The worker's own login bills it; its own default cwd (the repo
# its installer was pointed at) applies unless the caller chose.
account = (deps.worker_store.get(run_on) or {}).get("account") or account
elif harness == "claude":
cwd = cwd or accounts_mod.sidecar_payload(account)["cwd"]
out = _sidecar_call(deps, "/spawn", {"prompt": prompt, "sessionId": sid,
"model": model, "harness": harness,
"thinking": thinking, "effort": effort,
"account": account, "cwd": cwd},
worker=run_on)
# A fresh id can't be in the exit watcher's baseline — no shield needed.
deps.meta_store.update(sid, {"state": "running", "harness": harness,
"thinking": thinking, "effort": effort,
"account": account, "worker": run_on})
deps.hub.publish({"type": "meta", "id": sid})
return {"sessionId": sid, "pid": out.get("pid")}