# Conflicts: # backend/main.py # backend/schemas.py # sidecar/sidecar.py # sidecar/test_claude_args.py
280 lines
12 KiB
Python
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"}
|