Files
ai-agent/sidecar/routes_runs.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

280 lines
12 KiB
Python

"""
The run routes: ``/sessions``, ``/spawn``, ``/resume``, ``/interrupt`` and
``/message`` (a follow-up into a live run through its inbox socket).
"""
import json
import time
import os
import pathlib
import signal
import socket
import uuid as uuidlib
from fastapi import APIRouter, Header, HTTPException
from pydantic import BaseModel
import claude_cli
import pids
from claude_accounts import _require_login, _valid_account
from config import (ACCOUNTS, CLAUDE_PROJECTS, DEFAULT_ACCOUNT, DEFAULT_MODEL,
INBOX_PREFIX, LOG_DIR, PI_DEFAULT_MODEL, RUNNER_BIN)
from launch import _claude_args, _launch, _remote_control, _run_cwd, _stop_run
from validate import (_auth, _effort_args, _thinking_args, _valid_harness,
_valid_model, _valid_uuid)
router = APIRouter()
class SpawnBody(BaseModel):
prompt: str
sessionId: str | None = None # caller may pre-pick the id (ai-agent does)
model: str | None = None
cwd: str | None = None
harness: str | None = None # "claude" (default) | "pi"
thinking: bool | None = None # False ⇒ run with thinking off; None = on
effort: str | None = None # claude only; None ⇒ the CLI's default
account: str | None = None # claude only; None ⇒ DEFAULT_ACCOUNT
class ResumeBody(BaseModel):
sessionId: str # the existing session to continue
prompt: str # the new user turn to send
model: str | None = None
cwd: str | None = None # must match the session's original cwd
harness: str | None = None # "claude" (default) | "pi"
thinking: bool | None = None # False ⇒ run with thinking off; None = on
effort: str | None = None # claude only; None ⇒ what it last ran with
# The account the session started on: its transcript only exists in that
# account's projects/, so a resume on any other one can't find it.
account: str | None = None
# force=True: a run that's still live doesn't answer `reused` — it is
# stopped first (SIGINT its group, wait for it to die, escalate to SIGKILL)
# and the resume then launches with the new prompt. This is how "send a
# message to a running conversation" restarts the session on the new turn.
force: bool = False
@router.get("/sessions")
def sessions(authorization: str | None = Header(default=None)) -> dict:
"""List sessions still running (survivors of any sidecar restart included)."""
_auth(authorization)
out = []
if LOG_DIR.is_dir():
for p in sorted(LOG_DIR.glob("*.pid")):
sid = p.stem
rec = pids._live_record(sid)
if rec:
out.append({
"sessionId": sid,
"pid": rec.get("pid"),
"startedAt": rec.get("startedAt"),
"model": rec.get("model"),
"cwd": rec.get("cwd"),
"account": rec.get("account"),
})
return {"sessions": out, "count": len(out)}
@router.post("/spawn")
def spawn(body: SpawnBody, authorization: str | None = Header(default=None)) -> dict:
_auth(authorization)
prompt = (body.prompt or "").strip()
if not prompt:
raise HTTPException(400, "empty prompt")
sid = _valid_uuid(body.sessionId or str(uuidlib.uuid4()))
cwd = _run_cwd(body.cwd)
harness = _valid_harness(body.harness)
if harness == "pi":
# pi harness: same flag shape, different binary. The runner wraps
# `pi --mode json` and mirrors the session into the watched pi
# transcripts dir (permission scope = the run's cwd by default).
model = _valid_model(body.model, "pi") or PI_DEFAULT_MODEL
args = [RUNNER_BIN, "-p", prompt, "--session-id", sid, "--model", model]
args += _thinking_args(body.thinking, harness)
return _launch(args, sid, cwd, model=model)
# The viewer's composer sends the model tag the user picked (a CLI alias like
# `opus`, or a pinned id like `claude-opus-4-5-20251101`); with none picked we
# fall back to DEFAULT_MODEL. Either way `claude --model` resolves it.
model = _valid_model(body.model) or DEFAULT_MODEL
account = _valid_account(body.account)
_require_login(account)
args = _claude_args(prompt, "--session-id", sid) + ["--model", model]
args += _thinking_args(body.thinking, harness)
args += _effort_args(body.effort, harness)
args += _remote_control(account)
return _launch(args, sid, cwd, model=model, account=account)
@router.post("/resume")
def resume(body: ResumeBody, authorization: str | None = Header(default=None)) -> dict:
"""Continue an existing conversation with a new user turn.
Runs ``claude -p <prompt> --resume <sessionId>`` (default: reuses the
original session id, so it appends to the same
``~/.claude/projects/<proj>/<sid>.jsonl`` the viewer already watches — the
new turns stream straight into the open conversation). ``--resume`` is scoped
to the current directory, so the caller must pass the session's original
``cwd``.
``model`` switches the model the continuation runs on (the composer sends the
conversation's current model by default, so nothing changes unless the user
picks another tag). Omitted ⇒ no ``--model`` flag ⇒ the session keeps its own."""
_auth(authorization)
prompt = (body.prompt or "").strip()
if not prompt:
raise HTTPException(400, "empty prompt")
sid = _valid_uuid(body.sessionId)
harness = _valid_harness(body.harness)
model = _valid_model(body.model, harness)
cwd = _run_cwd(body.cwd)
account = _valid_account(body.account) if harness == "claude" else None
if account:
_require_login(account)
# Idempotency: if a run for this session is *really* still going, a rapid
# double-submit would launch a SECOND `claude -p --resume` appending duplicate
# turns to the same transcript (and clobbering the pidfile so only one stays
# interruptible). We take no new prompt in that case — and say so
# (`reused`), because the turn was NOT accepted: the caller has to surface
# that rather than show a message that will never be answered.
# With force=True the caller wants the new prompt to win instead: stop the
# live run (it lands its interrupt marker in the transcript), then resume.
rec = pids._live_record(sid)
if rec:
if not body.force:
return {"sessionId": sid, "pid": rec.get("pid"),
"log": str(LOG_DIR / f"{sid}.log"), "reused": True}
_stop_run(sid, rec)
if harness == "pi":
# The composer sends the conversation's current model by default; with
# nothing picked the runner's own default keeps the session coherent.
model = model or PI_DEFAULT_MODEL
args = [RUNNER_BIN, "-p", prompt, "--resume", sid, "--model", model]
args += _thinking_args(body.thinking, harness)
return _launch(args, sid, cwd, model=model, append=True)
args = _claude_args(prompt, "--resume", sid)
if model:
args += ["--model", model]
args += _thinking_args(body.thinking, harness)
args += _effort_args(body.effort, harness)
args += _remote_control(account)
return _launch(args, sid, cwd, model=model, append=True, account=account)
class InterruptBody(BaseModel):
sessionId: str
@router.post("/interrupt")
def interrupt(body: InterruptBody,
authorization: str | None = Header(default=None)) -> dict:
"""Stop a running session by sending its process group a SIGINT.
SIGINT is what Ctrl+C delivers to a terminal's foreground process group, so
a headless ``claude -p`` run handles it the same way: it aborts the current
turn, writes a ``[Request interrupted by user]`` marker to its transcript,
and exits. We signal the whole process group (pid == pgid, since the run was
started with ``start_new_session``) so in-flight tool subprocesses stop too.
Reads the run's identity from its on-disk pidfile, so this works even for a
session spawned by a *previous* sidecar instance that has since restarted."""
_auth(authorization)
sid = _valid_uuid(body.sessionId)
# Verify it's still our live process before signalling — never fire a signal
# at a PID that has already exited (possibly still a zombie), or been
# recycled into something else. _live_record drops the stale pidfile for us.
rec = pids._live_record(sid)
if rec is None:
raise HTTPException(404, "no running process for this session")
pid = rec["pid"]
pid_path = pids._pidfile(sid)
try:
os.killpg(pid, signal.SIGINT)
except ProcessLookupError:
# Raced us to exit between the liveness check and the signal.
pid_path.unlink(missing_ok=True)
raise HTTPException(404, "process already exited")
except OSError as e:
raise HTTPException(500, f"could not signal process: {e}")
pid_path.unlink(missing_ok=True)
return {"sessionId": sid, "pid": pid, "signal": "SIGINT", "ok": True}
class MessageBody(BaseModel):
sessionId: str
prompt: str
def _inbox_socket(rec: dict) -> pathlib.Path | None:
"""A live run's cross-session inbox socket, from the CLI's session registry
(``<config dir>/sessions/<pid>.json`` → ``messagingSocketPath``). None when
the run has none (a pi run, an older CLI, a registry not written yet)."""
cfg = ACCOUNTS.get(rec.get("account") or DEFAULT_ACCOUNT) \
or CLAUDE_PROJECTS.parent
try:
info = json.loads((cfg / "sessions" / f"{rec['pid']}.json").read_text())
except (OSError, ValueError, KeyError):
return None
p = info.get("messagingSocketPath") if isinstance(info, dict) else None
return pathlib.Path(p) if p else None
@router.post("/message")
def message(body: MessageBody,
authorization: str | None = Header(default=None)) -> dict:
"""Deliver a follow-up into a *live* run without stopping it.
The run keeps its background work (a build, a subagent, a Monitor on a
form) — the message lands as a user turn the model reads between tool
calls, or at once when it is idle-waiting. 404 when no process is live
for the session (the caller then does a plain resume), 409
``inbox_unavailable`` when the run can't take it (no registry entry, no
socket, refused) — the caller then falls back to the force-resume."""
_auth(authorization)
sid = _valid_uuid(body.sessionId)
prompt = (body.prompt or "").strip()
if not prompt:
raise HTTPException(400, "empty prompt")
rec = pids._live_record(sid)
if rec is None:
raise HTTPException(404, claude_cli.error_detail(
"no running process for this session", kind="not_running"))
# Inside a Stop-guard hold the CLI's inbox socket accepts and drops: the
# hook holds the turn (``<sid>.hold`` in LOG_DIR), so the text is left in
# its mailbox and the hook relays it in its block reason instead.
if (LOG_DIR / f"{sid}.hold").exists():
box = LOG_DIR / f"{sid}.inbox"
try:
box.mkdir(exist_ok=True)
(box / f"{time.time_ns()}.json").write_text(json.dumps(
{"prompt": prompt, "at": time.time()}))
except OSError as e:
raise HTTPException(409, claude_cli.error_detail(
f"could not leave the message for the held run: {e}",
kind="inbox_unavailable"))
return {"sessionId": sid, "pid": rec["pid"], "delivered": "hold"}
sock = _inbox_socket(rec)
if sock is None or not sock.exists():
raise HTTPException(409, claude_cli.error_detail(
"the running session has no inbox socket", kind="inbox_unavailable"))
line = json.dumps({"type": "user", "message": {
"role": "user", "content": INBOX_PREFIX + prompt}}) + "\n"
try:
with socket.socket(socket.AF_UNIX, socket.SOCK_STREAM) as s:
s.settimeout(5)
s.connect(str(sock))
s.sendall(line.encode())
except OSError as e:
raise HTTPException(409, claude_cli.error_detail(
f"the inbox socket refused the message: {e}", kind="inbox_unavailable"))
return {"sessionId": sid, "pid": rec["pid"], "delivered": "inbox"}