Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions PILOT.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
4 changes: 4 additions & 0 deletions TASKS.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
17 changes: 13 additions & 4 deletions scripts/pilot_cli_core.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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.

Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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
Expand Down
46 changes: 46 additions & 0 deletions tests/test_pilot_event_observer.py
Original file line number Diff line number Diff line change
@@ -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)