Files
ai-agent/backend/runs/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

182 lines
8.6 KiB
Python

"""Spawn / resume / fork / interrupt routes."""
import uuid as uuidlib
from fastapi import APIRouter, HTTPException
from pydantic import BaseModel
import accounts as accounts_mod
import conversations
import schemas
from conversations.catalog import _resolve_transcript, _summary_and_meta
from core.http import _r
from core.state import State
from runs.service import _send_message, _session_worker
from runs.sidecar_client import _sidecar_busy, _sidecar_call, _valid_sid
router = APIRouter()
class SpawnBody(BaseModel):
prompt: str
model: str | None = None
# Which agent harness runs the session: "claude" (the claude CLI, default)
# or "pi" (the pi.dev runner — see services/ai-agent/runner/).
harness: str | None = None
# Whether the session thinks. False ⇒ the sidecar runs the harness with
# extended thinking off (`--thinking disabled` / `--thinking off`); unset ⇒
# each CLI keeps its own default. Both harnesses support it.
thinking: bool | None = None
# How hard a claude run works a turn (`--effort`: low|medium|high|xhigh|max).
# Claude-only — it's what the composer shows instead of the thinking toggle
# for claude models; unset ⇒ the CLI's own default.
effort: str | None = None
# Which Claude account the run bills ("personal"/"work" — accounts.py).
# The composer's account chip; unset ⇒ routed from ``projects``.
account: str | None = None
# The location tags the composer had on (project slugs). A project owned
# by a non-default account routes the spawn there and starts it in that
# account's repo instead of the lab checkout.
projects: list[str] | None = None
# Host working dir for the run; unset ⇒ the account's own (else the
# sidecar's default, the lab checkout).
cwd: str | None = None
# Which runner launches it: a paired worker's id, "local" for the host
# sidecar, unset ⇒ the account's online worker if it has one, else local.
worker: str | None = None
class ResumeBody(BaseModel):
sessionId: str
prompt: str
model: str | None = None
cwd: str | None = None
harness: str | None = None # unset ⇒ the harness the session started on
thinking: bool | None = None # unset ⇒ what the session last ran with
effort: str | None = None # unset ⇒ what the session last ran with
# Accepted for symmetry but never switches a session: its transcript only
# exists in the account it started on (see _session_account).
account: str | None = None
# Resume a worker session even though its terminal looks live (the
# confirm behind a 409 ``terminal_live``).
takeover: bool = False
@router.post("/api/spawn", responses=_r(schemas.SpawnResult))
def spawn_conversation(deps: State, body: SpawnBody):
"""Kick off a new headless agent session on the host (see _send_message)."""
return _send_message(deps, body.prompt, model=body.model, harness=body.harness,
thinking=body.thinking, effort=body.effort,
account=body.account, projects=body.projects,
cwd=body.cwd, worker=body.worker)
@router.post("/api/resume", responses=_r(schemas.SpawnResult))
def resume_conversation(deps: State, body: ResumeBody):
"""Continue an existing conversation with a new user turn (see _send_message).
Runs ``claude -p <prompt> --resume <sessionId>`` in the session's original
cwd; resume reuses the session id, so Claude appends the new turns to the
same transcript the viewer is already watching and the open conversation
streams the continuation live. A still-running session is stopped first —
sending a message always wins."""
return _send_message(deps, body.prompt, session_id=body.sessionId,
model=body.model, cwd=body.cwd, harness=body.harness,
thinking=body.thinking, effort=body.effort,
takeover=body.takeover)
class ForkBody(BaseModel):
id: str # source conversation id (<slug>/<sid>.jsonl)
mi: int # ordinal of the user message being edited (fork point)
prompt: str # the edited message — becomes the fork's next turn
model: str | None = None
harness: str | None = None # unset ⇒ the harness the source ran on
@router.post("/api/fork", responses=_r(schemas.SpawnResult))
def fork_conversation(deps: State, body: ForkBody):
"""Branch a new conversation off an existing one at a user message.
The new session's history is the source transcript up to (excluding) user
message ``mi``; ``prompt`` (the edited message) is then sent as its next
turn. Neither CLI can truncate at a message (`claude --fork-session` only
clones whole sessions), so the host sidecar does transcript surgery: it
copies the first ``cutLine`` lines of the source session file under a
fresh session id and resumes that — see sidecar ``/fork``. ``cutText`` and
``cutUserOrd`` let the sidecar/runner verify they're slicing the same
message this viewer showed (the archive is an append-only copy of the
live file, so the line indices match; the canary catches drift)."""
prompt = (body.prompt or "").strip()
if not prompt:
raise HTTPException(400, "empty prompt")
ap = _resolve_transcript(body.id)
src_sid, summary, src_meta = _summary_and_meta(deps, body.id)
if not src_sid:
raise HTTPException(409, "source conversation has no session id")
cut = conversations.find_cut_line(ap, body.mi)
if cut is None:
raise HTTPException(404, "fork point not found — message ordinal "
"missing or not a user message")
if cut["line"] == 0:
raise HTTPException(400, "nothing before the first message — start a "
"new conversation instead")
harness = (body.harness or src_meta.get("harness") or "claude").lower()
# The fork is written next to the source transcript, in the source's
# account — it can't move to another one.
account = (accounts_mod.resolve(summary, src_meta) if harness == "claude"
else accounts_mod.DEFAULT)
# …and on the runner it lives on: the surgery happens next to the source.
run_on = _session_worker(deps, src_sid) if harness == "claude" else None
sid = str(uuidlib.uuid4())
out = _sidecar_call(deps, "/fork", {
"sessionId": sid, "srcSessionId": src_sid,
"cwd": summary.get("cwd"),
"cutLine": cut["line"], "cutText": cut["text"][:120],
"cutUserOrd": cut["userOrd"],
"prompt": prompt, "model": body.model, "harness": harness,
"account": account}, worker=run_on)
deps.meta_store.update(sid, {
"state": "running", "harness": harness, "account": account,
"worker": run_on,
"forkedFrom": {"id": body.id, "sessionId": src_sid, "mi": body.mi}})
deps.hub.publish({"type": "meta", "id": sid})
return {"sessionId": sid, "pid": out.get("pid")}
class InterruptBody(BaseModel):
sessionId: str
@router.post("/api/interrupt", responses=_r(schemas.InterruptResult))
def interrupt_conversation(deps: State, body: InterruptBody):
"""Stop a running Claude session this viewer spawned.
Proxies to the host sidecar, which sends the run's process group a SIGINT —
exactly like pressing Esc in a terminal. Claude aborts the turn and records a
``[Request interrupted by user]`` marker in the transcript, which the parser
then surfaces as the ``interrupted`` state. We also stamp the metadata
immediately so the UI flips before that marker has synced."""
sid = _valid_sid(body.sessionId)
# Shield the kill from the exit watcher: the vanish it causes must resolve
# to the "interrupted" stamped below, not a generic "finished".
_sidecar_busy.add(sid)
try:
try:
out = _sidecar_call(deps, "/interrupt", {"sessionId": sid},
worker=_session_worker(deps, sid))
except HTTPException as e:
# 404 from the sidecar = no live process for this session (it
# already exited, or was started from a terminal), passed through
# by _sidecar_call. A stale "running" badge over a dead run is
# exactly what the user is clearing here, so stamp it interrupted
# anyway.
if e.status_code != 404:
raise
out = {"note": "no live process — marked interrupted"}
deps.meta_store.update(sid, {"state": "interrupted"})
finally:
_sidecar_busy.discard(sid)
deps.hub.publish({"type": "meta", "id": sid})
return {"sessionId": sid, "ok": True, **out}