From 6e761b4495308fa8bee354b0a5e734a9ddc27714 Mon Sep 17 00:00:00 2001 From: Dat Date: Mon, 27 Jul 2026 05:53:37 +0700 Subject: [PATCH 1/7] fix(escalate): scope escalation to finished delegates with a failed gate maybe_escalate gated only on delegation_verdict != "failed", but that verdict is "failed" for ANY non-success terminal status, so a CANCELLED job (explicit user intent to stop) or a TIMED_OUT job was silently re-dispatched to a larger, costlier peer. Gate on JobState.SUCCEEDED first so only a declared-gate failure (check / scope / verify) escalates; a delegate that did not finish cleanly does not. TIMED_OUT does not escalate: it produced no graded artifact and a bigger, slower peer is at least as likely to time out again. Also guard the child-staging writes (create_job_dir, prompt/command writes, save_state) against OSError so the module's "never raises" contract holds on a read-only or full disk. --- src/crossagent/escalate.py | 104 +++++++++++++++++++++++-------------- tests/test_escalate.py | 71 +++++++++++++++++++++++++ 2 files changed, 137 insertions(+), 38 deletions(-) diff --git a/src/crossagent/escalate.py b/src/crossagent/escalate.py index c4e4ddc..92e201c 100644 --- a/src/crossagent/escalate.py +++ b/src/crossagent/escalate.py @@ -111,10 +111,30 @@ def maybe_escalate( """Re-dispatch a *failed* delegation to the next ladder rung, if any. Returns the spawned child's job id, or ``None`` when nothing was escalated - (the delegation did not fail, no rungs remain, the rung advisor is unknown, - or the depth cap was reached). Never raises: an escalation that cannot be - launched is recorded and skipped, never allowed to crash the worker. + (the delegate did not finish cleanly, no declared gate failed, no rungs + remain, the rung advisor is unknown, or the depth cap was reached). Never + raises: an escalation that cannot be launched is recorded and skipped, never + allowed to crash the worker. """ + # Escalation re-dispatches a delegate that FINISHED but whose work FAILED a + # declared gate — a failing check, a scope violation, or a failing + # verification (this module's stated scope). A job that did not finish + # cleanly is deliberately out of scope, so gate on SUCCEEDED first rather + # than on ``delegation_verdict != "failed"`` alone (which also returns + # "failed" for CANCELLED / TIMED_OUT / a crashed delegate): + # * CANCELLED is explicit user intent to stop; re-dispatching to a larger, + # costlier peer is the opposite of cancelling and spends real money. + # * TIMED_OUT (and any other non-success terminal status) produced no + # graded artifact — there is no gate failure to escalate, and a bigger + # model is generally slower, so it is at least as likely to time out + # again under the same budget. A hard task that needs a bigger model is a + # fresh dispatch decision, not an automatic ladder climb that silently + # burns budget. So TIMED_OUT does NOT escalate. + # Gating on SUCCEEDED means ``delegation_verdict == "failed"`` below can only + # be a declared-gate failure — mirroring cli._failed_reason's "did not finish + # cleanly" vs. gate-failure distinction. + if failed_job.status != jobs_mod.JobState.SUCCEEDED: + return None if delegation_verdict(failed_job) != "failed": return None @@ -145,41 +165,49 @@ def maybe_escalate( _audit_skip(job_dir, reason=f"escalation halted: {exc}") return None - child_dir = jobs_mod.create_job_dir(state_root, child_id) - _write_child_prompt(child_dir, prompt) - _write_child_command( - child_dir, - advisor=advisor, - model=model, - cwd=cwd, - registry_path=registry_path, - check=check, - check_timeout=check_timeout, - scope_paths=scope_paths, - pass_env=pass_env, - verify_with=verify_with, - verify_model=verify_model, - escalate_to=remaining, - ) - child = Job( - job_id=child_id, - status=jobs_mod.JobState.PENDING, - advisor=advisor.name, - name="", - cwd=cwd, - redacted_command="", - started_at=_now(), - updated_at=_now(), - last_activity_at=_now(), - last_event="escalation.created", - max_runtime_seconds=failed_job.max_runtime_seconds, - termination_grace_seconds=failed_job.termination_grace_seconds, - parent_job_id=parent_id, - trace_id=trace_id, - orchestrator_label=label, - nesting_depth=depth, - ) - jobs_mod.save_state(child_dir, child) + # Staging the child on disk touches the filesystem (mkdir, two file writes, + # a state save). A read-only or full disk raises OSError; catch it here so + # the "never raises" contract holds — a child that cannot be staged is + # recorded and skipped, exactly like a launch failure below. + try: + child_dir = jobs_mod.create_job_dir(state_root, child_id) + _write_child_prompt(child_dir, prompt) + _write_child_command( + child_dir, + advisor=advisor, + model=model, + cwd=cwd, + registry_path=registry_path, + check=check, + check_timeout=check_timeout, + scope_paths=scope_paths, + pass_env=pass_env, + verify_with=verify_with, + verify_model=verify_model, + escalate_to=remaining, + ) + child = Job( + job_id=child_id, + status=jobs_mod.JobState.PENDING, + advisor=advisor.name, + name="", + cwd=cwd, + redacted_command="", + started_at=_now(), + updated_at=_now(), + last_activity_at=_now(), + last_event="escalation.created", + max_runtime_seconds=failed_job.max_runtime_seconds, + termination_grace_seconds=failed_job.termination_grace_seconds, + parent_job_id=parent_id, + trace_id=trace_id, + orchestrator_label=label, + nesting_depth=depth, + ) + jobs_mod.save_state(child_dir, child) + except OSError as exc: + _audit_skip(job_dir, reason=f"escalation could not be staged: {exc}") + return None jobs_mod.append_event( job_dir, diff --git a/tests/test_escalate.py b/tests/test_escalate.py index 1bb57ca..d5363b8 100644 --- a/tests/test_escalate.py +++ b/tests/test_escalate.py @@ -255,6 +255,77 @@ def test_escalation_audit_event_on_parent(tmp_path): assert escalations[0]["trace_id"] == "trace_x" +# --------------------------------------------------------------------------- +# HIGH 1: escalation is scoped to a FINISHED delegate whose declared gate +# failed — never to a job that did not finish cleanly (CANCELLED / TIMED_OUT). +# --------------------------------------------------------------------------- + + +def _terminal_parent( + state_root: Path, *, status: JobState, job_id: str = "job_term" +) -> Job: + """Persist a parent in a non-success terminal state. + + Carries a *failing* check_result, mirroring a real cancel/timeout that + interrupts the delegate mid-run: under the old ``delegation_verdict != + "failed"`` gate (which returns "failed" for ANY non-success terminal status) + this would have been re-dispatched. + """ + job_dir = jobs_mod.create_job_dir(state_root, job_id) + now = datetime.now(timezone.utc).isoformat() + job = Job( + job_id=job_id, + status=status, + advisor="codex", + cwd=str(state_root), + started_at=now, + updated_at=now, + trace_id="trace_term", + nesting_depth=1, + check_result={ + "command": "pytest", + "exit_code": 1, + "stdout_tail": "", + "stderr_tail": "", + }, + ) + save_state(job_dir, job) + return job + + +def test_cancelled_parent_is_never_escalated(tmp_path): + """A user who cancels a job must not have it silently re-dispatched to a + larger, costlier peer — that is the opposite of cancelling (HIGH 1).""" + state_root = tmp_path / "state" + parent = _terminal_parent(state_root, status=JobState.CANCELLED) + spawned, launcher = _spawns() + assert _escalate(parent, state_root, ["claude:opus"], launcher) is None + assert spawned == [] + + +def test_timed_out_parent_is_never_escalated(tmp_path): + """Decision: a TIMED_OUT delegate produced no graded artifact, and a bigger, + slower peer is at least as likely to time out again under the same budget, so + it is NOT auto-escalated — the ladder is for gate failures on completed work, + not for work that never finished (HIGH 1).""" + state_root = tmp_path / "state" + parent = _terminal_parent(state_root, status=JobState.TIMED_OUT) + spawned, launcher = _spawns() + assert _escalate(parent, state_root, ["claude:opus"], launcher) is None + assert spawned == [] + + +def test_gate_failure_parent_still_escalates(tmp_path): + """The working path is not regressed: a SUCCEEDED delegate whose declared + gate (here, the check) failed is still escalated (HIGH 1).""" + state_root = tmp_path / "state" + parent = _failed_parent(state_root) # SUCCEEDED + failing check + spawned, launcher = _spawns() + child_id = _escalate(parent, state_root, ["claude:opus"], launcher) + assert child_id is not None + assert spawned == [(child_id, state_root)] + + def _events(job_dir: Path) -> list[dict]: path = job_dir / "events.jsonl" if not path.exists(): From 10df2ed952784c738a19341c6e75c90512196eca Mon Sep 17 00:00:00 2001 From: Dat Date: Mon, 27 Jul 2026 05:53:43 +0700 Subject: [PATCH 2/7] fix(verify): make schema temp-file lifecycle exception-safe MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit tempfile.mkstemp() ran before run_verification's try block, so a read-only /tmp, a full disk, or a restricted TMPDIR raised OSError out of the verifier — which both docstrings promise never happens — past the worker's unguarded caller, leaving the job non-terminal forever after the check and scope gates already ran and destroying the delegate's real work. Extract the schema file into a _schema_file context manager that degrades to non-structured mode when the temp file cannot be created and always unlinks on exit, so verification never propagates OSError. This also trims run_verification back under the 50-line guideline. --- src/crossagent/verify.py | 62 +++++++++++++++++++++++++++++----------- tests/test_verify.py | 21 ++++++++++++++ tests/test_worker.py | 40 ++++++++++++++++++++++++++ 3 files changed, 106 insertions(+), 17 deletions(-) diff --git a/src/crossagent/verify.py b/src/crossagent/verify.py index 7f28184..0131451 100644 --- a/src/crossagent/verify.py +++ b/src/crossagent/verify.py @@ -47,6 +47,8 @@ import os import subprocess import tempfile +from collections.abc import Iterator +from contextlib import contextmanager from dataclasses import dataclass from typing import Any, Optional @@ -279,6 +281,45 @@ def verifier_env(pass_env: Optional[list[str]] = None) -> dict[str, str]: } +@contextmanager +def _schema_file(structured: bool) -> Iterator[Optional[str]]: + """Yield a path to a temp file holding the verdict schema, or ``None``. + + Exception-safe by contract (the verifier "never raises", see + :func:`run_verification`): if the schema file cannot be created or written — + a read-only ``/tmp``, a full disk, a restricted ``TMPDIR`` in a sandbox/CI + container — this yields ``None`` so verification degrades to non-structured + mode instead of raising ``OSError`` out into the worker. The file is always + unlinked on exit. When *structured* is False no file is created. + + The mkstemp/write is in its own try/except (setup), separate from the + try/finally around the yield (cleanup), so an exception the caller raises + while the file is in use propagates untouched and is never mistaken for a + setup failure. + """ + if not structured: + yield None + return + schema_path: Optional[str] = None + try: + schema_fd, created_path = tempfile.mkstemp( + prefix="crossagent-verify-", suffix=".json" + ) + with os.fdopen(schema_fd, "w", encoding="utf-8") as handle: + json.dump(_VERDICT_SCHEMA, handle) + schema_path = created_path + except OSError: + schema_path = None + try: + yield schema_path + finally: + if schema_path is not None: + try: + os.unlink(schema_path) + except OSError: + pass + + def run_verification( advisor: Advisor, model: Optional[str], @@ -291,19 +332,12 @@ def run_verification( """Run one fresh, independent verification pass and return its outcome. Never raises: any launch or parse failure degrades to an ``error``/ - ``unverified`` verdict — a broken verifier must never crash the worker nor - silently green a delegation. + ``unverified`` verdict, and an inability to create the temp schema file + degrades to non-structured mode — a broken verifier must never crash the + worker nor silently green a delegation. """ structured = advisor.json_schema_flag is not None - schema_fd, schema_path = (-1, None) - if structured: - schema_fd, schema_path = tempfile.mkstemp( - prefix="crossagent-verify-", suffix=".json" - ) - try: - if schema_path is not None: - with os.fdopen(schema_fd, "w", encoding="utf-8") as handle: - json.dump(_VERDICT_SCHEMA, handle) + with _schema_file(structured) as schema_path: cmd = build_verifier_command(advisor, model, schema_path=schema_path) _append_prompt(cmd, advisor, artifact) parser = parsers_mod.get_parser(advisor.result_parser) @@ -323,12 +357,6 @@ def run_verification( structured, f"verifier failed to launch: {exc}", ) - finally: - if schema_path is not None: - try: - os.unlink(schema_path) - except OSError: - pass parsed = ( outcome.result diff --git a/tests/test_verify.py b/tests/test_verify.py index d7bdf3b..24f3129 100644 --- a/tests/test_verify.py +++ b/tests/test_verify.py @@ -187,6 +187,27 @@ def test_run_verification_no_result_is_error(tmp_path): assert outcome.verdict == "error" +def test_run_verification_survives_schema_temp_file_failure(tmp_path, monkeypatch): + """HIGH 2: a read-only /tmp, a full disk, or a restricted TMPDIR makes + ``tempfile.mkstemp`` raise OSError. That must degrade to non-structured mode, + never propagate out of ``run_verification`` (which promises never to raise) + and crash the worker after the delegate already did its work.""" + + def _boom(*args, **kwargs): + raise OSError("read-only file system") + + monkeypatch.setattr(verify_mod.tempfile, "mkstemp", _boom) + advisor = _structured_advisor(tmp_path, "pass") + + # Must not raise despite mkstemp failing: + outcome = run_verification(advisor, None, "ARTIFACT", cwd=str(tmp_path)) + + assert isinstance(outcome, VerifyOutcome) + # Degraded gracefully: without the schema flag the fake advisor still emitted + # a JSON verdict, which is parsed out of the answer. + assert outcome.verdict == "pass" + + def test_run_verification_never_raises_on_missing_executable(tmp_path): advisor = Advisor( name="ghost", diff --git a/tests/test_worker.py b/tests/test_worker.py index c404ed5..3a8ebfc 100644 --- a/tests/test_worker.py +++ b/tests/test_worker.py @@ -780,6 +780,46 @@ def test_worker_verify_audit_event_records_verdict_not_artifact(tmp_path, monkey assert verify_events[0]["advisor"] == "myverifier" +def test_worker_persists_terminal_record_when_verify_schema_temp_file_fails( + tmp_path, monkeypatch +): + """HIGH 2: if the verifier's schema temp file cannot be created — read-only + /tmp, full disk, restricted TMPDIR — the OSError must NOT propagate out of + run_verification into the worker after the check + scope gates already ran. + Otherwise the terminal transition never persists and the job is wedged + non-terminal forever with the delegate's real work destroyed. The worker must + still reach and persist a terminal record.""" + import tempfile + + from crossagent import verify as verify_mod + + verifier = _fake_verifier(tmp_path, "pass") # structured -> mkstemp attempted + monkeypatch.setattr( + "crossagent.advisors.resolve", lambda name, config_path=None: verifier + ) + + # Fail ONLY the verifier's schema temp file; leave the state-persistence + # temp files (a different prefix) working, so this isolates the verify path. + real_mkstemp = tempfile.mkstemp + + def _selective_mkstemp(*args, **kwargs): + if kwargs.get("prefix", "").startswith("crossagent-verify-"): + raise OSError("read-only file system") + return real_mkstemp(*args, **kwargs) + + monkeypatch.setattr(verify_mod.tempfile, "mkstemp", _selective_mkstemp) + + job = _run_job_through_worker(tmp_path, _RESULT_OK, verify_with="myverifier") + + # The worker reached a terminal state and recorded the verification, rather + # than crashing after the delegate's work with the job left non-terminal. + assert jobs_mod.is_terminal(job.status) + assert job.status == JobState.SUCCEEDED + assert job.verify_result is not None + # Degraded to non-structured mode: the verdict still parsed from the answer. + assert job.verify_result["verdict"] == "pass" + + def test_worker_escalates_failed_delegation_end_to_end(tmp_path, monkeypatch): """A failed delegation (failing check) with an escalation ladder spawns a same-trace child with parent_job_id set — the exact shape analytics counts.""" From e23e2c503cb72a86ea21b1af55cd781684ea774c Mon Sep 17 00:00:00 2001 From: Dat Date: Mon, 27 Jul 2026 05:53:55 +0700 Subject: [PATCH 3/7] fix(cli): scrub credentials on the foreground dispatch path Credential scrubbing was wired into the durable-job worker but not the default `crossagent --agent ... --prompt ...` invocation: _run_advisor ran the advisor subprocess with no env=, so it inherited the caller's full unscrubbed os.environ, and --pass-env was registered only for `start`. The security claim ("credential withholding for delegates") therefore exceeded the implementation for the original dispatch mode. Scrub via the same credentials.scrub_env helper the worker uses (one shared policy, no drift) and add the --pass-env escape hatch to the foreground parser. --- src/crossagent/cli.py | 32 +++++++++++++++-- tests/test_cli.py | 84 +++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 113 insertions(+), 3 deletions(-) diff --git a/src/crossagent/cli.py b/src/crossagent/cli.py index 379f892..57446c4 100644 --- a/src/crossagent/cli.py +++ b/src/crossagent/cli.py @@ -21,6 +21,7 @@ from . import __version__ from . import advisors as advisors_mod from . import check as check_mod +from . import credentials as credentials_mod from . import jobs as jobs_mod from . import parsers as parsers_mod from . import registry as reg @@ -109,10 +110,16 @@ def build_command( def _run_advisor( - cmd: list[str], cwd: str | None, parser_name: str + cmd: list[str], + cwd: str | None, + parser_name: str, + *, + env: dict[str, str] | None = None, ) -> tuple[int, parsers_mod.ParsedResult]: parser = parsers_mod.get_parser(parser_name) - outcome = runner_mod.run(cmd, cwd=cwd, consumer=parser, max_runtime_seconds=None) + outcome = runner_mod.run( + cmd, cwd=cwd, env=env, consumer=parser, max_runtime_seconds=None + ) parsed = ( outcome.result if isinstance(outcome.result, parsers_mod.ParsedResult) @@ -199,6 +206,18 @@ def parse_args(argv: list[str] | None = None) -> argparse.Namespace: parser.add_argument( "--registry", default=str(reg.DEFAULT_REGISTRY), help="Session registry path." ) + parser.add_argument( + "--pass-env", + action="append", + default=[], + dest="pass_env", + help=( + "Name of an environment variable to pass through to the advisor even " + "though it matches a credential pattern (e.g. the advisor's own API " + "key). Repeatable. By default all credential-bearing env vars are " + "withheld from the advisor, matching the durable-job path." + ), + ) parser.add_argument( "--list-advisors", action="store_true", help="Print known advisors and exit." ) @@ -278,7 +297,14 @@ def _dispatch( registry: dict[str, Any], registry_path: Path, ) -> int: - code, parsed = _run_advisor(cmd, args.cwd, advisor.result_parser) + # A foreground advisor is a delegate too: withhold the caller's ambient + # credentials (S4 policy), sharing the exact scrub the durable-job path uses + # (worker.build_advisor_env delegates to the same helper) so the two dispatch + # modes can never drift. ``--pass-env NAME`` is the caller's opt-out. + advisor_env = credentials_mod.scrub_env( + os.environ, pass_through=getattr(args, "pass_env", []) + ) + code, parsed = _run_advisor(cmd, args.cwd, advisor.result_parser, env=advisor_env) if parsed.failure: error = parsed.error or f"{advisor.name} exited with code {code}" diff --git a/tests/test_cli.py b/tests/test_cli.py index 07f7023..291bc83 100644 --- a/tests/test_cli.py +++ b/tests/test_cli.py @@ -1,3 +1,5 @@ +import sys + import pytest from crossagent import __version__, advisors @@ -138,6 +140,88 @@ def test_missing_advisor_cli_exits_cleanly_without_logging_prompt(monkeypatch, c assert captured.out == "" +# --------------------------------------------------------------------------- +# HIGH 3: the foreground (default) dispatch path scrubs credentials from the +# advisor env too — not only the durable-job path — with the same --pass-env +# escape hatch, so both dispatch modes share one policy. +# --------------------------------------------------------------------------- + +_FG_SECRET_NAME = "AWS_SECRET_ACCESS_KEY" +_FG_SECRET_VALUE = "fg-super-secret-value" + + +def _env_probe_advisor(tmp_path, probe_name): + """A fake claude-stream advisor that records one env var it was given.""" + script = tmp_path / "fake_fg_advisor.py" + script.write_text( + "import json, os\n" + f"open('fg_env_probe.txt', 'w').write(" + f"os.environ.get({probe_name!r}, 'ABSENT'))\n" + "print(json.dumps({'type': 'result', 'subtype': 'success', " + "'result': 'ok'}))\n", + encoding="utf-8", + ) + return Advisor( + name="fakefg", + executable=sys.executable, + base_args=(str(script),), + prompt_delivery="positional", + result_parser="claude-stream", + ) + + +def test_foreground_advisor_env_is_scrubbed(tmp_path, monkeypatch): + """The default `crossagent --agent ... --prompt ...` invocation must withhold + the caller's ambient credentials from the advisor, matching the durable-job + path (HIGH 3).""" + monkeypatch.setenv(_FG_SECRET_NAME, _FG_SECRET_VALUE) + advisor = _env_probe_advisor(tmp_path, _FG_SECRET_NAME) + monkeypatch.setattr(advisors, "resolve", lambda _name: advisor) + + code = main(["--agent", "fakefg", "--prompt", "hi", "--cwd", str(tmp_path)]) + + assert code == 0 + assert (tmp_path / "fg_env_probe.txt").read_text() == "ABSENT" + + +def test_foreground_pass_env_opts_a_named_var_back_in(tmp_path, monkeypatch): + """--pass-env NAME is the foreground escape hatch: the named credential var + reaches the advisor despite matching a credential pattern (HIGH 3).""" + monkeypatch.setenv(_FG_SECRET_NAME, _FG_SECRET_VALUE) + advisor = _env_probe_advisor(tmp_path, _FG_SECRET_NAME) + monkeypatch.setattr(advisors, "resolve", lambda _name: advisor) + + code = main( + [ + "--agent", + "fakefg", + "--prompt", + "hi", + "--cwd", + str(tmp_path), + "--pass-env", + _FG_SECRET_NAME, + ] + ) + + assert code == 0 + assert (tmp_path / "fg_env_probe.txt").read_text() == _FG_SECRET_VALUE + + +def test_foreground_secret_value_absent_from_output(tmp_path, monkeypatch, capsys): + """No credential VALUE appears in the foreground path's stdout/stderr — its + only output surface (the session registry stores no env) (HIGH 3).""" + monkeypatch.setenv(_FG_SECRET_NAME, _FG_SECRET_VALUE) + advisor = _env_probe_advisor(tmp_path, _FG_SECRET_NAME) + monkeypatch.setattr(advisors, "resolve", lambda _name: advisor) + + main(["--agent", "fakefg", "--prompt", "hi", "--cwd", str(tmp_path)]) + + captured = capsys.readouterr() + assert _FG_SECRET_VALUE not in captured.out + assert _FG_SECRET_VALUE not in captured.err + + def test_version_flag_prints_version_and_exits(capsys): with pytest.raises(SystemExit) as excinfo: main(["--version"]) From 37d7496dfea9f3d5b57c1f7941c9c4d0543ffaa0 Mon Sep 17 00:00:00 2001 From: Dat Date: Mon, 27 Jul 2026 06:07:56 +0700 Subject: [PATCH 4/7] style(verify,escalate): replace forbidden placeholder names obj and info --- src/crossagent/escalate.py | 4 ++-- src/crossagent/verify.py | 12 ++++++------ 2 files changed, 8 insertions(+), 8 deletions(-) diff --git a/src/crossagent/escalate.py b/src/crossagent/escalate.py index 92e201c..903b2f9 100644 --- a/src/crossagent/escalate.py +++ b/src/crossagent/escalate.py @@ -264,7 +264,7 @@ def _write_child_command( verify_model: Optional[str], escalate_to: list[str], ) -> None: - info = { + command_payload = { "command": _build_child_argv(advisor, model), "prompt_delivery": advisor.prompt_delivery, "cwd": cwd, @@ -286,7 +286,7 @@ def _write_child_command( "verify_model": verify_model, "escalate_to": escalate_to, } - jobs_mod.atomic_json_write(info, child_dir / "command.json") + jobs_mod.atomic_json_write(command_payload, child_dir / "command.json") def _audit_skip(job_dir: Path, *, reason: str) -> None: diff --git a/src/crossagent/verify.py b/src/crossagent/verify.py index 0131451..83b63d8 100644 --- a/src/crossagent/verify.py +++ b/src/crossagent/verify.py @@ -244,10 +244,10 @@ def _json_candidates(text: str) -> list[str]: return candidates -def _verdict_from_object(obj: dict[str, Any]) -> tuple[VerifyVerdict, str]: +def _verdict_from_object(verdict_object: dict[str, Any]) -> tuple[VerifyVerdict, str]: """Map a parsed verdict object to a ``(verdict, detail)`` pair.""" - raw = obj.get("verdict") - reason = obj.get("reason") + raw = verdict_object.get("verdict") + reason = verdict_object.get("reason") detail = str(reason) if isinstance(reason, str) and reason else "" if isinstance(raw, bool): return ("pass" if raw else "fail"), detail @@ -380,8 +380,8 @@ def run_verification( parsed.error or "verifier produced no answer", ) - obj = _parse_verdict_object(parsed.result) - if obj is None: + verdict_object = _parse_verdict_object(parsed.result) + if verdict_object is None: # Free prose with no machine-checkable verdict: inconclusive, not a pass. return VerifyOutcome( advisor.name, @@ -390,5 +390,5 @@ def run_verification( structured, "verifier returned no machine-checkable verdict (free prose)", ) - verdict, detail = _verdict_from_object(obj) + verdict, detail = _verdict_from_object(verdict_object) return VerifyOutcome(advisor.name, model, verdict, structured, detail) From 74596b6432496d98323180827ef4e2180d7b723c Mon Sep 17 00:00:00 2001 From: Dat Date: Mon, 27 Jul 2026 06:08:37 +0700 Subject: [PATCH 5/7] refactor(verify): extract _interpret_run_outcome from run_verification --- src/crossagent/verify.py | 23 +++++++++++++++++++---- 1 file changed, 19 insertions(+), 4 deletions(-) diff --git a/src/crossagent/verify.py b/src/crossagent/verify.py index 83b63d8..af08790 100644 --- a/src/crossagent/verify.py +++ b/src/crossagent/verify.py @@ -357,7 +357,22 @@ def run_verification( structured, f"verifier failed to launch: {exc}", ) + return _interpret_run_outcome(advisor.name, model, structured, outcome, timeout) + +def _interpret_run_outcome( + advisor_name: str, + model: Optional[str], + structured: bool, + outcome: runner_mod.RunOutcome, + timeout: float, +) -> VerifyOutcome: + """Map a completed runner outcome to a ``VerifyOutcome``. + + A timeout or a launch/parse failure degrades to ``error``; free prose with no + parseable verdict is ``unverified`` (inconclusive, never a pass); only a + machine-checkable object yields the graded ``pass``/``fail`` verdict. + """ parsed = ( outcome.result if isinstance(outcome.result, parsers_mod.ParsedResult) @@ -365,7 +380,7 @@ def run_verification( ) if outcome.timed_out: return VerifyOutcome( - advisor.name, + advisor_name, model, "error", structured, @@ -373,7 +388,7 @@ def run_verification( ) if parsed.failure or parsed.result is None: return VerifyOutcome( - advisor.name, + advisor_name, model, "error", structured, @@ -384,11 +399,11 @@ def run_verification( if verdict_object is None: # Free prose with no machine-checkable verdict: inconclusive, not a pass. return VerifyOutcome( - advisor.name, + advisor_name, model, "unverified", structured, "verifier returned no machine-checkable verdict (free prose)", ) verdict, detail = _verdict_from_object(verdict_object) - return VerifyOutcome(advisor.name, model, verdict, structured, detail) + return VerifyOutcome(advisor_name, model, verdict, structured, detail) From d25ce30105cc1b5cfcc8a435d65d635719bd388b Mon Sep 17 00:00:00 2001 From: Dat Date: Mon, 27 Jul 2026 06:13:29 +0700 Subject: [PATCH 6/7] refactor(escalate): extract eligibility guard, child Job build, and disk staging from maybe_escalate --- src/crossagent/escalate.py | 199 +++++++++++++++++++++++++------------ 1 file changed, 136 insertions(+), 63 deletions(-) diff --git a/src/crossagent/escalate.py b/src/crossagent/escalate.py index 903b2f9..7cf1f19 100644 --- a/src/crossagent/escalate.py +++ b/src/crossagent/escalate.py @@ -116,26 +116,7 @@ def maybe_escalate( raises: an escalation that cannot be launched is recorded and skipped, never allowed to crash the worker. """ - # Escalation re-dispatches a delegate that FINISHED but whose work FAILED a - # declared gate — a failing check, a scope violation, or a failing - # verification (this module's stated scope). A job that did not finish - # cleanly is deliberately out of scope, so gate on SUCCEEDED first rather - # than on ``delegation_verdict != "failed"`` alone (which also returns - # "failed" for CANCELLED / TIMED_OUT / a crashed delegate): - # * CANCELLED is explicit user intent to stop; re-dispatching to a larger, - # costlier peer is the opposite of cancelling and spends real money. - # * TIMED_OUT (and any other non-success terminal status) produced no - # graded artifact — there is no gate failure to escalate, and a bigger - # model is generally slower, so it is at least as likely to time out - # again under the same budget. A hard task that needs a bigger model is a - # fresh dispatch decision, not an automatic ladder climb that silently - # burns budget. So TIMED_OUT does NOT escalate. - # Gating on SUCCEEDED means ``delegation_verdict == "failed"`` below can only - # be a declared-gate failure — mirroring cli._failed_reason's "did not finish - # cleanly" vs. gate-failure distinction. - if failed_job.status != jobs_mod.JobState.SUCCEEDED: - return None - if delegation_verdict(failed_job) != "failed": + if not _is_escalatable_failure(failed_job): return None rungs = parse_rungs(escalate_to) @@ -154,7 +135,7 @@ def maybe_escalate( child_id = jobs_mod.generate_job_id() try: - parent_id, trace_id, label, depth = jobs_mod.resolve_lineage( + lineage = jobs_mod.resolve_lineage( parent_flag=failed_job.job_id, state_root=state_root, new_job_id=child_id, @@ -165,50 +146,28 @@ def maybe_escalate( _audit_skip(job_dir, reason=f"escalation halted: {exc}") return None - # Staging the child on disk touches the filesystem (mkdir, two file writes, - # a state save). A read-only or full disk raises OSError; catch it here so - # the "never raises" contract holds — a child that cannot be staged is - # recorded and skipped, exactly like a launch failure below. - try: - child_dir = jobs_mod.create_job_dir(state_root, child_id) - _write_child_prompt(child_dir, prompt) - _write_child_command( - child_dir, - advisor=advisor, - model=model, - cwd=cwd, - registry_path=registry_path, - check=check, - check_timeout=check_timeout, - scope_paths=scope_paths, - pass_env=pass_env, - verify_with=verify_with, - verify_model=verify_model, - escalate_to=remaining, - ) - child = Job( - job_id=child_id, - status=jobs_mod.JobState.PENDING, - advisor=advisor.name, - name="", - cwd=cwd, - redacted_command="", - started_at=_now(), - updated_at=_now(), - last_activity_at=_now(), - last_event="escalation.created", - max_runtime_seconds=failed_job.max_runtime_seconds, - termination_grace_seconds=failed_job.termination_grace_seconds, - parent_job_id=parent_id, - trace_id=trace_id, - orchestrator_label=label, - nesting_depth=depth, - ) - jobs_mod.save_state(child_dir, child) - except OSError as exc: - _audit_skip(job_dir, reason=f"escalation could not be staged: {exc}") + if not _create_and_save_child( + state_root=state_root, + job_dir=job_dir, + child_id=child_id, + prompt=prompt, + advisor=advisor, + model=model, + cwd=cwd, + registry_path=registry_path, + check=check, + check_timeout=check_timeout, + scope_paths=scope_paths, + pass_env=pass_env, + verify_with=verify_with, + verify_model=verify_model, + escalate_to=remaining, + lineage=lineage, + failed_job=failed_job, + ): return None + _, trace_id, _, depth = lineage jobs_mod.append_event( job_dir, "escalation", @@ -234,6 +193,33 @@ def maybe_escalate( # --------------------------------------------------------------------------- +def _is_escalatable_failure(failed_job: Job) -> bool: + """Return True only for a delegate that FINISHED but FAILED a declared gate. + + Escalation re-dispatches work that failed a check, a scope violation, or a + failing verification (this module's stated scope). A job that did not finish + cleanly is deliberately out of scope, so gate on SUCCEEDED first rather than + on ``delegation_verdict != "failed"`` alone (which also returns "failed" for + CANCELLED / TIMED_OUT / a crashed delegate): + + * CANCELLED is explicit user intent to stop; re-dispatching to a larger, + costlier peer is the opposite of cancelling and spends real money. + * TIMED_OUT (and any other non-success terminal status) produced no graded + artifact — there is no gate failure to escalate, and a bigger model is + generally slower, so it is at least as likely to time out again under the + same budget. A hard task that needs a bigger model is a fresh dispatch + decision, not an automatic ladder climb that silently burns budget. So + TIMED_OUT does NOT escalate. + + Gating on SUCCEEDED means ``delegation_verdict == "failed"`` can only be a + declared-gate failure — mirroring cli._failed_reason's "did not finish + cleanly" vs. gate-failure distinction. + """ + if failed_job.status != jobs_mod.JobState.SUCCEEDED: + return False + return delegation_verdict(failed_job) == "failed" + + def _drop_first_nonblank(rungs: Optional[list[str]]) -> list[str]: """Return the non-blank rungs with the first one (the spawned rung) removed.""" cleaned = [raw for raw in (rungs or []) if raw.strip()] @@ -289,6 +275,93 @@ def _write_child_command( jobs_mod.atomic_json_write(command_payload, child_dir / "command.json") +def _build_child_job( + child_id: str, + advisor: Advisor, + cwd: str, + lineage: tuple[Optional[str], str, Optional[str], Optional[int]], + failed_job: Job, +) -> Job: + """Assemble the PENDING child ``Job`` record for an escalation re-dispatch. + + *lineage* is the ``(parent_id, trace_id, label, depth)`` tuple returned by + :func:`~crossagent.jobs.resolve_lineage`; runtime bounds are inherited from + *failed_job* so the larger peer runs under the same limits. + """ + parent_id, trace_id, label, depth = lineage + return Job( + job_id=child_id, + status=jobs_mod.JobState.PENDING, + advisor=advisor.name, + name="", + cwd=cwd, + redacted_command="", + started_at=_now(), + updated_at=_now(), + last_activity_at=_now(), + last_event="escalation.created", + max_runtime_seconds=failed_job.max_runtime_seconds, + termination_grace_seconds=failed_job.termination_grace_seconds, + parent_job_id=parent_id, + trace_id=trace_id, + orchestrator_label=label, + nesting_depth=depth, + ) + + +def _create_and_save_child( + *, + state_root: Path, + job_dir: Path, + child_id: str, + prompt: str, + advisor: Advisor, + model: Optional[str], + cwd: str, + registry_path: str, + check: Optional[str], + check_timeout: float, + scope_paths: Optional[list[str]], + pass_env: list[str], + verify_with: Optional[str], + verify_model: Optional[str], + escalate_to: list[str], + lineage: tuple[Optional[str], str, Optional[str], Optional[int]], + failed_job: Job, +) -> bool: + """Stage the escalated child on disk: prompt, command, and state record. + + Staging touches the filesystem (mkdir, two file writes, a state save). A + read-only or full disk raises OSError; it is caught here so the caller's + "never raises" contract holds — a child that cannot be staged is recorded on + *job_dir* and skipped, exactly like a launch failure. Returns True on + success, False when staging was skipped. + """ + try: + child_dir = jobs_mod.create_job_dir(state_root, child_id) + _write_child_prompt(child_dir, prompt) + _write_child_command( + child_dir, + advisor=advisor, + model=model, + cwd=cwd, + registry_path=registry_path, + check=check, + check_timeout=check_timeout, + scope_paths=scope_paths, + pass_env=pass_env, + verify_with=verify_with, + verify_model=verify_model, + escalate_to=escalate_to, + ) + child = _build_child_job(child_id, advisor, cwd, lineage, failed_job) + jobs_mod.save_state(child_dir, child) + except OSError as exc: + _audit_skip(job_dir, reason=f"escalation could not be staged: {exc}") + return False + return True + + def _audit_skip(job_dir: Path, *, reason: str) -> None: jobs_mod.append_event( job_dir, "escalation_skipped", actor="system:escalate", reason=reason From 41ea2dd09ecdd1fe3171ee40a4c8a52eecb813c6 Mon Sep 17 00:00:00 2001 From: Dat Date: Mon, 27 Jul 2026 06:19:24 +0700 Subject: [PATCH 7/7] refactor(types): move gate-result shapes to dependency-free types module Relocates CheckResultDict/ScopeResultDict/ScopeStatus/VerifyResultDict/VerifyVerdict out of jobs.py into a new dependency-free types.py. This removes the inverted back-import where the gate producers (check/scope/verify) imported their own persisted-record shapes from their consumer (jobs). jobs re-exports the three it uses as Job field annotations for jobs_mod.* callers (worker). Behaviour unchanged; 513 tests green. --- src/crossagent/check.py | 2 +- src/crossagent/jobs.py | 92 ++++--------------------------------- src/crossagent/scope.py | 2 +- src/crossagent/types.py | 97 ++++++++++++++++++++++++++++++++++++++++ src/crossagent/verify.py | 2 +- 5 files changed, 108 insertions(+), 87 deletions(-) create mode 100644 src/crossagent/types.py diff --git a/src/crossagent/check.py b/src/crossagent/check.py index 8845b84..511aed9 100644 --- a/src/crossagent/check.py +++ b/src/crossagent/check.py @@ -23,7 +23,7 @@ from dataclasses import dataclass from typing import Optional -from .jobs import CheckResultDict +from .types import CheckResultDict # Keep only a bounded tail of the check's output — a full test run can emit # megabytes, and the whole Job record is loaded into memory (and the dashboard). diff --git a/src/crossagent/jobs.py b/src/crossagent/jobs.py index 358efe9..cc53b6e 100644 --- a/src/crossagent/jobs.py +++ b/src/crossagent/jobs.py @@ -16,7 +16,14 @@ from datetime import datetime, timezone from enum import Enum from pathlib import Path -from typing import Any, Callable, Literal, Optional, TypedDict +from typing import Any, Callable, Literal, Optional + +# The gate-result shapes now live in the dependency-free ``types`` module. +# ``jobs`` uses these three as ``Job`` field annotations, which also re-exports +# them for callers that still reference ``jobs_mod.CheckResultDict`` (e.g. +# ``worker``). ``ScopeStatus``/``VerifyVerdict`` are imported straight from +# ``types`` by ``scope``/``verify``. +from .types import CheckResultDict, ScopeResultDict, VerifyResultDict # Where a job's cost figure came from. ``advisor`` is the vendor-declared # estimate, ``computed`` is derived from a local price table, and ``unknown`` @@ -24,89 +31,6 @@ CostSource = Literal["advisor", "computed", "unknown"] -class CheckResultDict(TypedDict): - """Persisted outcome of the independent check-gate (slice S3). - - ``command`` is the caller-supplied check string; ``exit_code`` is the - deterministic delegation verdict (D5) — ``0`` means the work verified, any - non-zero value means the delegation *failed* regardless of whether the - delegate process itself exited 0. ``stdout_tail``/``stderr_tail`` are - bounded tails of the check's output. - - A ``Job.check_result`` of ``None`` means *no check ran* (unverified) — that - stays structurally distinct from a check that ran and failed (D7): absent is - never the same as false. - """ - - command: str - exit_code: int - stdout_tail: str - stderr_tail: str - - -# Outcome of the diff-scope assertion (slice S4). ``ok`` — every path the -# delegate modified was inside the declared allowlist; ``violated`` — it wrote -# outside its declared scope; ``undetermined`` — crossagent could not establish -# what changed (e.g. the cwd is not a git repo). ``undetermined`` is FAIL-CLOSED, -# never a pass: a scope check that silently passes when it cannot see the changes -# grants false assurance. -ScopeStatus = Literal["ok", "violated", "undetermined"] - - -class ScopeResultDict(TypedDict): - """Persisted outcome of the diff-scope assertion (slice S4). - - ``declared`` is the caller-supplied allowlist; ``violating_paths`` lists the - repo-relative paths the delegate modified outside it (empty unless - ``status == "violated"``); ``detail`` is a human-readable summary. - - A ``Job.scope_result`` of ``None`` means *no scope was declared* — scope - enforcement was off — which stays structurally distinct from a scope that - was declared and satisfied (D7: absent is never the same as ``ok``). - """ - - declared: list[str] - status: ScopeStatus - violating_paths: list[str] - detail: str - - -# Outcome of the independent verification pass (slice S5). A FRESH peer session -# grades the delegate's artifact supplied as user-turn input (D6), which removes -# the implicit-authorship channel that weakens self-grading — it does NOT claim -# to eliminate self-preference bias, which is a separate documented effect. -# ``pass`` — the verifier returned a machine-checkable verdict of correct. -# ``fail`` — the verifier returned a machine-checkable verdict of wrong. -# ``unverified`` — the verifier ran but produced no machine-checkable verdict -# (free prose, or the advisor lacks a structured-output -# contract). Treated as inconclusive, never as a pass (D4). -# ``error`` — the verifier could not run (no artifact, launch failure). -# ``unverified``/``error`` are inconclusive: they never green a delegation and -# never hard-fail it. Only ``fail`` blocks the green path. -VerifyVerdict = Literal["pass", "fail", "unverified", "error"] - - -class VerifyResultDict(TypedDict): - """Persisted outcome of the independent verification pass (slice S5). - - ``advisor``/``model`` identify the FRESH peer session that graded the work. - ``verdict`` is the deterministic pass/fail/inconclusive outcome; ``structured`` - records whether a machine-checkable contract (e.g. Claude ``--json-schema`` → - ``structured_output``) produced it, versus a JSON object parsed out of a prose - answer. ``detail`` is a human-readable summary (never the raw artifact). - - A ``Job.verify_result`` of ``None`` means *no verification was requested* — - structurally distinct from a verification that ran and failed, or ran and - could not decide (D7: absent is never the same as ``fail`` or ``unverified``). - """ - - advisor: str - model: Optional[str] - verdict: VerifyVerdict - structured: bool - detail: str - - # The delegation verdict keeps "the delegate finished" separate from "the work # was verified" (D5). See :func:`delegation_verdict`. DelegationVerdict = Literal["verified", "failed", "unverified", "incomplete"] diff --git a/src/crossagent/scope.py b/src/crossagent/scope.py index 35e53fd..70aa33a 100644 --- a/src/crossagent/scope.py +++ b/src/crossagent/scope.py @@ -47,7 +47,7 @@ from pathlib import Path from typing import Optional -from .jobs import ScopeResultDict, ScopeStatus +from .types import ScopeResultDict, ScopeStatus _GIT_TIMEOUT_SECONDS = 30.0 # porcelain -z entries are "XY ": two status chars, a space, then the path. diff --git a/src/crossagent/types.py b/src/crossagent/types.py new file mode 100644 index 0000000..e1e2789 --- /dev/null +++ b/src/crossagent/types.py @@ -0,0 +1,97 @@ +"""Shared persisted-record shapes for the delegation gates. + +These ``TypedDict``/``Literal`` definitions describe how the check (S3), scope +(S4), and verification (S5) gate outcomes are serialized into a ``Job`` record. +They live in their own dependency-free module so the gate modules (``check.py``, +``scope.py``, ``verify.py``) can name the persisted shape they produce WITHOUT +importing from ``jobs.py`` — their consumer — which would be an inverted +dependency. ``jobs.py`` imports these for its ``Job`` fields; this module imports +nothing from the package, so no import cycle is possible. +""" + +from __future__ import annotations + +from typing import Literal, Optional, TypedDict + + +class CheckResultDict(TypedDict): + """Persisted outcome of the independent check-gate (slice S3). + + ``command`` is the caller-supplied check string; ``exit_code`` is the + deterministic delegation verdict (D5) — ``0`` means the work verified, any + non-zero value means the delegation *failed* regardless of whether the + delegate process itself exited 0. ``stdout_tail``/``stderr_tail`` are + bounded tails of the check's output. + + A ``Job.check_result`` of ``None`` means *no check ran* (unverified) — that + stays structurally distinct from a check that ran and failed (D7): absent is + never the same as false. + """ + + command: str + exit_code: int + stdout_tail: str + stderr_tail: str + + +# Outcome of the diff-scope assertion (slice S4). ``ok`` — every path the +# delegate modified was inside the declared allowlist; ``violated`` — it wrote +# outside its declared scope; ``undetermined`` — crossagent could not establish +# what changed (e.g. the cwd is not a git repo). ``undetermined`` is FAIL-CLOSED, +# never a pass: a scope check that silently passes when it cannot see the changes +# grants false assurance. +ScopeStatus = Literal["ok", "violated", "undetermined"] + + +class ScopeResultDict(TypedDict): + """Persisted outcome of the diff-scope assertion (slice S4). + + ``declared`` is the caller-supplied allowlist; ``violating_paths`` lists the + repo-relative paths the delegate modified outside it (empty unless + ``status == "violated"``); ``detail`` is a human-readable summary. + + A ``Job.scope_result`` of ``None`` means *no scope was declared* — scope + enforcement was off — which stays structurally distinct from a scope that + was declared and satisfied (D7: absent is never the same as ``ok``). + """ + + declared: list[str] + status: ScopeStatus + violating_paths: list[str] + detail: str + + +# Outcome of the independent verification pass (slice S5). A FRESH peer session +# grades the delegate's artifact supplied as user-turn input (D6), which removes +# the implicit-authorship channel that weakens self-grading — it does NOT claim +# to eliminate self-preference bias, which is a separate documented effect. +# ``pass`` — the verifier returned a machine-checkable verdict of correct. +# ``fail`` — the verifier returned a machine-checkable verdict of wrong. +# ``unverified`` — the verifier ran but produced no machine-checkable verdict +# (free prose, or the advisor lacks a structured-output +# contract). Treated as inconclusive, never as a pass (D4). +# ``error`` — the verifier could not run (no artifact, launch failure). +# ``unverified``/``error`` are inconclusive: they never green a delegation and +# never hard-fail it. Only ``fail`` blocks the green path. +VerifyVerdict = Literal["pass", "fail", "unverified", "error"] + + +class VerifyResultDict(TypedDict): + """Persisted outcome of the independent verification pass (slice S5). + + ``advisor``/``model`` identify the FRESH peer session that graded the work. + ``verdict`` is the deterministic pass/fail/inconclusive outcome; ``structured`` + records whether a machine-checkable contract (e.g. Claude ``--json-schema`` → + ``structured_output``) produced it, versus a JSON object parsed out of a prose + answer. ``detail`` is a human-readable summary (never the raw artifact). + + A ``Job.verify_result`` of ``None`` means *no verification was requested* — + structurally distinct from a verification that ran and failed, or ran and + could not decide (D7: absent is never the same as ``fail`` or ``unverified``). + """ + + advisor: str + model: Optional[str] + verdict: VerifyVerdict + structured: bool + detail: str diff --git a/src/crossagent/verify.py b/src/crossagent/verify.py index af08790..4c269bd 100644 --- a/src/crossagent/verify.py +++ b/src/crossagent/verify.py @@ -56,7 +56,7 @@ from . import parsers as parsers_mod from . import runner as runner_mod from .advisors import Advisor -from .jobs import VerifyResultDict, VerifyVerdict +from .types import VerifyResultDict, VerifyVerdict # A verification that never terminates must not hang the worker forever. VERIFY_DEFAULT_TIMEOUT_SECONDS = 600.0