Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions corpus/skills/cat-mode/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
1 change: 1 addition & 0 deletions corpus/skills/cat-mode/references/execution-routing.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
23 changes: 23 additions & 0 deletions corpus/skills/cat-mode/references/subagents.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
55 changes: 46 additions & 9 deletions corpus/skills/cat-mode/scripts/route_execution.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"})

Expand All @@ -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",
Expand Down Expand Up @@ -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"):
Expand Down Expand Up @@ -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.

Expand All @@ -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"


Expand All @@ -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


Expand All @@ -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}))
Original file line number Diff line number Diff line change
Expand Up @@ -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]:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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))
Expand Down
114 changes: 114 additions & 0 deletions engine/hooks/_runner/report.py
Original file line number Diff line number Diff line change
Expand Up @@ -438,19 +438,133 @@ 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")
parser.add_argument("--json", action="store_true")
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()
Expand Down
57 changes: 57 additions & 0 deletions engine/hooks/_runner/tests/test_report.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
Loading
Loading