Completes all 3 model-divergence enforcement subtasks:
- scripts/detect_models.py: probes opencode.json + localhost endpoints,
builds models.json with --json/--write/--force
- scripts/status.py: CONFLICT_MATRIX, --model flag, --transition --model,
--claim --model, --audit Category 6, model-divergence brake gate in
--check-gate, helpers for manifest loading and conflict checking
- scripts/loop-runner.py: _role_model() helper + {model} passed via extras
dict to _invoke_harness for implement, verify, orchestrate roles
- tests/test_model_divergence.py: 33 tests covering all enforcement layers
- Single-LLM mode: record model advisory, no conflict check
- Multi-LLM mode (2+ models): conflict matrix enforced at transition, claim,
and loop brake gate
- Project-level models.json preferred over global ~/.automaton/models.json
996 lines
35 KiB
Python
996 lines
35 KiB
Python
#!/usr/bin/env python3
|
|
"""Automaton loop runner -- per-tick engine.
|
|
|
|
Invoked by the OS scheduler unit (automaton-loop-tick.sh / .bat generated by
|
|
`status.py --install-schedule`), or manually, or in --mode daemon.
|
|
|
|
Contract (design/loops/technical.md section 7):
|
|
1. load .state.loop + loop.json
|
|
2. status.py --check-gate <name> --json; SKIP on not-ok with exit 0
|
|
3. find_work(work_source dispatch: single/audit/backlog)
|
|
4. ensure worktree (D2 -- per-loop git worktree at <loop>/worktree,
|
|
branch loop/<name>; falls back to project root on non-git or failure)
|
|
5. spawn Implement role via loop.json harness.command
|
|
6. spawn Verify role via harness.command; output is graded JSON verdict
|
|
7. parse verdict (accepts raw JSON, fenced JSON blocks, line comments)
|
|
parse failure -> halt verifier_failed, exit 0
|
|
8. cap score_history at brakes.score_plateau_window
|
|
9. spawn Orchestrate role (the orchestrator calls status.py itself; runner
|
|
does not parse orchestrator output)
|
|
10. write iteration_count++, last_tick_at, last_verdict atomically
|
|
10.5. GC outputs/ (retain last N tick groups per outputs.retention, v1.1)
|
|
11. append TICK line to .state.log
|
|
|
|
Idempotent in failure: anything that fails before step 10 leaves .state.loop
|
|
unchanged. Harness subprocesses are not owned; runner does not kill process
|
|
groups in v1.
|
|
|
|
Stdlib only; no new pip deps.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import argparse
|
|
import contextlib
|
|
import json
|
|
import math
|
|
import os
|
|
import re
|
|
import subprocess
|
|
import sys
|
|
import time
|
|
from datetime import datetime, timezone
|
|
from pathlib import Path
|
|
from typing import Optional, Sequence
|
|
|
|
AUTOMATON_DIR = Path.home() / ".automaton"
|
|
STATUS_SCRIPT = AUTOMATON_DIR / "scripts" / "status.py"
|
|
VRAM_SCRIPT = AUTOMATON_DIR / "scripts" / "vram_detect.py"
|
|
|
|
LOOP_STATE_FILE = ".state.loop"
|
|
LOOP_CONFIG_FILE = "loop.json"
|
|
LOOP_TICK_LOG_NAME = ".state.log"
|
|
LOOP_OUTPUTS_DIR = "outputs"
|
|
LOOP_WORKTREE_DIR = "worktree"
|
|
LOOP_STATE_SCHEMA_VERSION = 1
|
|
CONTEXT_FLOOR_KB = 16_000
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Small duplicated helpers (kept local rather than imported across scripts;
|
|
# see technical.md "no cross-script imports")
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
def _find_project_dir(project: Optional[str]) -> Path:
|
|
if project:
|
|
p = Path(project).resolve()
|
|
return p
|
|
cwd = Path.cwd().resolve()
|
|
if cwd == AUTOMATON_DIR:
|
|
return AUTOMATON_DIR
|
|
if (cwd / ".automaton").exists():
|
|
return cwd
|
|
if cwd.parent == AUTOMATON_DIR:
|
|
return AUTOMATON_DIR
|
|
print(f"ERROR: not in an automaton project directory (cwd={cwd}). "
|
|
f"Use --project to specify the project path.", file=sys.stderr)
|
|
sys.exit(1)
|
|
|
|
|
|
def _loops_dir(project_dir: Path) -> Path:
|
|
if project_dir == AUTOMATON_DIR:
|
|
return AUTOMATON_DIR / "loops"
|
|
return project_dir / ".automaton" / "loops"
|
|
|
|
|
|
def _loop_dir(name: str, project_dir: Path) -> Path:
|
|
return _loops_dir(project_dir) / name
|
|
|
|
|
|
def _read_state_loop(loop_path: Path) -> Optional[dict]:
|
|
f = loop_path / LOOP_STATE_FILE
|
|
if not f.exists():
|
|
return None
|
|
try:
|
|
return json.loads(f.read_text())
|
|
except (OSError, json.JSONDecodeError):
|
|
return None
|
|
|
|
|
|
def _write_state_loop(loop_path: Path, state: dict) -> None:
|
|
f = loop_path / LOOP_STATE_FILE
|
|
tmp = loop_path / ".state.loop.tmp"
|
|
tmp.write_text(json.dumps(state, indent=2, sort_keys=True) + "\n")
|
|
tmp.replace(f)
|
|
|
|
|
|
def _read_loop_config(loop_path: Path) -> Optional[dict]:
|
|
f = loop_path / LOOP_CONFIG_FILE
|
|
if not f.exists():
|
|
return None
|
|
try:
|
|
return json.loads(f.read_text())
|
|
except (OSError, json.JSONDecodeError):
|
|
return None
|
|
|
|
|
|
def _append_tick_log(loop_path: Path, line: str) -> None:
|
|
log = loop_path / LOOP_TICK_LOG_NAME
|
|
ts = datetime.now(timezone.utc).isoformat()
|
|
with log.open("a", encoding="utf-8") as fh:
|
|
fh.write(f"[{ts}] {line}\n")
|
|
|
|
|
|
def _halt_loop(loop_path: Path, state: dict, reason: str) -> None:
|
|
state["status"] = "halted"
|
|
state["halt_reason"] = reason
|
|
_write_state_loop(loop_path, state)
|
|
_append_tick_log(loop_path, f"HALT reason={reason}")
|
|
|
|
|
|
_LOOP_LOCK_ENV_BYPASS = "AUTOMATON_NO_LOOP_LOCK"
|
|
|
|
|
|
@contextlib.contextmanager
|
|
def _loop_lock(loop_path: Path, exclusive: bool = True):
|
|
"""Cross-process file lock on <loop_path>/.state.lock held for the full tick.
|
|
|
|
Serializes the runner's read-modify-write cycle on `.state.loop` against
|
|
concurrent ticks (two scheduler firings on the same loop) and concurrent
|
|
`status.py --pause-loop` / --approve --loop writes. Blocking acquire;
|
|
no timeout in v1.1 (operators notice a wedged tick via `--loop-list`
|
|
stale `last_tick_at`).
|
|
|
|
Per-loop granularity: lock file lives in the loop's own dir, not the
|
|
framework root. A lock on loop A's tick does not block loop B.
|
|
|
|
The runner holds this lock across `_gate` (subprocess), the harness
|
|
subprocess (Implement/Verify/Orchestrate), and the state write (step 10).
|
|
`_gate` spawns `status.py --check-gate` which would otherwise deadlock
|
|
waiting on the same flock; to avoid this the runner passes
|
|
$AUTOMATON_NO_LOOP_LOCK=1 in that subprocess env, and `status.py`'s own
|
|
`_loop_lock` becomes a no-op that trusts the parent's outer lock.
|
|
|
|
NOT re-entrant: do not nest `_loop_lock` within itself. POSIX `flock` is
|
|
per-fd-per-process; a second runner process blocks cleanly until the
|
|
first releases.
|
|
|
|
NFS caveat: `flock` semantics differ on NFS-mounted loop dirs. The loop
|
|
dir is documented to be local (project root or `~/.automaton`).
|
|
"""
|
|
lock_file = loop_path / ".state.lock"
|
|
fd = os.open(str(lock_file), os.O_RDWR | os.O_CREAT, 0o644)
|
|
acquired = False
|
|
try:
|
|
if sys.platform == "win32":
|
|
import msvcrt
|
|
msvcrt.locking(fd, msvcrt.LK_LOCK if exclusive else msvcrt.LK_NBLCK, 1)
|
|
else:
|
|
import fcntl
|
|
fcntl.flock(fd, fcntl.LOCK_EX if exclusive else fcntl.LOCK_SH)
|
|
acquired = True
|
|
yield
|
|
finally:
|
|
if acquired:
|
|
if sys.platform == "win32":
|
|
import msvcrt
|
|
try:
|
|
msvcrt.locking(fd, msvcrt.LK_UNLCK, 1)
|
|
except OSError:
|
|
pass
|
|
else:
|
|
import fcntl
|
|
fcntl.flock(fd, fcntl.LOCK_UN)
|
|
os.close(fd)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# subprocess plumbing (testable via monkeypatch of subprocess.run)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
def _run_json(args: Sequence[str], env: Optional[dict] = None) -> Optional[dict]:
|
|
"""Run a subprocess, return parsed JSON or None."""
|
|
try:
|
|
res = subprocess.run(list(args), capture_output=True, text=True,
|
|
timeout=30, env=env)
|
|
except (OSError, subprocess.SubprocessError) as exc:
|
|
print(f"ERROR: subprocess {args[0]!r} failed: {exc}", file=sys.stderr)
|
|
return None
|
|
if res.returncode != 0:
|
|
return None
|
|
out = res.stdout.strip()
|
|
if not out:
|
|
return None
|
|
try:
|
|
return json.loads(out.splitlines()[-1])
|
|
except json.JSONDecodeError:
|
|
return None
|
|
|
|
|
|
def _substitute(template: str, mapping: dict) -> str:
|
|
out = template
|
|
for key, value in mapping.items():
|
|
out = out.replace("{" + key + "}", str(value))
|
|
return out
|
|
|
|
|
|
_TRUNCATE_MARKER = " …[truncated]"
|
|
|
|
|
|
def _truncate_tokens(text: str, max_tokens: int) -> str:
|
|
"""Approximate token cap (4 chars/token heuristic, stdlib only)."""
|
|
if not text:
|
|
return ""
|
|
if max_tokens <= 0:
|
|
return ""
|
|
char_budget = max_tokens * 4
|
|
if len(text) <= char_budget:
|
|
return text
|
|
return text[: char_budget - len(_TRUNCATE_MARKER)] + _TRUNCATE_MARKER
|
|
|
|
|
|
def _read_task_brief(task_dir: Path) -> str:
|
|
"""Pick the first existing phase doc to feed the verifier as {task_brief}."""
|
|
for name in ("RESEARCH.md", "DESIGN.md", "SPEC.md"):
|
|
f = task_dir / name
|
|
if f.exists():
|
|
try:
|
|
return f.read_text()
|
|
except OSError:
|
|
return ""
|
|
return ""
|
|
|
|
|
|
def _acceptance_criteria_text(cfg: dict) -> str:
|
|
raw = cfg.get("acceptance_criteria")
|
|
if raw is None:
|
|
return ""
|
|
if isinstance(raw, list):
|
|
return "\n".join(str(x) for x in raw)
|
|
return str(raw)
|
|
|
|
|
|
def _next_hint_text(state: dict) -> str:
|
|
last = state.get("last_verdict")
|
|
if not isinstance(last, dict):
|
|
return ""
|
|
return str(last.get("next_hint") or "")
|
|
|
|
|
|
def _git_run(args_list: list, cwd: str, timeout: int = 15) -> tuple[int, str, str]:
|
|
try:
|
|
res = subprocess.run(["git"] + args_list, cwd=cwd,
|
|
capture_output=True, text=True, timeout=timeout, check=False)
|
|
return res.returncode, res.stdout.strip(), res.stderr.strip()
|
|
except (OSError, subprocess.SubprocessError) as exc:
|
|
return -1, "", str(exc)
|
|
|
|
|
|
def _ensure_worktree(state: dict, cfg: dict, loop_path: Path, project_dir: Path) -> str:
|
|
blast = cfg.get("blast_radius") or {}
|
|
use_worktree = blast.get("use_worktree", True)
|
|
if not use_worktree:
|
|
return str(project_dir)
|
|
|
|
existing = state.get("worktree_path")
|
|
if existing and Path(existing).exists():
|
|
return existing
|
|
|
|
if existing and not Path(existing).exists():
|
|
state["worktree_path"] = None
|
|
state["worktree_branch"] = None
|
|
|
|
wt_path = loop_path / LOOP_WORKTREE_DIR
|
|
loop_name = state.get("name") or loop_path.name
|
|
branch = f"loop/{loop_name}"
|
|
|
|
rc, out, err = _git_run(["rev-parse", "--is-inside-work-tree"], str(project_dir))
|
|
if rc != 0 or out != "true":
|
|
_append_tick_log(loop_path, f"WARNING worktree skipped: not a git repo ({err})")
|
|
return str(project_dir)
|
|
|
|
rc, out, err = _git_run(["worktree", "add", str(wt_path), "-b", branch], str(project_dir))
|
|
if rc != 0:
|
|
if "already exists" in err or "exists" in err:
|
|
rc, out, err = _git_run(["worktree", "add", str(wt_path), branch], str(project_dir))
|
|
if rc != 0:
|
|
_append_tick_log(loop_path, f"WARNING worktree add failed: {err}")
|
|
return str(project_dir)
|
|
|
|
state["worktree_path"] = str(wt_path)
|
|
state["worktree_branch"] = branch
|
|
_write_state_loop(loop_path, state)
|
|
return str(wt_path)
|
|
|
|
|
|
def _resolve_prompt(prompt_ref: Optional[str], extras: Optional[dict],
|
|
loop_path: Path, tick_num: int, role: str) -> str:
|
|
"""Resolve a prompt reference to a file path with tokens substituted.
|
|
|
|
Searches <loop_path>/<prompt_ref> then ~/.automaton/prompts/<prompt_ref>.
|
|
Reads the file, substitutes content-level tokens ({task_brief},
|
|
{acceptance_criteria}, {next_hint}, {current_task}, {current_phase},
|
|
{verdict}, {artifact_content}), writes to a temp file in outputs/, and
|
|
returns the temp file path. Falls back to the raw prompt_ref if the file
|
|
is not found.
|
|
"""
|
|
if not prompt_ref:
|
|
return prompt_ref or ""
|
|
|
|
candidates = [
|
|
loop_path / prompt_ref,
|
|
AUTOMATON_DIR / "prompts" / prompt_ref,
|
|
]
|
|
src_path = None
|
|
for c in candidates:
|
|
if c.exists():
|
|
src_path = c
|
|
break
|
|
if src_path is None:
|
|
return prompt_ref
|
|
|
|
try:
|
|
content = src_path.read_text()
|
|
except OSError:
|
|
return prompt_ref
|
|
|
|
content_extras = {}
|
|
if extras:
|
|
for k in ("task_brief", "acceptance_criteria", "next_hint",
|
|
"current_task", "current_phase", "verdict"):
|
|
if k in extras:
|
|
content_extras[k] = extras[k]
|
|
|
|
if "{artifact_content}" in content:
|
|
artifact_path = extras.get("artifact") if extras else None
|
|
artifact_content = ""
|
|
if artifact_path:
|
|
try:
|
|
artifact_content = Path(artifact_path).read_text()
|
|
except OSError:
|
|
artifact_content = ""
|
|
content_extras["artifact_content"] = artifact_content
|
|
|
|
for key, value in content_extras.items():
|
|
content = content.replace("{" + key + "}", str(value))
|
|
|
|
out_dir = loop_path / LOOP_OUTPUTS_DIR
|
|
out_dir.mkdir(parents=True, exist_ok=True)
|
|
tmp_prompt = out_dir / f"tick{tick_num}-{role}-prompt.md"
|
|
tmp_prompt.write_text(content)
|
|
return str(tmp_prompt)
|
|
|
|
|
|
def _invoke_harness(
|
|
harness_cfg: Optional[dict],
|
|
role: str,
|
|
prompt_path: str,
|
|
cwd: str,
|
|
extras: Optional[dict] = None,
|
|
loop_path: Optional[Path] = None,
|
|
tick_num: int = 0,
|
|
) -> str:
|
|
"""Build the harness command from loop.json and invoke it. Returns stdout.
|
|
|
|
extras: substitution tokens specific to this role ({artifact}, {verdict}, etc).
|
|
If extras contains a "model" key, the ``{model}`` token in the harness
|
|
command is substituted. The caller is responsible for passing the model
|
|
via extras (extracted from loop.json role config or manifest default).
|
|
"""
|
|
resolved_prompt = prompt_path
|
|
if loop_path is not None:
|
|
resolved_prompt = _resolve_prompt(prompt_path, extras, loop_path, tick_num, role)
|
|
|
|
prompt_content = ""
|
|
try:
|
|
prompt_content = Path(resolved_prompt).read_text()
|
|
except (OSError, UnicodeDecodeError):
|
|
prompt_content = ""
|
|
|
|
if harness_cfg is None:
|
|
command = ["opencode", "run", "--dir", "{cwd}", "{prompt_content}"]
|
|
else:
|
|
command = list(harness_cfg.get("command") or [])
|
|
if not command:
|
|
command = ["opencode", "run", "--dir", "{cwd}", "{prompt_content}"]
|
|
mapping = {"prompt": resolved_prompt, "cwd": cwd, "prompt_content": prompt_content}
|
|
if extras:
|
|
mapping.update(extras)
|
|
final_argv = [_substitute(tok, mapping) for tok in command]
|
|
try:
|
|
res = subprocess.run(final_argv, capture_output=True, text=True, cwd=cwd)
|
|
except (OSError, subprocess.SubprocessError) as exc:
|
|
print(f"ERROR: harness invocation failed for role {role!r}: {exc}", file=sys.stderr)
|
|
return ""
|
|
return res.stdout
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Verdict parsing
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
_FENCE_RE = re.compile(r"```(?:json)?\s*(.*?)```", re.DOTALL)
|
|
|
|
|
|
def _strip_comments(text: str) -> str:
|
|
"""Remove // and # leading-comment lines (cheap, sufficient for v1)."""
|
|
kept = []
|
|
for raw in text.splitlines():
|
|
s = raw.lstrip()
|
|
if s.startswith("//") or s.startswith("#"):
|
|
continue
|
|
kept.append(raw)
|
|
return "\n".join(kept)
|
|
|
|
|
|
def parse_verdict(text: str) -> Optional[dict]:
|
|
"""Parse verifier JSON verdict. Accepts raw, fenced, or commented JSON.
|
|
|
|
Required keys: pass (bool — also accepts "true"/"false" strings,
|
|
case-insensitive), score (float — clamped to [0, 1]; NaN / non-finite
|
|
values default to 0.5; non-numeric values default to 0.5). Optional:
|
|
reasons (list[str]), next_hint (str). Returns None on parse failure.
|
|
"""
|
|
if not text or not text.strip():
|
|
return None
|
|
candidates = []
|
|
fence_match = _FENCE_RE.search(text)
|
|
if fence_match:
|
|
candidates.append(fence_match.group(1))
|
|
candidates.append(text)
|
|
for body in candidates:
|
|
body = _strip_comments(body).strip()
|
|
if not body:
|
|
continue
|
|
try:
|
|
data = json.loads(body)
|
|
except json.JSONDecodeError:
|
|
continue
|
|
if not isinstance(data, dict):
|
|
continue
|
|
if "pass" not in data:
|
|
continue
|
|
raw_pass = data.get("pass")
|
|
if isinstance(raw_pass, str):
|
|
lower = raw_pass.strip().lower()
|
|
if lower == "true":
|
|
verdict_pass = True
|
|
elif lower == "false":
|
|
verdict_pass = False
|
|
else:
|
|
verdict_pass = bool(raw_pass.strip())
|
|
else:
|
|
verdict_pass = bool(raw_pass)
|
|
try:
|
|
score = float(data.get("score", 0.0))
|
|
except (TypeError, ValueError):
|
|
score = 0.5
|
|
if not math.isfinite(score):
|
|
score = 0.5
|
|
score = max(0.0, min(1.0, score))
|
|
verdict = {
|
|
"pass": verdict_pass,
|
|
"score": score,
|
|
}
|
|
if "reasons" in data and isinstance(data["reasons"], list):
|
|
verdict["reasons"] = [str(r) for r in data["reasons"]]
|
|
else:
|
|
verdict["reasons"] = []
|
|
if "next_hint" in data and isinstance(data["next_hint"], str):
|
|
verdict["next_hint"] = data["next_hint"]
|
|
return verdict
|
|
return None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Tick
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
def _gate(loop_path: Path, loop_name: str, project_dir: Path) -> Optional[dict]:
|
|
"""Call status.py --check-gate; return the parsed JSON dict or None on subprocess error.
|
|
|
|
Sets $AUTOMATON_NO_LOOP_LOCK=1 in the subprocess env so status.py's
|
|
`_loop_lock` becomes a no-op. The runner is responsible for holding an
|
|
outer `_loop_lock` around the entire tick; if status.py also tried to
|
|
flock the same `.state.lock` it would deadlock waiting on the parent's
|
|
flock. The env var is scoped to this subprocess only — harness
|
|
subprocesses (Implement/Verify/Orchestrate) do NOT inherit it, so any
|
|
`status.py --transition` calls the harness makes lock normally.
|
|
"""
|
|
args = [sys.executable, str(STATUS_SCRIPT),
|
|
"--check-gate", loop_name,
|
|
"--project", str(project_dir),
|
|
"--json"]
|
|
env = {**os.environ, _LOOP_LOCK_ENV_BYPASS: "1"}
|
|
out = _run_json(args, env=env)
|
|
if out is None:
|
|
return None
|
|
return out
|
|
|
|
|
|
def _context_floor_ok() -> bool:
|
|
"""Return True iff vram_detect.py reports loop-mode eligibility."""
|
|
args = [sys.executable, str(VRAM_SCRIPT), "--loop-mode", "--json"]
|
|
out = _run_json(args)
|
|
if out is None:
|
|
return True # best-effort: if the tool is unavailable, allow the tick
|
|
return bool(out.get("loop_mode_eligible", True))
|
|
|
|
|
|
def _role_prompt(cfg: dict, role: str) ->Optional[str]:
|
|
roles = cfg.get("roles") or {}
|
|
role_cfg = roles.get(role) or {}
|
|
return role_cfg.get("prompt")
|
|
|
|
|
|
def _role_model(cfg: dict, role: str) -> Optional[str]:
|
|
"""Get the model configured for a role in loop.json, or the manifest default."""
|
|
roles = cfg.get("roles") or {}
|
|
role_cfg = roles.get(role) or {}
|
|
model = role_cfg.get("model")
|
|
if model:
|
|
return model
|
|
models_file = AUTOMATON_DIR / "models.json"
|
|
if models_file.exists():
|
|
try:
|
|
import json as _mj
|
|
manifest = _mj.loads(models_file.read_text())
|
|
return manifest.get("default")
|
|
except (OSError, _mj.JSONDecodeError):
|
|
pass
|
|
return None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Work sources (task add-goal-mode)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
_SEVERITY_RANK = {"high": 3, "med": 2, "low": 1}
|
|
|
|
|
|
def _task_dir_for(name: str, project_dir: Path) -> Path:
|
|
"""Local mirror of status.py _task_dir (no cross-script imports)."""
|
|
if project_dir == AUTOMATON_DIR:
|
|
base = AUTOMATON_DIR / "tasks"
|
|
else:
|
|
base = project_dir / ".automaton" / "tasks"
|
|
return base / name
|
|
|
|
|
|
def _slugify(text: str) -> str:
|
|
s = re.sub(r"[^A-Za-z0-9._-]+", "-", text.strip().lower())
|
|
s = re.sub(r"-+", "-", s).strip("-")
|
|
return s or "task"
|
|
|
|
|
|
def _find_work_single(state: dict, cfg: dict, loop_path: Path, project_dir: Path) -> Optional[str]:
|
|
"""Return current_task or None; unchanged from task 3 behaviour."""
|
|
return state.get("current_task")
|
|
|
|
|
|
def _find_work_audit(state: dict, cfg: dict, loop_path: Path, project_dir: Path) -> Optional[str]:
|
|
"""Run status.py --audit --json; pick the highest-severity unresolved violation."""
|
|
ws = cfg.get("work_source") or {}
|
|
audit_project = ws.get("project") or str(project_dir)
|
|
args = [sys.executable, str(STATUS_SCRIPT),
|
|
"--audit",
|
|
"--project", audit_project,
|
|
"--json"]
|
|
out = _run_json(args)
|
|
if not out:
|
|
return None
|
|
violations = out.get("violations") or []
|
|
unresolved = [v for v in violations if not v.get("resolved", False)]
|
|
if not unresolved:
|
|
return None
|
|
unresolved.sort(key=lambda v: _SEVERITY_RANK.get(v.get("severity", ""), 0), reverse=True)
|
|
top = unresolved[0]
|
|
task = top.get("task")
|
|
if task:
|
|
return task
|
|
msg = top.get("message") or "audit-violation"
|
|
slug = _slugify(msg)
|
|
if _task_dir_for(slug, project_dir).exists():
|
|
return slug
|
|
create_args = [sys.executable, str(STATUS_SCRIPT),
|
|
"--create-task", slug,
|
|
"--project", str(project_dir)]
|
|
try:
|
|
subprocess.run(create_args, capture_output=True, text=True, timeout=15)
|
|
except (OSError, subprocess.SubprocessError):
|
|
pass
|
|
return slug
|
|
|
|
|
|
def _find_work_backlog(state: dict, cfg: dict, loop_path: Path, project_dir: Path) -> Optional[str]:
|
|
"""Read design/<area>/BACKLOG.md; pick the topmost `- [ ]` item."""
|
|
ws = cfg.get("work_source") or {}
|
|
area = ws.get("area") or "loops"
|
|
if project_dir == AUTOMATON_DIR:
|
|
backlog = AUTOMATON_DIR / "design" / area / "BACKLOG.md"
|
|
else:
|
|
backlog = project_dir / "design" / area / "BACKLOG.md"
|
|
if not backlog.exists():
|
|
return None
|
|
try:
|
|
text = backlog.read_text()
|
|
except OSError:
|
|
return None
|
|
for line in text.splitlines():
|
|
stripped = line.strip()
|
|
if stripped.startswith("- [ ]"):
|
|
m = re.search(r"\*\*([A-Za-z0-9._-]+)\*\*", stripped)
|
|
if m:
|
|
return m.group(1)
|
|
body = stripped.replace("- [ ]", "", 1).strip()
|
|
return _slugify(body)
|
|
return None
|
|
|
|
|
|
_FIND_WORK_DISPATCH = {
|
|
"single": _find_work_single,
|
|
"audit": _find_work_audit,
|
|
"backlog": _find_work_backlog,
|
|
}
|
|
|
|
|
|
def _find_work(state: dict, cfg: dict, loop_path: Path, project_dir: Path) -> tuple[Optional[str], Optional[str]]:
|
|
"""Return (current_task, skip_reason). skip_reason is None when work was found.
|
|
|
|
Unknown/missing work_source falls back to 'single' with a WARNING log line.
|
|
"""
|
|
ws = cfg.get("work_source") or {}
|
|
kind = ws.get("kind") if isinstance(ws, dict) else None
|
|
if not kind:
|
|
kind = "single"
|
|
handler = _FIND_WORK_DISPATCH.get(kind)
|
|
if handler is None:
|
|
_append_tick_log(loop_path, f"WARNING unknown work_source.kind={kind!r}; falling back to single")
|
|
handler = _find_work_single
|
|
kind = "single"
|
|
task = handler(state, cfg, loop_path, project_dir)
|
|
if task is None:
|
|
if kind == "single":
|
|
return None, "no_current_task"
|
|
return None, "no_work"
|
|
return task, None
|
|
|
|
|
|
def _score_window(cfg: dict) -> int:
|
|
return int((cfg.get("brakes") or {}).get("score_plateau_window", 0) or 0)
|
|
|
|
|
|
def _loop_max_iterations(cfg: dict) -> int:
|
|
return int((cfg.get("brakes") or {}).get("max_iterations", 0) or 0)
|
|
|
|
|
|
def _outputs_dir(loop_path: Path) -> Path:
|
|
d = loop_path / LOOP_OUTPUTS_DIR
|
|
d.mkdir(parents=True, exist_ok=True)
|
|
return d
|
|
|
|
|
|
def _get_retention(cfg: Optional[dict]) -> int:
|
|
if cfg is None:
|
|
return 20
|
|
try:
|
|
raw = (cfg.get("outputs") or {}).get("retention", 20)
|
|
ret = int(raw)
|
|
except (TypeError, ValueError):
|
|
print(f"WARNING: outputs.retention={raw!r} is not an int; falling back to 20",
|
|
file=sys.stderr)
|
|
return 20
|
|
if ret < 0:
|
|
print(f"WARNING: outputs.retention={ret} is negative; treating as 0 (unlimited)",
|
|
file=sys.stderr)
|
|
return 0
|
|
return ret
|
|
|
|
|
|
def _gc_outputs(loop_path: Path, retention: int) -> None:
|
|
if retention <= 0:
|
|
return
|
|
out_dir = loop_path / LOOP_OUTPUTS_DIR
|
|
if not out_dir.exists():
|
|
return
|
|
try:
|
|
names = os.listdir(str(out_dir))
|
|
except OSError:
|
|
return
|
|
tick_re = re.compile(r"^tick(\d+)-")
|
|
max_seen = 0
|
|
for name in names:
|
|
m = tick_re.match(name)
|
|
if m:
|
|
idx = int(m.group(1))
|
|
if idx > max_seen:
|
|
max_seen = idx
|
|
if max_seen == 0:
|
|
return
|
|
cutoff = max_seen - retention + 1
|
|
for name in names:
|
|
m = tick_re.match(name)
|
|
if m:
|
|
idx = int(m.group(1))
|
|
if idx < cutoff:
|
|
try:
|
|
(out_dir / name).unlink()
|
|
except OSError as exc:
|
|
_append_tick_log(loop_path,
|
|
f"WARNING GC failed to remove {name}: {exc}")
|
|
|
|
|
|
def cmd_tick(args, runner_state: Optional[dict] = None) -> dict:
|
|
"""Execute one tick. Returns a summary dict (used both for --json output
|
|
and for daemon-mode bookkeeping).
|
|
|
|
`runner_state` is reserved for daemon mode to accumulate state across ticks.
|
|
|
|
The full read-modify-write cycle on `.state.loop` (from `_gate` through
|
|
step 10's state write) is wrapped in `_loop_lock(loop_path)` so that
|
|
concurrent ticks (two scheduler firings on the same loop) and concurrent
|
|
`status.py --pause-loop` / --approve --loop writes serialize rather than
|
|
overwriting each other. See `_loop_lock` docstring for the env-bypass
|
|
mechanism used to avoid deadlock with the `--check-gate` subprocess.
|
|
"""
|
|
project_dir = _find_project_dir(args.project)
|
|
loop_path = _loop_dir(args.loop, project_dir)
|
|
state = _read_state_loop(loop_path)
|
|
summary = {"loop": args.loop, "skipped": False, "halted": False, "reason": None,
|
|
"iter": None, "verdict": None}
|
|
|
|
if state is None:
|
|
_append_tick_log(loop_path, "SKIP untracked")
|
|
summary["reason"] = "untracked"
|
|
summary["skipped"] = True
|
|
return summary
|
|
|
|
cfg = _read_loop_config(loop_path) or {}
|
|
|
|
with _loop_lock(loop_path):
|
|
# Re-read fresh state under the lock; a concurrent --pause / --approve
|
|
# may have mutated it between the unlocked read above and here.
|
|
state = _read_state_loop(loop_path)
|
|
if state is None:
|
|
_append_tick_log(loop_path, "SKIP untracked")
|
|
summary["reason"] = "untracked"
|
|
summary["skipped"] = True
|
|
return summary
|
|
|
|
# Step 2: gate
|
|
gate = _gate(loop_path, args.loop, project_dir)
|
|
if gate is None:
|
|
_append_tick_log(loop_path, "SKIP gate_subprocess_failed")
|
|
summary["reason"] = "gate_subprocess_failed"
|
|
summary["skipped"] = True
|
|
return summary
|
|
if not gate.get("ok"):
|
|
reason = gate.get("reason") or "not_ok"
|
|
_append_tick_log(loop_path, f"SKIP reason={reason}")
|
|
summary["reason"] = reason
|
|
summary["skipped"] = True
|
|
return summary
|
|
|
|
# Step 3: find work (R1 -- dispatch on work_source.kind)
|
|
current_task, skip_reason = _find_work(state, cfg, loop_path, project_dir)
|
|
if current_task is None:
|
|
_append_tick_log(loop_path, f"SKIP {skip_reason}")
|
|
summary["reason"] = skip_reason
|
|
summary["skipped"] = True
|
|
return summary
|
|
|
|
# Step 3.5: claim task (R2 -- cross-loop ownership check)
|
|
if current_task != state.get("current_task"):
|
|
claim_env = {**os.environ, _LOOP_LOCK_ENV_BYPASS: "1"}
|
|
claim_args = [sys.executable, str(STATUS_SCRIPT),
|
|
"--claim-loop-task", args.loop,
|
|
"--task", current_task,
|
|
"--project", str(project_dir)]
|
|
try:
|
|
claim_res = subprocess.run(claim_args, capture_output=True, text=True,
|
|
timeout=15, env=claim_env)
|
|
except (OSError, subprocess.SubprocessError) as exc:
|
|
_append_tick_log(loop_path, f"SKIP claim_subprocess_failed:{exc}")
|
|
summary["reason"] = "claim_subprocess_failed"
|
|
summary["skipped"] = True
|
|
return summary
|
|
if claim_res.returncode != 0:
|
|
msg = claim_res.stderr.strip() or claim_res.stdout.strip() or "denied"
|
|
_append_tick_log(loop_path, f"SKIP {msg}")
|
|
summary["reason"] = msg
|
|
summary["skipped"] = True
|
|
return summary
|
|
state["current_task"] = current_task
|
|
|
|
# Step 4: ensure worktree exists (D2 -- per-loop git worktree)
|
|
cwd = _ensure_worktree(state, cfg, loop_path, project_dir)
|
|
|
|
# Step R5: context-floor guard (D13)
|
|
if not _context_floor_ok():
|
|
_halt_loop(loop_path, state, "human_intervention")
|
|
_append_tick_log(loop_path, "HALT human_intervention:context_below_floor")
|
|
summary["reason"] = "context_below_floor"
|
|
summary["halted"] = True
|
|
return summary
|
|
|
|
# Goal-oriented substitution tokens (R4, R5, R6).
|
|
task_dir = _task_dir_for(current_task, project_dir)
|
|
task_brief = _truncate_tokens(_read_task_brief(task_dir), 4000)
|
|
acceptance = _truncate_tokens(_acceptance_criteria_text(cfg), 2000)
|
|
next_hint = _truncate_tokens(_next_hint_text(state), 1000)
|
|
|
|
# Step 5: spawn Implement
|
|
implement_prompt = _role_prompt(cfg, "implement") or ""
|
|
harness_cfg = cfg.get("harness")
|
|
out_dir = _outputs_dir(loop_path)
|
|
tick_num = state.get('iteration_count', 0) + 1
|
|
impl_output = str(out_dir / f"tick{tick_num}-implement.json")
|
|
impl_model = _role_model(cfg, "implement")
|
|
implement_extras = {
|
|
"output": impl_output,
|
|
"current_task": current_task,
|
|
"task_brief": task_brief,
|
|
"acceptance_criteria": acceptance,
|
|
"next_hint": next_hint,
|
|
}
|
|
if impl_model:
|
|
implement_extras["model"] = impl_model
|
|
implement_stdout = _invoke_harness(
|
|
harness_cfg, "implement", implement_prompt, cwd,
|
|
extras=implement_extras,
|
|
loop_path=loop_path, tick_num=tick_num)
|
|
(Path(impl_output)).write_text(implement_stdout)
|
|
|
|
# Step 6: spawn Verify
|
|
verify_prompt = _role_prompt(cfg, "verify") or ""
|
|
verify_output = str(out_dir / f"tick{tick_num}-verify.json")
|
|
verify_model = _role_model(cfg, "verify")
|
|
verify_extras = {
|
|
"output": verify_output,
|
|
"artifact": impl_output,
|
|
"current_task": current_task,
|
|
"task_brief": task_brief,
|
|
"acceptance_criteria": acceptance,
|
|
"next_hint": next_hint,
|
|
}
|
|
if verify_model:
|
|
verify_extras["model"] = verify_model
|
|
verify_stdout = _invoke_harness(
|
|
harness_cfg, "verify", verify_prompt, cwd,
|
|
extras=verify_extras,
|
|
loop_path=loop_path, tick_num=tick_num)
|
|
(Path(verify_output)).write_text(verify_stdout)
|
|
|
|
# Step 7: parse verdict
|
|
verdict = parse_verdict(verify_stdout)
|
|
if verdict is None:
|
|
_halt_loop(loop_path, state, "verifier_failed")
|
|
_append_tick_log(loop_path, "HALT verifier_failed:unparseable")
|
|
summary["reason"] = "verifier_failed:unparseable"
|
|
summary["halted"] = True
|
|
return summary
|
|
|
|
# Step 8: cap score_history
|
|
window = _score_window(cfg)
|
|
history = list(state.get("score_history", []))
|
|
history.append(float(verdict.get("score", 0.0)))
|
|
if window > 0 and len(history) > window:
|
|
history = history[-window:]
|
|
state["score_history"] = history
|
|
state["last_verdict"] = verdict
|
|
|
|
# Step 9: spawn Orchestrate
|
|
orch_prompt = _role_prompt(cfg, "orchestrate") or ""
|
|
orch_output = str(out_dir / f"tick{tick_num}-orchestrate.json")
|
|
orch_model = _role_model(cfg, "orchestrate")
|
|
orch_extras = {
|
|
"output": orch_output,
|
|
"verdict": json.dumps(verdict),
|
|
"current_task": current_task,
|
|
"current_phase": state.get("current_phase", ""),
|
|
}
|
|
if orch_model:
|
|
orch_extras["model"] = orch_model
|
|
orch_stdout = _invoke_harness(
|
|
harness_cfg, "orchestrate", orch_prompt, cwd,
|
|
extras=orch_extras,
|
|
loop_path=loop_path, tick_num=tick_num)
|
|
(Path(orch_output)).write_text(orch_stdout)
|
|
|
|
# Step 9.5: release on terminal phase (R3)
|
|
task_state_path = _task_dir_for(current_task, project_dir) / ".state"
|
|
if task_state_path.exists():
|
|
try:
|
|
task_phase = task_state_path.read_text().strip()
|
|
except OSError:
|
|
task_phase = ""
|
|
if task_phase in ("complete", "human_intervention"):
|
|
state["current_task"] = None
|
|
_append_tick_log(loop_path,
|
|
f"RELEASE current_task={current_task} phase={task_phase}")
|
|
|
|
# Step 10: advance state (atomic)
|
|
state["iteration_count"] = int(state.get("iteration_count", 0)) + 1
|
|
state["last_tick_at"] = datetime.now(timezone.utc).isoformat()
|
|
_write_state_loop(loop_path, state)
|
|
|
|
# Step 10.5: GC outputs
|
|
retention = _get_retention(cfg)
|
|
_gc_outputs(loop_path, retention)
|
|
|
|
# Step 11: log
|
|
_append_tick_log(loop_path,
|
|
f"TICK pass={verdict['pass']} score={verdict['score']} "
|
|
f"iter={state['iteration_count']}")
|
|
|
|
summary["iter"] = state["iteration_count"]
|
|
summary["verdict"] = verdict
|
|
return summary
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Daemon mode
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
def cmd_daemon(args) -> int:
|
|
interval = args.interval
|
|
if interval is None:
|
|
# default from loop.json schedule.interval_seconds, else 3600
|
|
project_dir = _find_project_dir(args.project)
|
|
loop_path = _loop_dir(args.loop, project_dir)
|
|
cfg = _read_loop_config(loop_path) or {}
|
|
interval = int(((cfg.get("schedule") or {}).get("interval_seconds")) or 3600)
|
|
count = 0
|
|
max_iter = args.max_iterations or 0
|
|
try:
|
|
while max_iter == 0 or count < max_iter:
|
|
cmd_tick(args)
|
|
count += 1
|
|
if max_iter == 0 or count < max_iter:
|
|
time.sleep(interval)
|
|
except KeyboardInterrupt:
|
|
project_dir = _find_project_dir(args.project)
|
|
loop_path = _loop_dir(args.loop, project_dir)
|
|
if loop_path.exists():
|
|
_append_tick_log(loop_path, "DAEMON_STOPPED")
|
|
return 0
|
|
return 0
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# CLI
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
def main() -> int:
|
|
parser = argparse.ArgumentParser(description="Automaton loop runner")
|
|
parser.add_argument("--mode", required=True, choices=["tick", "daemon"],
|
|
help="Loop mode: tick (single tick) or daemon (sleep loop)")
|
|
parser.add_argument("--loop", required=True, help="Loop name")
|
|
parser.add_argument("--project", help="Project root directory (defaults to CWD)")
|
|
parser.add_argument("--interval", type=int, help="Daemon tick interval (seconds)")
|
|
parser.add_argument("--max-iterations", type=int, default=0,
|
|
help="Daemon max ticks (0 = unbounded)")
|
|
parser.add_argument("--json", action="store_true", dest="json_output",
|
|
help="Print machine-readable tick summary as last line")
|
|
args = parser.parse_args()
|
|
|
|
if args.mode == "tick":
|
|
summary = cmd_tick(args)
|
|
if args.json_output:
|
|
print(json.dumps(summary))
|
|
return 0
|
|
if args.mode == "daemon":
|
|
return cmd_daemon(args)
|
|
print(f"ERROR: unknown mode '{args.mode}'", file=sys.stderr)
|
|
return 2
|
|
|
|
|
|
if __name__ == "__main__":
|
|
sys.exit(main()) |