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:
2026-07-21 23:20:39 -04:00
parent 6c2e73d4ec
commit 1517e4dca8
17 changed files with 337 additions and 354 deletions
+39 -12
View File
@@ -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
+4 -5
View File
@@ -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",
+3 -4
View File
@@ -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
View File
@@ -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)
+13 -16
View File
@@ -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()
+12 -13
View File
@@ -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
+37 -76
View File
@@ -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:
+4 -23
View File
@@ -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)
+14 -60
View File
@@ -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
View File
@@ -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
View File
@@ -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
+14 -14
View File
@@ -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"]
-7
View File
@@ -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
+2 -6
View File
@@ -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(
+78 -15
View File
@@ -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"
-39
View File
@@ -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})
-5
View File
@@ -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") == ()