From 3492d92c55b2ce4b6f19feeca1367364f396abe9 Mon Sep 17 00:00:00 2001 From: Toni Nowak Date: Sun, 27 Sep 2026 14:38:23 +0200 Subject: [PATCH] feat(pilot): observe bounded events without losing failure evidence --- PILOT.md | 4 +++ TASKS.md | 4 +++ scripts/pilot_cli_core.py | 17 ++++++++--- tests/test_pilot_event_observer.py | 46 ++++++++++++++++++++++++++++++ 4 files changed, 67 insertions(+), 4 deletions(-) create mode 100644 tests/test_pilot_event_observer.py diff --git a/PILOT.md b/PILOT.md index 634b80e..ccda1a3 100644 --- a/PILOT.md +++ b/PILOT.md @@ -1751,3 +1751,7 @@ A valid diagnostic next step is now scored separately from complete causal order 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. + +## 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. diff --git a/TASKS.md b/TASKS.md index 42a926b..d5e7827 100644 --- a/TASKS.md +++ b/TASKS.md @@ -1155,3 +1155,7 @@ A valid diagnostic next step is now scored separately from complete causal order 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. + +## 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. diff --git a/scripts/pilot_cli_core.py b/scripts/pilot_cli_core.py index 8c83497..4b27437 100644 --- a/scripts/pilot_cli_core.py +++ b/scripts/pilot_cli_core.py @@ -26,7 +26,7 @@ import sys import tempfile import time -from typing import Any, Iterable +from typing import Any, Callable, Iterable from pilot_receipts import ACTION_TYPES, parse_codex_json_events @@ -818,6 +818,7 @@ def _isolated_environment( def _collect_events( process: subprocess.Popen[bytes], *, started: float, timeout: int, preserve_on_failure: bool = True, max_tokens: int | None = None, + event_observer: Callable[[dict[str, Any]], None] | None = None, ) -> tuple[list[str], list[float], str | None]: """Collect bounded JSONL, optionally stopping after a completed turn exceeds a token budget. @@ -849,13 +850,21 @@ def append_line(raw: bytes) -> None: line = raw.decode("utf-8", errors="replace") lines.append(line) times.append(time.monotonic()) - if max_tokens is None: + if max_tokens is None and event_observer is None: return try: event = json.loads(line) except (TypeError, json.JSONDecodeError): return - if not isinstance(event, dict) or event.get("type") != "turn.completed": + if not isinstance(event, dict): + return + if event_observer is not None: + try: + event_observer(event) + except Exception: + failure = "event_observer_error" + return + if max_tokens is None or event.get("type") != "turn.completed": return usage_status, usage = _safe_turn_usage(event) if usage_status != "available" or usage is None: @@ -887,7 +896,7 @@ def append_line(raw: bytes) -> None: raw_line = bytes(pending[:newline]) del pending[:newline + 1] append_line(raw_line) - if failure == "token_budget_exceeded": + if failure: break if failure: break diff --git a/tests/test_pilot_event_observer.py b/tests/test_pilot_event_observer.py new file mode 100644 index 0000000..d2f5e31 --- /dev/null +++ b/tests/test_pilot_event_observer.py @@ -0,0 +1,46 @@ +"""Real subprocess coverage for the bounded collector's optional observer.""" +import json +import subprocess +import sys +import time +import unittest +from pathlib import Path + +sys.path.insert(0, str(Path(__file__).resolve().parents[1] / 'scripts')) +import pilot_cli_core as core + + +class EventObserverTests(unittest.TestCase): + def collect(self, callback, *, max_tokens=None): + events = ['not-json', json.dumps({'type': 'item.completed', 'item': {'exit_code': 1}}), + json.dumps({'type': 'turn.completed', 'usage': {'input_tokens': 12, 'cached_input_tokens': 5, 'output_tokens': 3, 'reasoning_output_tokens': 0, 'cache_write_input_tokens': 0}})] + code = 'import sys; sys.stdout.write(' + repr('\n'.join(events) + '\n') + '); sys.stdout.flush()' + process = subprocess.Popen([sys.executable, '-c', code], stdout=subprocess.PIPE) + return core._collect_events(process, started=time.monotonic(), timeout=5, + event_observer=callback, max_tokens=max_tokens) + + def test_observer_receives_valid_events_and_malformed_line_is_retained(self): + seen = [] + lines, timestamps, failure = self.collect(seen.append) + self.assertIsNone(failure) + self.assertEqual(len(lines), 3) + self.assertEqual(len(timestamps), 3) + self.assertEqual([event['type'] for event in seen], ['item.completed', 'turn.completed']) + + def test_callback_error_retains_observed_lines_without_exception_text(self): + def broken(event): + raise RuntimeError('private callback detail') + lines, timestamps, failure = self.collect(broken) + self.assertEqual(failure, 'event_observer_error') + self.assertEqual(len(lines), 2) + self.assertEqual(len(timestamps), 2) + self.assertNotIn('private callback detail', '\n'.join(lines)) + + def test_observer_and_token_cap_coexist_without_double_counting_cached_input(self): + seen = [] + lines, _, failure = self.collect(seen.append, max_tokens=14) + self.assertEqual(failure, 'token_budget_exceeded') + self.assertEqual(len(lines), 3) + self.assertEqual(len(seen), 2) + _, _, failure = self.collect(lambda event: None, max_tokens=15) + self.assertIsNone(failure)