diff --git a/corpus/skills/cat-mode/SKILL.md b/corpus/skills/cat-mode/SKILL.md index 97a3caba..c826356b 100644 --- a/corpus/skills/cat-mode/SKILL.md +++ b/corpus/skills/cat-mode/SKILL.md @@ -138,6 +138,7 @@ isolated subagents and report back async rather than blocking on each one. durable artifact.** Separable and parallel is not authorization to fan out; a fan-out default cannot hand a subagent publishing authority the routing table never granted. Route that work through Execution routing. +- **Many PR stacks: one parallel unit per stack, never serial** (Invoker, else a worktree subagent each). - **A fork/subagent told to touch files must run in its own worktree, not the live checkout** — even when told "read-only." Scope wording is not filesystem isolation. diff --git a/corpus/skills/cat-mode/references/execution-routing.md b/corpus/skills/cat-mode/references/execution-routing.md index b5994260..dd105c7c 100644 --- a/corpus/skills/cat-mode/references/execution-routing.md +++ b/corpus/skills/cat-mode/references/execution-routing.md @@ -6,6 +6,7 @@ Catstack owns judgment and local fallback. Invoker owns durable plan submission, ## Decision +0. **More than one independent publishing unit** (several PR stacks to land or repair, several workflows): never serial in the parent chat. Invoker first, one workflow per unit; `subagent_worktree_per_unit` when Invoker is unavailable or the user directs subagents. `route_execution(units=N)` returns it; see [subagents.md](subagents.md). 1. **Invoker unavailable** (no `invoker_prepare_plan_review` / `invoker_submit_plan` tools): stay local — subagents, `loop-generator`, `land-stack`, current chat execution. 2. **Small local work** (one-file fix, short edit, read-only question): stay local even if Invoker is installed. Post-land wait until `MERGED`, merge-queue babysit, and already-named execution Backlog are **not** this bucket — they are `durable_parallel`. 3. **Approved plan or durable/parallel work** and Invoker MCP is available: delegate. If Invoker is missing, use a separate git worktree + PR stack. Do not park that work in the parent chat. diff --git a/corpus/skills/cat-mode/references/subagents.md b/corpus/skills/cat-mode/references/subagents.md index 8180b19c..5ed7b83d 100644 --- a/corpus/skills/cat-mode/references/subagents.md +++ b/corpus/skills/cat-mode/references/subagents.md @@ -36,6 +36,29 @@ definition whatever `produces` claims; and an empty or unrecognized `produces` raises rather than falling through to fan-out, because an output nobody declared is unchecked, not clean. +## Many stacks: parallel per unit, never serial + +Routing picks *who* runs publishing work; it never licenses running several +independent units one after another in the parent thread. When the work is N +independent PR stacks (landing, conflict repair, review fixes), each stack is +its own unit with its own worktree, and the units run in parallel: an Invoker +workflow per stack first, and one worktree-isolated subagent per stack as the +fallback when Invoker is unavailable or the user directs it +(`subagent_worktree_per_unit`). The per-unit subagent still inherits only the +scope the parent names, and its transcript is still grepped for writes. + +The failure shape: asked to land dozens of admin-bypass PRs across two repos, +the parent recommended working the ~20 rebases and review fixes "one at a +time" and started serially in one worktree, until the user asked for a +worktree subagent per stack. The routing table allowed it: `route_execution` +returned `local` for publishing work without Invoker at any unit count. + +Prior art: Amdahl's law — Gene M. Amdahl, "Validity of the single processor +approach to achieving large scale computing capabilities", AFIPS 1967, +https://doi.org/10.1145/1465482.1465560 — the serial fraction bounds the +whole job, so independent units forced through one thread set the finish +time. + ## Defer to the harness's routing skill The precedence above is catstack's fallback, not the owner. When a harness diff --git a/corpus/skills/cat-mode/scripts/route_execution.py b/corpus/skills/cat-mode/scripts/route_execution.py index 4fbf7028..2ad1bc08 100644 --- a/corpus/skills/cat-mode/scripts/route_execution.py +++ b/corpus/skills/cat-mode/scripts/route_execution.py @@ -11,8 +11,8 @@ from typing import Literal WorkKind = Literal["readonly", "small_local", "approved_plan", "durable_parallel"] -Route = Literal["local", "delegate_invoker"] -Delegation = Literal["local", "delegate_invoker", "subagent_fanout"] +Route = Literal["local", "delegate_invoker", "subagent_worktree_per_unit"] +Delegation = Literal["local", "delegate_invoker", "subagent_fanout", "subagent_worktree_per_unit"] DURABLE_ALIASES = frozenset({"post_land_babysit", "named_execution_backlog"}) @@ -36,6 +36,13 @@ "invoker_wait_for_workflow_or_status", ) +SUBAGENT_PER_UNIT_STEPS = ( + "one_worktree_per_unit", + "spawn_one_subagent_per_unit_in_parallel", + "collect_reports_async", + "grep_transcripts_for_writes", +) + SUBAGENT_FANOUT_STEPS = ( "spawn_worktree_isolated_subagents", "collect_reports_async", @@ -72,14 +79,32 @@ def normalize_work_kind(work_kind: str) -> WorkKind: raise ValueError(f"unknown work_kind: {work_kind!r}") -def route_execution(*, tools: set[str] | frozenset[str] | list[str], work_kind: str) -> Route: +def route_execution( + *, + tools: set[str] | frozenset[str] | list[str], + work_kind: str, + units: int = 1, + user_directed_subagents: bool = False, +) -> Route: """Return where execution should run for this request. - 1. Invoker MCP missing → local - 2. Small / read-only work → local even if Invoker exists - 3. Approved plan or durable/parallel → delegate_invoker + `units` counts independent publishing units (PR stacks, workflows). + More than one never runs serially in the parent thread: + Invoker first, else one worktree-isolated subagent per unit. + + 1. units > 1 → delegate_invoker, or subagent_worktree_per_unit when + Invoker is missing or the user directed subagents + 2. Invoker MCP missing → local + 3. Small / read-only work → local even if Invoker exists + 4. Approved plan or durable/parallel → delegate_invoker """ kind = normalize_work_kind(work_kind) + if units < 1: + raise ValueError(f"units must be >= 1, got {units!r}") + if units > 1 and kind != "readonly": + if invoker_mcp_available(tools) and not user_directed_subagents: + return "delegate_invoker" + return "subagent_worktree_per_unit" if not invoker_mcp_available(tools): return "local" if kind in ("readonly", "small_local"): @@ -112,6 +137,8 @@ def route_delegation( tools: set[str] | frozenset[str] | list[str], work_kind: str, produces: set[str] | frozenset[str] | list[str] | tuple[str, ...], + units: int = 1, + user_directed_subagents: bool = False, ) -> Delegation: """Resolve the Subagents default against the execution-routing table. @@ -122,7 +149,9 @@ def route_delegation( """ normalize_work_kind(work_kind) if publishes(work_kind=work_kind, produces=produces): - return route_execution(tools=tools, work_kind=work_kind) + return route_execution( + tools=tools, work_kind=work_kind, units=units, user_directed_subagents=user_directed_subagents, + ) return "subagent_fanout" @@ -131,6 +160,8 @@ def handoff_steps_for(route: Route | Delegation) -> tuple[str, ...]: return ("stay_local",) if route == "subagent_fanout": return SUBAGENT_FANOUT_STEPS + if route == "subagent_worktree_per_unit": + return SUBAGENT_PER_UNIT_STEPS return DELEGATE_HANDOFF_STEPS @@ -142,9 +173,15 @@ def handoff_steps_for(route: Route | Delegation) -> tuple[str, ...]: tools = payload.get("tools", []) work_kind = payload.get("work_kind", "small_local") produces = payload.get("produces") + units = int(payload.get("units", 1)) + directed = bool(payload.get("user_directed_subagents", False)) if produces is None: - route: Route | Delegation = route_execution(tools=tools, work_kind=work_kind) + route: Route | Delegation = route_execution( + tools=tools, work_kind=work_kind, units=units, user_directed_subagents=directed, + ) else: - route = route_delegation(tools=tools, work_kind=work_kind, produces=produces) + route = route_delegation( + tools=tools, work_kind=work_kind, produces=produces, units=units, user_directed_subagents=directed, + ) defer_to = installed_harness_routing_skill(payload.get("home")) print(json.dumps({"route": route, "steps": list(handoff_steps_for(route)), "defer_to": defer_to})) diff --git a/corpus/skills/principle-guard-the-context-window/scripts/ci_logs.py b/corpus/skills/principle-guard-the-context-window/scripts/ci_logs.py index cc2a69e3..69aece91 100755 --- a/corpus/skills/principle-guard-the-context-window/scripts/ci_logs.py +++ b/corpus/skills/principle-guard-the-context-window/scripts/ci_logs.py @@ -342,7 +342,7 @@ def capture_github(args: argparse.Namespace, root: Path) -> dict[str, Any]: "notes": notes, } write_manifest(directory, manifest) - return artifact_response(manifest, cached=False, downloaded=True, notes=notes) + return artifact_response(manifest, cached=False, downloaded=True, notes=[]) def judge_completeness(producer_status: int, job_status: str, log_bytes: int) -> tuple[str, str]: diff --git a/corpus/skills/principle-guard-the-context-window/tests/test_ci_logs.py b/corpus/skills/principle-guard-the-context-window/tests/test_ci_logs.py index 1c402bbc..905781a5 100644 --- a/corpus/skills/principle-guard-the-context-window/tests/test_ci_logs.py +++ b/corpus/skills/principle-guard-the-context-window/tests/test_ci_logs.py @@ -377,6 +377,26 @@ def test_in_progress_job_is_never_marked_complete(self): self.assertFalse(payload["artifact"]["complete"]) + def test_first_capture_lists_each_note_once(self): + env = self.gh_env(failing_log(), mode="in_progress") + payload, _ = self.run_helper( + "capture", + "--artifact-root", + str(self.artifacts), + "--repo", + "owner/name", + "--run", + "12345", + "--job", + "41", + "--gh-path", + str(self.gh), + env=env, + ) + notes = payload["notes"] + self.assertTrue(notes, payload) + self.assertEqual(len(notes), len(set(notes)), notes) + class TestSnippet(CiLogsCase): def test_structural_blocks_carry_the_failure_and_the_final_status(self): payload, _ = self.import_log(failing_log(blocks=1)) diff --git a/engine/hooks/_runner/report.py b/engine/hooks/_runner/report.py index afc5504f..71c99f33 100644 --- a/engine/hooks/_runner/report.py +++ b/engine/hooks/_runner/report.py @@ -438,6 +438,101 @@ def format_event_table(report: dict[str, Any]) -> str: return "\n".join(lines) + "\n" +JUDGE_STAGES = ("judge_skipped", "judge_queued", "judge_finished") +VERDICTS = ("hit", "clean", "unchecked") +LEAKS = ("no_transcript", "stuck", "undelivered", "undelivered_hits") + + +def is_delivery(row: dict[str, Any]) -> bool: + return row.get("harness") == "judge" and row.get("mode_source") == "judge" and row.get("action") in VERDICTS + + +def _judge_summary(hook: str) -> dict[str, Any]: + return { + "hook": hook, + "skipped": {}, + "queued": 0, + "finished": {verdict: 0 for verdict in VERDICTS}, + "delivered": {verdict: 0 for verdict in VERDICTS}, + **{leak: 0 for leak in LEAKS}, + } + + +def build_judge_report( + rows: list[dict[str, Any]], + malformed: int, + warnings: list[str], + now: datetime, + grace: timedelta, +) -> dict[str, Any]: + summaries: dict[str, dict[str, Any]] = {} + jobs: dict[str, dict[str, Any]] = {} + for row in rows: + action = row.get("action") + if action not in JUDGE_STAGES and not is_delivery(row): + continue + hook = str(row.get("hook") or "") + summary = summaries.setdefault(hook, _judge_summary(hook)) + reason = str(row.get("reason") or "") + ts = parse_ts(row.get("ts")) + job = jobs.setdefault(str(row.get("finding_id") or ""), {"hook": hook}) + if action == "judge_skipped": + summary["skipped"][reason] = summary["skipped"].get(reason, 0) + 1 + elif action == "judge_queued": + summary["queued"] += 1 + if reason == "no_transcript": + summary["no_transcript"] += 1 + job["queued"] = ts + elif action == "judge_finished": + verdict = reason if reason in VERDICTS else "unchecked" + summary["finished"][verdict] += 1 + job["finished"] = ts + job["verdict"] = verdict + else: + summary["delivered"][str(action)] += 1 + job["delivered"] = ts + for job in jobs.values(): + summary = summaries[job["hook"]] + queued, finished = job.get("queued"), job.get("finished") + if queued is not None and finished is None and now - queued > grace: + summary["stuck"] += 1 + if finished is not None and "delivered" not in job and now - finished > grace: + summary["undelivered"] += 1 + if job.get("verdict") == "hit": + summary["undelivered_hits"] += 1 + ordered = [summaries[hook] for hook in sorted(summaries)] + return { + "window_rows": len(rows), + "malformed_rows": malformed, + "warnings": warnings, + "grace_seconds": int(grace.total_seconds()), + "hooks": ordered, + "leaks": sum(summary[leak] for summary in ordered for leak in ("no_transcript", "stuck", "undelivered")), + } + + +def _counts(values: dict[str, int]) -> str: + return ",".join(f"{name}={count}" for name, count in sorted(values.items())) or "-" + + +def format_judge_table(report: dict[str, Any]) -> str: + lines = list(report["warnings"]) + if report["malformed_rows"]: + lines.append(f"skipped {report['malformed_rows']} malformed event row(s)") + lines.append("hook skipped queued finished delivered " + " ".join(LEAKS)) + for row in report["hooks"]: + lines.append( + f"{row['hook']} {_counts(row['skipped'])} {row['queued']} {_counts(row['finished'])} " + f"{_counts(row['delivered'])} " + " ".join(str(row[leak]) for leak in LEAKS) + ) + if report["leaks"]: + lines.append( + f"LEAK: {report['leaks']} judge job(s) queued with no transcript, stuck, or never delivered " + f"after {report['grace_seconds']}s" + ) + return "\n".join(lines) + "\n" + + def main(argv: list[str] | None = None) -> int: parser = argparse.ArgumentParser() parser.add_argument("--since", default="7d") @@ -445,12 +540,31 @@ def main(argv: list[str] | None = None) -> int: mode = parser.add_mutually_exclusive_group() mode.add_argument("--events", action="store_true", default=True) mode.add_argument("--runs", action="store_true") + mode.add_argument("--judge", action="store_true") + parser.add_argument("--grace", default="1h") + parser.add_argument("--check", action="store_true") args = parser.parse_args(argv) try: since = parse_since(args.since) + grace = parse_since(args.grace) except ValueError as exc: print(str(exc), file=sys.stderr) return 2 + if args.judge: + now = datetime.now(timezone.utc) + rows, malformed, warnings = read_event_rows(metrics_dir(), now - since) + if rows is None: + for warning in warnings: + print(warning) + return 2 + report = build_judge_report(rows, malformed, warnings, now, grace) + if args.json: + print(json.dumps(report, sort_keys=True)) + else: + print(format_judge_table(report), end="") + if warnings: + return 2 + return 1 if args.check and report["leaks"] else 0 if not args.runs: try: registry_data = registry.load_registry() diff --git a/engine/hooks/_runner/tests/test_report.py b/engine/hooks/_runner/tests/test_report.py index 71964d0b..43baa123 100644 --- a/engine/hooks/_runner/tests/test_report.py +++ b/engine/hooks/_runner/tests/test_report.py @@ -186,6 +186,63 @@ def write_events(self, rows: list[dict[str, object]], malformed: bool = False) - handle.write("{bad\n") return path + def stage(self, action: str, reason: str, job: str, hours_ago: float = 3, hook: str = "wrong-check-reflect") -> dict[str, object]: + ts = datetime.now(timezone.utc) - timedelta(hours=hours_ago) + return {"schema": "catstack.hook_event.v1", "ts": ts.isoformat(), "harness": "claude", "hook": hook, + "rule_id": "", "mode_source": "stage", "action": action, "reason": reason, "finding_id": job} + + def delivered(self, verdict: str, job: str, hook: str = "wrong-check-reflect") -> dict[str, object]: + return {"schema": "catstack.hook_event.v1", "ts": datetime.now(timezone.utc).isoformat(), "harness": "judge", + "hook": hook, "rule_id": hook, "mode_source": "judge", "action": verdict, "finding_id": job} + + def test_judge_report_counts_every_stage_and_each_leak(self) -> None: + self.write_events( + [ + self.stage("judge_skipped", "already_prompted", "s1"), + self.stage("judge_skipped", "already_prompted", "s2"), + self.stage("judge_skipped", "gate_off", "s3"), + self.stage("judge_queued", "transcript", "delivered-hit"), + self.stage("judge_finished", "hit", "delivered-hit"), + self.delivered("hit", "delivered-hit"), + self.stage("judge_queued", "transcript", "lost-hit"), + self.stage("judge_finished", "hit", "lost-hit"), + self.stage("judge_queued", "no_transcript", "no-path"), + self.stage("judge_finished", "clean", "no-path"), + self.stage("judge_queued", "transcript", "hung"), + self.stage("judge_queued", "transcript", "fresh", hours_ago=0.1), + ] + ) + result = self.run_report("--judge", "--json", "--since", "1d") + self.assertEqual(result.returncode, 0, result.stderr) + [row] = json.loads(result.stdout)["hooks"] + self.assertEqual(row["skipped"], {"already_prompted": 2, "gate_off": 1}) + self.assertEqual(row["queued"], 5) + self.assertEqual(row["finished"], {"hit": 2, "clean": 1, "unchecked": 0}) + self.assertEqual(row["delivered"], {"hit": 1, "clean": 0, "unchecked": 0}) + self.assertEqual( + {key: row[key] for key in ("no_transcript", "stuck", "undelivered", "undelivered_hits")}, + {"no_transcript": 1, "stuck": 1, "undelivered": 2, "undelivered_hits": 1}, + ) + + def test_judge_check_fails_on_a_leak_and_passes_when_every_verdict_is_delivered(self) -> None: + self.write_events([self.stage("judge_queued", "transcript", "j"), self.stage("judge_finished", "hit", "j")]) + leaked = self.run_report("--judge", "--check", "--since", "1d") + self.assertEqual(leaked.returncode, 1, leaked.stdout + leaked.stderr) + self.assertIn("LEAK: 1 judge job(s)", leaked.stdout) + self.write_events( + [self.stage("judge_queued", "transcript", "j"), self.stage("judge_finished", "hit", "j"), self.delivered("hit", "j")] + ) + clean = self.run_report("--judge", "--check", "--since", "1d") + self.assertEqual(clean.returncode, 0, clean.stdout + clean.stderr) + self.assertNotIn("LEAK", clean.stdout) + + def test_judge_rows_stay_out_of_the_rule_table(self) -> None: + self.seed_configs() + self.write_events([self.stage("judge_skipped", "gate_off", "s1")]) + result = self.run_report("--since", "1d") + self.assertNotIn("judge_skipped", result.stdout) + self.assertNotIn("wrong-check-reflect", result.stdout) + def test_seeded_rows_include_no_record_and_unregistered(self) -> None: self.seed_configs() self.write_rows( diff --git a/engine/hooks/_sdk/events.py b/engine/hooks/_sdk/events.py index 5e88bf00..8ee1a717 100644 --- a/engine/hooks/_sdk/events.py +++ b/engine/hooks/_sdk/events.py @@ -65,6 +65,21 @@ def once_per_session_or_compaction(hook: str, event: dict[str, object]) -> bool: return True +def write_stage_event( + hook: str, + harness: str, + session_id: str, + action: str, + reason: str, + finding_id: str | None = None, + stderr: TextIO | None = None, +) -> bool: + err = stderr if stderr is not None else sys.stderr + row = _row(hook, harness, {"session_id": session_id}, None, "", "stage", action, 0, finding_id) + row["reason"] = reason + return _append_rows(hook, [row], err) + + def prune_old_event_files(days: int = 30, stderr: TextIO | None = None) -> None: err = stderr if stderr is not None else sys.stderr root = _metrics_dir() diff --git a/engine/hooks/_sdk/tests/test_events.py b/engine/hooks/_sdk/tests/test_events.py index 38254f62..a994d5ec 100644 --- a/engine/hooks/_sdk/tests/test_events.py +++ b/engine/hooks/_sdk/tests/test_events.py @@ -14,7 +14,7 @@ sys.path.insert(0, str(SDK_DIR)) import runtime -from events import is_human_prompt, once_per_session_or_compaction, write_events +from events import is_human_prompt, once_per_session_or_compaction, write_events, write_stage_event from finding import Finding from modes import effective_mode @@ -177,6 +177,22 @@ def _stdio(self, stdin_text: str): sys.stdout = old_stdout sys.stderr = old_stderr + def test_stage_event_carries_its_reason_and_no_rule_id(self) -> None: + with tempfile.TemporaryDirectory() as tmp, mock.patch.dict( + os.environ, {"CATSTACK_HOOK_METRICS_DIR": tmp}, clear=False + ): + self.assertTrue( + write_stage_event("demo", "claude", "session-1", "judge_skipped", "already_prompted", "job-1") + ) + rows = self._rows(tmp) + + self.assertEqual(1, len(rows)) + self.assertEqual("judge_skipped", rows[0]["action"]) + self.assertEqual("already_prompted", rows[0]["reason"]) + self.assertEqual("stage", rows[0]["mode_source"]) + self.assertEqual("job-1", rows[0]["finding_id"]) + self.assertEqual("", rows[0]["rule_id"]) + def _rows(self, directory: str) -> list[dict[str, object]]: files = list(Path(directory).glob("events-*.jsonl")) self.assertEqual(1, len(files)) diff --git a/engine/hooks/llm-judge/judge.py b/engine/hooks/llm-judge/judge.py index de894513..1a50cc74 100644 --- a/engine/hooks/llm-judge/judge.py +++ b/engine/hooks/llm-judge/judge.py @@ -17,7 +17,7 @@ SDK_DIR = os.path.join(os.path.dirname(os.path.dirname(os.path.abspath(__file__))), "_sdk") sys.path.insert(0, SDK_DIR) -from events import write_events # noqa: E402 +from events import write_events, write_stage_event # noqa: E402 from finding import Finding # noqa: E402 TIMEOUT_SECONDS = 60 @@ -345,9 +345,47 @@ def enqueue(job: dict) -> str | None: stderr=log_handle, cwd=root, ) + record_queued(job) return job_id +def job_harness(job: dict) -> str: + harness = job.get("harness") + return harness if isinstance(harness, str) and harness else "judge" + + +def job_hook(job: dict) -> str: + hook = job.get("hook") + return hook if isinstance(hook, str) and hook else "llm-judge" + + +def record_queued(job: dict) -> None: + transcript = str(job.get("transcript") or "") + reason = "transcript" if transcript else "no_transcript" + write_stage_event(job_hook(job), job_harness(job), transcript, "judge_queued", reason, job["id"]) + if not transcript: + print( + f"catstack-hook-error {job_hook(job)}: judge job {job['id']} was queued with no transcript " + f"path, so its verdict can never be delivered", + file=sys.stderr, + ) + + +def record_finished(job: dict, result: dict) -> None: + errors = io.StringIO() + written = write_stage_event( + job_hook(job), + job_harness(job), + str(job.get("transcript") or ""), + "judge_finished", + str(result.get("outcome") or "unchecked"), + str(job.get("id") or "unknown"), + stderr=errors, + ) + if not written: + log(f"job {job.get('id')}: finished event not written: {errors.getvalue().strip()}") + + def run_job(path: str) -> dict: stem = os.path.splitext(os.path.basename(path))[0] job: dict = {"id": stem} @@ -365,6 +403,7 @@ def run_job(path: str) -> dict: result = verdict(job, {"outcome": "unchecked", "attempts": []}) result["reason"] = clip(f"judge error {type(exc).__name__}", str(exc)) write_json_atomic(os.path.join(verdict_dir(str(job.get("transcript") or "")), f"{stem}.json"), result) + record_finished(job, result) try: os.remove(path) except OSError as exc: diff --git a/engine/hooks/llm-judge/tests/test_judge.py b/engine/hooks/llm-judge/tests/test_judge.py index 33cce814..5d083d3d 100644 --- a/engine/hooks/llm-judge/tests/test_judge.py +++ b/engine/hooks/llm-judge/tests/test_judge.py @@ -3,6 +3,7 @@ Run: python3 -m unittest discover -s engine/hooks/llm-judge/tests -v """ +import io import json import os import sys @@ -10,6 +11,7 @@ import time import unittest import warnings +from contextlib import redirect_stderr from unittest.mock import patch LIB_DIR = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) @@ -407,6 +409,37 @@ def test_run_job_writes_verdict_under_transcript_hash_and_deletes_job(self): with open(written, encoding="utf-8") as handle: self.assertEqual(json.load(handle)["outcome"], "hit") + def stage_rows(self): + folder = os.path.join(self.state.name, "metrics") + rows = [] + for name in sorted(os.listdir(folder)) if os.path.isdir(folder) else []: + if name.startswith("events-"): + with open(os.path.join(folder, name), encoding="utf-8") as handle: + rows.extend(json.loads(line) for line in handle) + return [row for row in rows if row.get("mode_source") == "stage"] + + def test_enqueue_records_a_queued_stage_event_keyed_by_job_id(self): + with patch.object(judge.subprocess, "Popen"): + judge.enqueue(self.job(id="queued-job", harness="claude")) + rows = self.stage_rows() + self.assertEqual([(r["action"], r["reason"], r["finding_id"], r["harness"]) for r in rows], + [("judge_queued", "transcript", "queued-job", "claude")]) + + def test_enqueue_with_no_transcript_is_recorded_and_reported_as_an_error(self): + err = io.StringIO() + with patch.object(judge.subprocess, "Popen"), redirect_stderr(err): + judge.enqueue(self.job(id="lost-job", transcript="")) + self.assertEqual([(r["action"], r["reason"]) for r in self.stage_rows()], [("judge_queued", "no_transcript")]) + self.assertTrue(err.getvalue().startswith("catstack-hook-error demo-hook: judge job lost-job")) + + def test_run_job_records_a_finished_stage_event_with_the_verdict(self): + self.use_runners(ANSWER_MATCH) + path = os.path.join(self.state.name, "jobs", "job-1.json") + judge.write_json_atomic(path, self.job()) + judge.run_job(path) + self.assertEqual([(r["action"], r["reason"], r["finding_id"]) for r in self.stage_rows()], + [("judge_finished", "hit", "job-1")]) + def test_run_job_error_is_logged_and_still_writes_unchecked_verdict(self): os.environ[judge.RUNNERS_ENV] = "not json" path = os.path.join(self.state.name, "jobs", "broken-job.json") diff --git a/engine/hooks/wrong-check-reflect/claude_stop_check.py b/engine/hooks/wrong-check-reflect/claude_stop_check.py index f83368bb..925f7c2a 100644 --- a/engine/hooks/wrong-check-reflect/claude_stop_check.py +++ b/engine/hooks/wrong-check-reflect/claude_stop_check.py @@ -14,7 +14,7 @@ def main() -> None: except (json.JSONDecodeError, OSError): return payload = payload if isinstance(payload, dict) else {} - try_enqueue_judge(payload) + try_enqueue_judge(payload, "claude") if __name__ == "__main__": diff --git a/engine/hooks/wrong-check-reflect/codex_notify.py b/engine/hooks/wrong-check-reflect/codex_notify.py index 366b53c7..04950760 100644 --- a/engine/hooks/wrong-check-reflect/codex_notify.py +++ b/engine/hooks/wrong-check-reflect/codex_notify.py @@ -29,7 +29,7 @@ def main() -> None: if payload.get("type") != "agent-turn-complete": return - try_enqueue_judge(payload) + try_enqueue_judge(payload, "codex") if __name__ == "__main__": diff --git a/engine/hooks/wrong-check-reflect/cursor_session.py b/engine/hooks/wrong-check-reflect/cursor_session.py index 8c57541c..9a2f5136 100644 --- a/engine/hooks/wrong-check-reflect/cursor_session.py +++ b/engine/hooks/wrong-check-reflect/cursor_session.py @@ -15,7 +15,7 @@ def main() -> None: print(json.dumps({"followup_message": ""})) return payload = payload if isinstance(payload, dict) else {} - try_enqueue_judge(payload) + try_enqueue_judge(payload, "cursor") print(json.dumps({"followup_message": ""})) diff --git a/engine/hooks/wrong-check-reflect/detect.py b/engine/hooks/wrong-check-reflect/detect.py index af7d394a..0bb9339b 100644 --- a/engine/hooks/wrong-check-reflect/detect.py +++ b/engine/hooks/wrong-check-reflect/detect.py @@ -11,7 +11,10 @@ sys.path.insert(0, os.path.join( os.path.dirname(os.path.dirname(os.path.abspath(__file__))), "_flags")) +sys.path.insert(0, os.path.join( + os.path.dirname(os.path.dirname(os.path.abspath(__file__))), "_sdk")) +from events import write_stage_event # noqa: E402 from flags import enforcement_gate # noqa: E402 HOOKS_DIR = os.path.dirname(os.path.abspath(__file__)) @@ -322,30 +325,49 @@ def _phrases(): return module -def enqueue_judge(payload: dict) -> str | None: - if not isinstance(payload, dict) or payload.get("stop_hook_active"): - return None +def _session(payload: dict, transcript: str) -> str: + for key in ("session_id", "conversation_id", "conversationId"): + value = payload.get(key) + if isinstance(value, str) and value: + return value + return transcript + + +def _skipped(payload: dict, harness: str, transcript: str, reason: str) -> None: + write_stage_event("wrong-check-reflect", harness, _session(payload, transcript), "judge_skipped", reason) + return None + + +def enqueue_judge(payload: dict, harness: str = "unknown") -> str | None: + if not isinstance(payload, dict): + return _skipped({}, harness, "", "bad_payload") + if payload.get("stop_hook_active"): + return _skipped(payload, harness, "", "stop_hook_active") if not enforcement_gate("wrong-check-reflect", payload.get("cwd")): - return None + return _skipped(payload, harness, "", "gate_off") path = resolve_transcript(payload) text = last_assistant_text(payload, path) key = reply_key(path, text) - if not text.strip() or already_prompted(key): - return None + if not text.strip(): + return _skipped(payload, harness, path, "empty_reply") + if already_prompted(key): + return _skipped(payload, harness, path, "already_prompted") if path and user_already_asked_reflect(path): - return None + return _skipped(payload, harness, path, "user_asked_reflect") dictionary = _phrases().load("wrong-check-reflect") job = _phrases().job(dictionary, path, text) job["id"] = uuid.uuid4().hex + job["harness"] = harness job_id = _judge().enqueue(job) - if job_id is not None: - mark_prompted(key) + if job_id is None: + return _skipped(payload, harness, path, "judge_child") + mark_prompted(key) return job_id -def try_enqueue_judge(payload: dict) -> None: +def try_enqueue_judge(payload: dict, harness: str = "unknown") -> None: try: - enqueue_judge(payload) + enqueue_judge(payload, harness) except Exception as exc: print(f"catstack-hook-error wrong-check-reflect: {type(exc).__name__}: {exc}", file=sys.stderr) return diff --git a/engine/hooks/wrong-check-reflect/tests/test_hooks.py b/engine/hooks/wrong-check-reflect/tests/test_hooks.py index a76867a8..27a62d0c 100644 --- a/engine/hooks/wrong-check-reflect/tests/test_hooks.py +++ b/engine/hooks/wrong-check-reflect/tests/test_hooks.py @@ -505,6 +505,46 @@ def test_unreadable_transcript_reports_unchecked_instead_of_going_quiet(self): self.assertIsNone(rows) self.assertIn("unchecked", err.getvalue()) + def stage_rows(self): + folder = os.path.join(self.state.name, "metrics") + rows = [] + for name in sorted(os.listdir(folder)) if os.path.isdir(folder) else []: + if name.startswith("events-"): + with open(os.path.join(folder, name), encoding="utf-8") as handle: + rows.extend(json.loads(line) for line in handle) + return [row for row in rows if row.get("mode_source") == "stage"] + + def test_every_skip_records_its_reason(self): + path = self.write_transcript(("user", REFLECT_COMMAND), ("assistant", HIT_TEXT)) + detect.enqueue_judge({"transcript_path": path, "stop_hook_active": True}, "claude") + detect.enqueue_judge({"transcript_path": path}, "claude") + other = self.write_transcript(("assistant", HIT_TEXT), name="other.jsonl") + detect.mark_prompted(detect.reply_key(other, HIT_TEXT)) + detect.enqueue_judge({"transcript_path": other}, "cursor") + with patch.dict(os.environ, {flags.REFLECT_ENFORCEMENT: "0"}): + detect.enqueue_judge({"transcript_path": other}, "codex") + detect.enqueue_judge({"last_assistant_message": " "}, "claude") + self.assertEqual( + [(r["action"], r["reason"], r["harness"]) for r in self.stage_rows()], + [ + ("judge_skipped", "stop_hook_active", "claude"), + ("judge_skipped", "user_asked_reflect", "claude"), + ("judge_skipped", "already_prompted", "cursor"), + ("judge_skipped", "gate_off", "codex"), + ("judge_skipped", "empty_reply", "claude"), + ], + ) + self.assertEqual(self.jobs(), []) + + def test_claude_stop_records_the_queued_job_under_the_claude_harness(self): + self.use_runners(SLOW_CLEAN) + path = self.write_transcript(("assistant", HIT_TEXT)) + run_claude({"transcript_path": path, "session_id": "s-1"}) + queued = [r for r in self.stage_rows() if r["action"] == "judge_queued"] + self.assertEqual([(r["hook"], r["harness"], r["reason"]) for r in queued], + [("wrong-check-reflect", "claude", "transcript")]) + self.assertEqual(queued[0]["finding_id"], self.wait_for_jobs(1)[0][: -len(".json")]) + def test_claude_malformed_stdin_fail_open(self): err = io.StringIO() with patch.object(sys, "stdin", io.StringIO("not-json")): diff --git a/engine/skills/reflect/scripts/tests/test_token_audit.py b/engine/skills/reflect/scripts/tests/test_token_audit.py index 5f15c1fb..ef179928 100644 --- a/engine/skills/reflect/scripts/tests/test_token_audit.py +++ b/engine/skills/reflect/scripts/tests/test_token_audit.py @@ -1248,6 +1248,54 @@ def test_redundant_reads_flag_yes(self): os.unlink(path) os.unlink(out.name) + def test_bash_only_session_reports_tool_shape_detectors_unchecked(self): + u = {"input_tokens": 1, "output_tokens": 1, "cache_read_input_tokens": 0, "cache_creation_input_tokens": 0} + lines = [ + claude_assistant_line("m1", "u1", [{"type": "tool_use", "id": "t1", "name": "Bash", "input": {"command": "sed -n 1,50p /a.py"}}], u), + claude_assistant_line("m2", "u2", [{"type": "tool_use", "id": "t2", "name": "Bash", "input": {"command": "sed -n 1,50p /a.py"}}], u), + claude_assistant_line("m3", "u3", [{"type": "tool_use", "id": "t3", "name": "Bash", "input": {"command": "sed -i s/a/b/ /a.py"}}], u), + ] + path = write_jsonl(lines) + out = tempfile.NamedTemporaryFile(mode="w", suffix=".json", delete=False) + out.close() + try: + with redirect_stdout(io.StringIO()): + token_audit.audit_claude(path, out_path=out.name) + with open(out.name) as f: + report = json.load(f) + for name in ("model-tier-candidates", "redundant-reads", "no-verify-edit-streak"): + fl = self._flag_by_name(report, name) + self.assertEqual(fl["value"], "unchecked", name) + self.assertIsNone(fl["count"], name) + self.assertIn("went through Bash", fl["rationale"]) + finally: + os.unlink(path) + os.unlink(out.name) + + def test_structured_tool_session_still_measures_tool_shape_detectors(self): + u = {"input_tokens": 1, "output_tokens": 1, "cache_read_input_tokens": 0, "cache_creation_input_tokens": 0} + lines = [ + claude_assistant_line("m1", "u1", [{"type": "tool_use", "id": "t1", "name": "Read", "input": {"file_path": "/a.py", "offset": 1, "limit": 10}}], u), + claude_assistant_line("m2", "u2", [{"type": "tool_use", "id": "t2", "name": "Read", "input": {"file_path": "/a.py", "offset": 1, "limit": 10}}], u), + claude_assistant_line("m3", "u3", [{"type": "tool_use", "id": "t3", "name": "Bash", "input": {"command": "echo hi"}}], u), + ] + path = write_jsonl(lines) + out = tempfile.NamedTemporaryFile(mode="w", suffix=".json", delete=False) + out.close() + try: + with redirect_stdout(io.StringIO()): + token_audit.audit_claude(path, out_path=out.name) + with open(out.name) as f: + report = json.load(f) + for name in ("model-tier-candidates", "redundant-reads", "no-verify-edit-streak"): + fl = self._flag_by_name(report, name) + self.assertIn(fl["value"], ("yes", "no"), name) + self.assertIsInstance(fl["count"], int, name) + self.assertEqual(self._flag_by_name(report, "redundant-reads")["count"], 1) + finally: + os.unlink(path) + os.unlink(out.name) + def test_recurring_failure_signatures_flag_yes(self): u = {"input_tokens": 1, "output_tokens": 1, "cache_read_input_tokens": 0, "cache_creation_input_tokens": 0} lines = [] diff --git a/engine/skills/reflect/scripts/token_audit.py b/engine/skills/reflect/scripts/token_audit.py index 474f9c0f..40620564 100644 --- a/engine/skills/reflect/scripts/token_audit.py +++ b/engine/skills/reflect/scripts/token_audit.py @@ -792,7 +792,7 @@ def audit_subagents(path, audit_started_at=None): fired = [name for name in SUBAGENT_THRASH_FLAGS if flags.get(name, {}).get("value") == "yes"] if fired: thrash["by_agent"][fname] = fired - thrash["redundant_reads"] += flags["redundant-reads"]["count"] + thrash["redundant_reads"] += flags["redundant-reads"]["count"] or 0 for fp in stats.get("redundant_read_files", []): thrash["redundant_read_files"].append(f"{fname}:{fp}") thrash["tool_errors"] += stats["n_errors"] @@ -959,6 +959,15 @@ def audit_claude(path, out_path=None, include_subagents=True): simple_turns += 1 simple_turn_output_tokens += msg_usage[mid]["output_tokens"] + FILE_TOOLS = set(LOOKUP_TOOLS) | {"Edit", "Write"} + used_tools = {n for names in msg_tool_names.values() for n in names if isinstance(n, str)} + bash_only = "Bash" in used_tools and not (used_tools & FILE_TOOLS) + lookup_blind = read_blind = edit_blind = bash_only + BLIND_RATIONALE = ( + "all file access in this session went through Bash; this detector keys on the " + "{} tool name(s) and cannot observe it" + ) + grand = total_input + total_output + total_cache_read + total_cache_creation from_model = models.most_common(1)[0][0] if models else "claude-sonnet-5" tier_backtest = None @@ -1047,8 +1056,9 @@ def audit_claude(path, out_path=None, include_subagents=True): flags = [ _flag( "model-tier-candidates", - "yes" if simple_turns else "no", - simple_turns, + "unchecked" if lookup_blind else ("yes" if simple_turns else "no"), + None if lookup_blind else simple_turns, + BLIND_RATIONALE.format("Read/Grep/Glob") if lookup_blind else ( f"{simple_turns}/{n_assistant} turns called only Read/Grep/Glob " f"({simple_turn_output_tokens:,} output tokens on those turns)" @@ -1062,8 +1072,9 @@ def audit_claude(path, out_path=None, include_subagents=True): ), _flag( "redundant-reads", - "yes" if redundant else "no", - len(redundant), + "unchecked" if read_blind else ("yes" if redundant else "no"), + None if read_blind else len(redundant), + BLIND_RATIONALE.format("Read") if read_blind else f"{len(redundant)} redundant re-read(s) of an identical file+offset/limit window with no edit in between", ), _flag( @@ -1074,8 +1085,9 @@ def audit_claude(path, out_path=None, include_subagents=True): ), _flag( "no-verify-edit-streak", - "yes" if flagged_files or global_streak_max >= THRESH else "no", - global_streak_max, + "unchecked" if edit_blind else ("yes" if flagged_files or global_streak_max >= THRESH else "no"), + None if edit_blind else global_streak_max, + BLIND_RATIONALE.format("Edit/Write") if edit_blind else ( f"longest edit streak with zero verification: {global_streak_max}; " f"{len(flagged_files)} file(s) at or above threshold {THRESH}; " diff --git a/tests/test_execution_routing.py b/tests/test_execution_routing.py index 49b3a80f..8fb7bddb 100644 --- a/tests/test_execution_routing.py +++ b/tests/test_execution_routing.py @@ -40,6 +40,47 @@ def test_partial_tools_still_local(self): ) self.assertEqual(route, "local") + def test_many_publishing_units_never_run_serially(self): + tools = list(self.router.INVOKER_REQUIRED_TOOLS) + for kind in ("durable_parallel", "approved_plan", "post_land_babysit", "small_local"): + for available in (tools, []): + with self.subTest(kind=kind, invoker=bool(available)): + route = self.router.route_execution(tools=available, work_kind=kind, units=13) + self.assertNotEqual(route, "local") + + def test_many_units_prefer_invoker_then_worktree_subagent_per_unit(self): + tools = list(self.router.INVOKER_REQUIRED_TOOLS) + self.assertEqual( + self.router.route_execution(tools=tools, work_kind="post_land_babysit", units=13), + "delegate_invoker", + ) + self.assertEqual( + self.router.route_execution(tools=[], work_kind="post_land_babysit", units=13), + "subagent_worktree_per_unit", + ) + self.assertEqual( + self.router.route_execution( + tools=tools, work_kind="post_land_babysit", units=13, user_directed_subagents=True, + ), + "subagent_worktree_per_unit", + ) + steps = self.router.handoff_steps_for("subagent_worktree_per_unit") + self.assertEqual(steps[0], "one_worktree_per_unit") + self.assertIn("grep_transcripts_for_writes", steps) + + def test_many_units_route_delegation_reaches_per_unit_route(self): + route = self.router.route_delegation( + tools=[], work_kind="durable_parallel", produces=["pull_request"], units=4, + ) + self.assertEqual(route, "subagent_worktree_per_unit") + + def test_units_must_be_positive(self): + with self.assertRaises(ValueError): + self.router.route_execution(tools=[], work_kind="durable_parallel", units=0) + + def test_single_unit_keeps_existing_routes(self): + self.assertEqual(self.router.route_execution(tools=[], work_kind="approved_plan", units=1), "local") + def test_small_local_stays_local_even_with_invoker(self): tools = list(self.router.INVOKER_REQUIRED_TOOLS) for kind in ("small_local", "readonly"):