From bc33bead6f81ee0ab2d3107048ed1dc71dc723d4 Mon Sep 17 00:00:00 2001 From: Dat Date: Sun, 26 Jul 2026 17:07:50 +0700 Subject: [PATCH 1/4] feat(jobs): fold independent-verification gate into the delegation verdict (S5) Add VerifyResultDict + Job.verify_result (schema v3, None = not requested, distinct from a ran-and-failed verify per D7) and refactor delegation_verdict into a gate combiner: any failed gate blocks the green path, a passing gate greens it, an inconclusive/absent gate is neutral. A structured verify 'fail' therefore blocks green even when the shell check passed (D6); a prose-only or errored verifier degrades to inconclusive, never a pass (D4). Surface verify_result in runtime_status. --- src/crossagent/jobs.py | 96 ++++++++++++++++++++++++++++---- tests/test_jobs.py | 121 +++++++++++++++++++++++++++++++++++++++++ 2 files changed, 207 insertions(+), 10 deletions(-) diff --git a/src/crossagent/jobs.py b/src/crossagent/jobs.py index a7fdb6e..358efe9 100644 --- a/src/crossagent/jobs.py +++ b/src/crossagent/jobs.py @@ -71,6 +71,42 @@ class ScopeResultDict(TypedDict): 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"] @@ -228,6 +264,13 @@ class Job: # advisor child. ``None`` means the scrub did not run (a pre-S4 record); an # empty list means it ran and withheld nothing (D7: absent != empty). withheld_env: Optional[list[str]] = None + # --- Independent verification pass (schema v3, slice S5) ------------- + # ``verify_result`` is the outcome of grading the delegate's artifact in a + # FRESH peer session (D6). ``None`` means no verification was requested — + # distinct from a verification that ran and failed or could not decide (D7). + # A ``verdict`` of ``fail`` blocks the green path; ``pass`` can green it; + # ``unverified``/``error`` are inconclusive. See :func:`delegation_verdict`. + verify_result: Optional[VerifyResultDict] = None def delegation_verdict(job: Job) -> DelegationVerdict: @@ -244,25 +287,57 @@ def delegation_verdict(job: Job) -> DelegationVerdict: is NOT a pass: a missing gate is never green. - ``verified`` — the delegate finished cleanly and the check exited 0. - Slice S4 folds the diff-scope assertion into this same verdict rather than - adding a fifth state: a delegate that wrote outside its declared allowlist, - or whose adherence to that allowlist could not be determined, has not - produced trustworthy work — that is semantically a *failed* delegation, so - ``scope_result.status`` other than ``ok`` yields ``failed`` (fail closed). - When no scope was declared (``scope_result is None``) this gate is skipped, - preserving the pre-S4 four-state behaviour and its tests unchanged. + Slices S4 and S5 fold their gates into this same verdict rather than adding + new states. Every gate is combined the same way: **any failed gate blocks + the green path** (fail closed), at least one *passed* gate with no failures + greens the delegation, and a delegation with no decisive gate is + ``unverified`` — a missing gate is never a pass. + + Per-gate mapping (a gate that was not requested contributes nothing): + + - **Scope (S4, security gate):** ``violated``/``undetermined`` → fail (a + write outside the allowlist, or an inability to tell what changed, is never + trustworthy). ``ok`` is neutral — an in-bounds scope does not by itself + verify the *work*, so scope alone never greens a delegation. + - **Check (S3):** exit ``0`` → pass; any non-zero → fail. + - **Verify (S5):** ``pass`` → pass; ``fail`` → fail; ``unverified``/``error`` + → neutral (inconclusive; graceful degradation per D4 — a verifier that + could only produce prose, or could not run, never greens and never + hard-fails). + + A failing independent verification therefore blocks green even when the + shell check passed — which is the entire point of the fresh-session verifier + (D6). When no gate was declared this reduces to the pre-S3 behaviour and its + tests are unchanged. """ if not is_terminal(job.status): return "incomplete" if job.status != JobState.SUCCEEDED: return "failed" + + any_pass = False + scope = job.scope_result if scope is not None and scope.get("status") != "ok": return "failed" + check = job.check_result - if check is None: - return "unverified" - return "verified" if check.get("exit_code") == 0 else "failed" + if check is not None: + if check.get("exit_code") == 0: + any_pass = True + else: + return "failed" + + verify = job.verify_result + if verify is not None: + verdict = verify.get("verdict") + if verdict == "pass": + any_pass = True + elif verdict == "fail": + return "failed" + # "unverified"/"error" are inconclusive: neither green nor a hard fail. + + return "verified" if any_pass else "unverified" # --------------------------------------------------------------------------- @@ -594,6 +669,7 @@ def runtime_status(job: Job) -> dict[str, Any]: "check_result": job.check_result, "scope_result": job.scope_result, "withheld_env": job.withheld_env, + "verify_result": job.verify_result, "delegation_verdict": delegation_verdict(job), } diff --git a/tests/test_jobs.py b/tests/test_jobs.py index d4b1b87..27460d0 100644 --- a/tests/test_jobs.py +++ b/tests/test_jobs.py @@ -1839,3 +1839,124 @@ def test_v3_scope_result_round_trips_on_disk(tmp_path): loaded = load_state(job_dir) assert loaded.scope_result == _scope("violated", violating=["evil.py"]) assert loaded.withheld_env == ["GITHUB_TOKEN"] + + +# ========================================================================= +# Delegation verdict interaction with the independent verification pass +# (slice S5). A fresh-session verifier grades the delegate's artifact; a +# structured "fail" blocks green even when the shell check passed (D6), +# while a prose-only "unverified" degrades gracefully and never greens. +# ========================================================================= + + +def _verify(verdict, *, structured=True): + return { + "advisor": "codex", + "model": "gpt-5.6-sol", + "verdict": verdict, + "structured": structured, + "detail": "", + } + + +def _succeeded_verify(verify_verdict=None, *, check_exit=None, scope_status=None): + check = ( + None + if check_exit is None + else { + "command": "pytest", + "exit_code": check_exit, + "stdout_tail": "", + "stderr_tail": "", + } + ) + return Job( + job_id="job_v", + status=JobState.SUCCEEDED, + check_result=check, + scope_result=None if scope_status is None else _scope(scope_status), + verify_result=None if verify_verdict is None else _verify(verify_verdict), + ) + + +def test_delegation_verdict_verified_when_verification_passes(): + """A structured pass from the fresh verifier greens the delegation even with + no shell check declared.""" + assert delegation_verdict(_succeeded_verify("pass")) == "verified" + + +def test_delegation_verdict_failed_when_verification_fails(): + assert delegation_verdict(_succeeded_verify("fail")) == "failed" + + +def test_delegation_verdict_verify_fail_blocks_green_despite_passing_check(): + """The load-bearing D6 case: the shell check passed, but the independent + verifier says the work is wrong -> failed. A failing verification blocks the + green path a same-context self-grade might have waved through.""" + job = _succeeded_verify("fail", check_exit=0) + assert delegation_verdict(job) == "failed" + + +def test_delegation_verdict_prose_verification_is_unverified_not_pass(): + """An advisor that produced no machine-checkable verdict (prose) degrades to + 'unverified' and never greens: inconclusive is not a pass (D4).""" + job = _succeeded_verify("unverified", check_exit=None) + assert delegation_verdict(job) == "unverified" + + +def test_delegation_verdict_errored_verification_is_inconclusive(): + """A verifier that could not run leaves the delegation unverified, never a + hard fail (the verifier broke, not the work).""" + assert delegation_verdict(_succeeded_verify("error")) == "unverified" + + +def test_delegation_verdict_prose_verify_does_not_downgrade_passing_check(): + """A passing shell check greens the delegation; an inconclusive verify does + not drag it back to unverified.""" + job = _succeeded_verify("unverified", check_exit=0) + assert delegation_verdict(job) == "verified" + + +def test_delegation_verdict_scope_violation_dominates_passing_verification(): + """Fail closed still wins: a scope violation fails the delegation even when + the independent verifier passed.""" + job = _succeeded_verify("pass", scope_status="violated") + assert delegation_verdict(job) == "failed" + + +def test_delegation_verdict_absent_verify_preserves_pre_s5_behaviour(): + job = Job(job_id="job_v", status=JobState.SUCCEEDED, verify_result=None) + assert delegation_verdict(job) == "unverified" + + +def test_runtime_status_surfaces_verify_result(): + job = _succeeded_verify("fail", check_exit=0) + result = runtime_status(job) + assert result["verify_result"]["verdict"] == "fail" + assert result["delegation_verdict"] == "failed" + + +def test_v3_verify_result_round_trips_on_disk(tmp_path): + job_dir = create_job_dir(tmp_path, "job_verify_rt") + job = Job( + job_id="job_verify_rt", + status=JobState.SUCCEEDED, + verify_result=_verify("pass", structured=True), + ) + save_state(job_dir, job) + loaded = load_state(job_dir) + assert loaded.verify_result == _verify("pass", structured=True) + assert delegation_verdict(loaded) == "verified" + + +def test_v2_record_without_verify_result_still_loads(tmp_path): + """A pre-S5 record on disk (no verify_result key) loads with it absent, not + crashing — absent (not requested) stays distinct from failed (D7).""" + job_dir = create_job_dir(tmp_path, "job_v2_verify") + atomic_json_write( + {"schema_version": 2, "job_id": "job_v2_verify", "status": "succeeded"}, + job_dir / "state.json", + ) + loaded = load_state(job_dir) + assert loaded.verify_result is None + assert delegation_verdict(loaded) == "unverified" From c53aa342dc45b6c2b0ebc7e14e1544047744c1f3 Mon Sep 17 00:00:00 2001 From: Dat Date: Sun, 26 Jul 2026 17:07:58 +0700 Subject: [PATCH 2/4] feat(verify): fresh independent verification pass with graceful degradation (S5) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Spawn a FRESH peer session that grades the delegate's artifact supplied as user-turn input (D6) — build_verifier_command adds no resume/fork/name flag and never reads the session registry, removing the implicit-authorship channel that weakens self-grading. Prefer a machine-checkable contract where the advisor exposes one (Claude --json-schema -> structured_output; new Advisor.json_schema_flag), else parse a JSON verdict out of the answer; free prose -> unverified, never a pass. The verifier is a delegate too: credentials scrubbed and CROSSAGENT_* lineage stripped. Never raises — a broken verifier degrades to error/unverified. Scope-honest: removes the implicit-authorship channel; does NOT claim to eliminate self-preference bias (a separate documented effect), and the source paper did not study many-turn agentic settings — crossagent's own setting. --- src/crossagent/advisors.py | 8 + src/crossagent/verify.py | 366 +++++++++++++++++++++++++++++++++++++ tests/test_verify.py | 274 +++++++++++++++++++++++++++ 3 files changed, 648 insertions(+) create mode 100644 src/crossagent/verify.py create mode 100644 tests/test_verify.py diff --git a/src/crossagent/advisors.py b/src/crossagent/advisors.py index ad27043..18926b5 100644 --- a/src/crossagent/advisors.py +++ b/src/crossagent/advisors.py @@ -46,6 +46,13 @@ class Advisor: result_parser: str = "text" resume_command: tuple[str, ...] | None = None session_event_field: str | None = None + # Flag that requests a machine-checkable JSON output contract for the + # independent verification pass (slice S5). ``None`` means the advisor has no + # such contract, so a verifier built on it degrades to parsing a JSON verdict + # out of the answer text (D4 graceful degradation — never a hard failure). + # Claude exposes ``--json-schema`` (research finding [6]: the payload lands in + # ``structured_output``); no other built-in advisor has a verified equivalent. + json_schema_flag: str | None = None experimental: bool = False notes: str = "" @@ -84,6 +91,7 @@ def supports_stream(self) -> bool: session_name_flag="--name", fork_flag="--fork-session", result_parser="claude-stream", + json_schema_flag="--json-schema", ), "codex": Advisor( name="codex", diff --git a/src/crossagent/verify.py b/src/crossagent/verify.py new file mode 100644 index 0000000..7f28184 --- /dev/null +++ b/src/crossagent/verify.py @@ -0,0 +1,366 @@ +"""Independent verification pass for delegated work (slice S5). + +Research finding [5]: asking the producing agent to grade its own work is a +measurably weaker gate than an independent check, and the degradation is +triggered by *implicit* authorship — the artifact sitting in the model's own +prior/current turn — not by being told it authored the work. Decision D6 follows +from that mechanism: the verifier must be a **FRESH peer session** with the +artifact supplied as **user-turn input**. A resumed session, or one where the +artifact arrives as prior assistant context, reintroduces exactly the +implicit-authorship channel this pass exists to remove. + +Scope honesty: supplying the artifact as user-turn input to a fresh session +*removes the implicit-authorship channel*. It does **not** eliminate +self-preference bias — residual self-recognition preference is a separate +documented effect — and the paper behind this design explicitly did not study +many-turn agentic settings, which is precisely crossagent's setting. Treat the +verdict as a stronger-but-not-infallible gate, not a proof of correctness. + +How freshness is guaranteed +--------------------------- +:func:`build_verifier_command` builds the advisor argv itself and never adds a +resume, fork, or session-name flag and never consults the session registry, so +there is no channel by which a prior conversation (the delegate's own, or any +other) can be attached. The artifact is appended as the advisor's ordinary +prompt argument — which every supported CLI treats as a user turn — so it can +never arrive as assistant context. Both properties are asserted by the tests. + +Machine-checkable contract, with graceful degradation (D4) +---------------------------------------------------------- +Where the advisor exposes a JSON-schema contract (Claude ``--json-schema`` → +``structured_output``), the verifier requests it and reads the structured +verdict. Where it does not, the verifier still asks for a JSON verdict in the +prompt and parses one out of the answer; an answer with no parseable verdict is +recorded as ``unverified`` (free prose), never as a pass. A verifier that cannot +run at all is ``error``. Neither ``unverified`` nor ``error`` hard-fails the +delegation — only a machine-readable ``fail`` blocks the green path. + +Security: the verifier is a delegate too. It runs with credential-bearing env +vars scrubbed (the caller's ``--pass-env`` opt-ins excepted) and with +crossagent's own ``CROSSAGENT_*`` lineage variables stripped, so a spawned +verifier can neither inherit ambient secrets nor silently attach to a job tree. +""" + +from __future__ import annotations + +import json +import os +import subprocess +import tempfile +from dataclasses import dataclass +from typing import Any, Optional + +from . import credentials as credentials_mod +from . import parsers as parsers_mod +from . import runner as runner_mod +from .advisors import Advisor +from .jobs import VerifyResultDict, VerifyVerdict + +# A verification that never terminates must not hang the worker forever. +VERIFY_DEFAULT_TIMEOUT_SECONDS = 600.0 +# Keep only a bounded head of the delegate's answer/diff in the artifact — an +# LLM prompt does not need (and should not pay for) an unbounded transcript. +_ARTIFACT_MAX_CHARS = 24000 +_DIFF_MAX_CHARS = 16000 +_GIT_TIMEOUT_SECONDS = 30.0 +_LINEAGE_ENV_PREFIX = "CROSSAGENT_" + +# The machine-checkable contract we ask the verifier to satisfy. Kept tiny on +# purpose: a single verdict token plus a short reason. +_VERDICT_SCHEMA: dict[str, Any] = { + "type": "object", + "properties": { + "verdict": {"type": "string", "enum": ["pass", "fail"]}, + "reason": {"type": "string"}, + }, + "required": ["verdict"], + "additionalProperties": False, +} + +_PROMPT_INSTRUCTIONS = ( + "You are an INDEPENDENT reviewer. You did not write the work below; judge it " + "on its merits. Decide whether the delegated work correctly and completely " + "satisfies the stated task. Do not modify any files. Respond with a single " + 'JSON object and nothing else: {"verdict": "pass" | "fail", "reason": ' + '""}. Use "pass" only if the work is correct and complete; ' + 'otherwise "fail".' +) + + +@dataclass(frozen=True) +class VerifyOutcome: + """The result of one independent verification pass.""" + + advisor: str + model: Optional[str] + verdict: VerifyVerdict + structured: bool + detail: str + + def to_dict(self) -> VerifyResultDict: + """Return the persisted-record shape for ``Job.verify_result``.""" + return { + "advisor": self.advisor, + "model": self.model, + "verdict": self.verdict, + "structured": self.structured, + "detail": self.detail, + } + + +# --------------------------------------------------------------------------- +# Fresh command construction (no session flags — this is the freshness contract) +# --------------------------------------------------------------------------- + + +def build_verifier_command( + advisor: Advisor, model: Optional[str], *, schema_path: Optional[str] = None +) -> list[str]: + """Build the verifier argv for *advisor*, WITHOUT the prompt. + + Deliberately omits every session-attachment flag (resume, fork, name) and + never reads the session registry: a fresh session is the whole point (D6). + The prompt is appended separately by :func:`_append_prompt` as a user turn. + When the advisor supports a JSON-schema contract and *schema_path* is given, + single-shot JSON output plus the schema flag are added so the structured + verdict lands in ``structured_output``. + """ + cmd = [advisor.executable, *advisor.base_args, *advisor.invoke_args] + if model and advisor.model_flag: + cmd.extend([advisor.model_flag, model]) + if advisor.supports_stream: + # Single-shot JSON (not streaming): the terminal event is the whole + # payload, which is the cleanest carrier for a structured verdict. + cmd.extend(advisor.json_args) + if advisor.json_schema_flag and schema_path is not None: + cmd.extend([advisor.json_schema_flag, schema_path]) + return cmd + + +def _append_prompt(cmd: list[str], advisor: Advisor, prompt: str) -> None: + """Append *prompt* as the advisor's user-turn argument. + + Mirrors the delivery rule the runner uses so the artifact is delivered + exactly as a normal user prompt — never as assistant/system context. + """ + delivery = advisor.prompt_delivery + if delivery == "dashdash": + cmd.extend(["--", prompt]) + elif delivery.startswith("flag:"): + cmd.extend([delivery.split(":", 1)[1], prompt]) + else: # "positional" + cmd.append(prompt) + + +# --------------------------------------------------------------------------- +# Artifact assembly (what the verifier grades — supplied as user-turn input) +# --------------------------------------------------------------------------- + + +def build_artifact(task_prompt: str, result_text: Optional[str], cwd: str) -> str: + """Assemble the artifact the verifier grades, as a user-turn string. + + Combines the original task, the delegate's answer, and — when *cwd* is a git + repo — a best-effort working-tree diff of what the delegate changed. Bounded + so the verification prompt stays a reasonable size. + """ + sections = [ + _PROMPT_INSTRUCTIONS, + "\n\n=== TASK GIVEN TO THE DELEGATE ===\n" + task_prompt.strip(), + ] + answer = (result_text or "").strip() + sections.append( + "\n\n=== DELEGATE'S ANSWER ===\n" + + (answer or "(the delegate produced no answer)") + ) + diff = _git_diff(cwd) + if diff: + sections.append("\n\n=== WORKING-TREE DIFF (git diff HEAD) ===\n" + diff) + artifact = "".join(sections) + if len(artifact) > _ARTIFACT_MAX_CHARS: + artifact = artifact[:_ARTIFACT_MAX_CHARS] + "\n...[artifact truncated]" + return artifact + + +def _git_diff(cwd: str) -> Optional[str]: + """Return a bounded ``git diff HEAD`` for *cwd*, or ``None`` on any failure. + + Best-effort only (never raises): the answer is still gradable without a diff, + and a non-repo cwd or a git error must not break verification. git is run as + an argument list with ``shell=False`` (matching ``check.py``/``scope.py``); + no delegate output is interpolated into the command line. + """ + try: + completed = subprocess.run( + ["git", "-c", "core.pager=cat", "diff", "HEAD", "--"], + cwd=cwd, + capture_output=True, + text=True, + timeout=_GIT_TIMEOUT_SECONDS, + check=False, + ) + except (OSError, subprocess.SubprocessError): + return None + if completed.returncode != 0: + return None + diff = completed.stdout + if not diff.strip(): + return None + if len(diff) > _DIFF_MAX_CHARS: + diff = diff[:_DIFF_MAX_CHARS] + "\n...[diff truncated]" + return diff + + +# --------------------------------------------------------------------------- +# Verdict extraction +# --------------------------------------------------------------------------- + + +def _parse_verdict_object(text: str) -> Optional[dict[str, Any]]: + """Extract the verdict JSON object from *text*, tolerant of surrounding prose. + + Tries the whole string first, then the widest ``{...}`` span. Returns ``None`` + when nothing parses to a mapping — i.e. the answer was free prose. + """ + stripped = text.strip() + for candidate in _json_candidates(stripped): + try: + parsed = json.loads(candidate) + except (json.JSONDecodeError, ValueError): + continue + if isinstance(parsed, dict): + return parsed + return None + + +def _json_candidates(text: str) -> list[str]: + candidates = [text] + start = text.find("{") + end = text.rfind("}") + if 0 <= start < end: + candidates.append(text[start : end + 1]) + return candidates + + +def _verdict_from_object(obj: dict[str, Any]) -> tuple[VerifyVerdict, str]: + """Map a parsed verdict object to a ``(verdict, detail)`` pair.""" + raw = obj.get("verdict") + reason = obj.get("reason") + detail = str(reason) if isinstance(reason, str) and reason else "" + if isinstance(raw, bool): + return ("pass" if raw else "fail"), detail + token = str(raw).strip().lower() if raw is not None else "" + if token in ("pass", "passed", "true", "ok", "correct"): + return "pass", detail + if token in ("fail", "failed", "false", "incorrect", "reject", "rejected"): + return "fail", detail + # A JSON object with an unrecognised verdict token is inconclusive, not a + # pass — fall back to prose semantics. + return "unverified", detail or f"unrecognised verdict token: {raw!r}" + + +# --------------------------------------------------------------------------- +# Public entry point +# --------------------------------------------------------------------------- + + +def verifier_env(pass_env: Optional[list[str]] = None) -> dict[str, str]: + """Return the environment for the verifier subprocess. + + Credential-bearing vars are scrubbed (the caller's ``--pass-env`` opt-ins + excepted), and crossagent's own ``CROSSAGENT_*`` lineage vars are stripped so + a nested crossagent inside the verifier cannot attach to a stale job tree. + """ + scrubbed = credentials_mod.scrub_env(os.environ, pass_through=pass_env or []) + return { + key: value + for key, value in scrubbed.items() + if not key.startswith(_LINEAGE_ENV_PREFIX) + } + + +def run_verification( + advisor: Advisor, + model: Optional[str], + artifact: str, + *, + cwd: str, + pass_env: Optional[list[str]] = None, + timeout: float = VERIFY_DEFAULT_TIMEOUT_SECONDS, +) -> VerifyOutcome: + """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. + """ + 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) + cmd = build_verifier_command(advisor, model, schema_path=schema_path) + _append_prompt(cmd, advisor, artifact) + parser = parsers_mod.get_parser(advisor.result_parser) + try: + outcome = runner_mod.run( + cmd, + cwd=cwd, + env=verifier_env(pass_env), + consumer=parser, + max_runtime_seconds=timeout, + ) + except (OSError, ValueError) as exc: + return VerifyOutcome( + advisor.name, + model, + "error", + 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 + if isinstance(outcome.result, parsers_mod.ParsedResult) + else parsers_mod.ParsedResult() + ) + if outcome.timed_out: + return VerifyOutcome( + advisor.name, + model, + "error", + structured, + f"verifier timed out after {timeout:g}s", + ) + if parsed.failure or parsed.result is None: + return VerifyOutcome( + advisor.name, + model, + "error", + structured, + parsed.error or "verifier produced no answer", + ) + + obj = _parse_verdict_object(parsed.result) + if obj is None: + # Free prose with no machine-checkable verdict: inconclusive, not a pass. + return VerifyOutcome( + advisor.name, + model, + "unverified", + structured, + "verifier returned no machine-checkable verdict (free prose)", + ) + verdict, detail = _verdict_from_object(obj) + return VerifyOutcome(advisor.name, model, verdict, structured, detail) diff --git a/tests/test_verify.py b/tests/test_verify.py new file mode 100644 index 0000000..d7bdf3b --- /dev/null +++ b/tests/test_verify.py @@ -0,0 +1,274 @@ +"""Tests for the independent verification pass (slice S5).""" + +from __future__ import annotations + +import subprocess +import sys +import textwrap +from pathlib import Path + +from crossagent import verify as verify_mod +from crossagent.advisors import Advisor, resolve +from crossagent.verify import ( + VerifyOutcome, + _parse_verdict_object, + _verdict_from_object, + build_artifact, + build_verifier_command, + run_verification, + verifier_env, +) + +# --------------------------------------------------------------------------- +# Freshness contract: the verifier command must carry no session attachment. +# --------------------------------------------------------------------------- + +_SESSION_FLAGS = ("--resume", "--fork-session", "--name") + + +def test_verifier_command_has_no_session_flags(): + """A fresh session is the whole point (D6): no resume/fork/name flag may + appear, so no prior conversation can be attached to the verifier.""" + claude = resolve("claude") + cmd = build_verifier_command(claude, "sonnet", schema_path="/tmp/schema.json") + for flag in _SESSION_FLAGS: + assert flag not in cmd, flag + + +def test_verifier_command_requests_structured_output_when_supported(): + claude = resolve("claude") + cmd = build_verifier_command(claude, None, schema_path="/tmp/schema.json") + assert "--json-schema" in cmd + assert "/tmp/schema.json" in cmd + + +def test_verifier_command_omits_schema_flag_for_unsupported_advisor(): + """An advisor with no json_schema_flag never gets a schema flag — it will + degrade to prose parsing, not crash.""" + codex = resolve("codex") + cmd = build_verifier_command(codex, None, schema_path="/tmp/schema.json") + assert "--json-schema" not in cmd + + +def test_artifact_is_delivered_as_the_prompt_argument(): + """The artifact must arrive as the advisor's user-turn prompt, never as + assistant/system context. For claude (dashdash delivery) it is the final + argument after ``--``.""" + claude = resolve("claude") + cmd = build_verifier_command(claude, None) + verify_mod._append_prompt(cmd, claude, "THE-ARTIFACT") + assert cmd[-1] == "THE-ARTIFACT" + assert cmd[-2] == "--" + + +# --------------------------------------------------------------------------- +# Verdict extraction +# --------------------------------------------------------------------------- + + +def test_parse_verdict_object_from_raw_json(): + obj = _parse_verdict_object('{"verdict": "pass", "reason": "ok"}') + assert obj == {"verdict": "pass", "reason": "ok"} + + +def test_parse_verdict_object_extracts_from_surrounding_prose(): + text = 'Here is my review.\n\n{"verdict": "fail", "reason": "bug"}\n\nThanks!' + obj = _parse_verdict_object(text) + assert obj["verdict"] == "fail" + + +def test_parse_verdict_object_returns_none_for_prose(): + assert _parse_verdict_object("The code looks fine to me.") is None + + +def test_verdict_from_object_maps_tokens(): + assert _verdict_from_object({"verdict": "pass"})[0] == "pass" + assert _verdict_from_object({"verdict": "FAILED"})[0] == "fail" + assert _verdict_from_object({"verdict": True})[0] == "pass" + assert _verdict_from_object({"verdict": False})[0] == "fail" + # An unrecognised token is inconclusive, never a pass. + assert _verdict_from_object({"verdict": "maybe"})[0] == "unverified" + + +# --------------------------------------------------------------------------- +# run_verification end-to-end via a fake advisor script +# --------------------------------------------------------------------------- + + +def _fake_advisor( + tmp_path: Path, + body: str, + *, + result_parser: str = "claude-stream", + json_schema_flag: str | None = "--json-schema", +) -> Advisor: + script = tmp_path / "fake_verifier.py" + script.write_text(textwrap.dedent(body), encoding="utf-8") + return Advisor( + name="fakeverify", + executable=sys.executable, + base_args=(str(script),), + prompt_delivery="positional", + result_parser=result_parser, + json_args=(), + json_schema_flag=json_schema_flag, + ) + + +# A fake claude-style advisor: emits a structured_output verdict and records the +# prompt (its last argv) so a test can prove the artifact arrived as user input. +def _structured_advisor(tmp_path: Path, verdict: str) -> Advisor: + return _fake_advisor( + tmp_path, + f""" + import json, sys + with open('prompt_seen.txt', 'w') as handle: + handle.write(sys.argv[-1]) + print(json.dumps({{ + "type": "result", "subtype": "success", + "structured_output": {{"verdict": {verdict!r}, "reason": "because"}}, + }})) + """, + ) + + +def test_run_verification_structured_pass(tmp_path): + advisor = _structured_advisor(tmp_path, "pass") + outcome = run_verification(advisor, None, "ARTIFACT-TEXT", cwd=str(tmp_path)) + assert outcome.verdict == "pass" + assert outcome.structured is True + # The artifact reached the advisor as its user-turn prompt argument. + assert (tmp_path / "prompt_seen.txt").read_text() == "ARTIFACT-TEXT" + + +def test_run_verification_structured_fail(tmp_path): + advisor = _structured_advisor(tmp_path, "fail") + outcome = run_verification(advisor, None, "ARTIFACT", cwd=str(tmp_path)) + assert outcome.verdict == "fail" + + +def test_run_verification_prose_degrades_to_unverified(tmp_path): + """An advisor with no structured-output contract that returns prose degrades + to 'unverified' — never a pass, never a crash (D4).""" + advisor = _fake_advisor( + tmp_path, + """ + print("The delegate's work looks reasonable to me.") + """, + result_parser="text", + json_schema_flag=None, + ) + outcome = run_verification(advisor, None, "ARTIFACT", cwd=str(tmp_path)) + assert outcome.verdict == "unverified" + assert outcome.structured is False + + +def test_run_verification_prose_advisor_can_still_parse_json_answer(tmp_path): + """A text advisor that happens to answer in clean JSON yields a usable + verdict (structured=False but decisive).""" + advisor = _fake_advisor( + tmp_path, + """ + print('{"verdict": "pass", "reason": "all good"}') + """, + result_parser="text", + json_schema_flag=None, + ) + outcome = run_verification(advisor, None, "ARTIFACT", cwd=str(tmp_path)) + assert outcome.verdict == "pass" + assert outcome.structured is False + + +def test_run_verification_no_result_is_error(tmp_path): + """A verifier that emits no result event is recorded as 'error', which is + inconclusive — it must not green or hard-fail the delegation.""" + advisor = _fake_advisor(tmp_path, "pass\n") + outcome = run_verification(advisor, None, "ARTIFACT", cwd=str(tmp_path)) + assert outcome.verdict == "error" + + +def test_run_verification_never_raises_on_missing_executable(tmp_path): + advisor = Advisor( + name="ghost", + executable="/nonexistent/verifier-binary", + result_parser="text", + ) + outcome = run_verification(advisor, None, "ARTIFACT", cwd=str(tmp_path)) + assert isinstance(outcome, VerifyOutcome) + assert outcome.verdict == "error" + + +# --------------------------------------------------------------------------- +# Credential + lineage scrubbing for the verifier (it is a delegate too) +# --------------------------------------------------------------------------- + + +def test_verifier_env_scrubs_credentials(monkeypatch): + monkeypatch.setenv("AWS_SECRET_ACCESS_KEY", "super-secret") + env = verifier_env() + assert "AWS_SECRET_ACCESS_KEY" not in env + + +def test_verifier_env_honours_pass_env_opt_in(monkeypatch): + monkeypatch.setenv("ANTHROPIC_API_KEY", "sk-xxx") + env = verifier_env(["ANTHROPIC_API_KEY"]) + assert env["ANTHROPIC_API_KEY"] == "sk-xxx" + + +def test_verifier_env_strips_lineage_vars(monkeypatch): + monkeypatch.setenv("CROSSAGENT_PARENT_JOB_ID", "job_parent") + monkeypatch.setenv("CROSSAGENT_TRACE_ID", "trace_x") + env = verifier_env() + assert "CROSSAGENT_PARENT_JOB_ID" not in env + assert "CROSSAGENT_TRACE_ID" not in env + + +def test_run_verification_secret_does_not_reach_verifier(tmp_path, monkeypatch): + """End-to-end: a credential env var is withheld from the verifier child.""" + monkeypatch.setenv("AWS_SECRET_ACCESS_KEY", "leak-me") + advisor = _fake_advisor( + tmp_path, + """ + import json, os + with open('env_probe.txt', 'w') as handle: + handle.write(os.environ.get('AWS_SECRET_ACCESS_KEY', 'ABSENT')) + print(json.dumps({"type": "result", "subtype": "success", + "structured_output": {"verdict": "pass"}})) + """, + ) + run_verification(advisor, None, "ARTIFACT", cwd=str(tmp_path)) + assert (tmp_path / "env_probe.txt").read_text() == "ABSENT" + + +# --------------------------------------------------------------------------- +# Artifact assembly includes a git diff when the cwd is a repo +# --------------------------------------------------------------------------- + + +def _git(args, cwd): + subprocess.run( + ["git", *args], cwd=str(cwd), check=True, capture_output=True, text=True + ) + + +def test_build_artifact_includes_git_diff(tmp_path): + repo = tmp_path / "repo" + repo.mkdir() + _git(["init"], repo) + _git(["config", "user.email", "t@e.com"], repo) + _git(["config", "user.name", "T"], repo) + (repo / "app.py").write_text("original\n", encoding="utf-8") + _git(["add", "app.py"], repo) + _git(["commit", "-m", "init"], repo) + (repo / "app.py").write_text("delegate changed this\n", encoding="utf-8") + + artifact = build_artifact("do the task", "I edited app.py", str(repo)) + assert "TASK GIVEN TO THE DELEGATE" in artifact + assert "I edited app.py" in artifact + assert "delegate changed this" in artifact # the diff is present + + +def test_build_artifact_tolerates_non_git_cwd(tmp_path): + artifact = build_artifact("do the task", "the answer", str(tmp_path)) + assert "the answer" in artifact + assert "WORKING-TREE DIFF" not in artifact From d6cdbdcec22009713bbfce88b7bc650093dfa2d3 Mon Sep 17 00:00:00 2001 From: Dat Date: Sun, 26 Jul 2026 17:08:06 +0700 Subject: [PATCH 3/4] feat(escalate): re-dispatch failed delegations up a ladder as same-trace children (S5) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit On a failed delegation (failing check, scope violation, or failing verification -> delegation_verdict == failed), re-dispatch the task to the next larger peer. Recorded as option (a): a same-trace child job with parent_job_id set and the trace_id preserved, resolved through resolve_lineage — exactly what the shipped analytics.py escalation-rate definition counts, so no analytics change is needed. Ladder = ordered advisor[:model] rungs; each hop hands the remainder down, and the check/verify/scope/pass_env posture is propagated so the larger peer is held to the same write boundary and credential withholding. Bounded twice over: the rung list shrinks each hop, and MAX_NESTING_DEPTH is enforced (a LineageError is caught and audited as a skipped escalation, never a crash or a runaway loop). --- src/crossagent/escalate.py | 280 +++++++++++++++++++++++++++++++++++++ tests/test_escalate.py | 266 +++++++++++++++++++++++++++++++++++ 2 files changed, 546 insertions(+) create mode 100644 src/crossagent/escalate.py create mode 100644 tests/test_escalate.py diff --git a/src/crossagent/escalate.py b/src/crossagent/escalate.py new file mode 100644 index 0000000..c4e4ddc --- /dev/null +++ b/src/crossagent/escalate.py @@ -0,0 +1,280 @@ +"""Escalate-on-failure ladder for delegated work (slice S5). + +When a delegation fails — a failing check, a scope violation, or a failing +independent verification, all of which resolve to +:func:`~crossagent.jobs.delegation_verdict` == ``"failed"`` — the caller may +declare an *escalation ladder*: an ordered list of larger peers to re-dispatch +the same task to. On failure the first rung is spawned; the remaining rungs are +handed to that child so a further failure climbs the next rung. + +Recording (option (a) — satisfies the shipped analytics definition) +------------------------------------------------------------------- +``analytics.py`` (merged before this slice) computes the escalation rate as +"a failed delegation counts as escalated if it has a **same-trace child**". So a +re-dispatch is recorded as exactly that: a child job with the failed job as its +``parent_job_id`` and the **same ``trace_id``**, resolved through the existing +:func:`~crossagent.jobs.resolve_lineage` machinery. Nothing in ``analytics.py`` +changes; the escalation column starts reflecting reality the moment this ships. + +Runaway protection +------------------ +An escalation ladder is a recursion source. It is bounded twice over: + +1. The ladder is a finite list that shrinks by one rung each hop, so it + self-terminates even if every rung fails. +2. Each child is one level deeper, and lineage resolution enforces + :data:`~crossagent.jobs.MAX_NESTING_DEPTH`; a rung that would exceed the cap + raises :class:`~crossagent.jobs.LineageError`, which is caught and recorded as + a skipped escalation rather than crashing or looping. + +Security posture carried to the child +------------------------------------- +An escalated child is a delegate too. It goes through the ordinary worker path, +so credential scrubbing (``credentials.py``) and the diff-scope assertion apply +to it unchanged: the declared ``scope_paths`` allowlist and ``pass_env`` opt-ins +are propagated so the larger peer is held to the *same* write boundary and the +*same* credential withholding as the original delegate. The check and +verification gates are propagated too, so each rung's output is judged the same +way before the ladder climbs again. +""" + +from __future__ import annotations + +from pathlib import Path +from typing import Any, Callable, Optional + +from . import advisors as advisors_mod +from . import jobs as jobs_mod +from .advisors import Advisor +from .jobs import Job, delegation_verdict + +# Launches a detached worker for a child job. Injectable so escalation can be +# unit-tested without spawning a real process. +Launcher = Callable[[str, Path], Any] + + +def parse_rungs(rungs: Optional[list[str]]) -> list[tuple[str, Optional[str]]]: + """Parse ``advisor[:model]`` ladder rungs into ``(advisor, model)`` pairs. + + Splits on the first colon only, so a model containing no colon (the usual + case: ``opus``, ``gpt-5.6-sol``) is preserved intact. Blank entries are + dropped so a stray empty ``--escalate-to`` cannot spawn a nameless job. + """ + parsed: list[tuple[str, Optional[str]]] = [] + for raw in rungs or []: + entry = raw.strip() + if not entry: + continue + advisor_name, sep, model = entry.partition(":") + advisor_name = advisor_name.strip() + if not advisor_name: + continue + parsed.append((advisor_name, model.strip() if sep else None)) + return parsed + + +def _build_child_argv(advisor: Advisor, model: Optional[str]) -> list[str]: + """Build the escalated child's advisor argv — a FRESH delegation. + + Mirrors the advisor-invocation core of ``crossagent start`` (default stream + mode) but adds no session-attachment flag: an escalation re-dispatches the + task to a new, larger peer, so there is no prior session to resume. The + child does not inherit the parent's fine-grained invocation flags + (``--tools``, ``--safe-mode``, ``--permission-mode``); its write boundary is + enforced structurally by the propagated scope allowlist instead. + """ + cmd = [advisor.executable, *advisor.base_args, *advisor.invoke_args] + if model and advisor.model_flag: + cmd.extend([advisor.model_flag, model]) + if advisor.supports_stream: + cmd.extend(advisor.stream_args) + return cmd + + +def maybe_escalate( + failed_job: Job, + *, + prompt: str, + state_root: Path, + job_dir: Path, + cwd: str, + registry_path: str, + escalate_to: Optional[list[str]], + check: Optional[str], + check_timeout: float, + scope_paths: Optional[list[str]], + pass_env: list[str], + verify_with: Optional[str], + verify_model: Optional[str], + launcher: Optional[Launcher] = None, +) -> Optional[str]: + """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. + """ + if delegation_verdict(failed_job) != "failed": + return None + + rungs = parse_rungs(escalate_to) + if not rungs: + return None + advisor_name, model = rungs[0] + # ``remaining`` keeps the raw ``advisor[:model]`` strings for the child, minus + # the rung being spawned now — the ladder that a further failure will climb. + remaining = _drop_first_nonblank(escalate_to) + + try: + advisor = advisors_mod.resolve(advisor_name) + except KeyError as exc: + _audit_skip(job_dir, reason=f"unknown escalation advisor: {exc}") + return None + + child_id = jobs_mod.generate_job_id() + try: + parent_id, trace_id, label, depth = jobs_mod.resolve_lineage( + parent_flag=failed_job.job_id, + state_root=state_root, + new_job_id=child_id, + ) + except jobs_mod.LineageError as exc: + # Depth cap reached (or a corrupt chain): the ladder stops here. This is + # the runaway guard doing its job, not an error. + _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) + + jobs_mod.append_event( + job_dir, + "escalation", + actor="system:escalate", + child_job_id=child_id, + advisor=advisor.name, + model=model, + trace_id=trace_id, + depth=depth, + ) + + launch = launcher if launcher is not None else _default_launcher + try: + launch(child_id, state_root) + except OSError as exc: + _audit_skip(job_dir, reason=f"escalation worker failed to launch: {exc}") + return None + return child_id + + +# --------------------------------------------------------------------------- +# Helpers +# --------------------------------------------------------------------------- + + +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()] + return cleaned[1:] + + +def _write_child_prompt(child_dir: Path, prompt: str) -> None: + prompt_path = child_dir / "prompt" + prompt_path.write_text(prompt, encoding="utf-8") + try: + prompt_path.chmod(0o600) + except OSError: + pass + + +def _write_child_command( + child_dir: Path, + *, + 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], +) -> None: + info = { + "command": _build_child_argv(advisor, model), + "prompt_delivery": advisor.prompt_delivery, + "cwd": cwd, + "result_parser": advisor.result_parser, + "registry_path": registry_path, + "key": "", + "name": None, + "model": model or "", + "advisor": advisor.name, + "check": check, + "check_timeout": check_timeout, + # Security posture propagated to the larger peer (same write boundary, + # same credential withholding) — see the module docstring. + "scope_paths": scope_paths, + "pass_env": pass_env, + # The verification gate is re-run on the escalated output, and the + # remaining ladder lets a further failure climb the next rung. + "verify_with": verify_with, + "verify_model": verify_model, + "escalate_to": escalate_to, + } + jobs_mod.atomic_json_write(info, child_dir / "command.json") + + +def _audit_skip(job_dir: Path, *, reason: str) -> None: + jobs_mod.append_event( + job_dir, "escalation_skipped", actor="system:escalate", reason=reason + ) + + +def _default_launcher(child_id: str, state_root: Path) -> Any: + # Lazy import breaks the worker <-> escalate import cycle. + from .worker import start_worker + + return start_worker(child_id, state_root) + + +def _now() -> str: + from datetime import datetime, timezone + + return datetime.now(timezone.utc).isoformat() diff --git a/tests/test_escalate.py b/tests/test_escalate.py new file mode 100644 index 0000000..1bb57ca --- /dev/null +++ b/tests/test_escalate.py @@ -0,0 +1,266 @@ +"""Tests for the escalate-on-failure ladder (slice S5).""" + +from __future__ import annotations + +import json +from datetime import datetime, timezone +from pathlib import Path + +from crossagent import jobs as jobs_mod +from crossagent.escalate import _drop_first_nonblank, maybe_escalate, parse_rungs +from crossagent.jobs import ( + MAX_NESTING_DEPTH, + Job, + JobState, + load_state, + save_state, +) + + +# --------------------------------------------------------------------------- +# Rung parsing +# --------------------------------------------------------------------------- + + +def test_parse_rungs_advisor_only(): + assert parse_rungs(["claude"]) == [("claude", None)] + + +def test_parse_rungs_advisor_and_model(): + assert parse_rungs(["claude:opus", "codex:gpt-5.6-sol"]) == [ + ("claude", "opus"), + ("codex", "gpt-5.6-sol"), + ] + + +def test_parse_rungs_splits_on_first_colon_only(): + assert parse_rungs(["claude:some:model"]) == [("claude", "some:model")] + + +def test_parse_rungs_drops_blank_entries(): + assert parse_rungs(["", " ", "claude"]) == [("claude", None)] + + +def test_parse_rungs_none_is_empty(): + assert parse_rungs(None) == [] + + +def test_drop_first_nonblank_removes_only_the_spawned_rung(): + assert _drop_first_nonblank(["claude:opus", "codex"]) == ["codex"] + assert _drop_first_nonblank(["", "claude", "codex"]) == ["codex"] + assert _drop_first_nonblank(["only"]) == [] + + +# --------------------------------------------------------------------------- +# maybe_escalate +# --------------------------------------------------------------------------- + + +def _failed_parent( + state_root: Path, + *, + job_id="job_parent", + depth=1, + trace="trace_x", + parent_job_id=None, +) -> Job: + """Persist a failed delegation (failing check) as the escalation parent.""" + job_dir = jobs_mod.create_job_dir(state_root, job_id) + now = datetime.now(timezone.utc).isoformat() + job = Job( + job_id=job_id, + status=JobState.SUCCEEDED, # the delegate FINISHED... + advisor="codex", + cwd=str(state_root), + started_at=now, + updated_at=now, + trace_id=trace, + parent_job_id=parent_job_id, + nesting_depth=depth, + # ...but the check failed -> delegation_verdict == "failed". + check_result={ + "command": "pytest", + "exit_code": 1, + "stdout_tail": "", + "stderr_tail": "", + }, + ) + save_state(job_dir, job) + return job + + +def _spawns() -> tuple[list[tuple[str, Path]], object]: + spawned: list[tuple[str, Path]] = [] + + def launcher(child_id: str, state_root: Path) -> None: + spawned.append((child_id, state_root)) + + return spawned, launcher + + +def _escalate(job, state_root, escalate_to, launcher, **overrides): + kwargs = dict( + prompt="do the task", + state_root=state_root, + job_dir=state_root / job.job_id, + cwd=str(state_root), + registry_path=str(state_root / "sessions.json"), + escalate_to=escalate_to, + check="pytest", + check_timeout=30.0, + scope_paths=["src"], + pass_env=["ANTHROPIC_API_KEY"], + verify_with="claude", + verify_model=None, + launcher=launcher, + ) + kwargs.update(overrides) + return maybe_escalate(job, **kwargs) + + +def test_escalation_creates_same_trace_child_with_parent_link(tmp_path): + """Option (a): the re-dispatch is a same-trace child with parent_job_id set, + which is exactly what analytics.py counts as an escalation.""" + state_root = tmp_path / "state" + parent = _failed_parent(state_root) + spawned, launcher = _spawns() + + child_id = _escalate(parent, state_root, ["claude:opus"], launcher) + + assert child_id is not None + child = load_state(state_root / child_id) + assert child.parent_job_id == "job_parent" + assert child.trace_id == "trace_x" # SAME trace + assert child.nesting_depth == 2 # one deeper than the parent + assert child.advisor == "claude" + # The worker was actually launched for the child. + assert spawned == [(child_id, state_root)] + + +def test_escalation_is_recognised_by_analytics(tmp_path): + """The whole point of option (a): the shipped analytics escalation rate must + now see the re-dispatch, without any change to analytics.py.""" + from crossagent.analytics import build_analytics + + state_root = tmp_path / "state" + parent = _failed_parent(state_root) + _spawned, launcher = _spawns() + child_id = _escalate(parent, state_root, ["claude:opus"], launcher) + + parent = load_state(state_root / "job_parent") + child = load_state(state_root / child_id) + escalation = build_analytics([parent, child])["totals"]["escalation"] + assert escalation["failed"] == 1 + assert escalation["escalated"] == 1 + assert escalation["rate"] == 1.0 + + +def test_escalation_propagates_remaining_ladder_and_gates(tmp_path): + state_root = tmp_path / "state" + parent = _failed_parent(state_root) + _spawned, launcher = _spawns() + child_id = _escalate( + parent, state_root, ["claude:opus", "codex:gpt-5.6-sol"], launcher + ) + + command = json.loads((state_root / child_id / "command.json").read_text()) + # The spawned rung is dropped; the rest is handed to the child. + assert command["escalate_to"] == ["codex:gpt-5.6-sol"] + # Security posture + gates carried to the larger peer. + assert command["scope_paths"] == ["src"] + assert command["pass_env"] == ["ANTHROPIC_API_KEY"] + assert command["check"] == "pytest" + assert command["verify_with"] == "claude" + + +def test_no_escalation_when_delegation_did_not_fail(tmp_path): + """A verified/unverified delegation is never escalated.""" + state_root = tmp_path / "state" + job_dir = jobs_mod.create_job_dir(state_root, "job_ok") + now = datetime.now(timezone.utc).isoformat() + job = Job( + job_id="job_ok", + status=JobState.SUCCEEDED, + cwd=str(state_root), + started_at=now, + updated_at=now, + trace_id="trace_ok", + nesting_depth=1, + check_result={ + "command": "pytest", + "exit_code": 0, + "stdout_tail": "", + "stderr_tail": "", + }, + ) + save_state(job_dir, job) + spawned, launcher = _spawns() + assert _escalate(job, state_root, ["claude"], launcher) is None + assert spawned == [] + + +def test_no_escalation_when_no_rungs(tmp_path): + state_root = tmp_path / "state" + parent = _failed_parent(state_root) + spawned, launcher = _spawns() + assert _escalate(parent, state_root, None, launcher) is None + assert _escalate(parent, state_root, [], launcher) is None + assert spawned == [] + + +def test_escalation_halts_at_depth_cap(tmp_path): + """The runaway guard: a parent already at MAX_NESTING_DEPTH cannot spawn a + deeper child. The ladder stops and the skip is audited, not crashed.""" + state_root = tmp_path / "state" + # A parent already at the cap whose ancestor chain is broken: lineage + # resolution falls back to parent.nesting_depth + 1, which exceeds the cap. + parent = _failed_parent( + state_root, depth=MAX_NESTING_DEPTH, parent_job_id="job_missing_ancestor" + ) + spawned, launcher = _spawns() + + result = _escalate(parent, state_root, ["claude:opus"], launcher) + + assert result is None + assert spawned == [] + events = _events(state_root / "job_parent") + skips = [e for e in events if e.get("event") == "escalation_skipped"] + assert len(skips) == 1 + assert ( + "depth" in skips[0]["reason"].lower() or "nesting" in skips[0]["reason"].lower() + ) + + +def test_escalation_unknown_advisor_is_skipped_not_crashed(tmp_path): + state_root = tmp_path / "state" + parent = _failed_parent(state_root) + spawned, launcher = _spawns() + assert _escalate(parent, state_root, ["nosuchadvisor"], launcher) is None + assert spawned == [] + events = _events(state_root / "job_parent") + assert any(e.get("event") == "escalation_skipped" for e in events) + + +def test_escalation_audit_event_on_parent(tmp_path): + state_root = tmp_path / "state" + parent = _failed_parent(state_root) + _spawned, launcher = _spawns() + child_id = _escalate(parent, state_root, ["claude:opus"], launcher) + + events = _events(state_root / "job_parent") + escalations = [e for e in events if e.get("event") == "escalation"] + assert len(escalations) == 1 + assert escalations[0]["child_job_id"] == child_id + assert escalations[0]["advisor"] == "claude" + assert escalations[0]["trace_id"] == "trace_x" + + +def _events(job_dir: Path) -> list[dict]: + path = job_dir / "events.jsonl" + if not path.exists(): + return [] + return [ + json.loads(line) + for line in path.read_text(encoding="utf-8").splitlines() + if line.strip() + ] From f4ff13e0b14f674ce5e14f57caec86aef3ace9a5 Mon Sep 17 00:00:00 2001 From: Dat Date: Sun, 26 Jul 2026 17:08:15 +0700 Subject: [PATCH 4/4] feat(delegate): wire --verify-with and --escalate-to through start and the worker (S5) Add --verify-with/--verify-model/--escalate-to to 'start', persist them in command.json, and load them in the worker. The worker runs the fresh verification pass after the check/scope gates and persists verify_result on the SAME terminal transition (the only writer of terminal state); after that transition it re-dispatches up the escalation ladder when the delegation failed. Refine the result-command verdict line to name the actual failing gate. E2E worker tests drive real jobs through worker_main and reload from disk: a structured verify 'fail' blocks green while the delegate exited 0; a pass yields verified; the verifier argv carries no session flag and the artifact is its final user-turn arg; a prose verifier stays unverified; a failed delegation spawns a same-trace child with parent_job_id set at depth+1. --- src/crossagent/cli.py | 63 ++++++++++++-- src/crossagent/worker.py | 102 +++++++++++++++++++++- tests/test_worker.py | 183 ++++++++++++++++++++++++++++++++++++++- 3 files changed, 340 insertions(+), 8 deletions(-) diff --git a/src/crossagent/cli.py b/src/crossagent/cli.py index 1f0a7ae..379f892 100644 --- a/src/crossagent/cli.py +++ b/src/crossagent/cli.py @@ -387,6 +387,37 @@ def _parse_job_args(subcommand: str, argv: list[str]) -> argparse.Namespace: "vars are withheld from the delegate." ), ) + parser.add_argument( + "--verify-with", + dest="verify_with", + help=( + "Advisor that independently verifies the delegate's work in a " + "FRESH peer session, with the diff/answer supplied as user-turn " + "input (removes the implicit-authorship channel that weakens " + "self-grading). A failing verdict blocks the green path; a " + "prose-only or errored verifier degrades to unverified, never a " + "pass. Omit to leave verification off." + ), + ) + parser.add_argument( + "--verify-model", + dest="verify_model", + help="Model/alias for the --verify-with advisor (advisor default if omitted).", + ) + parser.add_argument( + "--escalate-to", + action="append", + default=None, + dest="escalate_to", + metavar="ADVISOR[:MODEL]", + help=( + "Re-dispatch a FAILED delegation (failing check, scope violation, " + "or failing verification) to this larger peer as a same-trace " + "child job. Repeatable to form an escalation ladder; each rung is " + "tried in turn as the previous fails. Bounded by the nesting-depth " + "cap. Omit to leave escalation off." + ), + ) parser.add_argument("--json", action="store_true") parser.set_defaults(stream=True) elif subcommand == "wait": @@ -657,23 +688,38 @@ def _cmd_result(args: argparse.Namespace) -> int: def _print_verdict(job: jobs_mod.Job, verdict: str, *, file: Any = sys.stderr) -> None: """Print a one-line delegation verdict to *file* (stderr by default).""" - check = job.check_result if verdict == "verified": - print("[crossagent] delegation verified — check passed", file=file) + print("[crossagent] delegation verified — all declared gates passed", file=file) elif verdict == "unverified": print( - "[crossagent] delegation UNVERIFIED — no --check ran; the delegate " - "finished but its work was not checked", + "[crossagent] delegation UNVERIFIED — the delegate finished but no " + "gate confirmed its work (no --check/--verify-with, or an " + "inconclusive verifier)", file=file, ) elif verdict == "failed": - exit_code = check.get("exit_code") if check else None print( - f"[crossagent] delegation FAILED verification — check exited {exit_code}", + f"[crossagent] delegation FAILED verification — {_failed_reason(job)}", file=file, ) +def _failed_reason(job: jobs_mod.Job) -> str: + """Describe why a delegation failed, naming the actual failing gate.""" + if job.status != jobs_mod.JobState.SUCCEEDED: + return f"delegate did not finish cleanly (status {job.status.value})" + scope = job.scope_result + if scope is not None and scope.get("status") != "ok": + return f"scope {scope.get('status')} ({len(scope.get('violating_paths') or [])} path(s))" + check = job.check_result + if check is not None and check.get("exit_code") != 0: + return f"check exited {check.get('exit_code')}" + verify = job.verify_result + if verify is not None and verify.get("verdict") == "fail": + return f"independent verification by {verify.get('advisor')} returned fail" + return "a declared gate did not pass" + + def _metrics_summary(job: jobs_mod.Job) -> str: """Return a one-line advisor-metrics summary, or ``""`` when nothing was measured. Cost/token/duration go to stderr so piped stdout stays the result. @@ -899,6 +945,11 @@ def _write_command_info( # --allow-path was given (enforcement off), distinct from an empty list. "scope_paths": getattr(args, "allow_path", None), "pass_env": getattr(args, "pass_env", []), + # Independent verification + escalation ladder (S5). ``verify_with`` is + # ``None`` when off; ``escalate_to`` is the (possibly empty) rung list. + "verify_with": getattr(args, "verify_with", None), + "verify_model": getattr(args, "verify_model", None), + "escalate_to": getattr(args, "escalate_to", None) or [], } jobs_mod.atomic_json_write(info, job_dir / "command.json") diff --git a/src/crossagent/worker.py b/src/crossagent/worker.py index 079a3df..5d3ce79 100644 --- a/src/crossagent/worker.py +++ b/src/crossagent/worker.py @@ -17,13 +17,16 @@ from pathlib import Path from typing import Any, Optional +from . import advisors as advisors_mod from . import check as check_mod from . import credentials as credentials_mod +from . import escalate as escalate_mod from . import jobs as jobs_mod from . import parsers as parsers_mod from . import registry as reg_mod from . import runner as runner_mod from . import scope as scope_mod +from . import verify as verify_mod # --------------------------------------------------------------------------- @@ -50,6 +53,13 @@ class _JobCommand: # credential env vars the caller opted to pass through to the delegate. scope_paths: Optional[list[str]] pass_env: list[str] + # Independent verification + escalation (slice S5). ``verify_with`` names the + # advisor that grades the delegate's artifact in a fresh session (``None`` = + # off). ``escalate_to`` is the ordered ladder of ``advisor[:model]`` rungs a + # failed delegation is re-dispatched up ([] = off). + verify_with: Optional[str] + verify_model: Optional[str] + escalate_to: list[str] # --------------------------------------------------------------------------- @@ -234,6 +244,12 @@ def _should_cancel() -> bool: # distinctly and drives the delegation verdict to failed, never a pass. scope_result = _run_scope_assertion(command, scope_baseline, job_dir) + # Run the independent verification pass (S5) if a verifier was declared. A + # FRESH peer session grades the delegate's artifact supplied as user-turn + # input (D6). A structured ``fail`` verdict blocks the green path; a + # prose-only or errored verifier degrades to inconclusive, never a pass. + verify_result = _run_verification(command, prompt, parsed.result, job_dir) + now = datetime.now(timezone.utc).isoformat() # Persist the advisor telemetry the parser extracted (S1) and the check-gate # outcome (S3) on the SAME terminal transition — the worker is the only @@ -241,7 +257,7 @@ def _should_cancel() -> bool: # disk. ``parsed`` fields and ``check_result`` already default to unknown / # None when nothing was measured (D4/D7), so this never fails the job and # never turns an unmeasured metric into a zero or an unrun check into a pass. - jobs_mod.transition_to( + job = jobs_mod.transition_to( job, final_state, job_dir=job_dir, @@ -257,6 +273,27 @@ def _should_cancel() -> bool: check_result=check_result, scope_result=scope_result, withheld_env=withheld, + verify_result=verify_result, + ) + + # Escalate-on-failure (S5). Runs AFTER the terminal state is persisted, so + # the failed parent is complete on disk before its same-trace child is + # spawned. maybe_escalate is a no-op unless the delegation FAILED and an + # escalation ladder remains; it never raises and respects MAX_NESTING_DEPTH. + escalate_mod.maybe_escalate( + job, + prompt=prompt, + state_root=state_dir, + job_dir=job_dir, + cwd=command.cwd, + registry_path=command.registry_path, + escalate_to=command.escalate_to, + check=command.check, + check_timeout=command.check_timeout, + scope_paths=command.scope_paths, + pass_env=command.pass_env, + verify_with=command.verify_with, + verify_model=command.verify_model, ) return 0 @@ -318,6 +355,66 @@ def _run_scope_assertion( return outcome.to_dict() +def _run_verification( + command: _JobCommand, + prompt: str, + result_text: Optional[str], + job_dir: Path, +) -> Optional[jobs_mod.VerifyResultDict]: + """Run the independent verification pass (S5), if a verifier was declared. + + Returns ``None`` when no verifier was requested — distinct from a + verification that ran and failed (D7). The verifier is a FRESH peer session + grading the delegate's artifact as user-turn input (D6); it is a delegate + too, so it runs with credentials scrubbed. An unknown verifier advisor or a + delegate that produced no artifact is recorded as an ``error`` outcome + (inconclusive), never a crash and never a silent pass. + """ + if not command.verify_with: + return None + try: + advisor = advisors_mod.resolve(command.verify_with) + except KeyError as exc: + outcome = verify_mod.VerifyOutcome( + command.verify_with, + command.verify_model, + "error", + False, + f"unknown verifier advisor: {exc}", + ) + else: + structured = advisor.json_schema_flag is not None + if result_text is None: + outcome = verify_mod.VerifyOutcome( + advisor.name, + command.verify_model, + "error", + structured, + "delegate produced no artifact to verify", + ) + else: + artifact = verify_mod.build_artifact(prompt, result_text, command.cwd) + outcome = verify_mod.run_verification( + advisor, + command.verify_model, + artifact, + cwd=command.cwd, + pass_env=command.pass_env, + timeout=verify_mod.VERIFY_DEFAULT_TIMEOUT_SECONDS, + ) + # Audit the verdict, never the artifact (which can contain repo content). + jobs_mod.append_event( + job_dir, + "verify", + actor="system:verify", + advisor=outcome.advisor, + model=outcome.model, + verdict=outcome.verdict, + structured=outcome.structured, + ) + return outcome.to_dict() + + # --------------------------------------------------------------------------- # Helpers # --------------------------------------------------------------------------- @@ -342,6 +439,9 @@ def _load_command(job_dir: Path) -> _JobCommand: ), scope_paths=_load_scope_paths(data.get("scope_paths")), pass_env=[str(name) for name in data.get("pass_env", [])], + verify_with=data.get("verify_with"), + verify_model=data.get("verify_model"), + escalate_to=[str(rung) for rung in data.get("escalate_to", [])], ) diff --git a/tests/test_worker.py b/tests/test_worker.py index 74600a1..c404ed5 100644 --- a/tests/test_worker.py +++ b/tests/test_worker.py @@ -131,12 +131,16 @@ def _run_job_through_worker( check: str | None = None, check_timeout: float = 30.0, pass_env: list[str] | None = None, + verify_with: str | None = None, + verify_model: str | None = None, + escalate_to: list[str] | None = None, ) -> Job: """Set up a job whose advisor is a fake claude-stream script, run the worker synchronously, and return the Job reloaded from disk. When *check* is given it is written into command.json so the worker runs - the S3 check-gate after the delegate finishes.""" + the S3 check-gate after the delegate finishes. ``verify_with`` / + ``escalate_to`` drive the S5 verification pass and escalation ladder.""" state_dir = tmp_path / "state" job_dir = jobs_mod.create_job_dir(state_dir, job_id) @@ -153,6 +157,8 @@ def _run_job_through_worker( cwd=str(tmp_path), started_at=now, updated_at=now, + trace_id=f"trace_{job_id}", + nesting_depth=1, ), ) (job_dir / "prompt").write_text("hello", encoding="utf-8") @@ -172,6 +178,11 @@ def _run_job_through_worker( command_info["check_timeout"] = check_timeout if pass_env is not None: command_info["pass_env"] = pass_env + if verify_with is not None: + command_info["verify_with"] = verify_with + command_info["verify_model"] = verify_model + if escalate_to is not None: + command_info["escalate_to"] = escalate_to jobs_mod.atomic_json_write(command_info, job_dir / "command.json") exit_code = worker_main(job_id, state_dir) @@ -635,3 +646,173 @@ def test_scope_module_importable_without_error(): # Guard: the module and its git timeout constant are wired. assert scope_mod._GIT_TIMEOUT_SECONDS > 0 assert pytest is not None + + +# ========================================================================= +# End-to-end independent verification + escalation (slice S5): drive a real +# job through worker_main and reload state from disk. A unit test over the +# combiner alone would not catch the worker forgetting to forward the verify +# result or to spawn the escalation child — this drives the whole path. +# ========================================================================= + +from crossagent.advisors import Advisor # noqa: E402 + + +def _fake_verifier( + tmp_path: Path, verdict: str, *, record_argv: bool = False +) -> Advisor: + """A fake claude-style verifier advisor emitting a structured verdict.""" + record = ( + "import sys\n" + "open('verifier_argv.json','w').write(__import__('json').dumps(sys.argv))\n" + if record_argv + else "" + ) + script = tmp_path / "fake_verifier.py" + script.write_text( + record + + "import json\n" + + "print(json.dumps({'type':'result','subtype':'success'," + + f"'structured_output': {{'verdict': {verdict!r}, 'reason': 'because'}}}}))\n", + encoding="utf-8", + ) + return Advisor( + name="myverifier", + executable=sys.executable, + base_args=(str(script),), + prompt_delivery="positional", + result_parser="claude-stream", + json_args=(), + json_schema_flag="--json-schema", + ) + + +def test_worker_verification_fail_blocks_green_end_to_end(tmp_path, monkeypatch): + """A structured 'fail' from the fresh verifier is persisted and drives the + delegation verdict to failed even though the delegate exited 0 (D6).""" + verifier = _fake_verifier(tmp_path, "fail") + monkeypatch.setattr( + "crossagent.advisors.resolve", lambda name, config_path=None: verifier + ) + + job = _run_job_through_worker(tmp_path, _RESULT_OK, verify_with="myverifier") + + assert job.status == JobState.SUCCEEDED # the delegate DID finish + assert job.verify_result is not None + assert job.verify_result["verdict"] == "fail" + assert job.verify_result["advisor"] == "myverifier" + assert jobs_mod.delegation_verdict(job) == "failed" + + +def test_worker_verification_pass_yields_verified_end_to_end(tmp_path, monkeypatch): + verifier = _fake_verifier(tmp_path, "pass") + monkeypatch.setattr( + "crossagent.advisors.resolve", lambda name, config_path=None: verifier + ) + + job = _run_job_through_worker(tmp_path, _RESULT_OK, verify_with="myverifier") + + assert job.verify_result["verdict"] == "pass" + assert job.verify_result["structured"] is True + assert jobs_mod.delegation_verdict(job) == "verified" + + +def test_worker_verifier_session_is_fresh_and_artifact_is_user_turn( + tmp_path, monkeypatch +): + """The verifier subprocess receives NO session-attachment flag (fresh + session) and the artifact as its final positional argument (user turn).""" + verifier = _fake_verifier(tmp_path, "pass", record_argv=True) + monkeypatch.setattr( + "crossagent.advisors.resolve", lambda name, config_path=None: verifier + ) + + _run_job_through_worker(tmp_path, _RESULT_OK, verify_with="myverifier") + + argv = json.loads((tmp_path / "verifier_argv.json").read_text()) + for flag in ("--resume", "--fork-session", "--name"): + assert flag not in argv, flag + # The artifact is the LAST argument (user-turn input), and it embeds the + # delegate's answer rather than arriving as prior assistant context. + assert "DELEGATE'S ANSWER" in argv[-1] + assert "ok" in argv[-1] + + +def test_worker_unverified_verifier_does_not_green_end_to_end(tmp_path, monkeypatch): + """A prose-only verifier (no structured verdict) degrades to unverified: it + neither greens nor hard-fails the delegation.""" + prose = tmp_path / "fake_prose_verifier.py" + prose.write_text("print('looks fine to me')\n", encoding="utf-8") + verifier = Advisor( + name="prosever", + executable=sys.executable, + base_args=(str(prose),), + prompt_delivery="positional", + result_parser="text", + json_schema_flag=None, + ) + monkeypatch.setattr( + "crossagent.advisors.resolve", lambda name, config_path=None: verifier + ) + + job = _run_job_through_worker(tmp_path, _RESULT_OK, verify_with="prosever") + + assert job.verify_result["verdict"] == "unverified" + assert jobs_mod.delegation_verdict(job) == "unverified" + + +def test_worker_verify_audit_event_records_verdict_not_artifact(tmp_path, monkeypatch): + verifier = _fake_verifier(tmp_path, "fail") + monkeypatch.setattr( + "crossagent.advisors.resolve", lambda name, config_path=None: verifier + ) + _run_job_through_worker(tmp_path, _RESULT_OK, verify_with="myverifier") + + events_path = tmp_path / "state" / "job_e2e" / "events.jsonl" + lines = [ + json.loads(line) + for line in events_path.read_text(encoding="utf-8").splitlines() + if line.strip() + ] + verify_events = [event for event in lines if event.get("event") == "verify"] + assert len(verify_events) == 1 + assert verify_events[0]["verdict"] == "fail" + assert verify_events[0]["advisor"] == "myverifier" + + +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.""" + spawned: list[tuple[str, Path]] = [] + monkeypatch.setattr( + "crossagent.escalate._default_launcher", + lambda child_id, state_root: spawned.append((child_id, state_root)), + ) + + job = _run_job_through_worker( + tmp_path, _RESULT_OK, check=_check_cmd(1), escalate_to=["codex:gpt-5.6-sol"] + ) + + assert jobs_mod.delegation_verdict(job) == "failed" + assert len(spawned) == 1 + child_id, _ = spawned[0] + child = jobs_mod.load_state(tmp_path / "state" / child_id) + assert child.parent_job_id == job.job_id + assert child.trace_id == job.trace_id # SAME trace (option a) + assert child.nesting_depth == 2 + assert child.advisor == "codex" + + +def test_worker_does_not_escalate_a_passing_delegation(tmp_path, monkeypatch): + spawned: list[tuple[str, Path]] = [] + monkeypatch.setattr( + "crossagent.escalate._default_launcher", + lambda child_id, state_root: spawned.append((child_id, state_root)), + ) + + job = _run_job_through_worker( + tmp_path, _RESULT_OK, check=_check_cmd(0), escalate_to=["codex"] + ) + + assert jobs_mod.delegation_verdict(job) == "verified" + assert spawned == []