diff --git a/docs/guides/custom-agent-runner-integration.md b/docs/guides/custom-agent-runner-integration.md index 36477710ba..5f046fe9d9 100644 --- a/docs/guides/custom-agent-runner-integration.md +++ b/docs/guides/custom-agent-runner-integration.md @@ -207,6 +207,23 @@ the outer runner applies and acknowledges scheduler state correctly. Those are useful extension and contribution surfaces for making Turn more mature; an Agent process exit code or scraped transcript is not a substitute for them. +### Watch a running Turn + +In a second terminal, inspect the same Turn while `run-once` is active: + +```bash +loopx --registry --runtime-root turn inspect-journal \ + --goal-id --agent-id --turn-key \ + --watch --watch-interval 1 --format json +``` + +Watch mode emits one newline-delimited JSON progress event whenever the journal +status or completed phase list changes. The event contains only Turn identity, +journal status, completed phases, and an empty effects list; it never includes +host output or session content. It stops at `committed`, `stopped`, or `failed`. +The command is read-only; use Ctrl-C to stop watching early. Start it after the +Turn journal exists, and use the same registry and runtime root as the runner. + ## Acceptance Checklist Before calling the integration autonomous, prove that: diff --git a/docs/guides/custom-agent-runner-integration.zh-CN.md b/docs/guides/custom-agent-runner-integration.zh-CN.md index 59355f0e41..b930fc6d05 100644 --- a/docs/guides/custom-agent-runner-integration.zh-CN.md +++ b/docs/guides/custom-agent-runner-integration.zh-CN.md @@ -183,6 +183,22 @@ effect,并且外层 runner 能正确应用和 ACK scheduler state。这些正 成熟度的 extension / contribution surface;Agent 进程退出码或从 transcript 猜结果不能 替代这些证明。 +### 观察正在运行的 Turn + +在第二个终端中,可以在 `run-once` 执行期间查看同一 Turn: + +```bash +loopx --registry --runtime-root turn inspect-journal \ + --goal-id --agent-id --turn-key \ + --watch --watch-interval 1 --format json +``` + +Watch 模式只在 journal 状态或已完成阶段列表变化时输出一条 JSON Lines 进度事件。 +事件仅包含 Turn 身份、journal 状态、已完成阶段和空 effects 列表,不包含 host 输出或 +session 内容。状态变为 `committed`、`stopped` 或 `failed` 后命令退出。该命令只读; +提前停止可按 Ctrl-C。请在 Turn journal 创建后启动,并使用与 runner 相同的 registry 和 +runtime root。 + ## 验收清单 在把集成称为“自主运行”前,至少证明: diff --git a/loopx/cli_commands/turn_inspection.py b/loopx/cli_commands/turn_inspection.py index eb3f47ba59..1f74b1d17b 100644 --- a/loopx/cli_commands/turn_inspection.py +++ b/loopx/cli_commands/turn_inspection.py @@ -2,7 +2,10 @@ import argparse from collections.abc import Callable +import json +import math from pathlib import Path +import time from ..control_plane.runtime.status_projection_cache import ( resolve_status_projection_cache_runtime_root, @@ -14,6 +17,8 @@ None, ] FormatSelector = Callable[..., str] +TURN_PROGRESS_EVENT_SCHEMA_VERSION = "loopx_turn_progress_event_v0" +_TERMINAL_TURN_STATUSES = {"committed", "stopped", "failed"} def handle_turn_journal_inspection( @@ -44,6 +49,99 @@ def handle_turn_journal_inspection( runtime_root, goal_id=args.goal_id, turn_key=args.turn_key, ) scope = str(((plan.get("session") or {}).get("context_policy") or {}).get("binding_scope") or "todo") + fmt = output_format(args) + if args.watch: + interval = args.watch_interval + if not math.isfinite(interval) or interval <= 0: + payload = { + "ok": False, + "schema_version": LOOPX_TURN_JOURNAL_INSPECTION_SCHEMA_VERSION, + "error": "watch interval must be a finite positive number", + "effects": [], + } + print_payload( + payload, + fmt, + render_loopx_turn_journal_inspection_markdown, + ) + return 1 + + previous: tuple[str, tuple[str, ...]] | None = None + try: + while True: + payload = inspect_loopx_turn_journal( + runtime_root, + goal_id=args.goal_id, + agent_id=args.agent_id, + turn_key=args.turn_key, + retry_failed=bool(args.retry_failed_turn), + session_binding_resolver=( + lambda turn_envelope: codex_cli_session_binding( + runtime_root, + turn_envelope, + session_scope=scope, + ) + ), + ) + if payload.get("ok") is not True: + print_payload( + payload, + fmt, + render_loopx_turn_journal_inspection_markdown, + ) + return 1 + + if payload.get("journal_consistent") is not True: + print_payload( + payload, + fmt, + render_loopx_turn_journal_inspection_markdown, + ) + return 1 + + status = str(payload.get("journal_status") or "") + phases = payload.get("completed_phases") + completed_phases = ( + [str(phase) for phase in phases] + if isinstance(phases, list) + else [] + ) + snapshot = (status, tuple(completed_phases)) + if snapshot != previous: + event: dict[str, object] = { + "schema_version": TURN_PROGRESS_EVENT_SCHEMA_VERSION, + "goal_id": args.goal_id, + "agent_id": args.agent_id, + "turn_key": args.turn_key, + "journal_status": status, + "completed_phases": completed_phases, + "effects": [], + } + if fmt == "json": + print( + json.dumps( + event, + ensure_ascii=False, + separators=(",", ":"), + ), + flush=True, + ) + else: + from .turn_rendering import ( + render_loopx_turn_journal_progress_markdown, + ) + + print( + render_loopx_turn_journal_progress_markdown(event), + flush=True, + ) + previous = snapshot + if status in _TERMINAL_TURN_STATUSES: + return 0 + time.sleep(interval) + except KeyboardInterrupt: + return 130 + payload = inspect_loopx_turn_journal( runtime_root, goal_id=args.goal_id, diff --git a/loopx/cli_commands/turn_registration.py b/loopx/cli_commands/turn_registration.py index ff582534b2..a3147bbb4b 100644 --- a/loopx/cli_commands/turn_registration.py +++ b/loopx/cli_commands/turn_registration.py @@ -48,6 +48,21 @@ def register_turn_commands( "Host Session binding check." ), ) + inspect_journal.add_argument( + "--watch", + action="store_true", + help=( + "Poll the read-only journal projection and emit one progress event " + "per changed phase/status until the Turn is terminal. With --format " + "json, events are newline-delimited JSON." + ), + ) + inspect_journal.add_argument( + "--watch-interval", + type=float, + default=1.0, + help="Seconds between journal reads while --watch is active (default: 1).", + ) plan = command_sub.add_parser( "plan", diff --git a/loopx/cli_commands/turn_rendering.py b/loopx/cli_commands/turn_rendering.py index 110476a4b6..136deda1f7 100644 --- a/loopx/cli_commands/turn_rendering.py +++ b/loopx/cli_commands/turn_rendering.py @@ -274,6 +274,26 @@ def render_loopx_turn_journal_inspection_markdown( ) +def render_loopx_turn_journal_progress_markdown( + payload: dict[str, object], +) -> str: + """Render one allowlisted progress event from a Turn journal watch.""" + phases = payload.get("completed_phases") + completed = ( + ", ".join(str(phase) for phase in phases) + if isinstance(phases, list) and phases + else "none" + ) + return "\n".join( + [ + "# LoopX Turn Progress", + f"- journal_status: {payload.get('journal_status')}", + f"- completed_phases: {completed}", + "- effects: none", + ] + ) + + def render_loopx_turn_managed_step_markdown(payload: dict[str, object]) -> str: if not payload.get("ok"): error = payload.get("error") or "Turn managed step failed" diff --git a/tests/test_loopx_turn_journal_inspection.py b/tests/test_loopx_turn_journal_inspection.py index cea505e151..7d77653359 100644 --- a/tests/test_loopx_turn_journal_inspection.py +++ b/tests/test_loopx_turn_journal_inspection.py @@ -4,6 +4,7 @@ import io import json from pathlib import Path +import threading from typing import Any import pytest @@ -11,9 +12,11 @@ from loopx.cli import main as cli_main from loopx.cli_commands import turn as turn_command from loopx.cli_commands import turn_decision +from loopx.cli_commands import turn_inspection from loopx.cli_commands import turn_rendering, turn_run_once, turn_todo_writeback from loopx.control_plane.turn_driver import executor from loopx.control_plane.turn_driver import turn_journal_runtime +from loopx.file_lock import exclusive_file_lock TURN_KEY = "sha256:" + "a" * 64 @@ -79,6 +82,8 @@ def _run_inspection_cli( goal_id: str = "fixture-goal", agent_id: str = "fixture-agent", turn_key: str = TURN_KEY, + watch: bool = False, + watch_interval: float = 0.01, ) -> tuple[int, str]: output = io.StringIO() with contextlib.redirect_stdout(output): @@ -96,6 +101,11 @@ def _run_inspection_cli( agent_id, "--turn-key", turn_key, + *( + ["--watch", "--watch-interval", str(watch_interval)] + if watch + else [] + ), "--format", output_format, ] @@ -255,6 +265,254 @@ def test_inspection_returns_versioned_allowlisted_projection_without_mutation( assert "do-not-expose" not in json.dumps(result) +def test_inspect_journal_watch_emits_only_changed_safe_progress_events( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + from loopx.control_plane.turn_driver.transaction import TRANSACTION_PHASES + + snapshots = iter([ + { + "ok": True, + "journal_consistent": True, + "journal_status": "in_progress", + "completed_phases": list(TRANSACTION_PHASES[:2]), + "effects": [], + "host_result": {"private": "must-not-leak"}, + }, + { + "ok": True, + "journal_consistent": True, + "journal_status": "in_progress", + "completed_phases": list(TRANSACTION_PHASES[:2]), + "effects": [], + }, + { + "ok": True, + "journal_consistent": True, + "journal_status": "in_progress", + "completed_phases": list(TRANSACTION_PHASES[:3]), + "effects": [], + }, + { + "ok": True, + "journal_consistent": True, + "journal_status": "committed", + "completed_phases": list(TRANSACTION_PHASES), + "effects": [], + }, + ]) + monkeypatch.setattr( + executor, + "inspect_loopx_turn_journal", + lambda *args, **kwargs: next(snapshots), + ) + monkeypatch.setattr( + "loopx.control_plane.turn_driver.inspect_loopx_turn_journal", + lambda *args, **kwargs: next(snapshots), + ) + monkeypatch.setattr(turn_inspection.time, "sleep", lambda _seconds: None) + + output = io.StringIO() + with contextlib.redirect_stdout(output): + exit_code = cli_main([ + "--registry", str(tmp_path / "registry.json"), + "--runtime-root", str(tmp_path), + "turn", "inspect-journal", + "--goal-id", "fixture-goal", + "--agent-id", "fixture-agent", + "--turn-key", TURN_KEY, + "--watch", + "--watch-interval", "0.01", + "--format", "json", + ]) + + lines = output.getvalue().splitlines() + events = [json.loads(line) for line in lines] + assert exit_code == 0 + assert len(events) == 3 + assert [event["journal_status"] for event in events] == [ + "in_progress", "in_progress", "committed", + ] + assert events[0]["completed_phases"] == list(TRANSACTION_PHASES[:2]) + assert events[1]["completed_phases"] == list(TRANSACTION_PHASES[:3]) + assert events[2]["completed_phases"] == list(TRANSACTION_PHASES) + assert all(event["effects"] == [] for event in events) + assert "must-not-leak" not in output.getvalue() + + +@pytest.mark.parametrize( + ("agent_id", "completed_phases", "expected_violation"), + [ + ("other-agent", COMPLETED_PHASES, "owner_mismatch"), + ("fixture-agent", ["host_execute", "not-a-transaction-phase"], + "completed_phases_not_ordered_prefix"), + ], +) +def test_inspect_journal_watch_rejects_inconsistent_real_journal( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, + agent_id: str, + completed_phases: list[str], + expected_violation: str, +) -> None: + journal = _journal() + journal["completed_phases"] = completed_phases + path = _write_journal(tmp_path, journal) + before = path.read_bytes() + monkeypatch.setattr(turn_inspection.time, "sleep", lambda _seconds: None) + + exit_code, raw_output = _run_inspection_cli( + tmp_path, output_format="json", agent_id=agent_id, watch=True, + ) + + assert exit_code == 1 + diagnostic = json.loads(raw_output) + assert diagnostic["schema_version"] == "loopx_turn_journal_inspection_v1" + assert diagnostic["journal_consistent"] is False + assert expected_violation in diagnostic["violations"] + assert diagnostic["schema_version"] != turn_inspection.TURN_PROGRESS_EVENT_SCHEMA_VERSION + assert path.read_bytes() == before + assert "do-not-expose" not in raw_output + + +def test_inspect_journal_watch_keeps_consistent_in_progress_observable( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + from loopx.control_plane.turn_driver.transaction import TRANSACTION_PHASES + + journal = _journal(status="in_progress") + journal["completed_phases"] = list(TRANSACTION_PHASES[:2]) + path = _write_journal(tmp_path, journal) + + def finish_committed_turn(_seconds: float) -> None: + committed = _journal() + path.write_text(json.dumps(committed, indent=2) + "\n", encoding="utf-8") + + monkeypatch.setattr(turn_inspection.time, "sleep", finish_committed_turn) + + exit_code, raw_output = _run_inspection_cli( + tmp_path, output_format="json", watch=True, + ) + + events = [json.loads(line) for line in raw_output.splitlines()] + assert exit_code == 0 + assert len(events) == 2 + assert events[0]["schema_version"] == turn_inspection.TURN_PROGRESS_EVENT_SCHEMA_VERSION + assert events[0]["journal_status"] == "in_progress" + assert events[0]["completed_phases"] == list(TRANSACTION_PHASES[:2]) + assert events[1]["journal_status"] == "committed" + assert events[1]["completed_phases"] == list(TRANSACTION_PHASES) + + +def test_inspect_journal_watch_rejects_non_positive_interval_before_read( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + def unexpected_read(*args: object, **kwargs: object) -> None: + raise AssertionError("invalid watch interval must fail before journal access") + + monkeypatch.setattr(executor, "inspect_loopx_turn_journal", unexpected_read) + monkeypatch.setattr( + "loopx.control_plane.turn_driver.inspect_loopx_turn_journal", + unexpected_read, + ) + output = io.StringIO() + with contextlib.redirect_stdout(output): + exit_code = cli_main([ + "--registry", str(tmp_path / "registry.json"), + "--runtime-root", str(tmp_path), + "turn", "inspect-journal", + "--goal-id", "fixture-goal", + "--agent-id", "fixture-agent", + "--turn-key", TURN_KEY, + "--watch", + "--watch-interval", "0", + "--format", "json", + ]) + + assert exit_code == 1 + assert "watch interval must be a finite positive number" in output.getvalue() + + +def test_inspect_journal_watch_observes_checkpoints_from_concurrent_writer( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + from loopx.control_plane.turn_driver.transaction import TRANSACTION_PHASES + + journal = _journal(status="in_progress") + journal["completed_phases"] = list(TRANSACTION_PHASES[:1]) + journal_path = _write_journal(tmp_path, journal) + updates = [ + (list(TRANSACTION_PHASES[:2]), "in_progress"), + (list(TRANSACTION_PHASES[:3]), "in_progress"), + (list(TRANSACTION_PHASES), "committed"), + ] + requested = threading.Event() + updated = threading.Event() + writer_errors: list[BaseException] = [] + + def write_checkpoints() -> None: + try: + for phases, status in updates: + assert requested.wait(timeout=5) + requested.clear() + with exclusive_file_lock(journal_path): + current = json.loads(journal_path.read_text(encoding="utf-8")) + current["completed_phases"] = phases + current["status"] = status + journal_path.write_text( + json.dumps(current, indent=2) + "\n", encoding="utf-8", + ) + updated.set() + except BaseException as exc: # surfaced in the command thread below + writer_errors.append(exc) + updated.set() + + writer = threading.Thread(target=write_checkpoints, daemon=True) + writer.start() + + def wait_for_checkpoint(_seconds: float) -> None: + requested.set() + assert updated.wait(timeout=5) + updated.clear() + if writer_errors: + raise writer_errors[0] + + monkeypatch.setattr(turn_inspection.time, "sleep", wait_for_checkpoint) + output = io.StringIO() + with contextlib.redirect_stdout(output): + exit_code = cli_main([ + "--registry", str(tmp_path / "registry.json"), + "--runtime-root", str(tmp_path), + "turn", "inspect-journal", + "--goal-id", "fixture-goal", + "--agent-id", "fixture-agent", + "--turn-key", TURN_KEY, + "--watch", + "--watch-interval", "0.01", + "--format", "json", + ]) + writer.join(timeout=5) + + assert not writer.is_alive() + assert not writer_errors + assert exit_code == 0 + events = [json.loads(line) for line in output.getvalue().splitlines()] + assert [event["completed_phases"] for event in events] == [ + list(TRANSACTION_PHASES[:1]), + list(TRANSACTION_PHASES[:2]), + list(TRANSACTION_PHASES[:3]), + list(TRANSACTION_PHASES), + ] + assert [event["journal_status"] for event in events] == [ + "in_progress", "in_progress", "in_progress", "committed", + ] + assert all(event["effects"] == [] for event in events) + + def test_inspection_reports_blocked_journal_as_successful_read(tmp_path: Path) -> None: journal = _journal(status="in_progress") journal["completed_phases"] = ["typed_result"]