From 8005720a8149c30e9bc0645d024432a5942eacaa Mon Sep 17 00:00:00 2001 From: Toni Nowak Date: Sun, 27 Sep 2026 14:26:49 +0200 Subject: [PATCH 1/4] fix(pilot): require symmetric observed Git validation --- PILOT.md | 4 + TASKS.md | 4 + scripts/pilot_contract_triage_pair.py | 106 +++++++++++++- tests/test_pilot_contract_case_profile.py | 4 + tests/test_pilot_profile_git_validation.py | 158 +++++++++++++++++++++ 5 files changed, 270 insertions(+), 6 deletions(-) create mode 100644 tests/test_pilot_profile_git_validation.py diff --git a/PILOT.md b/PILOT.md index 3b86a86..47ef5a1 100644 --- a/PILOT.md +++ b/PILOT.md @@ -1731,3 +1731,7 @@ This integration is an offline measurement capability, not delivery or efficacy ## VCR391 — Profile advice acknowledgment timing The profile event parser now observes acknowledgment after both legacy and configured valid triage results. A real synthetic event stream verifies acknowledgment before the next tool and distinguishes missing or late acknowledgment. This fixes delivery measurement only; it does not establish native delivery or efficacy. + +## VCR392 — Symmetric repair-arm Git validation + +Repair-profile arms now receive isolated Git baselines before execution. Setup failure prevents launch. The runner records the agent’s exact diff-check invocation and exit codes; missing or unsuccessful checks prevent task acceptance in both arms. Protected fixture digests exclude only internal Git metadata. Four offline tests cover setup failure, metadata handling, whitespace errors and missing/nonzero checks; existing profile tests remain green. No live delivery or efficacy result is established. diff --git a/TASKS.md b/TASKS.md index 1525bc3..cd5211d 100644 --- a/TASKS.md +++ b/TASKS.md @@ -1135,3 +1135,7 @@ The runner now writes the bounded redacted measurement to its private output dir ## VCR391 — Profile advice acknowledgment timing The profile event parser now observes acknowledgment after both legacy and configured valid triage results. A real synthetic event stream verifies acknowledgment before the next tool and distinguishes missing or late acknowledgment. This fixes delivery measurement only; it does not establish native delivery or efficacy. + +## VCR392 — Symmetric repair-arm Git validation + +Repair-profile arms now receive isolated Git baselines before execution. Setup failure prevents launch. The runner records the agent’s exact diff-check invocation and exit codes; missing or unsuccessful checks prevent task acceptance in both arms. Protected fixture digests exclude only internal Git metadata. Four offline tests cover setup failure, metadata handling, whitespace errors and missing/nonzero checks; existing profile tests remain green. No live delivery or efficacy result is established. diff --git a/scripts/pilot_contract_triage_pair.py b/scripts/pilot_contract_triage_pair.py index 5dccdaf..e768c4e 100644 --- a/scripts/pilot_contract_triage_pair.py +++ b/scripts/pilot_contract_triage_pair.py @@ -361,6 +361,7 @@ def _case_treatment_prompt(profile: CaseProfile, advice_policy: str) -> str: guidance += ( f"Complete the requested behavior change in {profile.source_file}. " f"Do not modify {profile.focused_test_file} or configured evidence files. " + "Run git diff --check and require exit code zero after the change. " "Run the focused test again and the required full test suite after the change. " "The supervisor runs a separate immutable oracle." ) @@ -537,6 +538,8 @@ def _event_receipts( evidence: set[str] = set() focused_exits: list[int] = [] full_exits: list[int] = [] + git_diff_check_exits: list[int] = [] + git_diff_check_invoked = False useful_failure_ms: float | None = None useful_failure_observed = False triage_seen = False @@ -624,6 +627,10 @@ def _event_receipts( if kind: pending[event_id] = (kind, None) argv = common._command_argv(item) + if (profile is not None and profile.outcome_mode == "repair" + and argv == ["git", "diff", "--check"]): + git_diff_check_invoked = True + pending[event_id] = ("git-diff-check", None) if argv and "jevcompass" in argv: triage_seen = True exact = _is_triage(item, focused_exits[0] if focused_exits else None, profile) @@ -643,7 +650,9 @@ def _event_receipts( elif event_type == "item.completed" and event_id in pending: kind, _ = pending.pop(event_id) code = item.get("exit_code") - if kind == "focused" and isinstance(code, int) and not isinstance(code, bool): + if kind == "git-diff-check" and isinstance(code, int) and not isinstance(code, bool): + git_diff_check_exits.append(code) + elif kind == "focused" and isinstance(code, int) and not isinstance(code, bool): focused_exits.append(code) output = item.get("aggregated_output") if not isinstance(output, str): @@ -718,6 +727,16 @@ def _event_receipts( "codex_billing_estimate": None, "event_count": len(list(lines)) if isinstance(lines, list) else None, } + 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 + result["agent_git_diff_check_exit_code"] = ( + git_diff_check_exits[-1] if git_diff_check_exits else None + ) + result["agent_git_diff_check_passed"] = bool( + git_diff_check_invoked and git_diff_check_exits + and all(code == 0 for code in git_diff_check_exits) + ) if advice_policy == "nonbinding": result["workflow_acknowledgment"] = workflow_acknowledgment result["triage_result_acknowledgment"] = triage_result_acknowledgment @@ -807,8 +826,58 @@ def _run_arm( return result, answer -def _fixture_digest(root: Path, *, exclude_relative: str | None = None) -> str | None: - """Hash fixture files while ignoring generated Python bytecode.""" +def _initialize_private_git_baseline(root: Path) -> bool: + """Initialize an isolated Git baseline for one copied repair fixture.""" + git = shutil.which("git") + if not git: + return False + env = { + "PATH": os.environ.get("PATH", ""), + "HOME": str(root), + "GIT_CONFIG_NOSYSTEM": "1", + "GIT_CONFIG_GLOBAL": os.devnull, + "GIT_TERMINAL_PROMPT": "0", + "GIT_AUTHOR_NAME": "JevCompass Fixture", + "GIT_AUTHOR_EMAIL": "fixture@example.invalid", + "GIT_COMMITTER_NAME": "JevCompass Fixture", + "GIT_COMMITTER_EMAIL": "fixture@example.invalid", + } + commands = ( + [git, "init", "--quiet"], + [git, "add", "-A"], + [git, "-c", "commit.gpgsign=false", "-c", "core.hooksPath=/dev/null", + "commit", "--quiet", "-m", "Immutable synthetic fixture baseline"], + ) + try: + for argv in commands: + completed = subprocess.run( + argv, cwd=root, env=env, stdout=subprocess.DEVNULL, + stderr=subprocess.DEVNULL, timeout=10, check=False, + ) + if completed.returncode != 0: + return False + except (OSError, subprocess.TimeoutExpired): + return False + return True + + +def _repair_git_diff_check_valid(arm: dict[str, Any]) -> bool: + """Require an observed, successful diff check for repair-profile completion.""" + exits = arm.get("agent_git_diff_check_exit_codes") + return bool( + arm.get("agent_git_diff_check_invocation_observed") is True + and isinstance(exits, list) and exits + and all(isinstance(code, int) and not isinstance(code, bool) and code == 0 + for code in exits) + and arm.get("agent_git_diff_check_passed") is True + ) + + +def _fixture_digest( + root: Path, *, exclude_relative: str | None = None, + exclude_internal_git: bool = False, +) -> str | None: + """Hash fixture files while ignoring bytecode and optional private .git metadata.""" digest = hashlib.sha256() try: for path in sorted(root.rglob("*")): @@ -817,6 +886,8 @@ def _fixture_digest(root: Path, *, exclude_relative: str | None = None) -> str | relative_text = path.relative_to(root).as_posix() if relative_text == exclude_relative: continue + if exclude_internal_git and (relative_text == ".git" or relative_text.startswith(".git/")): + continue relative = relative_text.encode("utf-8") digest.update(len(relative).to_bytes(4, "big")) digest.update(relative) @@ -963,9 +1034,11 @@ def run_pair( for name in (source_file_rel, test_file_rel) } initial_fixture_digest = _fixture_digest(fixture_source) + repair_profile = bool(case_profile and case_profile.outcome_mode == "repair") initial_protected_digest = _fixture_digest( fixture_source, exclude_relative=source_file_rel, - ) if case_profile and case_profile.outcome_mode == "repair" else initial_fixture_digest + exclude_internal_git=repair_profile, + ) if repair_profile else initial_fixture_digest for true_arm in ("baseline", "treatment"): fixtures[true_arm] = private_root / f"{true_arm}-fixture" homes[true_arm] = private_root / f"{true_arm}-home" @@ -982,6 +1055,19 @@ def run_pair( return receipt if digests["baseline"] != digests["treatment"]: raise RuntimeError("fixture parity verification failed") + if repair_profile and not all( + _initialize_private_git_baseline(fixtures[name]) + for name in ("baseline", "treatment") + ): + receipt = { + "schema_version": 1, "run_id": run_id, "status": "failed", + "failure": "repair_git_setup_failed", "arms": {}, + } + common._private_write( + output_dir / "receipt.json", + (json.dumps(receipt, sort_keys=True, indent=2) + "\n").encode("utf-8"), + ) + return receipt arms: dict[str, dict[str, Any]] = {} answers: dict[str, str | None] = {} @@ -1003,6 +1089,11 @@ def run_pair( measurement_path=output_dir / f"{label}-agent-measurement.json", profile=case_profile, ) + if repair_profile: + arms[label].setdefault("agent_git_diff_check_invocation_observed", False) + arms[label].setdefault("agent_git_diff_check_exit_codes", []) + arms[label].setdefault("agent_git_diff_check_exit_code", None) + arms[label].setdefault("agent_git_diff_check_passed", False) arms[label]["fixture_sha256_before"] = digests[true_arm] source_unchanged: dict[str, bool] = {} @@ -1026,8 +1117,10 @@ def run_pair( ) protected_fixture_unchanged[label] = ( initial_protected_digest is not None - and _fixture_digest(fixtures[true_arm], exclude_relative=source_file_rel) - == initial_protected_digest + and _fixture_digest( + fixtures[true_arm], exclude_relative=source_file_rel, + exclude_internal_git=repair_profile, + ) == initial_protected_digest ) source_bytes = _safe_artifact(fixtures[true_arm] / source_file_rel) if source_bytes is not None: @@ -1070,6 +1163,7 @@ def run_pair( and source_changed.get(label) is True and tests_unchanged.get(label) is True and protected_fixture_unchanged.get(label) is True + and _repair_git_diff_check_valid(arm) ) arm["validated_completion_ms"] = endpoint_ms if completion_valid else None arm["task_outcome_status"] = ( diff --git a/tests/test_pilot_contract_case_profile.py b/tests/test_pilot_contract_case_profile.py index ae4a850..c5077be 100644 --- a/tests/test_pilot_contract_case_profile.py +++ b/tests/test_pilot_contract_case_profile.py @@ -160,6 +160,10 @@ def fake_arm(**kwargs): "first_useful_failure_observed": True, "full_suite_invocation_observed": True, "full_suite_exit": 0 if repair_succeeds else 1, + "agent_git_diff_check_invocation_observed": repair_succeeds, + "agent_git_diff_check_exit_codes": [0] if repair_succeeds else [], + "agent_git_diff_check_exit_code": 0 if repair_succeeds else None, + "agent_git_diff_check_passed": repair_succeeds, } if advice_policy == "nonbinding" and kwargs["treatment"]: result.update({ diff --git a/tests/test_pilot_profile_git_validation.py b/tests/test_pilot_profile_git_validation.py new file mode 100644 index 0000000..90d1cac --- /dev/null +++ b/tests/test_pilot_profile_git_validation.py @@ -0,0 +1,158 @@ +from __future__ import annotations + +import hashlib +import json +import shlex +import subprocess +import sys +import tempfile +import types +import unittest +from pathlib import Path +from unittest.mock import patch + +ROOT = Path(__file__).resolve().parents[1] +sys.path.insert(0, str(ROOT / "scripts")) + +import pilot_contract_triage_pair as runner + + +def _event(kind: str, event_id: str, argv: list[str], **extra: object) -> str: + item = { + "id": event_id, + "type": "command_execution", + "command": shlex.join(argv), + **extra, + } + return json.dumps({"type": kind, "item": item}) + + +def _repair_diff_receipt(check_argv: list[str] | None, exit_code: int | None): + lines = [] + if check_argv is not None: + lines = [ + _event("item.started", "diff", check_argv), + _event("item.completed", "diff", check_argv, exit_code=exit_code), + ] + result, _ = runner._event_receipts( + lines, [1.0, 2.0], 0.0, + profile=types.SimpleNamespace(outcome_mode="repair", focused_command=(), evidence_files={}), + ) + return result + + +class PilotProfileGitValidationTests(unittest.TestCase): + def _profile(self, fixture: Path, oracle: Path) -> runner.CaseProfile: + oracle.write_text("pass\n", encoding="utf-8") + return runner.CaseProfile( + case_id="synthetic-git-validation", + fixture_source=fixture, + task_prompt="Repair the copied source.", + source_file="src/service.py", + focused_test_file="tests/test_service.py", + focused_command=("python", "-m", "unittest", "tests.test_service"), + evidence_files={}, + evidence_markers={}, + failure_markers=("AssertionError",), + triage_kinds=(), + triage_hypotheses=(), + triage_accepted_ids=(), + triage_accepted_statuses=(), + triage_observations={}, + rank_hypotheses=False, + outcome_mode="repair", + oracle_script=oracle, + oracle_sha256=hashlib.sha256(oracle.read_bytes()).hexdigest(), + ) + + def test_repair_pair_aborts_before_agent_when_git_is_unavailable(self): + with tempfile.TemporaryDirectory() as temp: + root = Path(temp) + fixture = root / "fixture" + (fixture / "src").mkdir(parents=True) + (fixture / "tests").mkdir() + (fixture / "src/service.py").write_text("VALUE = 1\n", encoding="utf-8") + (fixture / "tests/test_service.py").write_text("assert True\n", encoding="utf-8") + profile = self._profile(fixture, root / "oracle.py") + output = root / "output" + with ( + patch.object(runner, "_verify_codex_version", return_value=True), + patch.object(runner.core, "_copy_auth", return_value=True), + patch.object(runner.shutil, "which", return_value=None), + patch.object(runner, "_run_arm") as run_arm, + ): + receipt = runner.run_pair( + codex="codex", model="test-model", reasoning_effort="low", + timeout=2, seed=4, output_dir=output, case_profile=profile, + ) + self.assertEqual(receipt["failure"], "repair_git_setup_failed") + self.assertEqual(receipt["status"], "failed") + run_arm.assert_not_called() + self.assertEqual( + json.loads((output / "receipt.json").read_text())["failure"], + "repair_git_setup_failed", + ) + + def test_private_git_baseline_and_protected_digest_ignore_only_git_metadata(self): + with tempfile.TemporaryDirectory() as temp: + fixture = Path(temp) + (fixture / "src").mkdir() + (fixture / "src/service.py").write_text("VALUE = 1\n", encoding="utf-8") + before = runner._fixture_digest( + fixture, exclude_relative="src/service.py", exclude_internal_git=True, + ) + self.assertTrue(runner._initialize_private_git_baseline(fixture)) + after = runner._fixture_digest( + fixture, exclude_relative="src/service.py", exclude_internal_git=True, + ) + self.assertEqual(before, after) + self.assertNotEqual( + runner._fixture_digest(fixture, exclude_relative="src/service.py"), + before, + ) + + def test_whitespace_patch_fails_real_diff_check_and_repair_acceptance(self): + with tempfile.TemporaryDirectory() as temp: + fixture = Path(temp) + (fixture / "src").mkdir() + source = fixture / "src/service.py" + source.write_text("VALUE = 1\n", encoding="utf-8") + self.assertTrue(runner._initialize_private_git_baseline(fixture)) + source.write_text("VALUE = 2 \n", encoding="utf-8") + check = subprocess.run( + ["git", "diff", "--check"], cwd=fixture, + stdout=subprocess.PIPE, stderr=subprocess.PIPE, check=False, + ) + self.assertNotEqual(check.returncode, 0) + receipt, _ = runner._event_receipts( + [ + _event("item.started", "diff", ["git", "diff", "--check"]), + _event("item.completed", "diff", ["git", "diff", "--check"], + exit_code=check.returncode, aggregated_output=check.stderr.decode()), + ], + [1.0, 2.0], 0.0, + profile=types.SimpleNamespace(outcome_mode="repair", focused_command=(), evidence_files={}), + ) + self.assertTrue(receipt["agent_git_diff_check_invocation_observed"]) + self.assertEqual(receipt["agent_git_diff_check_exit_code"], check.returncode) + self.assertFalse(receipt["agent_git_diff_check_passed"]) + self.assertFalse(runner._repair_git_diff_check_valid(receipt)) + + def test_missing_or_nonzero_observed_diff_check_cannot_accept_repair(self): + missing = _repair_diff_receipt(None, None) + self.assertFalse(missing["agent_git_diff_check_invocation_observed"]) + self.assertEqual(missing["agent_git_diff_check_exit_codes"], []) + self.assertFalse(runner._repair_git_diff_check_valid(missing)) + + nonzero = _repair_diff_receipt(["git", "diff", "--check"], 2) + self.assertTrue(nonzero["agent_git_diff_check_invocation_observed"]) + self.assertEqual(nonzero["agent_git_diff_check_exit_codes"], [2]) + self.assertFalse(nonzero["agent_git_diff_check_passed"]) + self.assertFalse(runner._repair_git_diff_check_valid(nonzero)) + + accepted = _repair_diff_receipt(["git", "diff", "--check"], 0) + self.assertTrue(runner._repair_git_diff_check_valid(accepted)) + + +if __name__ == "__main__": + unittest.main() From deba15ad5f68a00ba9b07f70a9c1553bc8f0623b Mon Sep 17 00:00:00 2001 From: Toni Nowak Date: Sun, 27 Sep 2026 14:28:20 +0200 Subject: [PATCH 2/4] fix(pilot): separate diagnostic and causal ranking validity --- PILOT.md | 4 ++ TASKS.md | 4 ++ scripts/pilot_contract_triage_pair.py | 38 +++++++++----- tests/test_pilot_profile_rank_dimensions.py | 57 +++++++++++++++++++++ 4 files changed, 90 insertions(+), 13 deletions(-) create mode 100644 tests/test_pilot_profile_rank_dimensions.py diff --git a/PILOT.md b/PILOT.md index 6e96f9f..b13b4c2 100644 --- a/PILOT.md +++ b/PILOT.md @@ -1741,3 +1741,7 @@ Repair-profile arms now receive isolated Git baselines before execution. Setup f Staged a new synthetic URL query migration task with versioned contracts, historical evidence, immutable tests, and an external hash-pinned oracle. Offline checks establish the expected initial failure, a source-only reference repair, and rejection of protected-file tampering or unexpected files. The oracle independently runs focused/full checks and checks the source patch. This is preparation, not a launched pair or efficacy evidence. Launch remains gated on symmetric arm Git initialization and observed agent diff validation. The accepted diagnostic metadata includes contract confirmation, but the current runner does not yet request that candidate; causal ranking and diagnostic selection require separate integration. + +## VCR394 — Independent diagnostic and causal ranking receipts + +A valid diagnostic next step is now scored separately from complete causal ordering. Unestablished, partial or malformed rankings remain explicitly non-complete without erasing valid diagnostic delivery. Complete ranking requires a unique full causal permutation; contract confirmation remains a requested diagnostic candidate only. Three offline tests cover these boundaries. No live delivery or efficacy claim. diff --git a/TASKS.md b/TASKS.md index 5449beb..a35fd51 100644 --- a/TASKS.md +++ b/TASKS.md @@ -1145,3 +1145,7 @@ Repair-profile arms now receive isolated Git baselines before execution. Setup f Staged a new synthetic URL query migration task with versioned contracts, historical evidence, immutable tests, and an external hash-pinned oracle. Offline checks establish the expected initial failure, a source-only reference repair, and rejection of protected-file tampering or unexpected files. The oracle independently runs focused/full checks and checks the source patch. This is preparation, not a launched pair or efficacy evidence. Launch remains gated on symmetric arm Git initialization and observed agent diff validation. The accepted diagnostic metadata includes contract confirmation, but the current runner does not yet request that candidate; causal ranking and diagnostic selection require separate integration. + +## VCR394 — Independent diagnostic and causal ranking receipts + +A valid diagnostic next step is now scored separately from complete causal ordering. Unestablished, partial or malformed rankings remain explicitly non-complete without erasing valid diagnostic delivery. Complete ranking requires a unique full causal permutation; contract confirmation remains a requested diagnostic candidate only. Three offline tests cover these boundaries. No live delivery or efficacy claim. diff --git a/scripts/pilot_contract_triage_pair.py b/scripts/pilot_contract_triage_pair.py index e768c4e..120fd6d 100644 --- a/scripts/pilot_contract_triage_pair.py +++ b/scripts/pilot_contract_triage_pair.py @@ -414,7 +414,11 @@ def _triage_argv(exit_code: int, profile: CaseProfile | None = None) -> list[str suffix = [] for kind in profile.triage_kinds: suffix.extend(("--kind", kind)) - for hypothesis in profile.triage_hypotheses: + candidates = list(profile.triage_hypotheses) + if ("confirm_behavior_contract" in profile.triage_accepted_ids + and "confirm_behavior_contract" not in candidates): + candidates.append("confirm_behavior_contract") + for hypothesis in candidates: suffix.extend(("--hypothesis", hypothesis)) if profile.rank_hypotheses: suffix.append("--rank-hypotheses") @@ -499,20 +503,28 @@ def _validated_triage(payload: Any, focused_exit: int, or value.get("test_failed") is not True ): return {"status": "unscored", "candidate_ids": []} + parsed["hypothesis_order"] = [] + parsed["hypothesis_ranking_status"] = "not_established" if profile is not None and profile.rank_hypotheses: order = value.get("hypothesis_order") - if (value.get("hypothesis_ranking_status") != "complete" - or not isinstance(order, list) - or len(order) != len(profile.triage_hypotheses) - or any(not isinstance(item, str) for item in order) - or set(order) != set(profile.triage_hypotheses) - or len(set(order)) != len(order)): - return {"status": "unscored", "candidate_ids": []} - parsed["hypothesis_ranking_status"] = "complete" - parsed["hypothesis_order"] = list(order) - else: - parsed["hypothesis_ranking_status"] = "not_established" - parsed["hypothesis_order"] = [] + status = value.get("hypothesis_ranking_status") + causal_ids = { + item for item in profile.triage_hypotheses + if item != "confirm_behavior_contract" + } + if (status == "complete" and isinstance(order, list) + and all(isinstance(item, str) for item in order) + and len(order) == len(causal_ids) + and len(set(order)) == len(order) + and set(order) == causal_ids and len(causal_ids) >= 2): + parsed["hypothesis_ranking_status"] = "complete" + parsed["hypothesis_order"] = list(order) + elif status == "not_established" and order == []: + parsed["hypothesis_ranking_status"] = "not_established" + else: + # A valid next step remains valid independently of an incomplete + # or malformed causal ordering; never promote that order to success. + parsed["hypothesis_ranking_status"] = "incomplete" return parsed diff --git a/tests/test_pilot_profile_rank_dimensions.py b/tests/test_pilot_profile_rank_dimensions.py new file mode 100644 index 0000000..9f3a60d --- /dev/null +++ b/tests/test_pilot_profile_rank_dimensions.py @@ -0,0 +1,57 @@ +"""Diagnostic validity is independent of prospective causal ranking validity.""" +import json +import sys +import unittest +from pathlib import Path +from types import SimpleNamespace + +sys.path.insert(0, str(Path(__file__).resolve().parents[1] / 'scripts')) +import pilot_contract_triage_pair as runner + + +class ProfileRankDimensionsTests(unittest.TestCase): + def test_diagnostic_survives_unestablished_partial_and_invalid_rank(self): + profile = SimpleNamespace( + triage_accepted_ids=('assertion_behavior_regression', 'confirm_behavior_contract'), + triage_accepted_statuses=('no-remote-choice',), + triage_hypotheses=('assertion_behavior_regression', 'assertion_expectation_drift', 'confirm_behavior_contract'), + rank_hypotheses=True, + ) + for status, order, expected in ( + ('not_established', [], 'not_established'), + ('incomplete', [], 'incomplete'), + ('complete', ['confirm_behavior_contract', 'assertion_behavior_regression'], 'incomplete'), + ('complete', ['assertion_behavior_regression'] * 2, 'incomplete'), + ('complete', ['assertion_expectation_drift', 'assertion_behavior_regression'], 'complete'), + ): + with self.subTest(status=status, order=order): + payload = dict(observed_exit_status=1, test_failed=True, status='no-remote-choice', + steps=[{'id': 'confirm_behavior_contract'}], executed=False, + decision_usage=None, hypothesis_order=order, + hypothesis_ranking_status=status) + result = runner._validated_triage(json.dumps(payload), 1, profile) + self.assertEqual(result['candidate_ids'], ['confirm_behavior_contract']) + self.assertEqual(result['hypothesis_ranking_status'], expected) + self.assertEqual(result['hypothesis_order'], order if expected == 'complete' else []) + + def test_invalid_diagnostic_is_not_rescued_by_complete_rank(self): + profile = SimpleNamespace(triage_accepted_ids=('assertion_behavior_regression',), + triage_accepted_statuses=('no-remote-choice',), + triage_hypotheses=('assertion_behavior_regression', 'assertion_expectation_drift'), + rank_hypotheses=True) + payload = dict(observed_exit_status=0, test_failed=True, status='no-remote-choice', + steps=[{'id': 'assertion_behavior_regression'}], executed=False, + hypothesis_order=list(profile.triage_hypotheses), hypothesis_ranking_status='complete') + self.assertEqual(runner._validated_triage(json.dumps(payload), 1, profile)['status'], 'unscored') + + def test_diagnostic_candidate_requested_without_changing_causal_profile(self): + profile = SimpleNamespace( + triage_kinds=("assertion",), + triage_hypotheses=("assertion_behavior_regression", "assertion_expectation_drift"), + triage_accepted_ids=("assertion_behavior_regression", "confirm_behavior_contract"), + rank_hypotheses=True, triage_observations={}, + ) + argv = runner._triage_argv(1, profile) + self.assertIn("confirm_behavior_contract", argv) + self.assertEqual(argv.count("confirm_behavior_contract"), 1) + self.assertNotIn("confirm_behavior_contract", profile.triage_hypotheses) From 3d93d0c0a4072e3953220450eafc0e8ef043e0e4 Mon Sep 17 00:00:00 2001 From: Toni Nowak Date: Sun, 27 Sep 2026 14:30:17 +0200 Subject: [PATCH 3/4] test(pilot): align fixture with diagnostic candidate support --- tests/test_query_migration_fixture.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/tests/test_query_migration_fixture.py b/tests/test_query_migration_fixture.py index 7c56ce4..866ad60 100644 --- a/tests/test_query_migration_fixture.py +++ b/tests/test_query_migration_fixture.py @@ -48,7 +48,8 @@ def test_profile_loads_and_keeps_diagnostic_candidates_separate_from_causal_orde argv = pair._triage_argv(1, profile) self.assertIn("assertion_behavior_regression", argv) self.assertIn("assertion_expectation_drift", argv) - self.assertNotIn("confirm_behavior_contract", argv) + self.assertIn("confirm_behavior_contract", argv) + self.assertNotIn("confirm_behavior_contract", profile.triage_hypotheses) self.assertTrue(profile.rank_hypotheses) self.assertEqual(profile.triage_accepted_statuses, ("no-remote-choice",)) From 1ac409004597ee1975cd80fd42024a19c529bb16 Mon Sep 17 00:00:00 2001 From: Toni Nowak Date: Sun, 27 Sep 2026 14:31:35 +0200 Subject: [PATCH 4/4] feat(pilot): add credential-isolated profile triage bridge --- PILOT.md | 6 + TASKS.md | 6 + scripts/pilot_profile_triage_bridge.py | 614 ++++++++++++++++++++++ tests/test_pilot_profile_triage_bridge.py | 344 ++++++++++++ 4 files changed, 970 insertions(+) create mode 100644 scripts/pilot_profile_triage_bridge.py create mode 100644 tests/test_pilot_profile_triage_bridge.py diff --git a/PILOT.md b/PILOT.md index b13b4c2..634b80e 100644 --- a/PILOT.md +++ b/PILOT.md @@ -1745,3 +1745,9 @@ This is preparation, not a launched pair or efficacy evidence. Launch remains ga ## VCR394 — Independent diagnostic and causal ranking receipts A valid diagnostic next step is now scored separately from complete causal ordering. Unestablished, partial or malformed rankings remain explicitly non-complete without erasing valid diagnostic delivery. Complete ranking requires a unique full causal permutation; contract confirmation remains a requested diagnostic candidate only. Three offline tests cover these boundaries. No live delivery or efficacy claim. + +## VCR393 — Credential-isolated profile triage bridge + +Added a reusable one-shot supervisor bridge for validated enum-only profile triage. The child receives no OpenRouter key. Accepted requests persist a private typed receipt before returning catalog-authored advice, with diagnostic choice, causal order, usage and actual provider transport calls recorded separately. Nine offline tests cover isolation, complete/incomplete ranking, local abstention, fallback and invalid requests. + +The module is not yet integrated into the profile runner and proves neither native delivery nor benefit. Supplying the observed client disables internal typed decision caching. Broader child network egress is not restricted by this module; equivalent egress configuration is required for paired trials. diff --git a/TASKS.md b/TASKS.md index a35fd51..42a926b 100644 --- a/TASKS.md +++ b/TASKS.md @@ -1149,3 +1149,9 @@ This is preparation, not a launched pair or efficacy evidence. Launch remains ga ## VCR394 — Independent diagnostic and causal ranking receipts A valid diagnostic next step is now scored separately from complete causal ordering. Unestablished, partial or malformed rankings remain explicitly non-complete without erasing valid diagnostic delivery. Complete ranking requires a unique full causal permutation; contract confirmation remains a requested diagnostic candidate only. Three offline tests cover these boundaries. No live delivery or efficacy claim. + +## VCR393 — Credential-isolated profile triage bridge + +Added a reusable one-shot supervisor bridge for validated enum-only profile triage. The child receives no OpenRouter key. Accepted requests persist a private typed receipt before returning catalog-authored advice, with diagnostic choice, causal order, usage and actual provider transport calls recorded separately. Nine offline tests cover isolation, complete/incomplete ranking, local abstention, fallback and invalid requests. + +The module is not yet integrated into the profile runner and proves neither native delivery nor benefit. Supplying the observed client disables internal typed decision caching. Broader child network egress is not restricted by this module; equivalent egress configuration is required for paired trials. diff --git a/scripts/pilot_profile_triage_bridge.py b/scripts/pilot_profile_triage_bridge.py new file mode 100644 index 0000000..0d1b714 --- /dev/null +++ b/scripts/pilot_profile_triage_bridge.py @@ -0,0 +1,614 @@ +"""Credential-isolated supervisor bridge for profile-driven triage decisions. + +The agent subprocess gets no Jev/OpenRouter key. A runner observes the configured +focused-test failure, starts this one-shot loopback bridge, and installs the exact +Python shim built here. The supervisor validates only fixed enum metadata, calls +the production triage function, writes a redacted typed receipt, then returns a +catalog-authored result to the child. Bridged Codex network egress must be scoped +separately; this module does not claim localhost-only egress. +""" +from __future__ import annotations + +from dataclasses import dataclass +import json +import math +import os +from pathlib import Path +import sys +import threading +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +from typing import Any, Callable, Mapping, Sequence + + +MAX_REQUEST_BYTES = 4096 +MAX_RESPONSE_BYTES = 64 * 1024 +MAX_USAGE = 10**12 +_OBSERVATION_ENUMS = { + "import": ("ImportObservation", "import_observations"), + "assertion": ("AssertionObservation", "assertion_observations"), + "timeout": ("TimeoutObservation", "timeout_observations"), +} +_OBSERVATION_FLAGS = { + "import": "--import-observation", + "assertion": "--assertion-observation", + "timeout": "--timeout-observation", +} + + +def _strict_json(raw: str) -> Any: + def unique_pairs(pairs: list[tuple[str, Any]]) -> dict[str, Any]: + result: dict[str, Any] = {} + for key, value in pairs: + if key in result: + raise ValueError("duplicate JSON key") + result[key] = value + return result + + def finite_float(value: str) -> float: + number = float(value) + if not math.isfinite(number): + raise ValueError("non-finite JSON number") + return number + + return json.loads(raw, object_pairs_hook=unique_pairs, parse_float=finite_float, + parse_constant=lambda _value: (_ for _ in ()).throw( + ValueError("non-finite JSON"))) + + +@dataclass(frozen=True) +class ProfileTriageSpec: + """Supervisor-validated enum contract; values contain no source or free text.""" + + failure_kinds: tuple[str, ...] + hypotheses: tuple[str, ...] + allowed_observations: Mapping[str, tuple[str, ...]] + rank_hypotheses: bool = False + + def __post_init__(self) -> None: + from jevcompass.triage import FailureKind, HypothesisId + + if (not self.failure_kinds or not self.hypotheses + or any(not isinstance(value, str) for value in self.failure_kinds) + or any(not isinstance(value, str) for value in self.hypotheses) + or len(set(self.failure_kinds)) != len(self.failure_kinds) + or len(set(self.hypotheses)) != len(self.hypotheses) + or any(value not in {item.value for item in FailureKind} + for value in self.failure_kinds) + or any(value not in {item.value for item in HypothesisId} + for value in self.hypotheses) + or len(self.failure_kinds) > 4 or len(self.hypotheses) > 8): + raise ValueError("invalid profile triage enum contract") + if type(self.rank_hypotheses) is not bool: + raise ValueError("rank_hypotheses must be boolean") + if not isinstance(self.allowed_observations, Mapping): + raise ValueError("invalid profile observation allowlist") + for group, values in self.allowed_observations.items(): + if group not in _OBSERVATION_ENUMS: + raise ValueError("unknown observation group") + enum_name, _ = _OBSERVATION_ENUMS[group] + enum_type = getattr(__import__("jevcompass.triage", fromlist=[enum_name]), enum_name) + allowed = {item.value for item in enum_type} + if (not isinstance(values, tuple) or len(values) > len(allowed) + or any(not isinstance(value, str) or value not in allowed for value in values) + or len(set(values)) != len(values)): + raise ValueError("invalid profile observation enum allowlist") + + +class _ObservedDecisionClient: + """Lazy DecisionsClient subclass counting actual outbound transport calls.""" + + def __new__(cls, factory: Callable[[], Any], counter: list[int]): + from jevcompass.decisions import DecisionsClient + + class Observed(DecisionsClient): + def __init__(self) -> None: + self._factory = factory + self._counter = counter + self._delegate = None + + def decide_with_usage(self, state: dict[str, Any], questions: dict[str, Any]): + if self._delegate is None: + delegate = self._factory() + if not isinstance(delegate, DecisionsClient): + raise TypeError("decision client factory returned an unsupported client") + original = delegate.transport + + def counted_transport(url: str, body: bytes, api_key: str, timeout: float): + self._counter[0] += 1 + return original(url, body, api_key, timeout) + + delegate.transport = counted_transport + self._delegate = delegate + return self._delegate.decide_with_usage(state, questions) + + return Observed() + + +def _safe_usage(value: Any) -> dict[str, int | float | None] | None: + if value is None: + return None + inputs, outputs, cost = value.input_tokens, value.output_tokens, value.cost_usd + if (isinstance(inputs, bool) or not isinstance(inputs, int) or not 0 <= inputs <= MAX_USAGE + or isinstance(outputs, bool) or not isinstance(outputs, int) + or not 0 <= outputs <= MAX_USAGE): + raise ValueError("invalid decision usage") + if cost is not None and ( + isinstance(cost, bool) or not isinstance(cost, (int, float)) + or not math.isfinite(float(cost)) or not 0 <= cost <= MAX_USAGE + ): + raise ValueError("invalid decision usage") + return {"input_tokens": inputs, "output_tokens": outputs, + "cost_usd": float(cost) if cost is not None else None} + + +def _private_write(path: Path, payload: Mapping[str, Any]) -> None: + parent = path.parent + if parent.is_symlink() or not parent.is_dir() or parent.stat().st_mode & 0o077: + raise OSError("receipt parent must be a private real directory") + encoded = (json.dumps(payload, sort_keys=True, separators=(",", ":"), allow_nan=False) + + "\n").encode("utf-8") + if len(encoded) > MAX_RESPONSE_BYTES: + raise OSError("typed receipt exceeds size limit") + flags = os.O_WRONLY | os.O_CREAT | os.O_EXCL + flags |= getattr(os, "O_NOFOLLOW", 0) + fd = os.open(path, flags, 0o600) + try: + with os.fdopen(fd, "wb") as stream: + stream.write(encoded) + stream.flush() + os.fsync(stream.fileno()) + except BaseException: + try: + path.unlink() + except OSError: + pass + raise + + +class ProfileTriageBridge: + """One request after a runner-confirmed focused failure; response is catalog-only.""" + + def __init__( + self, spec: ProfileTriageSpec, *, receipt_path: Path, + client_factory: Callable[[], Any] | None = None, + ) -> None: + self.spec = spec + self.receipt_path = receipt_path + self.client_factory = client_factory + self._lock = threading.Lock() + self._request_count = 0 + self._observed_exit: int | None = None + self._focused_seen = False + self._state = "not_requested" + self._decision_receipt: dict[str, Any] | None = None + self._server = self._make_server() + self._thread = threading.Thread(target=self._server.serve_forever, + kwargs={"poll_interval": 0.05}, daemon=True) + + def _make_server(self) -> ThreadingHTTPServer: + bridge = self + + class Handler(BaseHTTPRequestHandler): + protocol_version = "HTTP/1.0" + server_version = "JevProfileTriageBridge" + sys_version = "" + + def log_message(self, _format: str, *_args: Any) -> None: + return + + def do_POST(self) -> None: + bridge._handle(self) + + def do_GET(self) -> None: + bridge._handle(self) + + def do_PUT(self) -> None: + bridge._handle(self) + + def do_DELETE(self) -> None: + bridge._handle(self) + + server = ThreadingHTTPServer(("127.0.0.1", 0), Handler) + server.daemon_threads = True + server.block_on_close = False + return server + + @property + def url(self) -> str: + host, port = self._server.server_address + return f"http://{host}:{port}/triage" + + def __enter__(self) -> "ProfileTriageBridge": + self._thread.start() + return self + + def __exit__(self, *_exc: object) -> None: + if self._thread.is_alive(): + self._server.shutdown() + self._thread.join(timeout=2.0) + self._server.server_close() + + def observe_focused_failure(self, exit_status: int, *, failure_confirmed: bool) -> bool: + """Called by the runner only for the first focused-command completion.""" + if type(failure_confirmed) is not bool: + raise ValueError("failure confirmation must be boolean") + with self._lock: + if self._focused_seen: + return False + self._focused_seen = True + if (failure_confirmed and isinstance(exit_status, int) + and not isinstance(exit_status, bool) and exit_status != 0): + self._observed_exit = exit_status + return True + self._state = "not_eligible" + return False + + def receipt(self) -> dict[str, Any]: + with self._lock: + return { + "request_count": min(self._request_count, 2), + "bridge_request_count": min(self._request_count, 2), + "state": self._state, + "observed_focused_failure": self._observed_exit is not None, + "decision": dict(self._decision_receipt) if self._decision_receipt else None, + } + + def _send(self, handler: BaseHTTPRequestHandler, status: int, payload: Mapping[str, Any]) -> None: + encoded = json.dumps(payload, sort_keys=True, separators=(",", ":"), + allow_nan=False).encode("utf-8") + if len(encoded) > MAX_RESPONSE_BYTES: + status, encoded = 502, b'{"status":"unavailable"}' + try: + handler.send_response(status) + handler.send_header("Content-Type", "application/json") + handler.send_header("Content-Length", str(len(encoded))) + handler.send_header("Connection", "close") + handler.end_headers() + handler.wfile.write(encoded) + handler.close_connection = True + except (BrokenPipeError, ConnectionResetError, OSError): + return + + def _handle(self, handler: BaseHTTPRequestHandler) -> None: + if handler.client_address[0] != "127.0.0.1": + self._send(handler, 403, {"status": "unavailable"}) + return + with self._lock: + self._request_count += 1 + number, observed_exit = self._request_count, self._observed_exit + if number > 1: + self._state = "duplicate" + elif observed_exit is None: + self._state = "premature" + if number > 1: + self._send(handler, 429, {"status": "unavailable"}) + return + if observed_exit is None: + self._write_rejection("premature") + self._send(handler, 409, {"status": "unavailable"}) + return + try: + request = self._read_request(handler) + args = self._validated_request(request, observed_exit) + payload, receipt = self._decide(args, observed_exit) + _private_write(self.receipt_path, receipt) + with self._lock: + self._state = "completed" + self._decision_receipt = receipt + self._send(handler, 200, payload) + except (ValueError, TypeError, OSError, UnicodeError, json.JSONDecodeError): + self._write_rejection("invalid_request") + self._send(handler, 400, {"status": "unavailable"}) + except Exception: + # Never forward provider, transport, or exception text to the child. + self._write_rejection("provider_or_bridge_error") + self._send(handler, 200, self._local_fallback(observed_exit)) + + def _read_request(self, handler: BaseHTTPRequestHandler) -> Any: + if (handler.command != "POST" or handler.path != "/triage" + or handler.headers.get("Transfer-Encoding") is not None): + raise ValueError("invalid request") + lengths = handler.headers.get_all("Content-Length", []) + if len(lengths) != 1 or handler.headers.get_content_type() != "application/json": + raise ValueError("invalid request") + length = int(lengths[0]) + if length < 1 or length > MAX_REQUEST_BYTES: + raise ValueError("invalid request") + raw = handler.rfile.read(length) + if len(raw) != length: + raise ValueError("short request") + return _strict_json(raw.decode("utf-8")) + + def _validated_request(self, request: Any, observed_exit: int) -> dict[str, Any]: + if not isinstance(request, dict) or set(request) != { + "observed_exit_status", "failure_kinds", "hypotheses", "observations", + "rank_hypotheses", + }: + raise ValueError("invalid enum-only request") + if (type(request["observed_exit_status"]) is not int + or request["observed_exit_status"] != observed_exit + or request["failure_kinds"] != list(self.spec.failure_kinds) + or request["hypotheses"] != list(self.spec.hypotheses) + or type(request["rank_hypotheses"]) is not bool + or request["rank_hypotheses"] != self.spec.rank_hypotheses): + raise ValueError("request does not match profile") + raw_observations = request["observations"] + if not isinstance(raw_observations, dict) or set(raw_observations) - set(self.spec.allowed_observations): + raise ValueError("invalid observation groups") + observations: dict[str, tuple[str, ...]] = {} + for group, values in raw_observations.items(): + allowed = set(self.spec.allowed_observations[group]) + if (not isinstance(values, list) or len(values) > len(allowed) + or any(not isinstance(value, str) or value not in allowed for value in values) + or len(set(values)) != len(values)): + raise ValueError("invalid observation enum") + observations[group] = tuple(values) + return {"observations": observations} + + def _decide(self, args: dict[str, Any], observed_exit: int): + from jevcompass.decisions import DecisionsClient + from jevcompass.triage import ( + CATALOG, AssertionObservation, FailureKind, HypothesisId, + ImportObservation, TriageDecisionReason, TimeoutObservation, + TriageResult, triage_failure, + ) + + counter = [0] + def factory(): + return self.client_factory() if self.client_factory is not None else DecisionsClient() + + client = _ObservedDecisionClient(factory, counter) + enum_groups = { + "import": (ImportObservation, "import_observations"), + "assertion": (AssertionObservation, "assertion_observations"), + "timeout": (TimeoutObservation, "timeout_observations"), + } + triage_observations = { + enum_groups[group][1]: tuple(enum_groups[group][0](value) for value in values) + for group, values in args["observations"].items() + } + result = triage_failure( + tuple(FailureKind(value) for value in self.spec.failure_kinds), + tuple(HypothesisId(value) for value in self.spec.hypotheses), + observed_exit, + client=client, + rank_hypotheses=self.spec.rank_hypotheses, + **triage_observations, + ) + if not isinstance(result, TriageResult) or result.observed_exit_status != observed_exit: + raise ValueError("invalid production triage result") + reason = result.decision_reason + if reason is not None and not isinstance(reason, TriageDecisionReason): + raise ValueError("invalid decision reason") + allowed_hypotheses = set(self.spec.hypotheses) + step_ids = [step.id.value for step in result.steps] + if (len(step_ids) > 2 or len(set(step_ids)) != len(step_ids) + or any(identifier not in allowed_hypotheses for identifier in step_ids)): + raise ValueError("invalid diagnostic choice") + order = [item.value for item in result.hypothesis_order] + if (len(order) > 4 or len(set(order)) != len(order) + or any(item not in allowed_hypotheses or item == "confirm_behavior_contract" + for item in order)): + raise ValueError("invalid causal hypothesis order") + ranking_status = result.hypothesis_ranking_status + if ranking_status not in {"complete", "incomplete", "not_established"}: + raise ValueError("invalid hypothesis ranking status") + if not self.spec.rank_hypotheses and (order or ranking_status != "not_established"): + raise ValueError("unexpected hypothesis ranking") + if ranking_status == "complete" and len(order) < 2: + raise ValueError("complete ranking must order multiple causes") + if result.status not in {"remote-choice", "no-remote-choice"}: + raise ValueError("invalid triage status") + if result.status == "remote-choice" and ( + not step_ids or reason is not TriageDecisionReason.ACCEPTED + ): + raise ValueError("invalid accepted choice") + if result.status != "remote-choice" and reason is TriageDecisionReason.ACCEPTED: + raise ValueError("accepted reason requires remote choice") + usage = _safe_usage(result.decision_usage) + source = ( + "cached_preferred_next_step" if result.cache_hit + else "remote_preferred_next_step" + ) if reason is TriageDecisionReason.ACCEPTED else ( + "locally_resolved_guidance" if reason is TriageDecisionReason.LOCAL_RESOLUTION + else "unranked_local_fallback" + ) + steps = [] + for index, identifier in enumerate(step_ids): + entry = CATALOG[HypothesisId(identifier)].step + steps.append({ + "id": identifier, "title": entry.title, "instruction": entry.instruction, + "selection_source": source if index == 0 else "unranked_local_fallback", + }) + usage_status = "reported" if usage is not None else ( + "not_invoked" if counter[0] == 0 else "not_reported" + ) + if type(result.cache_hit) is not bool: + raise ValueError("invalid cache status") + receipt = { + "schema_version": 1, + "status": result.status, + "bridge_request_count": 1, + "observed_exit_status": observed_exit, + "test_failed": True, + "executed": False, + "decision_reason": reason.value if reason is not None else None, + "diagnostic_step_ids": step_ids, + "diagnostic_choice_id": step_ids[0] if step_ids else None, + "diagnostic_selection_source": source if step_ids else "none", + "hypothesis_order": order, + "hypothesis_ranking_status": ranking_status, + "cache_hit": result.cache_hit is True, + "provider_transport_call_count": counter[0], + "decision_usage_status": usage_status, + "decision_usage": usage, + } + payload = { + "status": result.status, + "observed_exit_status": observed_exit, + "test_failed": True, + "executed": False, + "decision_reason": receipt["decision_reason"], + "cache_hit": receipt["cache_hit"], + "hypothesis_ranking_status": ranking_status, + "steps": steps, + "decision_usage": usage, + } + if self.spec.rank_hypotheses: + payload["hypothesis_order"] = order + return payload, receipt + + def _local_fallback(self, observed_exit: int) -> dict[str, Any]: + from jevcompass.triage import CATALOG, HypothesisId + + steps = [CATALOG[HypothesisId(value)].step for value in self.spec.hypotheses[:2]] + return { + "status": "no-remote-choice", "observed_exit_status": observed_exit, + "test_failed": True, "executed": False, "decision_reason": "provider_error", + "cache_hit": False, "hypothesis_ranking_status": "not_established", + "steps": [{ + "id": step.id.value, "title": step.title, "instruction": step.instruction, + "selection_source": "unranked_local_fallback", + } for step in steps], + "decision_usage": None, + } + + def _write_rejection(self, state: str) -> None: + with self._lock: + self._state = state + if self._decision_receipt is not None: + return + receipt = { + "schema_version": 1, "status": "rejected", "request_state": state, + "bridge_request_count": min(self._request_count, 2), + "observed_exit_status": self._observed_exit, + "provider_transport_call_count": 0, + "diagnostic_choice_id": None, "hypothesis_order": [], + "hypothesis_ranking_status": "not_established", "cache_hit": False, + "decision_usage_status": "not_invoked", "decision_usage": None, + } + try: + _private_write(self.receipt_path, receipt) + except FileExistsError: + pass + except OSError: + pass + with self._lock: + self._decision_receipt = receipt + + +def command_for(spec: ProfileTriageSpec, exit_status: int, + observations: Mapping[str, Sequence[str]] | None = None) -> tuple[str, ...]: + """Return the canonical safe CLI argv for the validated profile.""" + if isinstance(exit_status, bool) or not isinstance(exit_status, int) or exit_status == 0: + raise ValueError("triage bridge requires a nonzero observed exit") + observed = observations or {} + if set(observed) - set(spec.allowed_observations): + raise ValueError("unknown observation group") + args = ["python", "-m", "jevcompass", "triage", "--exit-code", str(exit_status)] + for kind in spec.failure_kinds: + args.extend(("--kind", kind)) + for hypothesis in spec.hypotheses: + args.extend(("--hypothesis", hypothesis)) + for group in ("import", "assertion", "timeout"): + values = tuple(observed.get(group, ())) + allowed = set(spec.allowed_observations.get(group, ())) + if (len(set(values)) != len(values) + or any(not isinstance(value, str) or value not in allowed for value in values)): + raise ValueError("invalid observation enum") + for value in values: + args.extend((_OBSERVATION_FLAGS[group], value)) + if spec.rank_hypotheses: + args.append("--rank-hypotheses") + args.append("--json") + return tuple(args) + + +def write_python_shim(directory: Path, bridge: ProfileTriageBridge, + real_python: str = sys.executable) -> Path: + """Write a private exact-command shim; unknown commands execute real Python.""" + if (not isinstance(real_python, str) or not real_python + or not bridge.url.startswith("http://127.0.0.1:")): + raise ValueError("invalid bridge shim configuration") + bindir = directory / "bin" + bindir.mkdir(mode=0o700, parents=True, exist_ok=True) + shim = bindir / "python" + groups = {key: list(value) for key, value in bridge.spec.allowed_observations.items()} + source = f'''#!/usr/bin/env python3 +import json, os, sys +from urllib.request import Request, urlopen +REAL = {real_python!r} +URL = {bridge.url!r} +KINDS = {list(bridge.spec.failure_kinds)!r} +HYPOTHESES = {list(bridge.spec.hypotheses)!r} +OBSERVATIONS = {groups!r} +RANK = {bridge.spec.rank_hypotheses!r} +FLAGS = {dict(_OBSERVATION_FLAGS)!r} +ARGS = sys.argv[1:] +def fallback(): + os.execv(REAL, [REAL, *ARGS]) +def pairs(flag, values): + result = [] + for value in values: + result.extend((flag, value)) + return result +prefix = ["-m", "jevcompass", "triage", "--exit-code"] +if len(ARGS) < 6 or ARGS[:4] != prefix: + fallback() +try: + code = int(ARGS[4]) + if code == 0: + fallback() + tail = ARGS[5:] + expected = pairs("--kind", KINDS) + pairs("--hypothesis", HYPOTHESES) + if tail[:len(expected)] != expected: + fallback() + tail = tail[len(expected):] + observations = {{group: [] for group in OBSERVATIONS}} + while tail and tail[0] != "--json" and tail[0] != "--rank-hypotheses": + if len(tail) < 2: + fallback() + flag, value = tail[0], tail[1] + group = next((name for name, item in FLAGS.items() if item == flag), None) + if group is None or value not in OBSERVATIONS.get(group, []) or value in observations[group]: + fallback() + observations[group].append(value) + tail = tail[2:] + if RANK: + if not tail or tail[0] != "--rank-hypotheses": + fallback() + tail = tail[1:] + elif tail and tail[0] == "--rank-hypotheses": + fallback() + if tail != ["--json"]: + fallback() + body = json.dumps({{ + "observed_exit_status": code, "failure_kinds": KINDS, + "hypotheses": HYPOTHESES, "observations": observations, + "rank_hypotheses": RANK, + }}, separators=(",", ":")).encode() + req = Request(URL, data=body, headers={{"Content-Type": "application/json"}}, method="POST") + with urlopen(req, timeout=2.0) as response: + data = response.read({MAX_RESPONSE_BYTES + 1}) + if len(data) <= {MAX_RESPONSE_BYTES}: + sys.stdout.buffer.write(data + (b"\\n" if not data.endswith(b"\\n") else b"")) + raise SystemExit(0) +except SystemExit: + raise +except Exception: + pass +fallback() +''' + fd = os.open(shim, os.O_WRONLY | os.O_CREAT | os.O_EXCL | getattr(os, "O_NOFOLLOW", 0), 0o700) + try: + with os.fdopen(fd, "w", encoding="utf-8") as stream: + stream.write(source) + stream.flush() + os.fsync(stream.fileno()) + except BaseException: + try: + shim.unlink() + except OSError: + pass + raise + return bindir diff --git a/tests/test_pilot_profile_triage_bridge.py b/tests/test_pilot_profile_triage_bridge.py new file mode 100644 index 0000000..72456a0 --- /dev/null +++ b/tests/test_pilot_profile_triage_bridge.py @@ -0,0 +1,344 @@ +"""Offline tests for profile-driven supervisor triage transport.""" +from __future__ import annotations + +import json +import os +from pathlib import Path +import subprocess +import sys +import tempfile +import unittest +from urllib.error import HTTPError +from urllib.request import Request, urlopen + +ROOT = Path(__file__).resolve().parents[1] +sys.path.insert(0, str(ROOT / "src")) +sys.path.insert(0, str(ROOT / "scripts")) + +from jevcompass.decisions import DecisionsClient +from pilot_profile_triage_bridge import ( + ProfileTriageBridge, ProfileTriageSpec, command_for, write_python_shim, +) + + +TIMEOUT_SPEC = ProfileTriageSpec( + failure_kinds=("timeout",), + hypotheses=("timeout_contention", "timeout_nonterminating"), + allowed_observations={"timeout": ( + "progress_observed", "no_progress_observed", "resource_contention_observed", + "no_resource_contention_observed", "wait_condition_satisfiable", + "wait_condition_unsatisfiable", + )}, + rank_hypotheses=True, +) + + +def reply_transport(answers=None, usage=None, *, error=None, seen=None): + def transport(url, body, api_key, timeout): + if seen is not None: + seen.append({"url": url, "body": json.loads(body), "key": api_key}) + if error is not None: + raise error + return json.dumps({"answers": answers or {}, "usage": usage}).encode() + return transport + + +def successful_answers(): + return { + "diagnostic": { + "type": "choice", "choice": "timeout_nonterminating", "confidence": 0.91, + }, + "hypothesis_pair_0_1": { + "type": "choice", "choice": "timeout_nonterminating", "confidence": 0.88, + }, + } + + +def post(url, payload): + data = json.dumps(payload, separators=(",", ":")).encode() + request = Request(url, data=data, headers={"Content-Type": "application/json"}, + method="POST") + try: + with urlopen(request, timeout=2) as response: + return response.status, response.read() + except HTTPError as error: + try: + return error.code, error.read() + finally: + error.close() + + +def request_for(**overrides): + value = { + "observed_exit_status": 1, + "failure_kinds": ["timeout"], + "hypotheses": ["timeout_contention", "timeout_nonterminating"], + "observations": {"timeout": []}, + "rank_hypotheses": True, + } + value.update(overrides) + return value + + +class ProfileTriageBridgeTests(unittest.TestCase): + def private_dir(self, temporary): + directory = Path(temporary) / "private" + directory.mkdir(mode=0o700) + return directory + + def test_exact_shim_keeps_key_parent_only_and_persists_choice_rank_and_usage(self): + with tempfile.TemporaryDirectory() as temporary: + private = self.private_dir(temporary) + seen = [] + factory = lambda: DecisionsClient( + api_key="SUPERVISOR_ONLY_SECRET", model="synthetic", + transport=reply_transport( + successful_answers(), + {"input_tokens": 14, "output_tokens": 5, "cost": 0.0012}, + seen=seen, + ), + ) + receipt_path = private / "decision.json" + bridge = ProfileTriageBridge( + TIMEOUT_SPEC, receipt_path=receipt_path, client_factory=factory, + ) + with bridge: + self.assertTrue(bridge.observe_focused_failure(1, failure_confirmed=True)) + bin_dir = write_python_shim(private, bridge) + child_env = { + "PATH": str(bin_dir) + os.pathsep + os.environ.get("PATH", ""), + "HOME": str(private), "PYTHONPATH": str(ROOT / "src"), + "JEVCOMPASS_TRIAGE_BRIDGE_URL": bridge.url, + } + self.assertNotIn("OPENROUTER_API_KEY", child_env) + command = command_for(TIMEOUT_SPEC, 1) + child = subprocess.run( + [str(bin_dir / "python"), *command[1:]], + cwd=ROOT, env=child_env, stdin=subprocess.DEVNULL, + stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, + timeout=4, check=False, + ) + self.assertEqual(child.returncode, 0, child.stderr) + response = json.loads(child.stdout) + receipt = json.loads(receipt_path.read_text()) + self.assertEqual(receipt["diagnostic_choice_id"], "timeout_nonterminating") + self.assertEqual(receipt["diagnostic_selection_source"], "remote_preferred_next_step") + self.assertEqual(receipt["hypothesis_order"], [ + "timeout_nonterminating", "timeout_contention", + ]) + self.assertEqual(receipt["hypothesis_ranking_status"], "complete") + self.assertEqual(receipt["provider_transport_call_count"], 1) + self.assertEqual(receipt["decision_usage"], { + "input_tokens": 14, "output_tokens": 5, "cost_usd": 0.0012, + }) + self.assertFalse(receipt["cache_hit"]) + self.assertEqual(response["hypothesis_order"], receipt["hypothesis_order"]) + self.assertEqual(response["steps"][0]["selection_source"], "remote_preferred_next_step") + self.assertEqual(seen[0]["key"], "SUPERVISOR_ONLY_SECRET") + self.assertEqual(seen[0]["body"]["state"], { + "test_outcome": "failed", "failure_kinds": ["timeout"], + "hypotheses": ["timeout_contention", "timeout_nonterminating"], + }) + self.assertNotIn("private/source.py", json.dumps(seen[0]["body"])) + self.assertNotIn("SUPERVISOR_ONLY_SECRET", child.stdout + child.stderr) + self.assertEqual(receipt_path.stat().st_mode & 0o777, 0o600) + self.assertEqual(bridge.receipt()["request_count"], 1) + + def test_local_resolution_is_saved_without_provider_call_or_fake_ranking(self): + with tempfile.TemporaryDirectory() as temporary: + private = self.private_dir(temporary) + calls = [] + factory = lambda: DecisionsClient( + api_key="unused", transport=reply_transport(successful_answers(), seen=calls), + ) + path = private / "decision.json" + bridge = ProfileTriageBridge(TIMEOUT_SPEC, receipt_path=path, client_factory=factory) + with bridge: + bridge.observe_focused_failure(1, failure_confirmed=True) + status, body = post(bridge.url, request_for( + observations={"timeout": ["wait_condition_unsatisfiable"]}, + )) + self.assertEqual(status, 200) + response = json.loads(body) + receipt = json.loads(path.read_text()) + self.assertEqual(receipt["decision_reason"], "local_resolution") + self.assertEqual(receipt["diagnostic_choice_id"], "timeout_nonterminating") + self.assertEqual(receipt["hypothesis_ranking_status"], "not_established") + self.assertEqual(receipt["hypothesis_order"], []) + self.assertEqual(receipt["provider_transport_call_count"], 0) + self.assertEqual(receipt["decision_usage_status"], "not_invoked") + self.assertEqual(calls, []) + self.assertEqual(response["status"], "no-remote-choice") + + def test_local_abstention_is_preserved_without_provider_call(self): + with tempfile.TemporaryDirectory() as temporary: + private = self.private_dir(temporary) + calls = [] + spec = ProfileTriageSpec( + failure_kinds=("timeout",), hypotheses=("timeout_nonterminating",), + allowed_observations={}, rank_hypotheses=True, + ) + factory = lambda: DecisionsClient( + api_key="unused", transport=reply_transport(successful_answers(), seen=calls), + ) + path = private / "decision.json" + bridge = ProfileTriageBridge(spec, receipt_path=path, client_factory=factory) + with bridge: + bridge.observe_focused_failure(1, failure_confirmed=True) + status, body = post(bridge.url, { + "observed_exit_status": 1, "failure_kinds": ["timeout"], + "hypotheses": ["timeout_nonterminating"], "observations": {}, + "rank_hypotheses": True, + }) + response, receipt = json.loads(body), json.loads(path.read_text()) + self.assertEqual(status, 200) + self.assertEqual(receipt["decision_reason"], "local_abstention") + self.assertEqual(receipt["hypothesis_ranking_status"], "not_established") + self.assertEqual(receipt["provider_transport_call_count"], 0) + self.assertEqual(receipt["decision_usage_status"], "not_invoked") + self.assertEqual(response["status"], "no-remote-choice") + self.assertEqual(calls, []) + + def test_low_confidence_keeps_local_next_step_and_usage_without_claiming_remote_choice(self): + with tempfile.TemporaryDirectory() as temporary: + private = self.private_dir(temporary) + factory = lambda: DecisionsClient( + api_key="synthetic", + transport=reply_transport({ + "diagnostic": { + "type": "choice", "choice": "timeout_contention", "confidence": 0.2, + }, + "hypothesis_pair_0_1": { + "type": "choice", "choice": "timeout_contention", "confidence": 0.9, + }, + }, {"input_tokens": 9, "output_tokens": 3, "cost": None}), + ) + spec = ProfileTriageSpec( + TIMEOUT_SPEC.failure_kinds, TIMEOUT_SPEC.hypotheses, {}, True, + ) + path = private / "decision.json" + bridge = ProfileTriageBridge(spec, receipt_path=path, client_factory=factory) + with bridge: + bridge.observe_focused_failure(1, failure_confirmed=True) + status, body = post(bridge.url, request_for( + observations={}, rank_hypotheses=True, + )) + self.assertEqual(status, 200) + response, receipt = json.loads(body), json.loads(path.read_text()) + self.assertEqual(response["status"], "no-remote-choice") + self.assertEqual(receipt["decision_reason"], "insufficient_confidence") + self.assertEqual(receipt["diagnostic_selection_source"], "unranked_local_fallback") + self.assertEqual(receipt["hypothesis_ranking_status"], "complete") + self.assertEqual(receipt["provider_transport_call_count"], 1) + self.assertEqual(receipt["decision_usage_status"], "reported") + + def test_accepted_diagnostic_choice_survives_incomplete_causal_ranking(self): + with tempfile.TemporaryDirectory() as temporary: + private = self.private_dir(temporary) + factory = lambda: DecisionsClient( + api_key="synthetic", + transport=reply_transport({ + "diagnostic": { + "type": "choice", "choice": "timeout_nonterminating", "confidence": 0.9, + }, + }, {"input_tokens": 6, "output_tokens": 2, "cost": None}), + ) + path = private / "decision.json" + bridge = ProfileTriageBridge(TIMEOUT_SPEC, receipt_path=path, + client_factory=factory) + with bridge: + bridge.observe_focused_failure(1, failure_confirmed=True) + status, body = post(bridge.url, request_for()) + response, receipt = json.loads(body), json.loads(path.read_text()) + self.assertEqual(status, 200) + self.assertEqual(response["status"], "remote-choice") + self.assertEqual(receipt["diagnostic_choice_id"], "timeout_nonterminating") + self.assertEqual(receipt["hypothesis_ranking_status"], "incomplete") + self.assertEqual(receipt["hypothesis_order"], []) + self.assertEqual(receipt["provider_transport_call_count"], 1) + self.assertEqual(receipt["decision_usage_status"], "reported") + + def test_provider_error_is_generic_fallback_and_counted_at_transport_boundary(self): + with tempfile.TemporaryDirectory() as temporary: + private = self.private_dir(temporary) + path = private / "decision.json" + factory = lambda: DecisionsClient( + api_key="synthetic", + transport=reply_transport(error=RuntimeError("PRIVATE_PROVIDER_DETAIL")), + ) + bridge = ProfileTriageBridge(TIMEOUT_SPEC, receipt_path=path, client_factory=factory) + with bridge: + bridge.observe_focused_failure(1, failure_confirmed=True) + status, body = post(bridge.url, request_for()) + receipt = json.loads(path.read_text()) + self.assertEqual(status, 200) + self.assertEqual(json.loads(body)["decision_reason"], "provider_error") + self.assertEqual(receipt["provider_transport_call_count"], 1) + self.assertNotIn("PRIVATE_PROVIDER_DETAIL", body.decode()) + self.assertNotIn("PRIVATE_PROVIDER_DETAIL", path.read_text()) + + def test_malformed_decision_is_fallback_and_retains_valid_usage(self): + with tempfile.TemporaryDirectory() as temporary: + private = self.private_dir(temporary) + path = private / "decision.json" + factory = lambda: DecisionsClient( + api_key="synthetic", + transport=reply_transport( + {"wrong-question": {}}, + {"input_tokens": 8, "output_tokens": 2, "cost": 0.0004}, + ), + ) + bridge = ProfileTriageBridge(TIMEOUT_SPEC, receipt_path=path, client_factory=factory) + with bridge: + bridge.observe_focused_failure(1, failure_confirmed=True) + status, body = post(bridge.url, request_for()) + receipt = json.loads(path.read_text()) + self.assertEqual(status, 200) + self.assertEqual(json.loads(body)["decision_reason"], "invalid_response") + self.assertEqual(receipt["provider_transport_call_count"], 1) + self.assertEqual(receipt["decision_usage"]["input_tokens"], 8) + self.assertEqual(receipt["decision_usage_status"], "reported") + + def test_premature_duplicate_and_invalid_enum_requests_never_call_provider(self): + with tempfile.TemporaryDirectory() as temporary: + private = self.private_dir(temporary) + calls = [] + factory = lambda: DecisionsClient( + api_key="synthetic", transport=reply_transport(successful_answers(), seen=calls), + ) + early = ProfileTriageBridge( + TIMEOUT_SPEC, receipt_path=private / "early.json", client_factory=factory, + ) + with early: + status, _ = post(early.url, request_for()) + self.assertEqual(status, 409) + self.assertEqual(json.loads((private / "early.json").read_text())["request_state"], + "premature") + self.assertEqual(calls, []) + + invalid = ProfileTriageBridge( + TIMEOUT_SPEC, receipt_path=private / "invalid.json", client_factory=factory, + ) + with invalid: + invalid.observe_focused_failure(1, failure_confirmed=True) + status, _ = post(invalid.url, request_for( + hypotheses=["timeout_contention", "PRIVATE_TEXT"], + )) + self.assertEqual(status, 400) + status, _ = post(invalid.url, request_for()) + self.assertEqual(status, 429) + self.assertEqual(calls, []) + self.assertEqual(json.loads((private / "invalid.json").read_text())["status"], + "rejected") + + def test_spec_and_command_reject_non_enum_and_zero_exit(self): + with self.assertRaises(ValueError): + ProfileTriageSpec(("timeout",), ("raw failure text",), {}) + with self.assertRaises(ValueError): + command_for(TIMEOUT_SPEC, 0) + with self.assertRaises(ValueError): + command_for(TIMEOUT_SPEC, 1, {"timeout": ["raw diagnostic"]}) + + +if __name__ == "__main__": + unittest.main()