mirror of
https://github.com/0xWheatyz/handler.git
synced 2026-08-29 19:21:40 +00:00
feat!: headless is the only runner - delete the tmux run path (phase 4)
Agent runs are now always worker-owned 'claude -p' subprocesses; tmux survives only for the interactive /login flow. - deleted: worker.capture_agent_output/_pane_tail + the capture loop arm (the empty-/log bug's home), spawn's tmux launch/_claude_command, the tmux resume/kill branches (the silent-send-keys bug's home), tmux.session_name/list_sessions, the CLI attach subcommand, the 'runner' setting - spawn: task is now a hard requirement (headless has no idle REPL) - enforced in spawn (SpawnError) and the API (400); onboarding seeding dropped (-p skips the trust dialog) - resume: single headless path; pre-headless agent rows (no session_id) degrade to the context-re-injection fresh run - settings_gen: permissions allowlist is always emitted - credsync: change-triggered uploads key on .claude/.credentials.json only (claude touches ~/.claude.json every run - keying on it would ping-pong uploads between workers); logins still publish explicitly - cli list: liveness from agent_runs in the DB, not tmux - tests: spawn/kill/resume re-pointed at the fake_launch seam (conftest); integration test now drives API -> worker -> real fake claude subprocess -> events endpoint; README documents the headless model + multi-worker deployment invariants Suite 295 green; frontend unchanged since phase 3.
This commit is contained in:
@@ -13,10 +13,11 @@ git remote, and your own network exposure.
|
||||
> control layer, HTTP API, database, migrations, and verification/approval hooks are
|
||||
> implemented and tested (106 tests, SQLite). Phase 2 adds credential resolution +
|
||||
> injection, role-based forge-workflow skills, a hard approval gate, and a CI-status
|
||||
> poller. Live end-to-end agent spawning against a real `claude` binary + tmux is stubbed
|
||||
> behind mockable seams (`tmux`, `verify`, `forge`, `gitops`, `spawn.resume`) and wired
|
||||
> but not yet exercised against production binaries. See [`docs/PLAN.md`](docs/PLAN.md)
|
||||
> for the full design and roadmap.
|
||||
> poller. Agent runs are headless (`claude -p --output-format stream-json`, supervised
|
||||
> by the worker, events persisted to the DB); the run/kill/resume paths are exercised
|
||||
> end-to-end against a scripted fake claude binary, with a manual validation script
|
||||
> (`scripts/validate_claude_headless.sh`) for the real one. See
|
||||
> [`docs/PLAN.md`](docs/PLAN.md) for the full design and roadmap.
|
||||
|
||||
---
|
||||
|
||||
@@ -45,19 +46,23 @@ the control layer and API are disposable compute that can restart or scale out f
|
||||
```
|
||||
writes reads (+ answer backfill)
|
||||
┌──────────────────┐ ┌──────────────┐ ┌──────────────────┐
|
||||
│ control layer │───────▶│ database │◀───────│ HTTP API │
|
||||
│ worker(s) │───────▶│ database │◀───────│ HTTP API │
|
||||
│ (CLI + hooks) │ │ PG / SQLite │ │ (FastAPI) │
|
||||
└──────────────────┘ └──────────────┘ └──────────────────┘
|
||||
│ ▲ ▲
|
||||
│ spawns │ Stop / PreToolUse / Notification hooks │ curl, UI, any client
|
||||
▼ │ write checkmark + log rows │ (bearer token)
|
||||
tmux + claude binary (one working dir / worktree per agent)
|
||||
▼ │ + streamed run events, checkmark, log │ (bearer token)
|
||||
claude -p --output-format stream-json (one working dir / worktree per agent)
|
||||
```
|
||||
|
||||
- **Control layer** (`handler.control`) — the only writer. Spawns/lists/attaches/kills
|
||||
agents as `tmux` sessions running the `claude` binary, one working directory or git
|
||||
worktree per agent, namespaced `project__agent`. Stateless; every write goes straight
|
||||
to the database.
|
||||
- **Control layer / workers** (`handler.control`) — the only writer. Runs each agent as
|
||||
a **headless** `claude -p --output-format stream-json` subprocess (one working
|
||||
directory or git worktree per agent), streams every stdout event into the database as
|
||||
it happens, and reconciles agent status from the process itself (exit code + EOF —
|
||||
positive liveness, no screen scraping). Stateless: repo state is pulled from git when
|
||||
a task is claimed, claude session transcripts are archived to / materialized from the
|
||||
DB for cross-worker `--resume`, and the claude login credential bundle is distributed
|
||||
encrypted through the DB. tmux survives only to drive the interactive `/login` flow.
|
||||
- **Hooks** (`handler.hooks`) — run inside each agent via a generated `settings.json`.
|
||||
They write the checkpoint/log rows and enforce the test and push gates.
|
||||
- **API** (`handler.api`) — a thin, read-mostly HTTP layer over the same database (the
|
||||
@@ -77,10 +82,32 @@ The data model is defined once (SQLAlchemy Core) and renders correctly on both:
|
||||
Portable column types bridge the two, and the checkmark upsert uses native
|
||||
`INSERT … ON CONFLICT DO UPDATE` on both dialects. Migrations are Alembic, dual-dialect.
|
||||
|
||||
### Scaling workers horizontally
|
||||
|
||||
Multiple worker containers can drain the same command queue concurrently (Postgres
|
||||
`FOR UPDATE SKIP LOCKED`); each supervises up to `MAX_CONCURRENT_RUNS` claude processes
|
||||
and skips claiming run-starting commands while full, leaving them for a less-loaded
|
||||
worker. Workers heartbeat into the DB; if one dies mid-run, any surviving worker's
|
||||
reaper marks its runs (and their agents) `crashed` — visible in the UI with the last
|
||||
output preserved — and the operator resumes explicitly on whichever worker picks it up.
|
||||
|
||||
Deployment invariants for multi-worker:
|
||||
|
||||
- **No shared filesystems.** Git carries repo state (workers clone/pull on claim);
|
||||
claude session transcripts live in `session_archives`; login credentials are
|
||||
Fernet-encrypted into `runtime_secrets` and materialized by every worker.
|
||||
- **Identical `PROJECTS_ROOT` on every worker** — claude keys its session storage to
|
||||
the absolute working-dir path, so cross-worker `--resume` needs the same layout.
|
||||
- **The same `HANDLER_SECRET_KEY` on every worker** (and the API) — without it, the
|
||||
credential bundle can't be distributed and only the worker that ran `/login` can run
|
||||
agents.
|
||||
- The two-step web login is automatically pinned to one worker
|
||||
(`commands.target_worker`), so it works unchanged with a fleet.
|
||||
|
||||
## Requirements
|
||||
|
||||
- Python 3.11+
|
||||
- `git` and `tmux` (for live spawning)
|
||||
- `git` (for live spawning) and `tmux` (only for the web `/login` flow)
|
||||
- A `claude` binary, authenticated (for live spawning)
|
||||
- `mise` in each managed project, with a `.mise.toml` defining at least a `test` task
|
||||
- Postgres (default) — or nothing but a file path for the SQLite fallback
|
||||
|
||||
@@ -12,7 +12,6 @@ from fastapi import APIRouter, Depends, HTTPException, Query, status
|
||||
from sqlalchemy import Connection
|
||||
from sqlalchemy.exc import IntegrityError
|
||||
|
||||
from ...config import get_settings
|
||||
from ...db import repository as repo
|
||||
from ..deps import db_conn, require_admin, require_auth
|
||||
from ..schemas import (
|
||||
@@ -79,10 +78,10 @@ def enqueue_spawn(project: str, body: SpawnIn, conn: Connection = Depends(db_con
|
||||
status.HTTP_409_CONFLICT,
|
||||
detail=f"agent '{body.name}' already exists in project '{project}'",
|
||||
)
|
||||
if get_settings().runner == "headless" and not body.task:
|
||||
# A tmux agent can idle at the REPL awaiting input; a headless `claude -p` run
|
||||
# with no prompt exits immediately. Reject here (400) instead of letting the
|
||||
# command fail asynchronously in the worker.
|
||||
if not body.task:
|
||||
# A headless `claude -p` run with no prompt exits immediately having done
|
||||
# nothing. Reject here (400) instead of letting the command fail asynchronously
|
||||
# in the worker.
|
||||
raise HTTPException(
|
||||
status.HTTP_400_BAD_REQUEST,
|
||||
detail="a task is required: the headless runner has no idle-REPL mode",
|
||||
|
||||
@@ -51,10 +51,9 @@ class Settings(BaseSettings):
|
||||
forge_bin: str = "forge"
|
||||
git_bin: str = "git"
|
||||
|
||||
# ---- Headless runner (claude -p --output-format stream-json). ``runner`` selects the
|
||||
# launch path: "tmux" (legacy interactive session, the default until the headless path
|
||||
# is validated) or "headless" (worker-owned subprocess streaming events to the DB).
|
||||
runner: str = "tmux"
|
||||
# ---- Headless runner (claude -p --output-format stream-json): worker-owned
|
||||
# subprocesses streaming events to the DB. Agent runs are always headless; tmux
|
||||
# remains only for the interactive /login flow.
|
||||
# How many concurrent claude runs one worker container supervises; commands that would
|
||||
# start a run are left queued (for another worker) while all slots are busy.
|
||||
max_concurrent_runs: int = 4
|
||||
|
||||
+11
-26
@@ -1,10 +1,11 @@
|
||||
"""``handler`` CLI — the control layer's write side.
|
||||
|
||||
Spawn/list/attach/kill manage agent processes. Phase 2 adds the forge-workflow control
|
||||
commands: ``approve``/``reject`` (the senior agent records its verdict, which the deploy
|
||||
gate checks), ``poll-ci`` (backfill CI verdicts), and ``forge-init`` (write the role
|
||||
skills into a managed repo). The DB is the source of truth for what agents exist; tmux is
|
||||
cross-checked for liveness. All commands are project-namespaced.
|
||||
Spawn/list/kill manage agent runs (headless ``claude -p`` processes supervised by the
|
||||
worker). Phase 2 adds the forge-workflow control commands: ``approve``/``reject`` (the
|
||||
senior agent records its verdict, which the deploy gate checks), ``poll-ci`` (backfill
|
||||
CI verdicts), and ``forge-init`` (write the role skills into a managed repo). The DB is
|
||||
the single source of truth: what agents exist AND whether their runs are live both come
|
||||
from it. All commands are project-namespaced.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -15,7 +16,7 @@ import sys
|
||||
|
||||
from ..db import repository as repo
|
||||
from ..db.engine import connection
|
||||
from . import poller, reposync, skills_gen, spawn, tmux, worker
|
||||
from . import poller, reposync, skills_gen, spawn, worker
|
||||
|
||||
|
||||
def _cmd_spawn(args: argparse.Namespace) -> int:
|
||||
@@ -37,38 +38,27 @@ def _cmd_spawn(args: argparse.Namespace) -> int:
|
||||
print(f" role: {args.role}")
|
||||
if agent.get("forge_note"):
|
||||
print(f" warning: {agent['forge_note']}", file=sys.stderr)
|
||||
print(f" tmux session: {tmux.session_name(args.project, args.name)}")
|
||||
return 0
|
||||
|
||||
|
||||
def _cmd_list(args: argparse.Namespace) -> int:
|
||||
live = set(tmux.list_sessions())
|
||||
with connection() as conn:
|
||||
live = {run["agent_id"] for run in repo.list_running_runs(conn)}
|
||||
if args.project:
|
||||
projects = [args.project] if repo.get_project(conn, args.project) else []
|
||||
else:
|
||||
projects = [p["id"] for p in repo.list_projects(conn)]
|
||||
for project_id in projects:
|
||||
for agent in repo.list_agents(conn, project_id):
|
||||
session = tmux.session_name(project_id, agent["name"])
|
||||
alive = "live" if session in live else "-"
|
||||
alive = "live" if agent["id"] in live else "-"
|
||||
role = agent.get("role") or "-"
|
||||
worker_id = agent.get("worker_id") or "-"
|
||||
print(
|
||||
f"{project_id}/{agent['name']}\t{role}\t{agent['status']}\t{alive}\t{session}"
|
||||
f"{project_id}/{agent['name']}\t{role}\t{agent['status']}\t{alive}\t{worker_id}"
|
||||
)
|
||||
return 0
|
||||
|
||||
|
||||
def _cmd_attach(args: argparse.Namespace) -> int:
|
||||
session = tmux.session_name(args.project, args.name)
|
||||
if not tmux.has_session(session):
|
||||
print(f"error: no live session '{session}'", file=sys.stderr)
|
||||
return 1
|
||||
# Replace this process with an interactive tmux attach.
|
||||
os.execvp("tmux", ["tmux", "attach", "-t", session])
|
||||
return 0 # pragma: no cover - execvp does not return
|
||||
|
||||
|
||||
def _cmd_kill(args: argparse.Namespace) -> int:
|
||||
try:
|
||||
spawn.kill(args.project, args.name)
|
||||
@@ -233,11 +223,6 @@ def build_parser() -> argparse.ArgumentParser:
|
||||
p_list.add_argument("--project", help="limit to one project")
|
||||
p_list.set_defaults(func=_cmd_list)
|
||||
|
||||
p_attach = sub.add_parser("attach", help="attach to an agent's tmux session")
|
||||
p_attach.add_argument("--project", required=True)
|
||||
p_attach.add_argument("--name", required=True)
|
||||
p_attach.set_defaults(func=_cmd_attach)
|
||||
|
||||
p_kill = sub.add_parser("kill", help="kill an agent's session")
|
||||
p_kill.add_argument("--project", required=True)
|
||||
p_kill.add_argument("--name", required=True)
|
||||
|
||||
@@ -43,15 +43,19 @@ def _credential_files() -> dict[str, str]:
|
||||
|
||||
|
||||
def fingerprint() -> tuple:
|
||||
"""(path, mtime_ns, size) of the on-disk credential files — cheap change detection."""
|
||||
fp = []
|
||||
for path in sorted(_credential_files().values()):
|
||||
try:
|
||||
st = os.stat(path)
|
||||
fp.append((path, st.st_mtime_ns, st.st_size))
|
||||
except OSError:
|
||||
continue
|
||||
return tuple(fp)
|
||||
"""(path, mtime_ns, size) of the OAuth token file — cheap change detection.
|
||||
|
||||
Deliberately only ``.claude/.credentials.json``: claude touches ``~/.claude.json``
|
||||
on every run (project entries, UI state), and treating those as "new credentials"
|
||||
would ping-pong uploads between workers forever. A login that only rewrites
|
||||
``.claude.json`` is still published — the login flow calls :func:`upload` directly.
|
||||
"""
|
||||
path = _credential_files()[".claude/.credentials.json"]
|
||||
try:
|
||||
st = os.stat(path)
|
||||
except OSError:
|
||||
return ()
|
||||
return ((path, st.st_mtime_ns, st.st_size),)
|
||||
|
||||
|
||||
def upload() -> bool:
|
||||
@@ -165,10 +169,3 @@ def refresh() -> str | None:
|
||||
_state.seen_updated_at = stored["updated_at"] if stored else None
|
||||
return "uploaded"
|
||||
return None
|
||||
|
||||
|
||||
def note_local_write() -> None:
|
||||
"""Record that this process just changed local credentials deliberately (e.g. the
|
||||
claude_config onboarding merge at spawn), so refresh() doesn't misread the mtime
|
||||
bump as a new login and ping-pong uploads between workers."""
|
||||
_state.last_fingerprint = fingerprint()
|
||||
|
||||
@@ -21,7 +21,7 @@ def _hook_command(event: str) -> str:
|
||||
return f"{sys.executable} -m handler.hooks {event}"
|
||||
|
||||
|
||||
def build_settings(headless: bool = False) -> dict:
|
||||
def build_settings() -> dict:
|
||||
settings = {
|
||||
"hooks": {
|
||||
"Stop": [{"hooks": [{"type": "command", "command": _hook_command("stop")}]}],
|
||||
@@ -39,24 +39,23 @@ def build_settings(headless: bool = False) -> dict:
|
||||
],
|
||||
}
|
||||
}
|
||||
if headless:
|
||||
# ``claude -p`` never prompts — anything that would ask for permission is
|
||||
# auto-denied. The allowlist is therefore what lets normal work (git, mise, the
|
||||
# project's own tooling) proceed; the PreToolUse/Stop hooks above remain the hard
|
||||
# gate either way, since a hook deny overrides any allow.
|
||||
s = get_settings()
|
||||
settings["permissions"] = {
|
||||
"defaultMode": s.headless_permission_mode,
|
||||
"allow": s.headless_allowed_tools_list,
|
||||
}
|
||||
# ``claude -p`` never prompts — anything that would ask for permission is
|
||||
# auto-denied. The allowlist is therefore what lets normal work (git, mise, the
|
||||
# project's own tooling) proceed; the PreToolUse/Stop hooks above remain the hard
|
||||
# gate either way, since a hook deny overrides any allow.
|
||||
s = get_settings()
|
||||
settings["permissions"] = {
|
||||
"defaultMode": s.headless_permission_mode,
|
||||
"allow": s.headless_allowed_tools_list,
|
||||
}
|
||||
return settings
|
||||
|
||||
|
||||
def write_settings(working_dir: str, headless: bool = False) -> str:
|
||||
def write_settings(working_dir: str) -> str:
|
||||
"""Write ``.claude/settings.json`` under the agent's working dir; return its path."""
|
||||
claude_dir = os.path.join(working_dir, ".claude")
|
||||
os.makedirs(claude_dir, exist_ok=True)
|
||||
path = os.path.join(claude_dir, "settings.json")
|
||||
with open(path, "w") as fh:
|
||||
json.dump(build_settings(headless=headless), fh, indent=2)
|
||||
json.dump(build_settings(), fh, indent=2)
|
||||
return path
|
||||
|
||||
@@ -15,16 +15,13 @@ from ..config import get_settings
|
||||
from ..db import repository as repo
|
||||
from ..db.engine import connection
|
||||
from . import (
|
||||
claude_config,
|
||||
credentials,
|
||||
credsync,
|
||||
forge,
|
||||
gitops,
|
||||
headless,
|
||||
mise,
|
||||
reposync,
|
||||
settings_gen,
|
||||
tmux,
|
||||
worktree,
|
||||
)
|
||||
|
||||
@@ -47,18 +44,6 @@ def require_test_task(working_dir: str) -> None:
|
||||
)
|
||||
|
||||
|
||||
def _claude_command(task: str | None, settings_path: str) -> str:
|
||||
claude = get_settings().claude_bin
|
||||
argv = [claude, "--settings", settings_path]
|
||||
if task:
|
||||
argv.append(_shell_quote(task))
|
||||
return " ".join(argv)
|
||||
|
||||
|
||||
def _shell_quote(value: str) -> str:
|
||||
return "'" + value.replace("'", "'\\''") + "'"
|
||||
|
||||
|
||||
def _install_git_credentials(
|
||||
working_dir: str, git_remote: str | None, conn=None
|
||||
) -> None:
|
||||
@@ -97,10 +82,10 @@ def spawn(
|
||||
test gate. ``worker_id`` identifies the calling worker container (headless runs record
|
||||
it on the run row; the CLI defaults to a pid-scoped id).
|
||||
"""
|
||||
if get_settings().runner == "headless" and not task:
|
||||
# A tmux agent without a task idles at the REPL waiting for input; ``claude -p``
|
||||
# has no such mode — an empty prompt would exit immediately having done nothing.
|
||||
raise SpawnError("a headless agent requires a task (the runner is 'headless')")
|
||||
if not task:
|
||||
# ``claude -p`` has no idle-REPL mode — an empty prompt would exit immediately
|
||||
# having done nothing, so a task is a hard requirement.
|
||||
raise SpawnError("an agent requires a task (headless claude has no idle mode)")
|
||||
sync_note = None
|
||||
with connection() as conn:
|
||||
project = repo.get_project(conn, project_id)
|
||||
@@ -154,8 +139,7 @@ def spawn(
|
||||
role=role,
|
||||
)
|
||||
|
||||
headless_run = get_settings().runner == "headless"
|
||||
settings_path = settings_gen.write_settings(working_dir, headless=headless_run)
|
||||
settings_path = settings_gen.write_settings(working_dir)
|
||||
env = _agent_env(project, agent, token, role=role, mise_init=mise_init)
|
||||
|
||||
# Verify the pinned forge version, if one is configured. Non-fatal: a version drift
|
||||
@@ -163,25 +147,14 @@ def spawn(
|
||||
# touches forge and the base image is the real pin (README 3.6, Phase 2).
|
||||
forge_note = _check_forge_version(working_dir)
|
||||
|
||||
# Mark Claude Code onboarding complete + trust the working dir before launching. The
|
||||
# tmux path needs both (no human at the TTY to answer the theme/trust screens);
|
||||
# ``-p`` skips the trust dialog but still reads onboarding state, so keep it for both.
|
||||
claude_config.ensure_onboarded(working_dir)
|
||||
credsync.note_local_write()
|
||||
|
||||
if headless_run:
|
||||
headless.launch(
|
||||
agent,
|
||||
kind="spawn",
|
||||
prompt=task,
|
||||
settings_path=settings_path,
|
||||
env=env,
|
||||
worker_id=worker_id or f"cli-{os.getpid()}",
|
||||
)
|
||||
else:
|
||||
session = tmux.session_name(project_id, name)
|
||||
command = _claude_command(task, settings_path)
|
||||
tmux.new_session(session, cwd=working_dir, command=command, env=env)
|
||||
headless.launch(
|
||||
agent,
|
||||
kind="spawn",
|
||||
prompt=task,
|
||||
settings_path=settings_path,
|
||||
env=env,
|
||||
worker_id=worker_id or f"cli-{os.getpid()}",
|
||||
)
|
||||
agent = {**agent, "forge_note": forge_note, "sync_note": sync_note}
|
||||
return agent
|
||||
|
||||
@@ -234,51 +207,35 @@ def _check_forge_version(working_dir: str) -> str | None:
|
||||
|
||||
|
||||
def kill(project_id: str, name: str) -> None:
|
||||
"""Stop an agent: flag its running run for cancel and mark the row done.
|
||||
|
||||
The owning worker's supervisor polls the cancel flag and SIGTERMs its own child
|
||||
(cross-worker safe — nobody signals a process they don't own). No running run means
|
||||
the process is already gone; the status update is all that's left to do.
|
||||
"""
|
||||
with connection() as conn:
|
||||
agent = repo.get_agent_by_name(conn, project_id, name)
|
||||
if agent is None:
|
||||
raise SpawnError(f"agent '{name}' not found in project '{project_id}'")
|
||||
if agent.get("session_id"):
|
||||
# Headless agent: flag the running run for cancel; the owning worker's
|
||||
# supervisor polls the flag and SIGTERMs its own child (cross-worker safe —
|
||||
# nobody signals a process they don't own). No running run = already dead.
|
||||
run = repo.get_latest_run(conn, agent["id"])
|
||||
if run is not None and run["status"] == "running":
|
||||
repo.request_run_cancel(conn, run["id"])
|
||||
repo.set_agent_status(conn, agent["id"], "done")
|
||||
return
|
||||
session = tmux.session_name(project_id, name)
|
||||
if tmux.has_session(session):
|
||||
tmux.kill_session(session)
|
||||
run = repo.get_latest_run(conn, agent["id"])
|
||||
if run is not None and run["status"] == "running":
|
||||
repo.request_run_cancel(conn, run["id"])
|
||||
repo.set_agent_status(conn, agent["id"], "done")
|
||||
|
||||
|
||||
def resume(agent: dict, answer: str, worker_id: str | None = None) -> tuple[bool, str]:
|
||||
"""Feed an operator's answer back to an agent.
|
||||
"""Feed an operator's answer back to an agent as a new ``claude -p --resume`` run.
|
||||
|
||||
The seam the API's ``/resume`` route calls (and the one tests mock). Legacy tmux
|
||||
agents (``session_id`` null) get the answer typed into their live session; headless
|
||||
agents get a brand-new ``claude -p --resume`` run on this worker, with the session
|
||||
transcript materialized from the DB archive first so any worker can serve the resume.
|
||||
"""
|
||||
if agent.get("session_id"):
|
||||
return _resume_headless(agent, answer, worker_id or f"cli-{os.getpid()}")
|
||||
session = tmux.session_name(agent["project_id"], agent["name"])
|
||||
if not tmux.has_session(session):
|
||||
return False, f"no live session '{session}' to resume"
|
||||
tmux.send_keys(session, answer)
|
||||
return True, f"answer delivered to session '{session}'"
|
||||
|
||||
|
||||
def _resume_headless(agent: dict, answer: str, worker_id: str) -> tuple[bool, str]:
|
||||
"""Launch a ``--resume`` run for a headless agent, materializing its session first.
|
||||
|
||||
Refuses while a run is still live (two concurrent processes on one session would
|
||||
corrupt it). When neither this worker nor the DB has the transcript — the owning
|
||||
worker died before its first archive — falls back to a *fresh* session whose prompt
|
||||
re-injects context from the DB (checkmark + open question + answer), recorded as a
|
||||
``worker`` event so the UI shows the degraded continuity.
|
||||
The seam the API's ``/resume`` route calls (and the one tests mock). The session
|
||||
transcript is materialized from the DB archive first, so ANY worker can serve the
|
||||
resume. Refuses while a run is still live (two concurrent processes on one session
|
||||
would corrupt it). When no transcript survives anywhere — the owning worker died
|
||||
before its first archive, or the row predates the headless runner — falls back to a
|
||||
*fresh* session whose prompt re-injects context from the DB (checkmark + open
|
||||
question + answer), recorded as a ``worker`` event so the UI shows the degraded
|
||||
continuity.
|
||||
"""
|
||||
worker_id = worker_id or f"cli-{os.getpid()}"
|
||||
with connection() as conn:
|
||||
run = repo.get_latest_run(conn, agent["id"])
|
||||
if run is not None and run["status"] == "running":
|
||||
@@ -289,7 +246,7 @@ def _resume_headless(agent: dict, answer: str, worker_id: str) -> tuple[bool, st
|
||||
return False, f"project '{agent['project_id']}' not registered"
|
||||
|
||||
working_dir = agent["working_dir"]
|
||||
settings_path = settings_gen.write_settings(working_dir, headless=True)
|
||||
settings_path = settings_gen.write_settings(working_dir)
|
||||
try:
|
||||
token = None
|
||||
with connection() as conn:
|
||||
@@ -298,6 +255,10 @@ def _resume_headless(agent: dict, answer: str, worker_id: str) -> tuple[bool, st
|
||||
return False, str(exc)
|
||||
env = _agent_env(project, agent, token)
|
||||
|
||||
if not agent.get("session_id"):
|
||||
# Pre-headless agent row (or a spawn that never launched): nothing to --resume.
|
||||
return _resume_reinjected(agent, answer, settings_path, env, worker_id)
|
||||
|
||||
transcript = headless.session_dir(working_dir) / f"{agent['session_id']}.jsonl"
|
||||
if archive is not None:
|
||||
try:
|
||||
|
||||
@@ -1,7 +1,8 @@
|
||||
"""Thin tmux wrapper — the single mock seam for spawning.
|
||||
"""Thin tmux wrapper — now used ONLY by the interactive ``/login`` flow.
|
||||
|
||||
Every tmux/claude invocation goes through these functions so tests can substitute a
|
||||
fake and never touch a real tmux server or ``claude`` binary.
|
||||
Agent runs are headless (``control.headless``); the one thing that still genuinely
|
||||
needs a TTY is driving claude's ``/login`` OAuth screens. Everything here goes through
|
||||
subprocess so the login tests can substitute a fake and never touch a real tmux server.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -11,14 +12,6 @@ import subprocess
|
||||
from ..config import get_settings
|
||||
|
||||
|
||||
def session_name(project_id: str, agent_name: str) -> str:
|
||||
"""``project__agent`` with tmux-illegal characters sanitized (README 3.4)."""
|
||||
safe = f"{project_id}__{agent_name}"
|
||||
for ch in (".", ":", " "):
|
||||
safe = safe.replace(ch, "-")
|
||||
return safe
|
||||
|
||||
|
||||
def new_session(
|
||||
name: str,
|
||||
cwd: str,
|
||||
@@ -57,18 +50,6 @@ def has_session(name: str) -> bool:
|
||||
return result.returncode == 0
|
||||
|
||||
|
||||
def list_sessions() -> list[str]:
|
||||
tmux = get_settings().tmux_bin
|
||||
result = subprocess.run(
|
||||
[tmux, "list-sessions", "-F", "#{session_name}"],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
)
|
||||
if result.returncode != 0:
|
||||
return []
|
||||
return [line for line in result.stdout.splitlines() if line]
|
||||
|
||||
|
||||
def kill_session(name: str) -> None:
|
||||
tmux = get_settings().tmux_bin
|
||||
subprocess.run([tmux, "kill-session", "-t", name], check=True)
|
||||
|
||||
@@ -1,11 +1,13 @@
|
||||
"""The control-container worker: executes commands the API enqueues.
|
||||
"""The control/worker container: executes commands the API enqueues and supervises
|
||||
headless claude runs.
|
||||
|
||||
The API (in its own container) has no ``git``/``tmux``/``claude`` and does not own the
|
||||
tmux sessions, so it cannot run control actions directly. Instead it writes a ``queued``
|
||||
row to the ``commands`` table; this worker — running in the control container — claims each
|
||||
row, dispatches it to the *same* control functions the CLI uses (``spawn``/``poller``/
|
||||
``skills_gen``/``repo.record_approval``), and writes the result or error back. It also runs
|
||||
the periodic CI sweep, subsuming the old ``poll-ci --watch`` loop.
|
||||
The API (in its own container) has no ``git``/``claude``, so it cannot run control
|
||||
actions directly. Instead it writes a ``queued`` row to the ``commands`` table; any
|
||||
worker claims each row (multi-worker safe — ``FOR UPDATE SKIP LOCKED`` + slot-aware
|
||||
claim filters), dispatches it to the *same* control functions the CLI uses (``spawn``/
|
||||
``poller``/``skills_gen``/``repo.record_approval``), and writes the result or error
|
||||
back. It also heartbeats + reaps dead workers' runs, syncs claude credentials, and runs
|
||||
the periodic CI sweep.
|
||||
|
||||
Every command runs in isolation: one bad command is recorded as ``failed`` and never stops
|
||||
the loop. ``execute_command`` is the pure dispatch seam (given a claimed command dict,
|
||||
@@ -23,7 +25,7 @@ from datetime import UTC, datetime, timedelta
|
||||
from ..config import get_settings
|
||||
from ..db import repository as repo
|
||||
from ..db.engine import connection
|
||||
from . import credsync, gitops, login, poller, reposync, skills_gen, spawn, tmux
|
||||
from . import credsync, gitops, login, poller, reposync, skills_gen, spawn
|
||||
|
||||
# Command types that launch a claude run and therefore need a free slot on this worker.
|
||||
# A worker with all slots busy leaves these queued for a less-loaded worker to claim.
|
||||
@@ -394,43 +396,6 @@ def fire_due_schedules(now: datetime | None = None) -> int:
|
||||
return fired
|
||||
|
||||
|
||||
# How many trailing pane lines to snapshot — enough to show the current screen (a menu, a
|
||||
# prompt, the tail of the last command) without bloating the row.
|
||||
_PANE_TAIL_LINES = 40
|
||||
|
||||
|
||||
def _pane_tail(pane: str, lines: int = _PANE_TAIL_LINES) -> str:
|
||||
"""The last ``lines`` of a captured pane, trailing blank lines trimmed so an idle
|
||||
screen doesn't store as a wall of whitespace."""
|
||||
rows = (pane or "").splitlines()
|
||||
while rows and not rows[-1].strip():
|
||||
rows.pop()
|
||||
return "\n".join(rows[-lines:])
|
||||
|
||||
|
||||
def capture_agent_output() -> int:
|
||||
"""Snapshot each working agent's live tmux pane tail into the DB.
|
||||
|
||||
The tmux socket lives only in the control container, so this is the one channel the
|
||||
API/UI have onto what a running — or wedged — agent is actually doing: an agent stuck
|
||||
on claude's first-run theme picker surfaces as that screen instead of a misleading
|
||||
green 'working'. A missing session is skipped (its process is gone). Returns the count
|
||||
updated.
|
||||
"""
|
||||
with connection() as conn:
|
||||
working = repo.list_agents_by_status(conn, "working")
|
||||
updated = 0
|
||||
for agent in working:
|
||||
session = tmux.session_name(agent["project_id"], agent["name"])
|
||||
if not tmux.has_session(session):
|
||||
continue
|
||||
tail = _pane_tail(tmux.capture_pane(session))
|
||||
with connection() as conn:
|
||||
repo.update_agent_output(conn, agent["id"], tail)
|
||||
updated += 1
|
||||
return updated
|
||||
|
||||
|
||||
def drain(worker_id: str, limit: int | None = None) -> int:
|
||||
"""Claim and run queued commands until the queue is empty (or ``limit`` reached).
|
||||
|
||||
@@ -454,15 +419,12 @@ def drain(worker_id: str, limit: int | None = None) -> int:
|
||||
def _full_slot_exclusions(worker_id: str) -> tuple[str, ...]:
|
||||
"""Command types this worker must not claim right now.
|
||||
|
||||
Headless runs are supervised in-process, so a worker at ``max_concurrent_runs`` skips
|
||||
Runs are supervised in-process, so a worker at ``max_concurrent_runs`` skips
|
||||
claiming run-starting commands — they stay queued for a worker with a free slot. Slot
|
||||
accounting is DB-driven (this worker's ``running`` runs), so it needs no in-memory
|
||||
registry and survives restarts (a fresh process gets a fresh worker id; stale rows
|
||||
belong to the old id and are the reaper's problem). Tmux runs are fire-and-forget and
|
||||
never consume a slot.
|
||||
belong to the old id and are the reaper's problem).
|
||||
"""
|
||||
if get_settings().runner != "headless":
|
||||
return ()
|
||||
with connection() as conn:
|
||||
active = len(repo.list_running_runs(conn, worker_id=worker_id))
|
||||
if active >= get_settings().max_concurrent_runs:
|
||||
@@ -474,20 +436,18 @@ def run(
|
||||
worker_id: str | None = None,
|
||||
poll_interval: float = 2.0,
|
||||
ci_interval: float = 30.0,
|
||||
capture_interval: float = 2.0,
|
||||
credsync_interval: float = 30.0,
|
||||
reap_interval: float = 15.0,
|
||||
iterations: int | None = None,
|
||||
) -> None:
|
||||
"""The control-container main loop: drain the command queue, snapshot live agent
|
||||
output, sync claude credentials, heartbeat + reap dead workers, and sweep CI.
|
||||
"""The control-container main loop: drain the command queue, sync claude
|
||||
credentials, heartbeat + reap dead workers, and sweep CI.
|
||||
|
||||
``iterations`` bounds the loop for tests; production runs unbounded. Sleeps
|
||||
``poll_interval`` only when a pass found no commands, so bursts drain promptly.
|
||||
"""
|
||||
worker_id = worker_id or make_worker_id()
|
||||
last_ci = 0.0
|
||||
last_capture = 0.0
|
||||
last_credsync = 0.0
|
||||
last_reap = 0.0
|
||||
count = 0
|
||||
@@ -508,12 +468,6 @@ def run(
|
||||
pass
|
||||
did_work = drain(worker_id) > 0
|
||||
now = time.monotonic()
|
||||
if capture_interval > 0 and now - last_capture >= capture_interval:
|
||||
try:
|
||||
capture_agent_output()
|
||||
except Exception: # noqa: BLE001 - a capture hiccup must not kill the worker
|
||||
pass
|
||||
last_capture = now
|
||||
if credsync_interval > 0 and (last_credsync == 0.0 or now - last_credsync >= credsync_interval):
|
||||
# First pass runs immediately: a fresh worker container must materialize the
|
||||
# claude credentials before it claims its first spawn.
|
||||
|
||||
+31
-6
@@ -1,7 +1,8 @@
|
||||
"""Shared fixtures. Everything runs on a fresh SQLite file per test, materialized via
|
||||
a *real* ``alembic upgrade head`` — so the migration path itself is under test, not
|
||||
just ``create_all``. No live claude/tmux/mise is ever touched: the three seams
|
||||
(``control.tmux``, ``hooks.verify``, ``control.spawn.resume``) are faked.
|
||||
just ``create_all``. No live claude/tmux/mise is ever touched: the seams
|
||||
(``control.headless.launch`` for runs, ``control.tmux`` for the login flow,
|
||||
``hooks.verify``, ``control.spawn.resume``) are faked.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -111,20 +112,44 @@ def fake_tmux(monkeypatch):
|
||||
def send_enter(name):
|
||||
calls["send_enter"].append({"name": name})
|
||||
|
||||
def list_sessions():
|
||||
return list(live)
|
||||
|
||||
monkeypatch.setattr(tmux, "new_session", new_session)
|
||||
monkeypatch.setattr(tmux, "has_session", has_session)
|
||||
monkeypatch.setattr(tmux, "kill_session", kill_session)
|
||||
monkeypatch.setattr(tmux, "send_keys", send_keys)
|
||||
monkeypatch.setattr(tmux, "send_text", send_text)
|
||||
monkeypatch.setattr(tmux, "send_enter", send_enter)
|
||||
monkeypatch.setattr(tmux, "list_sessions", list_sessions)
|
||||
|
||||
return {"calls": calls, "live": live}
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def fake_launch(monkeypatch):
|
||||
"""Record ``headless.launch`` calls instead of spawning a claude subprocess.
|
||||
|
||||
Mirrors the real launch's DB side effects (run row + agent session/worker) so kill/
|
||||
resume logic downstream of a fake spawn behaves like production, minus the process.
|
||||
"""
|
||||
from handler.control import headless
|
||||
from handler.db import repository as repo
|
||||
from handler.db.engine import connection
|
||||
|
||||
calls: list[dict] = []
|
||||
|
||||
def launch(agent, *, kind, prompt, settings_path, env, worker_id, on_exit=None):
|
||||
session_id = agent.get("session_id") if kind == "resume" else f"fake-sid-{len(calls) + 1}"
|
||||
with connection() as conn:
|
||||
run = repo.create_run(conn, agent["id"], session_id, worker_id, kind)
|
||||
repo.set_agent_session(conn, agent["id"], session_id, worker_id)
|
||||
calls.append(
|
||||
{"agent": agent, "kind": kind, "prompt": prompt, "settings_path": settings_path,
|
||||
"env": env, "worker_id": worker_id, "run": run}
|
||||
)
|
||||
return run
|
||||
|
||||
monkeypatch.setattr(headless, "launch", launch)
|
||||
return calls
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def fake_gitops(monkeypatch):
|
||||
"""Fake the git seam: record config/add/commit, return a controllable branch/sha."""
|
||||
|
||||
+75
-27
@@ -1,4 +1,7 @@
|
||||
"""Control-layer spawn: the hard test-task gate, settings generation, identity env."""
|
||||
"""Control-layer spawn: the hard test-task gate, settings generation, identity env.
|
||||
|
||||
Spawns go through the ``fake_launch`` seam (conftest) — the headless analogue of the old
|
||||
fake tmux: it records the launch and mirrors its DB side effects, no subprocess."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
@@ -24,36 +27,48 @@ def _write_mise(root, with_test=True):
|
||||
(root / ".mise.toml").write_text(body)
|
||||
|
||||
|
||||
def test_spawn_refuses_without_test_task(env, fake_tmux):
|
||||
def test_spawn_refuses_without_test_task(env, fake_launch):
|
||||
root = env["tmp"] / "proj"
|
||||
_write_mise(root, with_test=False)
|
||||
_register_project(root)
|
||||
with pytest.raises(spawn.SpawnError, match="no \\[tasks.test\\]"):
|
||||
spawn.spawn("proj", "api")
|
||||
assert fake_tmux["calls"]["new_session"] == []
|
||||
spawn.spawn("proj", "api", task="do it")
|
||||
assert fake_launch == []
|
||||
|
||||
|
||||
def test_spawn_refuses_without_mise_file(env, fake_tmux):
|
||||
def test_spawn_refuses_without_mise_file(env, fake_launch):
|
||||
root = env["tmp"] / "proj"
|
||||
root.mkdir(parents=True, exist_ok=True)
|
||||
_register_project(root)
|
||||
with pytest.raises(spawn.SpawnError, match="no mise config"):
|
||||
spawn.spawn("proj", "api", task="do it")
|
||||
|
||||
|
||||
def test_spawn_refuses_without_task(env, fake_launch):
|
||||
root = env["tmp"] / "proj"
|
||||
_write_mise(root)
|
||||
_register_project(root)
|
||||
with pytest.raises(spawn.SpawnError, match="requires a task"):
|
||||
spawn.spawn("proj", "api")
|
||||
# Fail-fast: no orphaned agent row behind the refused spawn.
|
||||
with get_engine().begin() as conn:
|
||||
assert repo.get_agent_by_name(conn, "proj", "api") is None
|
||||
assert fake_launch == []
|
||||
|
||||
|
||||
def test_spawn_accepts_dotless_mise_toml(env, fake_tmux):
|
||||
def test_spawn_accepts_dotless_mise_toml(env, fake_launch):
|
||||
# mise also reads `mise.toml` (no leading dot); the gate must honor it too.
|
||||
root = env["tmp"] / "proj"
|
||||
root.mkdir(parents=True, exist_ok=True)
|
||||
(root / "mise.toml").write_text("[tasks.test]\nrun = 'pytest'\n")
|
||||
_register_project(root)
|
||||
|
||||
agent = spawn.spawn("proj", "api")
|
||||
agent = spawn.spawn("proj", "api", task="do it")
|
||||
with get_engine().begin() as conn:
|
||||
assert repo.get_agent_by_name(conn, "proj", "api")["id"] == agent["id"]
|
||||
|
||||
|
||||
def test_spawn_creates_agent_settings_and_session(env, fake_tmux):
|
||||
def test_spawn_creates_agent_settings_and_run(env, fake_launch):
|
||||
root = env["tmp"] / "proj"
|
||||
_write_mise(root, with_test=True)
|
||||
_register_project(root)
|
||||
@@ -64,66 +79,99 @@ def test_spawn_creates_agent_settings_and_session(env, fake_tmux):
|
||||
with get_engine().begin() as conn:
|
||||
assert repo.get_agent_by_name(conn, "proj", "api")["id"] == agent["id"]
|
||||
|
||||
# settings.json wires all four hook events.
|
||||
# settings.json wires all four hook events AND the headless permission allowlist
|
||||
# (claude -p auto-denies anything that would prompt; the allowlist is what lets
|
||||
# normal work proceed — the hooks stay the hard gate).
|
||||
settings = json.loads((root / ".claude" / "settings.json").read_text())
|
||||
assert set(settings["hooks"]) == {"Stop", "SessionEnd", "PreToolUse", "Notification"}
|
||||
pre = settings["hooks"]["PreToolUse"][0]
|
||||
assert pre["matcher"] == "AskUserQuestion|Bash"
|
||||
assert "handler.hooks pre_tool_use" in pre["hooks"][0]["command"]
|
||||
assert settings["permissions"]["defaultMode"] == "acceptEdits"
|
||||
assert "Bash(git *)" in settings["permissions"]["allow"]
|
||||
|
||||
# tmux session named project__agent, with identity + DATABASE_URL in env.
|
||||
call = fake_tmux["calls"]["new_session"][0]
|
||||
assert call["name"] == "proj__api"
|
||||
# A headless run launched with identity + DATABASE_URL in env and the task as prompt.
|
||||
call = fake_launch[0]
|
||||
assert call["kind"] == "spawn"
|
||||
assert call["prompt"] == "build the thing"
|
||||
assert call["env"]["HANDLER_PROJECT_ID"] == "proj"
|
||||
assert call["env"]["HANDLER_AGENT_NAME"] == "api"
|
||||
assert call["env"]["HANDLER_AGENT_ID"] == str(agent["id"])
|
||||
assert call["env"]["DATABASE_URL"] == env["url"]
|
||||
# The run row + session id landed on the agent.
|
||||
with get_engine().begin() as conn:
|
||||
row = repo.get_agent_by_name(conn, "proj", "api")
|
||||
assert row["session_id"] == call["run"]["session_id"]
|
||||
assert repo.get_latest_run(conn, row["id"])["kind"] == "spawn"
|
||||
|
||||
|
||||
def test_spawn_mise_init_skips_test_gate_and_marks_env(env, fake_tmux):
|
||||
def test_spawn_mise_init_skips_test_gate_and_marks_env(env, fake_launch):
|
||||
# A repo with no .mise.toml at all: the normal gate would refuse, but the mise-init
|
||||
# bootstrap agent must launch anyway (creating that file is its whole job).
|
||||
root = env["tmp"] / "proj"
|
||||
root.mkdir(parents=True, exist_ok=True)
|
||||
_register_project(root)
|
||||
|
||||
agent = spawn.spawn("proj", "mise-init", require_tests=False, mise_init=True)
|
||||
agent = spawn.spawn(
|
||||
"proj", "mise-init", task="write the mise config", require_tests=False, mise_init=True
|
||||
)
|
||||
|
||||
with get_engine().begin() as conn:
|
||||
assert repo.get_agent_by_name(conn, "proj", "mise-init")["id"] == agent["id"]
|
||||
# The launched session carries HANDLER_MISE_INIT so its hooks enforce commit + push.
|
||||
call = fake_tmux["calls"]["new_session"][0]
|
||||
assert call["env"]["HANDLER_MISE_INIT"] == "1"
|
||||
# The launched run carries HANDLER_MISE_INIT so its hooks enforce commit + push.
|
||||
assert fake_launch[0]["env"]["HANDLER_MISE_INIT"] == "1"
|
||||
|
||||
|
||||
def test_spawn_still_gates_without_mise_init_flag(env, fake_tmux):
|
||||
def test_spawn_still_gates_without_mise_init_flag(env, fake_launch):
|
||||
root = env["tmp"] / "proj"
|
||||
root.mkdir(parents=True, exist_ok=True)
|
||||
_register_project(root)
|
||||
# require_tests defaults on, so a normal spawn against a mise-less repo still refuses.
|
||||
with pytest.raises(spawn.SpawnError, match="no mise config"):
|
||||
spawn.spawn("proj", "api")
|
||||
assert fake_tmux["calls"]["new_session"] == []
|
||||
spawn.spawn("proj", "api", task="do it")
|
||||
assert fake_launch == []
|
||||
|
||||
|
||||
def test_kill_sets_done_and_kills_session(env, fake_tmux):
|
||||
def test_kill_cancels_run_and_sets_done(env, fake_launch):
|
||||
root = env["tmp"] / "proj"
|
||||
_write_mise(root, with_test=True)
|
||||
_register_project(root)
|
||||
spawn.spawn("proj", "api")
|
||||
spawn.spawn("proj", "api", task="do it")
|
||||
|
||||
spawn.kill("proj", "api")
|
||||
assert "proj__api" in fake_tmux["calls"]["kill_session"]
|
||||
with get_engine().begin() as conn:
|
||||
assert repo.get_agent_by_name(conn, "proj", "api")["status"] == "done"
|
||||
agent = repo.get_agent_by_name(conn, "proj", "api")
|
||||
assert agent["status"] == "done"
|
||||
# The running run was flagged; the owning supervisor terminates its own child.
|
||||
assert repo.get_latest_run(conn, agent["id"])["cancel_requested"] is True
|
||||
|
||||
|
||||
def test_resume_sends_answer_to_live_session(env, fake_tmux):
|
||||
def test_resume_reinjects_when_no_transcript(env, fake_launch):
|
||||
"""A resume with no archive and no local transcript degrades to a fresh run whose
|
||||
prompt carries the operator's answer (context re-injection)."""
|
||||
root = env["tmp"] / "proj"
|
||||
_write_mise(root, with_test=True)
|
||||
_register_project(root)
|
||||
agent = spawn.spawn("proj", "api")
|
||||
spawn.spawn("proj", "api", task="do it")
|
||||
with get_engine().begin() as conn:
|
||||
agent = repo.get_agent_by_name(conn, "proj", "api")
|
||||
repo.finish_run(conn, repo.get_latest_run(conn, agent["id"])["id"], "completed")
|
||||
|
||||
ok, detail = spawn.resume(agent, "use Postgres")
|
||||
assert ok is True
|
||||
assert fake_tmux["calls"]["send_keys"][0] == {"name": "proj__api", "keys": "use Postgres"}
|
||||
assert "re-injected" in detail
|
||||
assert fake_launch[-1]["kind"] == "spawn"
|
||||
assert "use Postgres" in fake_launch[-1]["prompt"]
|
||||
|
||||
|
||||
def test_resume_refused_while_run_live(env, fake_launch):
|
||||
root = env["tmp"] / "proj"
|
||||
_write_mise(root, with_test=True)
|
||||
_register_project(root)
|
||||
spawn.spawn("proj", "api", task="do it") # fake run stays 'running'
|
||||
with get_engine().begin() as conn:
|
||||
agent = repo.get_agent_by_name(conn, "proj", "api")
|
||||
|
||||
ok, detail = spawn.resume(agent, "answer")
|
||||
assert ok is False
|
||||
assert "live run" in detail
|
||||
|
||||
@@ -19,15 +19,15 @@ def _register(root, **kw):
|
||||
repo.create_project(conn, "proj", str(root), **kw)
|
||||
|
||||
|
||||
def test_spawn_injects_credentials_and_installs_helper(env, fake_tmux, fake_gitops, monkeypatch):
|
||||
def test_spawn_injects_credentials_and_installs_helper(env, fake_launch, fake_gitops, monkeypatch):
|
||||
monkeypatch.setenv("PROJ_TOKEN", "s3cret")
|
||||
root = env["tmp"] / "proj"
|
||||
_write_mise(root)
|
||||
_register(root, git_remote="https://github.com/me/proj.git", credential_ref="env:PROJ_TOKEN")
|
||||
|
||||
spawn.spawn("proj", "junior", role="junior")
|
||||
spawn.spawn("proj", "junior", role="junior", task="do it")
|
||||
|
||||
call = fake_tmux["calls"]["new_session"][0]
|
||||
call = fake_launch[0]
|
||||
# Token injected under the generic + host-specific names, never the raw ref stored.
|
||||
assert call["env"]["FORGE_TOKEN"] == "s3cret"
|
||||
assert call["env"]["GITHUB_TOKEN"] == "s3cret"
|
||||
@@ -38,43 +38,43 @@ def test_spawn_injects_credentials_and_installs_helper(env, fake_tmux, fake_gito
|
||||
assert "$FORGE_TOKEN" in helper[0]["value"]
|
||||
|
||||
|
||||
def test_spawn_ssh_remote_installs_no_https_helper(env, fake_tmux, fake_gitops, monkeypatch):
|
||||
def test_spawn_ssh_remote_installs_no_https_helper(env, fake_launch, fake_gitops, monkeypatch):
|
||||
monkeypatch.setenv("PROJ_TOKEN", "s3cret")
|
||||
root = env["tmp"] / "proj"
|
||||
_write_mise(root)
|
||||
_register(root, git_remote="git@github.com:me/proj.git", credential_ref="env:PROJ_TOKEN")
|
||||
spawn.spawn("proj", "junior", role="junior")
|
||||
spawn.spawn("proj", "junior", role="junior", task="do it")
|
||||
# ssh remote -> token still injected, but no HTTPS credential helper installed.
|
||||
assert fake_tmux["calls"]["new_session"][0]["env"]["GITHUB_TOKEN"] == "s3cret"
|
||||
assert fake_launch[0]["env"]["GITHUB_TOKEN"] == "s3cret"
|
||||
assert fake_gitops["config"] == []
|
||||
|
||||
|
||||
def test_spawn_fails_fast_on_broken_credential_ref(env, fake_tmux, fake_gitops, monkeypatch):
|
||||
def test_spawn_fails_fast_on_broken_credential_ref(env, fake_launch, fake_gitops, monkeypatch):
|
||||
monkeypatch.delenv("ABSENT_TOKEN", raising=False)
|
||||
root = env["tmp"] / "proj"
|
||||
_write_mise(root)
|
||||
_register(root, credential_ref="env:ABSENT_TOKEN")
|
||||
|
||||
with pytest.raises(spawn.SpawnError, match="not set"):
|
||||
spawn.spawn("proj", "junior", role="junior")
|
||||
spawn.spawn("proj", "junior", role="junior", task="do it")
|
||||
# No agent row and no session left behind by the failed spawn.
|
||||
with get_engine().begin() as conn:
|
||||
assert repo.get_agent_by_name(conn, "proj", "junior") is None
|
||||
assert fake_tmux["calls"]["new_session"] == []
|
||||
assert fake_launch == []
|
||||
|
||||
|
||||
def test_spawn_without_credential_ref_injects_no_token(env, fake_tmux, fake_gitops):
|
||||
def test_spawn_without_credential_ref_injects_no_token(env, fake_launch, fake_gitops):
|
||||
root = env["tmp"] / "proj"
|
||||
_write_mise(root)
|
||||
_register(root)
|
||||
spawn.spawn("proj", "api")
|
||||
call = fake_tmux["calls"]["new_session"][0]
|
||||
spawn.spawn("proj", "api", task="do it")
|
||||
call = fake_launch[0]
|
||||
assert "FORGE_TOKEN" not in call["env"]
|
||||
# No token -> no credential helper installed.
|
||||
assert fake_gitops["config"] == []
|
||||
|
||||
|
||||
def test_spawn_reports_forge_version_mismatch(env, fake_tmux, fake_gitops, fake_forge, monkeypatch):
|
||||
def test_spawn_reports_forge_version_mismatch(env, fake_launch, fake_gitops, fake_forge, monkeypatch):
|
||||
monkeypatch.setenv("FORGE_VERSION", "9.9.9")
|
||||
from handler import config
|
||||
from handler.db import engine
|
||||
@@ -88,5 +88,5 @@ def test_spawn_reports_forge_version_mismatch(env, fake_tmux, fake_gitops, fake_
|
||||
fake_forge["version_ok"] = False
|
||||
fake_forge["version_out"] = "forge 1.2.3"
|
||||
|
||||
agent = spawn.spawn("proj", "api")
|
||||
agent = spawn.spawn("proj", "api", task="do it")
|
||||
assert "9.9.9" in agent["forge_note"]
|
||||
|
||||
@@ -101,10 +101,3 @@ def test_disabled_without_secret_key(env, tmp_path):
|
||||
_write_local_credentials(tmp_path)
|
||||
assert credsync.upload() is False
|
||||
assert credsync.refresh() is None
|
||||
|
||||
|
||||
def test_note_local_write_suppresses_upload(secret_env, tmp_path):
|
||||
_write_local_credentials(tmp_path)
|
||||
credsync.note_local_write()
|
||||
# The deliberate local write (e.g. ensure_onboarded at spawn) is not re-published.
|
||||
assert credsync.refresh() is None
|
||||
|
||||
@@ -26,7 +26,6 @@ def headless_env(env, monkeypatch):
|
||||
from handler import config
|
||||
|
||||
monkeypatch.setenv("CLAUDE_BIN", FAKE_CLAUDE)
|
||||
monkeypatch.setenv("RUNNER", "headless")
|
||||
config.get_settings.cache_clear()
|
||||
yield env
|
||||
config.get_settings.cache_clear()
|
||||
@@ -250,16 +249,13 @@ def test_resume_refused_while_run_live(headless_env, tmp_path, monkeypatch):
|
||||
_wait_for(_finished_run(run["id"]), timeout=30.0)
|
||||
|
||||
|
||||
def test_headless_settings_include_permissions(headless_env, tmp_path):
|
||||
path = settings_gen.write_settings(str(tmp_path / "wd"), headless=True)
|
||||
def test_settings_include_permissions_and_hooks(headless_env, tmp_path):
|
||||
path = settings_gen.write_settings(str(tmp_path / "wd"))
|
||||
data = json.loads(Path(path).read_text())
|
||||
assert data["permissions"]["defaultMode"] == "acceptEdits"
|
||||
assert "Bash(git *)" in data["permissions"]["allow"]
|
||||
assert "hooks" in data # the hard gate is untouched
|
||||
|
||||
tmux_path = settings_gen.write_settings(str(tmp_path / "wd2"))
|
||||
assert "permissions" not in json.loads(Path(tmux_path).read_text())
|
||||
|
||||
|
||||
def test_api_rejects_empty_task_headless_spawn(headless_env, client, auth):
|
||||
client.post(
|
||||
|
||||
@@ -1,13 +1,32 @@
|
||||
"""End-to-end web management: the dashboard's HTTP calls -> command queue -> worker ->
|
||||
real ``spawn.spawn`` -> tmux seam. Proves the full container-split flow works with only the
|
||||
tmux/claude boundary faked, not the control layer itself."""
|
||||
real ``spawn.spawn`` -> a real headless subprocess (the fake claude binary). Proves the
|
||||
full container-split flow works with only the claude binary faked, not the control
|
||||
layer: events stream into the DB, the run reconciles, kill cancels."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import time
|
||||
from pathlib import Path
|
||||
|
||||
import pytest
|
||||
|
||||
from handler.control import worker
|
||||
from handler.db import repository as repo
|
||||
from handler.db.engine import get_engine
|
||||
|
||||
REPO_ROOT = Path(__file__).resolve().parents[1]
|
||||
FAKE_CLAUDE = str(REPO_ROOT / "tests" / "fixtures" / "fake_claude.py")
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def headless_env(env, monkeypatch):
|
||||
from handler import config
|
||||
|
||||
monkeypatch.setenv("CLAUDE_BIN", FAKE_CLAUDE)
|
||||
config.get_settings.cache_clear()
|
||||
yield env
|
||||
config.get_settings.cache_clear()
|
||||
|
||||
|
||||
def _spawnable_project(root):
|
||||
root.mkdir(parents=True, exist_ok=True)
|
||||
@@ -16,10 +35,21 @@ def _spawnable_project(root):
|
||||
repo.create_project(conn, "proj", str(root))
|
||||
|
||||
|
||||
def test_spawn_via_api_then_worker_creates_agent_and_session(client, auth, env, fake_tmux):
|
||||
_spawnable_project(env["tmp"] / "proj")
|
||||
def _wait(predicate, timeout=20.0):
|
||||
deadline = time.monotonic() + timeout
|
||||
while time.monotonic() < deadline:
|
||||
result = predicate()
|
||||
if result:
|
||||
return result
|
||||
time.sleep(0.1)
|
||||
return None
|
||||
|
||||
# 1. The dashboard enqueues a spawn (202 + a queued command).
|
||||
|
||||
def test_spawn_via_api_then_worker_runs_headless_claude(client, auth, headless_env):
|
||||
_spawnable_project(headless_env["tmp"] / "proj")
|
||||
|
||||
# 1. The dashboard enqueues a spawn (202 + a queued command). A task is mandatory —
|
||||
# headless claude has no idle-REPL mode.
|
||||
r = client.post(
|
||||
"/projects/proj/agents/spawn",
|
||||
json={"name": "api", "task": "build the thing"},
|
||||
@@ -32,29 +62,62 @@ def test_spawn_via_api_then_worker_creates_agent_and_session(client, auth, env,
|
||||
# No agent yet — the worker hasn't run.
|
||||
assert client.get("/projects/proj/agents", headers=auth).json() == []
|
||||
|
||||
# 2. The control worker drains the queue (runs the real spawn.spawn).
|
||||
# 2. The control worker drains the queue (real spawn.spawn -> real subprocess).
|
||||
assert worker.drain("test-worker") == 1
|
||||
|
||||
# 3. The command is done and the agent + tmux session now exist.
|
||||
# 3. The command finished at launch (fire-and-forget)...
|
||||
got = client.get(f"/commands/{command_id}", headers=auth).json()
|
||||
assert got["status"] == "done"
|
||||
assert got["result"]["name"] == "api"
|
||||
|
||||
agents = client.get("/projects/proj/agents", headers=auth).json()
|
||||
assert [a["name"] for a in agents] == ["api"]
|
||||
assert fake_tmux["calls"]["new_session"][0]["name"] == "proj__api"
|
||||
# ...and the run's whole life shows up via the API: events stream in, the agent
|
||||
# reconciles to done, last_output is the assistant's text.
|
||||
def finished():
|
||||
agents = client.get("/projects/proj/agents", headers=auth).json()
|
||||
return agents if agents and agents[0]["status"] == "done" else None
|
||||
|
||||
agents = _wait(finished)
|
||||
assert agents is not None, "run never reconciled to done"
|
||||
agent = agents[0]
|
||||
assert agent["name"] == "api"
|
||||
assert agent["session_id"]
|
||||
assert agent["worker_id"] == "test-worker"
|
||||
assert agent["last_output"] == "working on: build the thing"
|
||||
|
||||
events = client.get("/projects/proj/agents/api/events", headers=auth).json()
|
||||
assert [e["type"] for e in events] == ["system", "assistant", "result"]
|
||||
|
||||
|
||||
def test_kill_via_api_then_worker(client, auth, env, fake_tmux):
|
||||
_spawnable_project(env["tmp"] / "proj")
|
||||
client.post("/projects/proj/agents/spawn", json={"name": "api"}, headers=auth)
|
||||
def test_spawn_without_task_is_rejected(client, auth, headless_env):
|
||||
_spawnable_project(headless_env["tmp"] / "proj")
|
||||
r = client.post("/projects/proj/agents/spawn", json={"name": "api"}, headers=auth)
|
||||
assert r.status_code == 400
|
||||
assert "task is required" in r.json()["detail"]
|
||||
|
||||
|
||||
def test_kill_via_api_then_worker(client, auth, headless_env, monkeypatch):
|
||||
monkeypatch.setenv("FAKE_CLAUDE_MODE", "hang")
|
||||
_spawnable_project(headless_env["tmp"] / "proj")
|
||||
client.post(
|
||||
"/projects/proj/agents/spawn", json={"name": "api", "task": "hang"}, headers=auth
|
||||
)
|
||||
worker.drain("w")
|
||||
|
||||
# The hanging run is live; kill flags it and the supervisor SIGTERMs its child.
|
||||
r = client.post("/projects/proj/agents/api/kill", headers=auth)
|
||||
assert r.status_code == 202
|
||||
worker.drain("w")
|
||||
|
||||
assert client.get(f"/commands/{r.json()['id']}", headers=auth).json()["status"] == "done"
|
||||
assert "proj__api" in fake_tmux["calls"]["kill_session"]
|
||||
with get_engine().begin() as conn:
|
||||
assert repo.get_agent_by_name(conn, "proj", "api")["status"] == "done"
|
||||
agent = repo.get_agent_by_name(conn, "proj", "api")
|
||||
assert agent["status"] == "done"
|
||||
|
||||
def canceled():
|
||||
with get_engine().begin() as conn:
|
||||
run = repo.get_latest_run(conn, agent["id"])
|
||||
return run if run["status"] != "running" else None
|
||||
|
||||
run = _wait(canceled, timeout=30.0)
|
||||
assert run is not None, "kill never terminated the hanging run"
|
||||
assert run["status"] == "canceled"
|
||||
|
||||
@@ -229,45 +229,6 @@ def test_bad_command_is_recorded_failed_not_raised(env):
|
||||
assert "agent name" in failed["error"]
|
||||
|
||||
|
||||
def test_capture_agent_output_snapshots_working_agents(env, monkeypatch):
|
||||
with get_engine().begin() as conn:
|
||||
repo.create_project(conn, "p", "/tmp/p")
|
||||
agent = repo.create_agent(conn, "p", "api", "/tmp/p/api", status="working")
|
||||
|
||||
monkeypatch.setattr(worker.tmux, "has_session", lambda name: True)
|
||||
monkeypatch.setattr(
|
||||
worker.tmux, "capture_pane", lambda name, escapes=False: "boot\nTheme picker\n\n\n"
|
||||
)
|
||||
|
||||
assert worker.capture_agent_output() == 1
|
||||
with get_engine().begin() as conn:
|
||||
row = repo.get_agent_by_id(conn, agent["id"])
|
||||
# The tail is stored with trailing blank lines trimmed.
|
||||
assert row["last_output"] == "boot\nTheme picker"
|
||||
assert row["output_at"] is not None
|
||||
|
||||
|
||||
def test_capture_agent_output_skips_dead_sessions_and_nonworking(env, monkeypatch):
|
||||
with get_engine().begin() as conn:
|
||||
repo.create_project(conn, "p", "/tmp/p")
|
||||
repo.create_agent(conn, "p", "gone", "/tmp/p/gone", status="working")
|
||||
done = repo.create_agent(conn, "p", "done", "/tmp/p/done", status="done")
|
||||
|
||||
captured = []
|
||||
monkeypatch.setattr(worker.tmux, "has_session", lambda name: False)
|
||||
monkeypatch.setattr(
|
||||
worker.tmux,
|
||||
"capture_pane",
|
||||
lambda name, escapes=False: captured.append(name) or "x",
|
||||
)
|
||||
|
||||
# The working agent's session is dead (skipped); the done agent isn't queried at all.
|
||||
assert worker.capture_agent_output() == 0
|
||||
assert captured == []
|
||||
with get_engine().begin() as conn:
|
||||
assert repo.get_agent_by_id(conn, done["id"])["last_output"] is None
|
||||
|
||||
|
||||
def test_drain_processes_multiple_then_stops(env, monkeypatch):
|
||||
_seed_project()
|
||||
monkeypatch.setattr(poller, "sweep", lambda project_id=None: {"checked": 0})
|
||||
|
||||
@@ -16,7 +16,6 @@ from handler.db.engine import get_engine
|
||||
def headless_env(env, monkeypatch):
|
||||
from handler import config
|
||||
|
||||
monkeypatch.setenv("RUNNER", "headless")
|
||||
monkeypatch.setenv("MAX_CONCURRENT_RUNS", "2")
|
||||
config.get_settings.cache_clear()
|
||||
yield env
|
||||
@@ -75,7 +74,3 @@ def test_slot_frees_when_run_finishes(headless_env, monkeypatch):
|
||||
with get_engine().begin() as conn:
|
||||
repo.finish_run(conn, run1["id"], "completed", exit_code=0)
|
||||
assert worker._full_slot_exclusions("w1") == ()
|
||||
|
||||
|
||||
def test_tmux_runner_never_excludes(env):
|
||||
assert worker._full_slot_exclusions("w") == ()
|
||||
|
||||
Reference in New Issue
Block a user