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>
182 lines
8.6 KiB
Python
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}
|