#!/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 --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 /worktree, branch loop/; 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 /.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 / then ~/.automaton/prompts/. 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//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())