From e3cd456b923372afae0b7293a7e03980e8e4d166 Mon Sep 17 00:00:00 2001 From: Toni Nowak Date: Sun, 27 Sep 2026 14:50:48 +0200 Subject: [PATCH] Integrate credential-isolated profile triage pilot --- PILOT.md | 6 + TASKS.md | 6 + scripts/pilot_contract_triage_pair.py | 290 ++++++++++- tests/test_pilot_contract_triage_pair.py | 8 +- .../test_pilot_profile_bridge_integration.py | 463 ++++++++++++++++++ 5 files changed, 747 insertions(+), 26 deletions(-) create mode 100644 tests/test_pilot_profile_bridge_integration.py diff --git a/PILOT.md b/PILOT.md index ccda1a3..67f3bdf 100644 --- a/PILOT.md +++ b/PILOT.md @@ -1755,3 +1755,9 @@ The module is not yet integrated into the profile runner and proves neither nati ## VCR397 — Bounded live event observer The shared collector supports an optional event observer without changing default collection. Callback failures retain collected events and return a generic failure status. Three real-subprocess tests verify malformed-line retention, callback-failure preservation and compatibility with the observed token cap without double-counting cached input. The existing 63 collector tests remain green. This is an integration prerequisite, not a live agent comparison. + +## VCR395 — Explicit profile bridge integration + +The profile runner now has explicit supervisor opt-in for remote triage and accepted result statuses. Both arms use equal network and token-budget settings; only treatment receives the credential-isolated bridge. The live event observer confirms the focused failure before the bridge can decide. Private agent measurement and typed bridge receipts survive optional parser failure. Diagnostic delivery, causal ranking and measured provider calls remain separate. + +Offline integration tests cover the arm pipeline, symmetric configuration, deterministic observer acknowledgment, and receipt preservation after timeout or parser failure. The complete local suite passes: 810 tests in 37.714 seconds. This is not native hook delivery or efficacy evidence. The optional completed-turn token cap is observational, not a provider-side spending limit, and enabled network access is not restricted to loopback. A prospective trial still requires frozen code/profile/oracle and green hosted CI. diff --git a/TASKS.md b/TASKS.md index d5e7827..1b8a40b 100644 --- a/TASKS.md +++ b/TASKS.md @@ -1159,3 +1159,9 @@ The module is not yet integrated into the profile runner and proves neither nati ## VCR397 — Bounded live event observer The shared collector supports an optional event observer without changing default collection. Callback failures retain collected events and return a generic failure status. Three real-subprocess tests verify malformed-line retention, callback-failure preservation and compatibility with the observed token cap without double-counting cached input. The existing 63 collector tests remain green. This is an integration prerequisite, not a live agent comparison. + +## VCR395 — Explicit profile bridge integration + +The profile runner now has explicit supervisor opt-in for remote triage and accepted result statuses. Both arms use equal network and token-budget settings; only treatment receives the credential-isolated bridge. The live event observer confirms the focused failure before the bridge can decide. Private agent measurement and typed bridge receipts survive optional parser failure. Diagnostic delivery, causal ranking and measured provider calls remain separate. + +Offline integration tests cover the arm pipeline, symmetric configuration, deterministic observer acknowledgment, and receipt preservation after timeout or parser failure. The complete local suite passes: 810 tests in 37.714 seconds. This is not native hook delivery or efficacy evidence. The optional completed-turn token cap is observational, not a provider-side spending limit, and enabled network access is not restricted to loopback. A prospective trial still requires frozen code/profile/oracle and green hosted CI. diff --git a/scripts/pilot_contract_triage_pair.py b/scripts/pilot_contract_triage_pair.py index 120fd6d..2d0aff7 100644 --- a/scripts/pilot_contract_triage_pair.py +++ b/scripts/pilot_contract_triage_pair.py @@ -34,6 +34,7 @@ build_agent_measurement_receipt, parse_codex_json_events, parse_choice_receipt, ) import pilot_triage_pair as triage_common +import pilot_profile_triage_bridge as profile_bridge ROOT = Path(__file__).resolve().parents[1] @@ -347,6 +348,19 @@ def load_case_profile(path: Path) -> CaseProfile: ) +def _profile_bridge_spec(profile: CaseProfile) -> profile_bridge.ProfileTriageSpec: + hypotheses = list(profile.triage_hypotheses) + for identifier in profile.triage_accepted_ids: + if identifier not in hypotheses: + hypotheses.append(identifier) + return profile_bridge.ProfileTriageSpec( + failure_kinds=profile.triage_kinds, + hypotheses=tuple(hypotheses), + allowed_observations=profile.triage_observations, + rank_hypotheses=profile.rank_hypotheses, + ) + + def _case_treatment_prompt(profile: CaseProfile, advice_policy: str) -> str: evidence_paths = ", ".join(profile.evidence_files.values()) triage_command = shlex.join(_triage_argv(1, profile)) @@ -479,8 +493,11 @@ def pairs(values: list[tuple[str, Any]]) -> dict[str, Any]: return json.loads(text, object_pairs_hook=pairs) -def _validated_triage(payload: Any, focused_exit: int, - profile: CaseProfile | None = None) -> dict[str, Any]: +def _validated_triage( + payload: Any, focused_exit: int, profile: CaseProfile | None = None, + accepted_statuses: tuple[str, ...] | None = None, + bridge_decision_receipt: dict[str, Any] | None = None, +) -> dict[str, Any]: if (not isinstance(payload, str) or len(payload.encode("utf-8")) > triage_common.MAX_TRIAGE_OUTPUT_BYTES): return {"status": "unscored", "candidate_ids": []} @@ -491,7 +508,10 @@ def _validated_triage(payload: Any, focused_exit: int, if not isinstance(value, dict) or value.get("executed") is not False: return {"status": "unscored", "candidate_ids": []} allowed_ids = profile.triage_accepted_ids if profile else TRIAGE_IDS - accepted_statuses = profile.triage_accepted_statuses if profile else ("no-remote-choice",) + accepted_statuses = ( + accepted_statuses if accepted_statuses is not None else + profile.triage_accepted_statuses if profile else ("no-remote-choice",) + ) parsed = parse_choice_receipt( value, choice_type="triage", candidate_ids=allowed_ids, ) @@ -503,6 +523,22 @@ def _validated_triage(payload: Any, focused_exit: int, or value.get("test_failed") is not True ): return {"status": "unscored", "candidate_ids": []} + if parsed.get("status") == "remote-choice": + remote_confirmed = ( + isinstance(bridge_decision_receipt, dict) + and bridge_decision_receipt.get("status") == "remote-choice" + and bridge_decision_receipt.get("observed_exit_status") == focused_exit + and bridge_decision_receipt.get("test_failed") is True + and bridge_decision_receipt.get("executed") is False + and bridge_decision_receipt.get("decision_reason") == "accepted" + and bridge_decision_receipt.get("diagnostic_choice_id") == parsed["candidate_ids"][0] + and bridge_decision_receipt.get("diagnostic_step_ids") == parsed["candidate_ids"] + and bridge_decision_receipt.get("diagnostic_selection_source") in { + "cached_preferred_next_step", "remote_preferred_next_step", + } + ) + if not remote_confirmed: + return {"status": "unscored", "candidate_ids": []} parsed["hypothesis_order"] = [] parsed["hypothesis_ranking_status"] = "not_established" if profile is not None and profile.rank_hypotheses: @@ -544,6 +580,8 @@ def _event_receipts( lines: Iterable[str], event_times: list[float], started: float, *, advice_policy: str = "legacy-required-step", profile: CaseProfile | None = None, + accepted_triage_statuses: tuple[str, ...] | None = None, + profile_bridge_receipt: dict[str, Any] | None = None, ) -> tuple[dict[str, Any], str | None]: lines = list(lines) pending: dict[str, tuple[str, str | None]] = {} @@ -696,7 +734,15 @@ def _event_receipts( if callable(usage_fn): triage_usage = usage_fn(output) if code == 0: - triage_choice = _validated_triage(output, focused_exits[0], profile) + bridge_decision = ( + profile_bridge_receipt.get("decision") + if isinstance(profile_bridge_receipt, dict) else None + ) + triage_choice = _validated_triage( + output, focused_exits[0], profile, + accepted_statuses=accepted_triage_statuses, + bridge_decision_receipt=bridge_decision, + ) if not triage_choice.get("candidate_ids"): triage_output_status = "invalid_output" elif profile is None: @@ -739,6 +785,26 @@ def _event_receipts( "codex_billing_estimate": None, "event_count": len(list(lines)) if isinstance(lines, list) else None, } + if profile_bridge_receipt is not None: + decision = profile_bridge_receipt.get("decision") + if not isinstance(decision, dict): + decision = {} + result["profile_triage_bridge"] = profile_bridge_receipt + result["provider_transport_call_count"] = decision.get( + "provider_transport_call_count", 0 + ) + result["decision_usage_status"] = decision.get( + "decision_usage_status", "not_invoked" + ) + result["decision_usage"] = decision.get("decision_usage") + result["diagnostic_choice_id"] = decision.get("diagnostic_choice_id") + result["diagnostic_selection_source"] = decision.get( + "diagnostic_selection_source", "none" + ) + result["hypothesis_ranking_status"] = decision.get( + "hypothesis_ranking_status", "not_established" + ) + result["hypothesis_order"] = decision.get("hypothesis_order", []) if profile is not None and profile.outcome_mode == "repair": result["agent_git_diff_check_invocation_observed"] = git_diff_check_invoked result["agent_git_diff_check_exit_codes"] = git_diff_check_exits @@ -781,6 +847,10 @@ def _run_arm( measurement_path: Path | None = None, advice_policy: str = "legacy-required-step", profile: CaseProfile | None = None, + profile_bridge_spec: profile_bridge.ProfileTriageSpec | None = None, + accepted_triage_statuses: tuple[str, ...] | None = None, + allow_network: bool = False, + max_tokens: int | None = None, ) -> tuple[dict[str, Any], str | None]: (home / ".codex").mkdir(mode=0o700, parents=True, exist_ok=True) measurement_path = measurement_path or (home / ".codex" / "agent-measurement.json") @@ -789,44 +859,166 @@ def _run_arm( ) env.pop("OPENROUTER_API_KEY", None) started = time.monotonic() + bridge = None + bridge_receipt_path: Path | None = None + bridge_summary: dict[str, Any] | None = None + typed_bridge_receipt: dict[str, Any] | None = None + process = None try: - process = subprocess.Popen( - common._cli_command( - codex, model, reasoning_effort, prompt, - allow_network=False, - ), - cwd=fixture, env=env, stdin=subprocess.DEVNULL, - stdout=subprocess.PIPE, stderr=subprocess.DEVNULL, - ) - lines, times, failure = core._collect_events( - process, started=started, timeout=timeout, preserve_on_failure=True, + if profile_bridge_spec is not None: + if not treatment or profile is None: + raise ValueError("profile bridge is treatment-only and requires a case profile") + bridge_receipt_path = measurement_path.with_name( + measurement_path.stem + "-profile-triage.json" + ) + bridge = profile_bridge.ProfileTriageBridge( + profile_bridge_spec, receipt_path=bridge_receipt_path, + ) + with bridge: + shim_path = profile_bridge.write_python_shim( + home / "profile-triage-bridge", bridge, + ) + env["PATH"] = os.pathsep.join( + item for item in (str(shim_path), env.get("PATH", "")) if item + ) + process = subprocess.Popen( + common._cli_command( + codex, model, reasoning_effort, prompt, + allow_network=allow_network, + ), + cwd=fixture, env=env, stdin=subprocess.DEVNULL, + stdout=subprocess.PIPE, stderr=subprocess.DEVNULL, + ) + focused_ids: set[str] = set() + focused_seen = False + + def observe(event: dict[str, Any]) -> None: + nonlocal focused_seen + event_type = event.get("type") + item = _item(event) + if item is None: + return + event_id = item.get("id") + if not isinstance(event_id, str): + return + if event_type == "item.started" and _focused(item, profile): + focused_ids.add(event_id) + elif (event_type == "item.completed" and event_id in focused_ids + and not focused_seen): + focused_seen = True + code = item.get("exit_code") + output = item.get("aggregated_output") + if not isinstance(output, str): + output = item.get("output") + confirmed = ( + isinstance(code, int) and not isinstance(code, bool) + and code != 0 and _matches_useful_failure(output, profile) + ) + bridge.observe_focused_failure( + code if isinstance(code, int) and not isinstance(code, bool) else 0, + failure_confirmed=confirmed, + ) + + lines, times, failure = core._collect_events( + process, started=started, timeout=timeout, + preserve_on_failure=True, max_tokens=max_tokens, + event_observer=observe, + ) + bridge_summary = bridge.receipt() + else: + process = subprocess.Popen( + common._cli_command( + codex, model, reasoning_effort, prompt, + allow_network=allow_network, + ), + cwd=fixture, env=env, stdin=subprocess.DEVNULL, + stdout=subprocess.PIPE, stderr=subprocess.DEVNULL, + ) + lines, times, failure = core._collect_events( + process, started=started, timeout=timeout, + preserve_on_failure=True, max_tokens=max_tokens, + ) + except Exception: + if process is not None: + try: + if process.poll() is None: + process.kill() + process.wait() + except (OSError, subprocess.SubprocessError): + pass + if bridge is not None: + bridge_summary = bridge.receipt() + failed = _empty_arm( + "profile_bridge_setup_failed" if profile_bridge_spec else "codex_unavailable" ) - except OSError: - return _empty_arm("codex_unavailable"), None + failed["profile_triage_bridge"] = bridge_summary + return failed, None ended = time.monotonic() wall_ms = round((ended - started) * 1000, 2) measurement = build_agent_measurement_receipt( lines, times, started, ended, process.returncode, failure, ) - # Write this bounded, redacted record before the optional task-specific - # event interpretation so a parser error cannot erase early measurements. + # Preserve the bounded event measurement before task-specific parsing. common._private_write( measurement_path, (json.dumps(measurement, sort_keys=True, separators=(",", ":")) + "\n").encode("utf-8"), ) + if bridge_receipt_path is not None and bridge_receipt_path.is_file(): + try: + if bridge_receipt_path.stat().st_size <= profile_bridge.MAX_RESPONSE_BYTES: + loaded = _strict_json(bridge_receipt_path.read_text(encoding="utf-8")) + if isinstance(loaded, dict): + typed_bridge_receipt = loaded + except (OSError, ValueError, TypeError, json.JSONDecodeError): + typed_bridge_receipt = None + + bridge_metadata = {} + if bridge_summary is not None: + decision = bridge_summary.get("decision") + if isinstance(decision, dict): + bridge_metadata = { + "provider_transport_call_count": decision.get("provider_transport_call_count", 0), + "decision_usage_status": decision.get("decision_usage_status", "not_invoked"), + "decision_usage": decision.get("decision_usage"), + "diagnostic_choice_id": decision.get("diagnostic_choice_id"), + "diagnostic_selection_source": decision.get("diagnostic_selection_source", "none"), + "hypothesis_order": decision.get("hypothesis_order", []), + "hypothesis_ranking_status": decision.get("hypothesis_ranking_status", "not_established"), + } + if failure is not None: failed = _empty_arm(failure) failed["completion_ms"] = wall_ms if wall_ms <= MAX_TIMEOUT * 1000 else None failed["cli_exit_code"] = process.returncode failed["agent_measurement"] = measurement failed["event_count"] = len(lines) + if bridge_summary is not None: + failed["profile_triage_bridge"] = bridge_summary + failed["profile_triage_typed_receipt"] = typed_bridge_receipt + failed.update(bridge_metadata) + return failed, None + + try: + observed, answer = _event_receipts( + lines, times, started, advice_policy=advice_policy, profile=profile, + accepted_triage_statuses=accepted_triage_statuses, + profile_bridge_receipt=bridge_summary, + ) + except Exception: + failed = _empty_arm("event_parser_error") + failed.update({ + "completion_ms": wall_ms if wall_ms <= MAX_TIMEOUT * 1000 else None, + "cli_exit_code": process.returncode, + "agent_measurement": measurement, + "event_count": len(lines), + "profile_triage_bridge": bridge_summary, + "profile_triage_typed_receipt": typed_bridge_receipt, + **bridge_metadata, + }) return failed, None - observed, answer = _event_receipts( - lines, times, started, advice_policy=advice_policy, profile=profile, - ) result = { "cli_status": "completed" if process.returncode == 0 else "failed", "failure": None if process.returncode == 0 else "cli_exit_nonzero", @@ -835,6 +1027,10 @@ def _run_arm( "agent_measurement": measurement, **observed, } + if bridge_summary is not None: + result["profile_triage_bridge"] = bridge_summary + result["profile_triage_typed_receipt"] = typed_bridge_receipt + result.update(bridge_metadata) return result, answer @@ -993,6 +1189,9 @@ def run_pair( output_dir: Path, fixture_source: Path = FIXTURE, advice_policy: str = "legacy-required-step", case_profile: CaseProfile | None = None, + remote_profile_triage: bool = False, + accepted_remote_statuses: tuple[str, ...] = (), + max_tokens: int | None = None, ) -> dict[str, Any]: """Run both local-only CLI arms and save private blind artifacts/receipts.""" if case_profile is not None: @@ -1008,6 +1207,24 @@ def run_pair( raise ValueError("model and reasoning effort must be simple identifiers") if not isinstance(advice_policy, str) or advice_policy not in ADVICE_POLICIES: raise ValueError("invalid advice policy") + if type(remote_profile_triage) is not bool: + raise ValueError("remote_profile_triage must be boolean") + if remote_profile_triage: + if case_profile is None: + raise ValueError("remote profile triage requires a validated case profile") + allowed_remote_statuses = {"remote-choice", "no-remote-choice"} + if (not isinstance(accepted_remote_statuses, tuple) + or not accepted_remote_statuses + or len(set(accepted_remote_statuses)) != len(accepted_remote_statuses) + or any(value not in allowed_remote_statuses for value in accepted_remote_statuses)): + raise ValueError("explicit supervisor acceptance statuses are required") + elif accepted_remote_statuses: + raise ValueError("accepted remote statuses require explicit remote profile triage") + if max_tokens is not None and ( + isinstance(max_tokens, bool) or not isinstance(max_tokens, int) + or not 1 <= max_tokens <= core.MAX_EVENT_TOKEN_BUDGET + ): + raise ValueError("max_tokens is outside the supported event-token budget") if not _verify_codex_version(codex): raise ValueError("Codex CLI 0.157.0 is required") @@ -1081,6 +1298,7 @@ def run_pair( ) return receipt + bridge_spec = _profile_bridge_spec(case_profile) if remote_profile_triage else None arms: dict[str, dict[str, Any]] = {} answers: dict[str, str | None] = {} for true_arm in order: @@ -1100,6 +1318,13 @@ def run_pair( else "legacy-required-step"), measurement_path=output_dir / f"{label}-agent-measurement.json", profile=case_profile, + profile_bridge_spec=bridge_spec if true_arm == "treatment" else None, + accepted_triage_statuses=( + accepted_remote_statuses + if remote_profile_triage and true_arm == "treatment" else None + ), + allow_network=remote_profile_triage, + max_tokens=max_tokens, ) if repair_profile: arms[label].setdefault("agent_git_diff_check_invocation_observed", False) @@ -1300,8 +1525,13 @@ def run_pair( "outcome_mode": case_profile.outcome_mode if case_profile else "legacy_contract_triage", "decision_scope": ( "local_fixture_triage_no_remote_choice" if case_profile is None - else "case_profile_local_triage_remote_unavailable" + else "case_profile_profile_bridge_remote_possible" + if remote_profile_triage else "case_profile_local_triage_remote_unavailable" ), + "remote_profile_triage_enabled": remote_profile_triage, + "accepted_remote_statuses": list(accepted_remote_statuses), + "network_access_enabled_for_both_arms": remote_profile_triage, + "observed_token_budget": max_tokens, "task_outcome_acceptance_status": ( "passed" if case_profile and case_profile.outcome_mode == "repair" and repair_task_correctness else "failed" if case_profile and case_profile.outcome_mode == "repair" @@ -1325,6 +1555,9 @@ def run_pair( "agent_measurement_captured": ( output_dir / f"{label}-agent-measurement.json" ).is_file(), + "profile_triage_typed_receipt_captured": ( + output_dir / f"{label}-agent-measurement-profile-triage.json" + ).is_file(), } for label in ("arm-a", "arm-b") }, @@ -1336,7 +1569,8 @@ def run_pair( "advice_policy": advice_policy, "decision_scope": ( "local_fixture_triage_no_remote_choice" if case_profile is None - else "case_profile_local_triage_remote_unavailable" + else "case_profile_profile_bridge_remote_possible" + if remote_profile_triage else "case_profile_local_triage_remote_unavailable" ), "task_correctness_status": "passed" if task_correctness else "failed", "protocol_delivery_status": delivery_status, @@ -1364,6 +1598,13 @@ def main(argv: list[str] | None = None) -> int: parser.add_argument("--advice-policy", choices=ADVICE_POLICIES, default="legacy-required-step", help="nonbinding is an opt-in profile; legacy remains the default") + parser.add_argument("--remote-profile-triage", action="store_true", + help="opt in to a treatment-only local bridge for a validated profile") + parser.add_argument("--accept-profile-triage-status", action="append", + choices=("remote-choice", "no-remote-choice"), default=[], + help="explicit supervisor acceptance status (repeat as needed)") + parser.add_argument("--max-tokens", type=int, + help="observed completed-turn token cap; not a provider-side limit") parser.add_argument("--output-dir", type=Path, required=True, help="new private directory for blind receipts and artifacts") args = parser.parse_args(argv) @@ -1380,6 +1621,9 @@ def main(argv: list[str] | None = None) -> int: timeout=args.timeout, seed=args.seed, output_dir=args.output_dir, advice_policy=args.advice_policy, case_profile=load_case_profile(args.case_profile) if args.case_profile else None, + remote_profile_triage=args.remote_profile_triage, + accepted_remote_statuses=tuple(args.accept_profile_triage_status), + max_tokens=args.max_tokens, ) except (OSError, ValueError, RuntimeError): print(json.dumps({"status": "failed", "failure": "runner_setup_failed"})) diff --git a/tests/test_pilot_contract_triage_pair.py b/tests/test_pilot_contract_triage_pair.py index df6f832..ce3b6d3 100644 --- a/tests/test_pilot_contract_triage_pair.py +++ b/tests/test_pilot_contract_triage_pair.py @@ -170,7 +170,7 @@ def fake_popen(command, *, cwd=None, env=None, **kwargs): }) return FakeProcess(command[-1]) - def collect(process, *, started, timeout, preserve_on_failure=False): + def collect(process, *, started, timeout, preserve_on_failure=False, max_tokens=None, event_observer=None): treatment = process.prompt != runner.BASE_PROMPT lines, times, _ = arm_events( treatment=treatment, include_evidence=evidence, @@ -323,8 +323,10 @@ def test_measurement_snapshot_survives_optional_receipt_parser_failure(self): with mock.patch.object( runner, "_event_receipts", side_effect=KeyError("synthetic optional parser"), ): - with self.assertRaisesRegex(KeyError, "synthetic optional parser"): - self._run_pair(temp) + receipt, _, _ = self._run_pair(temp) + self.assertNotEqual(receipt["status"], "completed") + self.assertTrue(all(arm["failure"] == "event_parser_error" + for arm in receipt["arms"].values())) snapshot = output / "arm-a-agent-measurement.json" self.assertTrue(snapshot.is_file()) self.assertEqual(stat.S_IMODE(snapshot.stat().st_mode), 0o600) diff --git a/tests/test_pilot_profile_bridge_integration.py b/tests/test_pilot_profile_bridge_integration.py new file mode 100644 index 0000000..f7a5601 --- /dev/null +++ b/tests/test_pilot_profile_bridge_integration.py @@ -0,0 +1,463 @@ +from __future__ import annotations + +import hashlib +import json +import os +from pathlib import Path +import shlex +import stat +import subprocess +import sys +import tempfile +import types +import unittest +from unittest import mock + +ROOT = Path(__file__).resolve().parents[1] +sys.path[:0] = [str(ROOT / "scripts"), str(ROOT / "src")] + +import pilot_contract_triage_pair as runner +import pilot_profile_triage_bridge as bridge_module + + +def _item_event(kind: str, event_id: str, command: list[str], **extra: object) -> str: + return json.dumps({ + "type": kind, + "item": { + "id": event_id, + "type": "command_execution", + "command": shlex.join(command), + **extra, + }, + }) + + +class ProfileBridgeArmIntegrationTests(unittest.TestCase): + def _profile(self, root: Path) -> runner.CaseProfile: + fixture = root / "fixture" + (fixture / "tests").mkdir(parents=True) + (fixture / "src").mkdir() + (fixture / "README.md").write_text("contract marker\n", encoding="utf-8") + (fixture / "src/service.py").write_text("VALUE = 0\n", encoding="utf-8") + (fixture / "tests/test_service.py").write_text("assert True\n", encoding="utf-8") + oracle = root / "oracle.py" + oracle.write_text("pass\n", encoding="utf-8") + return runner.CaseProfile( + case_id="bridge-integration", + fixture_source=fixture, + task_prompt="Inspect and classify the synthetic failing test.", + source_file="src/service.py", + focused_test_file="tests/test_service.py", + focused_command=("python", "-m", "unittest", "tests.test_service"), + evidence_files={"contract": "README.md"}, + evidence_markers={"contract": ("contract marker",)}, + failure_markers=("AssertionError: synthetic mismatch",), + triage_kinds=("assertion",), + triage_hypotheses=( + "assertion_behavior_regression", + "assertion_expectation_drift", + ), + triage_accepted_ids=("confirm_behavior_contract",), + triage_accepted_statuses=("no-remote-choice",), + triage_observations={}, + rank_hypotheses=False, + outcome_mode="contract_triage", + oracle_script=oracle, + oracle_sha256=hashlib.sha256(oracle.read_bytes()).hexdigest(), + ) + + def _fake_decision(self, _args, observed_exit): + identifier = "confirm_behavior_contract" + usage = {"input_tokens": 11, "output_tokens": 4} + step = { + "id": identifier, + "title": "Confirm behavior contract", + "instruction": "Confirm the expected behavior with the owner.", + "selection_source": "remote_preferred_next_step", + } + receipt = { + "schema_version": 1, + "status": "remote-choice", + "bridge_request_count": 1, + "observed_exit_status": observed_exit, + "test_failed": True, + "executed": False, + "decision_reason": "accepted", + "diagnostic_step_ids": [identifier], + "diagnostic_choice_id": identifier, + "diagnostic_selection_source": "remote_preferred_next_step", + "hypothesis_order": [], + "hypothesis_ranking_status": "not_established", + "cache_hit": False, + "provider_transport_call_count": 1, + "decision_usage_status": "reported", + "decision_usage": usage, + } + payload = { + "status": "remote-choice", + "observed_exit_status": observed_exit, + "test_failed": True, + "executed": False, + "decision_reason": "accepted", + "cache_hit": False, + "hypothesis_ranking_status": "not_established", + "steps": [step], + "decision_usage": usage, + } + return payload, receipt + + def _fake_cli_source(self, profile: runner.CaseProfile, marker: Path, sleep_after: float = 0) -> str: + triage = runner._triage_argv(1, profile) + focus = list(profile.focused_command) + evidence = ["cat", *profile.evidence_files.values()] + return f""" +import json, os, subprocess, sys, time +def emit(kind, event_id, command, **extra): + item = {{"id": event_id, "type": "command_execution", + "command": __import__("shlex").join(command), **extra}} + print(json.dumps({{"type": kind, "item": item}}), flush=True) +emit("item.started", "focus", {focus!r}) +emit("item.completed", "focus", {focus!r}, exit_code=1, + aggregated_output="AssertionError: synthetic mismatch") +emit("item.started", "evidence", {evidence!r}) +emit("item.completed", "evidence", {evidence!r}, exit_code=0, + aggregated_output="contract marker") +triage = {triage!r} +marker = {str(marker)!r} +deadline = time.monotonic() + 5 +while not os.path.exists(marker) and time.monotonic() < deadline: + time.sleep(0.005) +if not os.path.exists(marker): + raise RuntimeError("collector did not acknowledge focused failure") +emit("item.started", "triage", triage) +completed = subprocess.run(triage, text=True, capture_output=True, check=False) +emit("item.completed", "triage", triage, exit_code=completed.returncode, + aggregated_output=completed.stdout) +print(json.dumps({{"type": "item.completed", "item": {{ + "id": "final", "type": "agent_message", "text": "Done" +}}}}), flush=True) +time.sleep({sleep_after!r}) +""" + + def _collect_with_ack(self, marker: Path): + original = runner.core._collect_events + + def collect(*args, **kwargs): + observer = kwargs.get("event_observer") + + def observing(event): + if observer is not None: + observer(event) + item = event.get("item", {}) + if ( + event.get("type") == "item.completed" + and item.get("id") == "focus" + ): + marker.write_text("observed", encoding="utf-8") + + kwargs["event_observer"] = observing + return original(*args, **kwargs) + + return collect + + def test_timeout_preserves_measurement_and_typed_bridge_receipt(self): + with tempfile.TemporaryDirectory() as temp: + root = Path(temp) + profile = self._profile(root) + output = root / "timeout" + output.mkdir(mode=0o700) + measurement = output / "agent-measurement.json" + marker = measurement.with_suffix(".failure-observed") + source = self._fake_cli_source(profile, marker, sleep_after=2) + real_popen = subprocess.Popen + + with ( + mock.patch.object( + runner.core, "_collect_events", + side_effect=self._collect_with_ack(marker), + ), + mock.patch.object( + runner.common, "_cli_command", + return_value=[sys.executable, "-c", source], + ), + mock.patch.object(runner.subprocess, "Popen", wraps=real_popen), + mock.patch.object( + bridge_module.ProfileTriageBridge, "_decide", + new=self._fake_decision, + ), + ): + result, answer = runner._run_arm( + codex="fake", model="test-model", reasoning_effort="low", + prompt="synthetic", fixture=profile.fixture_source, + home=output / "home", timeout=1, treatment=True, + measurement_path=measurement, profile=profile, + profile_bridge_spec=runner._profile_bridge_spec(profile), + accepted_triage_statuses=("remote-choice",), + allow_network=True, max_tokens=300000, + ) + + self.assertIsNone(answer) + self.assertEqual(result["failure"], "timeout") + self.assertTrue(measurement.is_file()) + persisted_measurement = json.loads(measurement.read_text(encoding="utf-8")) + self.assertEqual(persisted_measurement["status"], "failed") + self.assertTrue(persisted_measurement["collector_failed"]) + typed_receipt = output / "agent-measurement-profile-triage.json" + self.assertTrue(typed_receipt.is_file()) + persisted = json.loads(typed_receipt.read_text(encoding="utf-8")) + self.assertEqual(persisted["status"], "remote-choice") + self.assertEqual(result["profile_triage_typed_receipt"], persisted) + self.assertEqual(result["provider_transport_call_count"], 1) + + def test_treatment_bridge_runs_after_observed_failure_and_receipt_survives_parser_error(self): + with tempfile.TemporaryDirectory() as temp: + root = Path(temp) + profile = self._profile(root) + real_popen = subprocess.Popen + spawn_envs = [] + network_flags = [] + cli_source = [""] + + def cli_command(_codex, _model, _effort, _prompt, *, allow_network=False): + network_flags.append(allow_network) + return [sys.executable, "-c", cli_source[0]] + + def popen(*args, **kwargs): + spawn_envs.append(dict(kwargs["env"])) + return real_popen(*args, **kwargs) + + def run_arm(label: str, parser_failure: bool): + output = root / label + output.mkdir(mode=0o700) + measurement = output / "agent-measurement.json" + marker = measurement.with_suffix(".failure-observed") + cli_source[0] = self._fake_cli_source(profile, marker) + with ( + mock.patch.object( + runner.core, "_collect_events", + side_effect=self._collect_with_ack(marker), + ), + mock.patch.dict(os.environ, {"OPENROUTER_API_KEY": "parent-only-secret"}), + mock.patch.object(runner.common, "_cli_command", side_effect=cli_command), + mock.patch.object(runner.subprocess, "Popen", side_effect=popen), + mock.patch.object( + bridge_module.ProfileTriageBridge, "_decide", + new=self._fake_decision, + ), + ): + if parser_failure: + with mock.patch.object( + runner, "_event_receipts", + side_effect=RuntimeError("synthetic parser failure"), + ): + return runner._run_arm( + codex="offline-fake-codex", model="test-model", + reasoning_effort="low", prompt="synthetic prompt", + fixture=profile.fixture_source, home=output / "home", + timeout=10, treatment=True, measurement_path=measurement, + profile=profile, profile_bridge_spec=runner._profile_bridge_spec(profile), + accepted_triage_statuses=("remote-choice",), + allow_network=True, max_tokens=300000, + ) + return runner._run_arm( + codex="offline-fake-codex", model="test-model", + reasoning_effort="low", prompt="synthetic prompt", + fixture=profile.fixture_source, home=output / "home", + timeout=10, treatment=True, measurement_path=measurement, + profile=profile, profile_bridge_spec=runner._profile_bridge_spec(profile), + accepted_triage_statuses=("remote-choice",), + allow_network=True, max_tokens=300000, + ) + + result, answer = run_arm("parsed", False) + self.assertEqual(result["triage_output_status"], "valid_configured_choice") + self.assertEqual(result["triage"]["status"], "remote-choice") + self.assertEqual(result["provider_transport_call_count"], 1) + self.assertEqual(result["diagnostic_choice_id"], "confirm_behavior_contract") + self.assertEqual(result["decision_usage_status"], "reported") + self.assertEqual(result["decision_usage"], {"input_tokens": 11, "output_tokens": 4}) + self.assertEqual(result["hypothesis_ranking_status"], "not_established") + self.assertEqual(answer, "Done") + self.assertTrue(result["profile_triage_bridge"]["observed_focused_failure"]) + self.assertEqual(result["profile_triage_bridge"]["state"], "completed") + + error_result, error_answer = run_arm("parser-error", True) + self.assertEqual(error_result["failure"], "event_parser_error") + self.assertIsNone(error_answer) + self.assertTrue((root / "parser-error/agent-measurement.json").is_file()) + typed_receipt = root / "parser-error/agent-measurement-profile-triage.json" + self.assertTrue(typed_receipt.is_file()) + persisted = json.loads(typed_receipt.read_text(encoding="utf-8")) + self.assertEqual(persisted["status"], "remote-choice") + self.assertEqual(error_result["profile_triage_typed_receipt"], persisted) + self.assertEqual(error_result["provider_transport_call_count"], 1) + + self.assertEqual(network_flags, [True, True]) + self.assertNotIn("OPENROUTER_API_KEY", spawn_envs[-1]) + self.assertEqual( + stat.S_IMODE((root / "parsed/agent-measurement-profile-triage.json").stat().st_mode), + 0o600, + ) + self.assertEqual( + json.loads((root / "parsed/agent-measurement.json").read_text())["status"], + "completed", + ) + + def test_explicit_remote_mode_gives_equal_egress_but_treatment_only_bridge(self): + with tempfile.TemporaryDirectory() as temp: + root = Path(temp) + profile = self._profile(root) + output = root / "pair-output" + received = [] + + def fake_arm(**kwargs): + received.append(kwargs) + return ({ + "cli_status": "completed", + "completion_ms": 5.0, + "initial_focused_exit": 1, + "focused_exit_codes": [1], + "first_useful_failure_observed": True, + "full_suite_invocation_observed": True, + "full_suite_exit": 1, + "triage_invocation_observed": kwargs["treatment"], + "triage_after_evidence": kwargs["treatment"], + "triage_invalid_invocation_observed": False, + "triage_exit_code": 0 if kwargs["treatment"] else None, + "triage_output_status": ( + "valid_configured_choice" if kwargs["treatment"] else "not_invoked" + ), + "evidence_complete_before_triage": kwargs["treatment"], + }, "answer") + + with ( + mock.patch.object(runner, "_verify_codex_version", return_value=True), + mock.patch.object(runner.core, "_copy_auth", return_value=True), + mock.patch.object(runner, "_run_arm", side_effect=fake_arm), + mock.patch.object(runner, "_run_independent_oracle", return_value={ + "status": "passed", "elapsed_ms": 1.0, + }), + ): + receipt = runner.run_pair( + codex="fake", model="test-model", reasoning_effort="low", + timeout=5, seed=1, output_dir=output, case_profile=profile, + remote_profile_triage=True, + accepted_remote_statuses=("remote-choice",), + max_tokens=300000, + ) + with self.assertRaisesRegex(ValueError, "explicit supervisor"): + runner.run_pair( + codex="fake", model="test-model", reasoning_effort="low", + timeout=5, seed=1, output_dir=root / "rejected", + case_profile=profile, remote_profile_triage=True, + ) + + self.assertEqual(len(received), 2) + self.assertEqual([kwargs["allow_network"] for kwargs in received], [True, True]) + self.assertEqual([kwargs["max_tokens"] for kwargs in received], [300000, 300000]) + treatment_call = next(kwargs for kwargs in received if kwargs["treatment"]) + baseline_call = next(kwargs for kwargs in received if not kwargs["treatment"]) + self.assertIsNotNone(treatment_call["profile_bridge_spec"]) + self.assertEqual(treatment_call["accepted_triage_statuses"], ("remote-choice",)) + self.assertIsNone(baseline_call["profile_bridge_spec"]) + self.assertIsNone(baseline_call["accepted_triage_statuses"]) + self.assertTrue(receipt["remote_profile_triage_enabled"]) + self.assertTrue(receipt["network_access_enabled_for_both_arms"]) + self.assertEqual(receipt["accepted_remote_statuses"], ["remote-choice"]) + self.assertEqual(receipt["observed_token_budget"], 300000) + + def test_explicit_local_abstention_status_remains_valid(self): + with tempfile.TemporaryDirectory() as temp: + root = Path(temp) + profile = self._profile(root) + profile = __import__("dataclasses").replace( + profile, + triage_accepted_ids=( + "assertion_behavior_regression", + "assertion_expectation_drift", + ), + ) + receipt_path = root / "abstention/agent-measurement-profile-triage.json" + measurement_path = root / "abstention/agent-measurement.json" + measurement_path.parent.mkdir(mode=0o700) + marker = measurement_path.with_suffix(".failure-observed") + source = self._fake_cli_source(profile, marker) + real_popen = subprocess.Popen + + def local_decision(_bridge, _args, observed_exit): + identifiers = [ + "assertion_behavior_regression", + "assertion_expectation_drift", + ] + steps = [{ + "id": identifier, + "title": identifier, + "instruction": "Compare the local contract and behavior.", + "selection_source": "unranked_local_fallback", + } for identifier in identifiers] + receipt = { + "schema_version": 1, + "status": "no-remote-choice", + "bridge_request_count": 1, + "observed_exit_status": observed_exit, + "test_failed": True, + "executed": False, + "decision_reason": "local_resolution", + "diagnostic_step_ids": identifiers, + "diagnostic_choice_id": identifiers[0], + "diagnostic_selection_source": "locally_resolved_guidance", + "hypothesis_order": [], + "hypothesis_ranking_status": "not_established", + "cache_hit": False, + "provider_transport_call_count": 0, + "decision_usage_status": "not_invoked", + "decision_usage": None, + } + return ({ + "status": "no-remote-choice", + "observed_exit_status": observed_exit, + "test_failed": True, + "executed": False, + "decision_reason": "local_resolution", + "cache_hit": False, + "hypothesis_ranking_status": "not_established", + "steps": steps, + "decision_usage": None, + }, receipt) + + with ( + mock.patch.object( + runner.core, "_collect_events", + side_effect=self._collect_with_ack(marker), + ), + mock.patch.object( + runner.common, "_cli_command", + return_value=[sys.executable, "-c", source], + ), + mock.patch.object(runner.subprocess, "Popen", wraps=real_popen), + mock.patch.object( + bridge_module.ProfileTriageBridge, "_decide", + new=local_decision, + ), + ): + result, _ = runner._run_arm( + codex="fake", model="test-model", reasoning_effort="low", + prompt="synthetic", fixture=profile.fixture_source, + home=root / "abstention/home", timeout=10, treatment=True, + measurement_path=measurement_path, profile=profile, + profile_bridge_spec=runner._profile_bridge_spec(profile), + accepted_triage_statuses=("no-remote-choice",), + allow_network=True, + ) + + self.assertEqual(result["triage"]["status"], "no-remote-choice") + self.assertEqual(result["triage_output_status"], "valid_configured_choice", result) + self.assertEqual(result["provider_transport_call_count"], 0) + self.assertEqual(result["decision_usage_status"], "not_invoked") + self.assertEqual(result["diagnostic_choice_id"], "assertion_behavior_regression") + self.assertTrue(receipt_path.is_file()) + self.assertEqual(result["profile_triage_typed_receipt"]["status"], "no-remote-choice") + + +if __name__ == "__main__": + unittest.main()