From 5efd5331c44c0054047470bee56479c4dfb82a09 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 02:37:18 +0800 Subject: [PATCH 1/4] fix(turn): preserve original journal effect readback on execution errors Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- loopx/cli_commands/turn.py | 18 ++++++- loopx/cli_commands/turn_rendering.py | 38 +++++++++++--- .../control_plane/turn_driver/turn_journal.ts | 5 ++ .../turn_journal_effect_readback.ts | 51 +++++++++++++++++++ .../turn_driver/turn_journal_runtime.py | 7 +++ 5 files changed, 112 insertions(+), 7 deletions(-) create mode 100644 loopx/control_plane/turn_driver/turn_journal_effect_readback.ts diff --git a/loopx/cli_commands/turn.py b/loopx/cli_commands/turn.py index 23eef2a5f8..a0cc0c7d88 100644 --- a/loopx/cli_commands/turn.py +++ b/loopx/cli_commands/turn.py @@ -50,6 +50,7 @@ load_loopx_turn_plan_from_journal, run_codex_cli_host, run_loopx_turn_once, + inspect_loopx_turn_journal, selected_turn_todo, ) from ..control_plane.turn_driver.host_binding import managed_executor_binding @@ -114,6 +115,7 @@ def handle_turn_command( output_format=output_format, print_payload=print_payload, ) payload: dict[str, Any] = {} + execution_started = False try: if getattr(args, "todo_id", None) is not None and ( getattr(args, "resume_turn_key", None) @@ -1056,6 +1058,7 @@ def on_managed_start_admitted() -> None: on_admitted=on_managed_start_admitted, ) + execution_started = bool(args.execute) payload = run_loopx_turn_once( payload, host_argv=raw_argv, @@ -1090,7 +1093,20 @@ def on_managed_start_admitted() -> None: else: raise ValueError("turn requires the `plan` or `run-once` subcommand") except Exception as exc: # noqa: BLE001 - CLI boundary renders typed JSON failure - payload = build_turn_error_payload(payload, exc, turn_command=args.turn_command) + journal_readback = None + if execution_started: + transaction = payload.get("transaction") or {} + try: + journal_readback = inspect_loopx_turn_journal( + runtime_root, goal_id=args.goal_id, agent_id=args.agent_id, + turn_key=str(transaction.get("turn_key") or ""), + ) + except Exception: # noqa: BLE001 - retain original error and unknown effects + pass + payload = build_turn_error_payload( + payload, exc, turn_command=args.turn_command, + execution_started=execution_started, journal_readback=journal_readback, + ) renderer = ( _render_loopx_turn_execution_markdown if args.turn_command == "run-once" diff --git a/loopx/cli_commands/turn_rendering.py b/loopx/cli_commands/turn_rendering.py index d22356e850..110476a4b6 100644 --- a/loopx/cli_commands/turn_rendering.py +++ b/loopx/cli_commands/turn_rendering.py @@ -15,6 +15,8 @@ def build_turn_error_payload( planned: dict[str, Any], exc: Exception, *, turn_command: str, + execution_started: bool = False, + journal_readback: Mapping[str, Any] | None = None, ) -> dict[str, Any]: """Preserve the failed Turn's identity and effect readback at the CLI edge.""" @@ -24,6 +26,15 @@ def build_turn_error_payload( run_once = turn_command == "run-once" error_code = getattr(exc, "code", None) error_payload = getattr(exc, "payload", None) + planned_effects = planned.get("effects") + planned_effects = planned_effects if isinstance(planned_effects, Mapping) else {} + # Once the executor was entered an exception can follow a protected effect. + # A missing reply is not proof of non-execution. The original journal is + # projected separately so a replay never looks like a second host launch. + effects = { + key: True if planned_effects.get(key) is True else None if execution_started else False + for key in ("host_invoked", "state_written", "scheduler_acknowledged", "quota_spent") + } return { **({"error_code": error_code, **(error_payload if isinstance(error_payload, Mapping) else {})} if isinstance(error_code, str) else {}), @@ -33,12 +44,16 @@ def build_turn_error_payload( ), "mode": "run_once" if run_once else "plan", "error": str(exc), - "effects": { - "host_invoked": False, - "state_written": False, - "scheduler_acknowledged": False, - "quota_spent": False, - }, + "effects": effects, + "effects_scope": "current_invocation", + **({"journal_observation": { + "scope": "original_turn", + "status": "observed" if journal_readback is not None else "unavailable", + **({key: journal_readback[key] for key in ( + "journal_consistent", "journal_status", "completed_phases", + "recorded_effects", "recovery_decision", + )} if journal_readback is not None else {}), + }} if execution_started else {}), **({ "resume_turn_key": turn_key, "journal_ref": f"turn:{turn_key.removeprefix('sha256:')[:16]}", @@ -114,6 +129,10 @@ def render_loopx_turn_execution_markdown(payload: dict[str, object]) -> str: if isinstance(payload.get("host_failure"), dict) else {} ) + journal_observation = ( + payload.get("journal_observation") + if isinstance(payload.get("journal_observation"), dict) else {} + ) return "\n".join( [ "# LoopX Turn Run Once", @@ -145,6 +164,12 @@ def render_loopx_turn_execution_markdown(payload: dict[str, object]) -> str: f"- host_invoked: {effects.get('host_invoked')}", f"- state_written: {effects.get('state_written')}", f"- quota_spent: {effects.get('quota_spent')}", + *([f"- error: {payload['error']}"] if payload.get("error") else []), + *([f"- journal_readback: {journal_observation.get('status')}", + f"- recorded_effects: {journal_observation.get('recorded_effects')}", + f"- recovery_from: {(journal_observation.get('recovery_decision') or {}).get('resume_from')}", + f"- recovery_reinvoke_host: {(journal_observation.get('recovery_decision') or {}).get('reinvoke_host')}"] + if journal_observation else []), *( [ f"- recovery_plan: {planned.get('action')}", @@ -209,6 +234,7 @@ def render_loopx_turn_journal_inspection_markdown( f"- recovery_reason: {recovery.get('reason')}", f"- recovery_checks: {rendered_checks or 'none'}", f"- journal_consistent: {payload.get('journal_consistent')}", + f"- original_turn_recorded_effects: {payload.get('recorded_effects')}", *( [ f"- last_recovery_plan: {last_planned.get('action')}", diff --git a/loopx/control_plane/turn_driver/turn_journal.ts b/loopx/control_plane/turn_driver/turn_journal.ts index 37a628998f..9cbf2e9938 100644 --- a/loopx/control_plane/turn_driver/turn_journal.ts +++ b/loopx/control_plane/turn_driver/turn_journal.ts @@ -11,6 +11,7 @@ import { type EffectTurn, } from "../effect_program.ts"; import { EffectRuntimeRequestError } from "../effect_runtime_errors.ts"; +import { recordedTurnEffects, type RecordedTurnEffects } from "./turn_journal_effect_readback.ts"; export const TURN_JOURNAL_INSPECTION_SCHEMA_VERSION = "loopx_turn_journal_inspection_v1"; @@ -76,6 +77,7 @@ export interface TurnJournalInspection { journal_consistent: boolean; recovery_decision: TurnRecoveryDecision; last_recovery: TurnRecoveryAudit | null; + recorded_effects: RecordedTurnEffects; effects: []; } @@ -92,6 +94,7 @@ export interface TurnJournalEffectContext { journal_consistent: boolean; recovery_decision: TurnRecoveryDecision; last_recovery: TurnRecoveryAudit | null; + recorded_effects: RecordedTurnEffects; } // Replay has its own verdict. It is not a quota decision and must not manufacture @@ -657,6 +660,7 @@ export function interpretTurnJournalEffect( journal_consistent: journalConsistent, recovery_decision: turnRecoveryDecision, last_recovery: projectRecoveryAudit(journal.recovery_audit), + recorded_effects: recordedTurnEffects(journal, completedPhases, journalConsistent), }, }, interpretation: { @@ -709,6 +713,7 @@ export function projectTurnJournalInspection( journal_consistent: context.journal_consistent, recovery_decision: context.recovery_decision, last_recovery: context.last_recovery, + recorded_effects: context.recorded_effects, effects: [], }; } diff --git a/loopx/control_plane/turn_driver/turn_journal_effect_readback.ts b/loopx/control_plane/turn_driver/turn_journal_effect_readback.ts new file mode 100644 index 0000000000..2e8a2ecd8b --- /dev/null +++ b/loopx/control_plane/turn_driver/turn_journal_effect_readback.ts @@ -0,0 +1,51 @@ +/** Lower-bound observations of the original Turn, never new execution effects. */ +export interface RecordedTurnEffects { + host_invoked: boolean | null; + state_written: boolean | null; + quota_spent: boolean | null; + scheduler_acknowledged: boolean | null; +} + +type JsonObject = Record; + +function object(value: unknown): JsonObject { + return typeof value === "object" && value !== null && !Array.isArray(value) + ? value as JsonObject : {}; +} + +export function recordedTurnEffects( + journal: JsonObject, + completedPhases: readonly string[], + journalConsistent: boolean, +): RecordedTurnEffects { + const unknown: RecordedTurnEffects = { + host_invoked: null, state_written: null, + quota_spent: null, scheduler_acknowledged: null, + }; + // Foreign, corrupt or contradictory lineage cannot supply effect facts. + if (!journalConsistent) return unknown; + const completed = new Set(completedPhases); + const attempts = object(journal.effect_attempts); + const attemptsMalformed = journal.effect_attempts !== undefined + && (journal.effect_attempts === null || Array.isArray(journal.effect_attempts) + || typeof journal.effect_attempts !== "object"); + // A malformed attempt is not evidence of absence either. Validation and + // recovery admission still belong to their original owners. + const pending = (step: string) => attemptsMalformed || Object.hasOwn(attempts, step); + const scheduler = object(journal.scheduler); + // An attempt is persisted BEFORE confirmation/host launch. It proves neither + // launch nor non-launch until the host checkpoint is durable. + const hostCount = journal.host_attempt_count; + const hostUncertain = hostCount !== undefined + && (!Number.isInteger(hostCount) || Number(hostCount) !== 0); + return { + host_invoked: completed.has("host_execute") ? true : hostUncertain ? null : false, + state_written: completed.has("durable_writeback") ? true + : pending("durable_writeback") || pending("terminal_closeout") ? null : false, + quota_spent: completed.has("quota_spend") ? true + : pending("quota_spend") ? null : false, + // A completed outer-controller phase need not acknowledge any host cadence. + scheduler_acknowledged: typeof scheduler.acknowledged === "boolean" + ? scheduler.acknowledged : completed.has("quota_spend") ? null : false, + }; +} diff --git a/loopx/control_plane/turn_driver/turn_journal_runtime.py b/loopx/control_plane/turn_driver/turn_journal_runtime.py index f6f09ddf85..0065faaf39 100644 --- a/loopx/control_plane/turn_driver/turn_journal_runtime.py +++ b/loopx/control_plane/turn_driver/turn_journal_runtime.py @@ -26,6 +26,7 @@ "journal_consistent", "recovery_decision", "last_recovery", + "recorded_effects", "effects", } _BOOLEAN_PROJECTION_KEYS = { @@ -145,6 +146,12 @@ def interpret_turn_journal_projection( or not all(isinstance(violation, str) for violation in payload["violations"]) or not _validate_recovery_decision(payload.get("recovery_decision")) or not _validate_recovery_audit(payload.get("last_recovery")) + or not isinstance(payload.get("recorded_effects"), dict) + or set(payload["recorded_effects"]) != { + "host_invoked", "state_written", "quota_spent", "scheduler_acknowledged", + } + or any(value is not None and not isinstance(value, bool) + for value in payload["recorded_effects"].values()) ): raise RuntimeError( "TypeScript Turn-journal inspection projection type mismatch" From cef790dec6ae31ed48c397f759104f7ed44c2bad Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 02:37:19 +0800 Subject: [PATCH 2/4] test(turn): cover uncertain effects and checkpoint recovery Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- tests/control_plane_ts/turn_journal.test.ts | 53 ++++++++++++++++++++ tests/test_loopx_turn_driver.py | 7 +++ tests/test_loopx_turn_error_readback.py | 55 +++++++++++++++++++++ tests/test_loopx_turn_journal_inspection.py | 9 ++++ 4 files changed, 124 insertions(+) create mode 100644 tests/test_loopx_turn_error_readback.py diff --git a/tests/control_plane_ts/turn_journal.test.ts b/tests/control_plane_ts/turn_journal.test.ts index d90ac878ad..39d9bf946b 100644 --- a/tests/control_plane_ts/turn_journal.test.ts +++ b/tests/control_plane_ts/turn_journal.test.ts @@ -157,6 +157,8 @@ test("journal replay uses its own decision without a quota action slot", () => { checks: [{ kind: "journal_consistency", outcome: "passed" }], }, last_recovery: null, + recorded_effects: {host_invoked: true, state_written: true, + quota_spent: true, scheduler_acknowledged: null}, effects: [], }); }); @@ -181,6 +183,57 @@ test("non-terminal replay blocking does not block executor recovery", () => { }); }); +test("recorded effects distinguish checkpoint, prepared uncertainty and no attempt", () => { + const input = request("in_progress"); + input.journal.completed_phases = ["host_execute", "typed_result", "validation"]; + input.journal.effect_attempts = {durable_writeback: {status: "prepared"}}; + const result = interpretTurnJournal(input); + assert.deepEqual(result.recorded_effects, { + host_invoked: true, state_written: null, quota_spent: false, + scheduler_acknowledged: false, + }); + assert.equal(result.recovery_decision.reinvoke_host, false); + + input.journal.completed_phases = []; + delete input.journal.effect_attempts; + assert.equal(interpretTurnJournal(input).recorded_effects.host_invoked, false); + input.journal.host_attempt_count = 1; + assert.equal(interpretTurnJournal(input).recorded_effects.host_invoked, null); +}); + +test("pending spend and foreign lineage cannot certify a quota charge", () => { + const input = request("in_progress"); + input.journal.completed_phases = ["host_execute", "typed_result", "validation", "durable_writeback"]; + input.journal.effect_attempts = {quota_spend: {status: "prepared"}}; + assert.deepEqual(interpretTurnJournal(input).recorded_effects, { + host_invoked: true, state_written: true, quota_spent: null, + scheduler_acknowledged: false, + }); + assert.deepEqual(interpretTurnJournal({...input, agent_id: "another-caller"}).recorded_effects, { + host_invoked: null, state_written: null, quota_spent: null, + scheduler_acknowledged: null, + }); + input.journal.completed_phases = ["typed_result"]; + assert.equal(interpretTurnJournal(input).recorded_effects.host_invoked, null); +}); + +test("outer-controller completion does not manufacture scheduler acknowledgement", () => { + const input = request(); + input.journal.scheduler = {completed: true, acknowledged: false}; + assert.equal(interpretTurnJournal(input).recorded_effects.scheduler_acknowledged, false); +}); + +test("malformed attempt facts never certify absence of an effect", () => { + const input = request("in_progress"); + input.journal.completed_phases = []; + input.journal.host_attempt_count = "1"; + input.journal.effect_attempts = {durable_writeback: {status: "invalid"}}; + assert.equal(interpretTurnJournal(input).recorded_effects.host_invoked, null); + assert.equal(interpretTurnJournal(input).recorded_effects.state_written, null); + input.journal.effect_attempts = []; + assert.equal(interpretTurnJournal(input).recorded_effects.quota_spent, null); +}); + test("scheduler recovery resumes only scheduler apply", () => { const input = request("scheduler_action_required"); input.journal.completed_phases = [ diff --git a/tests/test_loopx_turn_driver.py b/tests/test_loopx_turn_driver.py index 0bcb53e19e..9b2f56cbb9 100644 --- a/tests/test_loopx_turn_driver.py +++ b/tests/test_loopx_turn_driver.py @@ -2274,6 +2274,13 @@ def crash_after_canonical_completion(**kwargs: object) -> dict[str, object]: assert first_exit_code == 1, first assert first["error"] == "injected crash after canonical Todo commit" + assert first["effects"]["host_invoked"] is None + observation = first["journal_observation"] + assert observation["scope"] == "original_turn" + assert observation["recorded_effects"]["host_invoked"] is True + assert observation["recorded_effects"]["state_written"] is None + assert observation["recorded_effects"]["quota_spent"] is False + assert observation["recovery_decision"]["reinvoke_host"] is False assert committed_results[0]["status"] == "done" assert committed_results[0]["provider_status"] == "applied" assert committed_results[0]["idempotent_replay"] is False diff --git a/tests/test_loopx_turn_error_readback.py b/tests/test_loopx_turn_error_readback.py new file mode 100644 index 0000000000..c5ebb3592a --- /dev/null +++ b/tests/test_loopx_turn_error_readback.py @@ -0,0 +1,55 @@ +"""Error readback preserves uncertainty without inventing execution facts.""" +from __future__ import annotations + +from loopx.cli_commands.turn_rendering import ( + build_turn_error_payload, + render_loopx_turn_execution_markdown, +) + + +def test_pre_execution_error_preserves_hook_effects_without_claiming_a_host() -> None: + result = build_turn_error_payload( + {"effects": {"state_written": True}}, ValueError("invalid adapter"), + turn_command="run-once", + ) + assert result["effects"] == { + "host_invoked": False, "state_written": True, + "quota_spent": False, "scheduler_acknowledged": False, + } + assert "journal_observation" not in result + + +def test_interrupted_executor_retains_original_turn_observation_not_a_new_launch() -> None: + inspection = { + "journal_consistent": True, "journal_status": "in_progress", + "completed_phases": ["host_execute", "typed_result", "validation"], + "recorded_effects": {"host_invoked": True, "state_written": None, + "quota_spent": False, "scheduler_acknowledged": False}, + "recovery_decision": {"resume_from": "durable_writeback", "reinvoke_host": False}, + "private_body": "never expose this", + } + result = build_turn_error_payload( + {"transaction": {"turn_key": "sha256:" + "a" * 64}}, + OSError("lost writeback reply"), turn_command="run-once", + execution_started=True, journal_readback=inspection, + ) + assert result["effects_scope"] == "current_invocation" + assert all(value is None for value in result["effects"].values()) + assert result["journal_observation"]["recorded_effects"] == inspection["recorded_effects"] + assert result["resume_turn_key"] == "sha256:" + "a" * 64 + assert "never expose" not in str(result) + rendered = render_loopx_turn_execution_markdown(result) + assert "lost writeback reply" in rendered + assert "recovery_from: durable_writeback" in rendered + assert "recovery_reinvoke_host: False" in rendered + + +def test_unreadable_journal_does_not_replace_original_error_with_false_effects() -> None: + result = build_turn_error_payload( + {}, ValueError("completion validation required"), turn_command="run-once", + execution_started=True, + ) + assert result["error"] == "completion validation required" + assert result["effects"]["host_invoked"] is None + assert result["effects"]["quota_spent"] is None + assert result["journal_observation"] == {"scope": "original_turn", "status": "unavailable"} diff --git a/tests/test_loopx_turn_journal_inspection.py b/tests/test_loopx_turn_journal_inspection.py index a67e4cd123..2b4f5271d2 100644 --- a/tests/test_loopx_turn_journal_inspection.py +++ b/tests/test_loopx_turn_journal_inspection.py @@ -246,6 +246,8 @@ def test_inspection_returns_versioned_allowlisted_projection_without_mutation( "checks": [{"kind": "journal_consistency", "outcome": "passed"}], }, "last_recovery": None, + "recorded_effects": {"host_invoked": True, "state_written": True, + "quota_spent": True, "scheduler_acknowledged": None}, "effects": [], } assert journal_path.read_bytes() == before_bytes @@ -391,6 +393,7 @@ def test_inspect_journal_cli_json_and_markdown_share_allowlisted_projection( "journal_consistent", "recovery_decision", "last_recovery", + "recorded_effects", "effects", } assert markdown_output == ( @@ -402,6 +405,8 @@ def test_inspect_journal_cli_json_and_markdown_share_allowlisted_projection( "- recovery_reason: terminal_result_retained\n" "- recovery_checks: journal_consistency:passed\n" "- journal_consistent: True\n" + "- original_turn_recorded_effects: {'host_invoked': True, 'state_written': True, " + "'quota_spent': True, 'scheduler_acknowledged': None}\n" "- replay_decision: replay_legal\n" "- journal_status: committed\n" "- replay_legal: True\n" @@ -537,6 +542,8 @@ def rpc(method: str, params: dict[str, object]) -> dict[str, object]: ], }, "last_recovery": None, + "recorded_effects": {"host_invoked": True, "state_written": True, + "quota_spent": True, "scheduler_acknowledged": None}, "effects": [], } @@ -600,6 +607,8 @@ def test_typescript_runtime_rejects_malformed_projection_types( "checks": [{"kind": "journal_consistency", "outcome": "passed"}], }, "last_recovery": None, + "recorded_effects": {"host_invoked": True, "state_written": True, + "quota_spent": True, "scheduler_acknowledged": None}, "effects": [], } From 39684b7659a05d59b131d4e87a5cdf45e656faf3 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 02:37:19 +0800 Subject: [PATCH 3/4] docs(turn): clarify invocation versus original Turn effects Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../rfcs/loopx-overall-roadmap-v0.md | 6 +++++ docs/reference/protocols/loopx-turn-v0.md | 24 +++++++++++++++++++ 2 files changed, 30 insertions(+) diff --git a/docs/architecture/rfcs/loopx-overall-roadmap-v0.md b/docs/architecture/rfcs/loopx-overall-roadmap-v0.md index 2692441143..8f89f97933 100644 --- a/docs/architecture/rfcs/loopx-overall-roadmap-v0.md +++ b/docs/architecture/rfcs/loopx-overall-roadmap-v0.md @@ -512,6 +512,12 @@ preflight through the actual Turn dry-run and selected executor/profile. Task admission, current pinned acceptance and runtime availability remain distinct; unprobed generic/cloud availability stays unknown. The same inspection is available to CLI and enabled MCP/new Chat tools, without a new state store or launch effect. +Turn exception readback now distinguishes uncertain invocation effects from the +original journal's checkpointed effects and prepared-effect recovery. This is a +bounded R2 recovery slice, not runtime provisioning or G1 completion; caller +admission, fresh task derivation and frontend/Lark qualification remain open. +中文:异常读回区分本次调用未知副作用与原 Turn 的持久观察/待核对效果,只完成 +R2 的一个恢复切片;调用方准入、新任务生成与前端/Lark 验收仍未闭环。 The same owner-local panel now opens current validated artifact text and its version/source identifiers, accepts feedback through the original coordinator inbox and exposes coordinator pause with its actual scope. Stale reads clear prior diff --git a/docs/reference/protocols/loopx-turn-v0.md b/docs/reference/protocols/loopx-turn-v0.md index 9076421284..e46705c77c 100644 --- a/docs/reference/protocols/loopx-turn-v0.md +++ b/docs/reference/protocols/loopx-turn-v0.md @@ -383,6 +383,30 @@ quota construction, scheduler context, planning, host invocation, settlement, spend, or state writeback. Its `effects` field is therefore always an empty list. +`recorded_effects` separately reports lower-bound observations of the **original +Turn**, not effects performed by this inspection or a recovery invocation. +Boolean values mean checkpointed execution or no recorded attempt; `null` means +unknown. A saved Host attempt precedes launch confirmation; a `prepared` writeback +or spend is not a commit receipt. Inconsistent/foreign lineage supplies only +unknown effect facts. A scheduler phase does not imply host acknowledgement. + +If executing `run-once` raises unexpectedly, its error response retains the +original error and resume key, marks uncertain current-invocation `effects` as +`null`, and adds `journal_observation` from this same read-only TS owner. Its +`scope=original_turn` prevents a saved Host result being mistaken for another +launch. Unavailable inspection remains explicit; there is no Python fallback. +Use the original `recovery_decision` and provider readback to recover, never a +fresh task or repeated model call inferred from a failed CLI reply. This readback +does not bypass controller completion validation, lease conflicts or quota gates. +Pre-execution failures retain known Turn-start hook writes without claiming Host +execution. Normal successful/replayed `effects` remain invocation-scoped. + +中文:`recorded_effects` 是原 Turn 的持久观察,不是本次检查或恢复又发生了副作用。 +`null` 表示未知:已登记 Host attempt 不等于模型已启动,`prepared` 不等于写回或扣额 +已提交。异常返回保留原错误/恢复身份,以同一 TS owner 的只读 `journal_observation` +区分本次调用与原 Turn;读回失败不猜测“没有执行”。按原恢复判定和 provider 回执 +继续,不能因 CLI 报错新建任务重跑模型,也不放松验收、租约或扣额门禁。 + Exit zero means that inspection completed, including when `decision` is `replay_blocked`. A non-zero exit means the command could not inspect the requested journal because a selector, file, JSON document, or schema was From 49ac1fcccc5f372e70b99347bd25a7dc7dbf04c5 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 03:10:02 +0800 Subject: [PATCH 4/4] fix(turn): share prepared intent checks across journal reads and writes Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- docs/reference/protocols/loopx-turn-v0.md | 16 +++++- .../control_plane/turn_driver/turn_journal.ts | 16 +++++- .../turn_journal_attempt_contract.ts | 57 +++++++++++++++++++ .../turn_journal_effect_readback.ts | 22 ++++--- .../turn_driver/turn_journal_effects.ts | 49 +++------------- tests/control_plane_ts/turn_journal.test.ts | 37 +++++++++++- .../turn_journal_effects.test.ts | 57 +++++++++++++++++++ tests/test_loopx_turn_journal_inspection.py | 22 ++++++- 8 files changed, 216 insertions(+), 60 deletions(-) create mode 100644 loopx/control_plane/turn_driver/turn_journal_attempt_contract.ts diff --git a/docs/reference/protocols/loopx-turn-v0.md b/docs/reference/protocols/loopx-turn-v0.md index e46705c77c..5225a313cf 100644 --- a/docs/reference/protocols/loopx-turn-v0.md +++ b/docs/reference/protocols/loopx-turn-v0.md @@ -387,8 +387,16 @@ list. Turn**, not effects performed by this inspection or a recovery invocation. Boolean values mean checkpointed execution or no recorded attempt; `null` means unknown. A saved Host attempt precedes launch confirmation; a `prepared` writeback -or spend is not a commit receipt. Inconsistent/foreign lineage supplies only -unknown effect facts. A scheduler phase does not imply host acknowledgement. +or spend is not a commit receipt. Inconsistent/foreign identity or phase lineage +supplies only unknown effect facts. A scheduler phase does not imply host acknowledgement. + +The inspector and journal writer share the original TS prepared-intent contract: +one supported step, object shape, `prepared` status, exact settlement effect ref, +the next phase and a nonterminal status. Unknown or malformed intents block the +existing recovery decision, including Host reinvocation. Valid lineage's already +completed effects remain proved; other effects are unknown, not `false`. The +writer's empty-map rejection is retained: no pending intent uses an absent field; +committed history is carried by completed checkpoints, not retained intents. If executing `run-once` raises unexpectedly, its error response retains the original error and resume key, marks uncertain current-invocation `effects` as @@ -406,6 +414,10 @@ execution. Normal successful/replayed `effects` remain invocation-scoped. 已提交。异常返回保留原错误/恢复身份,以同一 TS owner 的只读 `journal_observation` 区分本次调用与原 Turn;读回失败不猜测“没有执行”。按原恢复判定和 provider 回执 继续,不能因 CLI 报错新建任务重跑模型,也不放松验收、租约或扣额门禁。 +检查与写入共用原 TS prepared-intent 合同,校验步骤、形状、状态、effect 身份和 +阶段绑定。未知或畸形 intent 阻断原恢复判定,不能建议重调 Host;保留合法身份与 +阶段已经证明的执行事实,其余返回未知而非 `false`。不放松原 writer 的空 map +拒绝语义:无待决 intent 应省略该字段,已提交历史由完成 checkpoint 表达。 Exit zero means that inspection completed, including when `decision` is `replay_blocked`. A non-zero exit means the command could not inspect the diff --git a/loopx/control_plane/turn_driver/turn_journal.ts b/loopx/control_plane/turn_driver/turn_journal.ts index 9cbf2e9938..3b10f9e0ba 100644 --- a/loopx/control_plane/turn_driver/turn_journal.ts +++ b/loopx/control_plane/turn_driver/turn_journal.ts @@ -11,6 +11,7 @@ import { type EffectTurn, } from "../effect_program.ts"; import { EffectRuntimeRequestError } from "../effect_runtime_errors.ts"; +import { preparedAttemptViolation } from "./turn_journal_attempt_contract.ts"; import { recordedTurnEffects, type RecordedTurnEffects } from "./turn_journal_effect_readback.ts"; export const TURN_JOURNAL_INSPECTION_SCHEMA_VERSION = @@ -622,8 +623,7 @@ export function interpretTurnJournalEffect( violations.push("journal_status_unsupported"); } - const replayLegal = violations.length === 0; - const journalConsistent = + const lineageConsistent = goalMatches && ownerMatches && settlementIdentityValid && @@ -632,6 +632,14 @@ export function interpretTurnJournalEffect( turnKeyMatches && phasesFormOrderedPrefix && supportedJournalStatuses.has(journalStatus); + const attemptViolation = preparedAttemptViolation(journal, { + status: journalStatus, + completedPhases, + effectId: settlementIdentityFromPlan(transaction).value?.effect_id ?? "", + }); + if (attemptViolation) violations.push(attemptViolation.code); + const replayLegal = violations.length === 0; + const journalConsistent = lineageConsistent && attemptViolation === null; const decision = replayLegal ? "replay_legal" : "replay_blocked"; const turnRecoveryDecision = recoveryDecision( request, @@ -660,7 +668,9 @@ export function interpretTurnJournalEffect( journal_consistent: journalConsistent, recovery_decision: turnRecoveryDecision, last_recovery: projectRecoveryAudit(journal.recovery_audit), - recorded_effects: recordedTurnEffects(journal, completedPhases, journalConsistent), + recorded_effects: recordedTurnEffects( + journal, completedPhases, lineageConsistent, attemptViolation === null, + ), }, }, interpretation: { diff --git a/loopx/control_plane/turn_driver/turn_journal_attempt_contract.ts b/loopx/control_plane/turn_driver/turn_journal_attempt_contract.ts new file mode 100644 index 0000000000..90e90a1b84 --- /dev/null +++ b/loopx/control_plane/turn_driver/turn_journal_attempt_contract.ts @@ -0,0 +1,57 @@ +/** The journal owner's prepared-intent contract, shared by writes and reads. */ +import transactionContract from "../turn_transaction_contract.json" with { type: "json" }; +import { SETTLEMENT_STEP_KINDS, type JsonObject } from "../effect_program.ts"; + +const preparedStepKinds: ReadonlySet = new Set( + SETTLEMENT_STEP_KINDS.filter((kind) => kind !== "validation"), +); + +interface AttemptViolation { + code: string; + message: string; +} + +function isObject(value: unknown): value is JsonObject { + return typeof value === "object" && value !== null && !Array.isArray(value); +} + +export function preparedAttemptViolation( + journal: JsonObject, + state: { status: string; completedPhases: readonly string[]; effectId: string }, +): AttemptViolation | null { + if (journal.effect_attempts === undefined) return null; + if (!isObject(journal.effect_attempts)) { + return { code: "prepared_effect_attempts_invalid", + message: "Turn journal prepared effects must be an object" }; + } + const entries = Object.entries(journal.effect_attempts); + // Preserve the writer contract: no pending intent is represented by an + // absent field, not by an empty map or a set of historical committed intents. + if (entries.length !== 1) { + return { code: "prepared_effect_count_invalid", + message: "Turn journal must carry at most one prepared effect" }; + } + const [stepKind, attempt] = entries[0]; + if (!preparedStepKinds.has(stepKind)) { + return { code: "prepared_effect_step_unsupported", + message: "Turn journal carries an unsupported prepared effect step" }; + } + if (!isObject(attempt) || attempt.status !== "prepared" + || attempt.effect_ref !== `${state.effectId}#${stepKind}`) { + return { code: "prepared_effect_identity_invalid", + message: "Turn journal prepared effect does not match settlement identity" }; + } + const phases: readonly string[] = transactionContract.phases; + const phaseIndex = phases.indexOf(stepKind === "terminal_closeout" + ? "scheduler_apply" : stepKind); + if (phaseIndex < 0 || state.completedPhases.length !== phaseIndex + || !state.completedPhases.every((phase, index) => phase === phases[index])) { + return { code: "prepared_effect_phase_invalid", + message: "Turn journal prepared effect is not the next settlement step" }; + } + if (!["in_progress", "failed"].includes(state.status)) { + return { code: "prepared_effect_status_invalid", + message: "Turn journal terminal state cannot retain a prepared effect" }; + } + return null; +} diff --git a/loopx/control_plane/turn_driver/turn_journal_effect_readback.ts b/loopx/control_plane/turn_driver/turn_journal_effect_readback.ts index 2e8a2ecd8b..b3867f9754 100644 --- a/loopx/control_plane/turn_driver/turn_journal_effect_readback.ts +++ b/loopx/control_plane/turn_driver/turn_journal_effect_readback.ts @@ -16,23 +16,27 @@ function object(value: unknown): JsonObject { export function recordedTurnEffects( journal: JsonObject, completedPhases: readonly string[], - journalConsistent: boolean, + lineageConsistent: boolean, + attemptsValid: boolean, ): RecordedTurnEffects { const unknown: RecordedTurnEffects = { host_invoked: null, state_written: null, quota_spent: null, scheduler_acknowledged: null, }; // Foreign, corrupt or contradictory lineage cannot supply effect facts. - if (!journalConsistent) return unknown; + if (!lineageConsistent) return unknown; const completed = new Set(completedPhases); - const attempts = object(journal.effect_attempts); - const attemptsMalformed = journal.effect_attempts !== undefined - && (journal.effect_attempts === null || Array.isArray(journal.effect_attempts) - || typeof journal.effect_attempts !== "object"); - // A malformed attempt is not evidence of absence either. Validation and - // recovery admission still belong to their original owners. - const pending = (step: string) => attemptsMalformed || Object.hasOwn(attempts, step); const scheduler = object(journal.scheduler); + // Unknown or contradictory intents cannot prove non-execution. Keep only + // facts proved by the valid lineage's already completed checkpoints. + if (!attemptsValid) return { + host_invoked: completed.has("host_execute") ? true : null, + state_written: completed.has("durable_writeback") ? true : null, + quota_spent: completed.has("quota_spend") ? true : null, + scheduler_acknowledged: scheduler.acknowledged === true ? true : null, + }; + const attempts = object(journal.effect_attempts); + const pending = (step: string) => Object.hasOwn(attempts, step); // An attempt is persisted BEFORE confirmation/host launch. It proves neither // launch nor non-launch until the host checkpoint is durable. const hostCount = journal.host_attempt_count; diff --git a/loopx/control_plane/turn_driver/turn_journal_effects.ts b/loopx/control_plane/turn_driver/turn_journal_effects.ts index 168337957c..c86da287dc 100644 --- a/loopx/control_plane/turn_driver/turn_journal_effects.ts +++ b/loopx/control_plane/turn_driver/turn_journal_effects.ts @@ -3,10 +3,7 @@ import { readFile } from "node:fs/promises"; import { isAbsolute } from "node:path"; import type { JsonObject } from "../effect_program.ts"; -import { - SETTLEMENT_STEP_KINDS, - settlementIdentityFromPlan, -} from "../effect_program.ts"; +import { settlementIdentityFromPlan } from "../effect_program.ts"; import { EffectRuntimeConflictError, EffectRuntimeRequestError, @@ -16,6 +13,7 @@ import { withFileMutationLock, } from "../effect_runtime_io.ts"; import { requireNonEmptyString as requiredString } from "../runtime_decode.ts"; +import { preparedAttemptViolation } from "./turn_journal_attempt_contract.ts"; import { interpretTurnJournalEffect, supportedJournalStatuses, @@ -23,9 +21,6 @@ import { } from "./turn_journal.ts"; const terminalStatuses = new Set(["committed", "stopped"]); -const preparedStepKinds: ReadonlySet = new Set( - SETTLEMENT_STEP_KINDS.filter((kind) => kind !== "validation"), -); const statusTransitions: Readonly>> = { in_progress: new Set([ "in_progress", @@ -91,39 +86,6 @@ function samePrefix(left: readonly string[], right: readonly string[]): boolean return left.length <= right.length && left.every((phase, index) => phase === right[index]); } -function requireValidPreparedAttempt( - journal: JsonObject, - state: Pick, -): void { - if (journal.effect_attempts === undefined) return; - const attempts = asObject(journal.effect_attempts); - if (Object.keys(attempts).length !== 1) { - conflict("Turn journal must carry at most one prepared effect"); - } - const [stepKind, rawAttempt] = Object.entries(attempts)[0] ?? []; - if (!stepKind || !preparedStepKinds.has(stepKind)) { - conflict("Turn journal carries an unsupported prepared effect step"); - } - const attempt = asObject(rawAttempt); - const effectRef = `${state.effectId}#${stepKind}`; - if (attempt.status !== "prepared" || attempt.effect_ref !== effectRef) { - conflict("Turn journal prepared effect does not match settlement identity"); - } - const phaseIndex = stepKind === "terminal_closeout" - ? transactionPhases.indexOf("scheduler_apply") - : transactionPhases.indexOf(stepKind); - if ( - phaseIndex < 0 || - state.completedPhases.length !== phaseIndex || - !samePrefix(state.completedPhases, transactionPhases) - ) { - conflict("Turn journal prepared effect is not the next settlement step"); - } - if (!["in_progress", "failed"].includes(state.status)) { - conflict("Turn journal terminal state cannot retain a prepared effect"); - } -} - function requireJournalState(journal: JsonObject): JournalState { if (journal.schema_version !== "loopx_turn_journal_v0") { throw new EffectRuntimeRequestError( @@ -145,6 +107,12 @@ function requireJournalState(journal: JsonObject): JournalState { }); const context = effect.request.context; const effectId = journalEffectId(journal); + const attemptViolation = preparedAttemptViolation(journal, { + status: context.journal_status, + completedPhases: context.completed_phases, + effectId: effectId ?? "", + }); + if (attemptViolation) conflict(attemptViolation.message); if (!context.journal_consistent || !effectId) { throw new EffectRuntimeRequestError( `Turn journal snapshot is inconsistent: ${context.violations.join(", ")}`, @@ -190,7 +158,6 @@ function requireJournalState(journal: JsonObject): JournalState { } } const state = { status, completedPhases, effectId, failedPhase }; - requireValidPreparedAttempt(journal, state); return state; } diff --git a/tests/control_plane_ts/turn_journal.test.ts b/tests/control_plane_ts/turn_journal.test.ts index 39d9bf946b..0745e61c18 100644 --- a/tests/control_plane_ts/turn_journal.test.ts +++ b/tests/control_plane_ts/turn_journal.test.ts @@ -186,7 +186,9 @@ test("non-terminal replay blocking does not block executor recovery", () => { test("recorded effects distinguish checkpoint, prepared uncertainty and no attempt", () => { const input = request("in_progress"); input.journal.completed_phases = ["host_execute", "typed_result", "validation"]; - input.journal.effect_attempts = {durable_writeback: {status: "prepared"}}; + input.journal.effect_attempts = {durable_writeback: { + status: "prepared", effect_ref: `${effectId("fixture-goal", "fixture-agent")}#durable_writeback`, + }}; const result = interpretTurnJournal(input); assert.deepEqual(result.recorded_effects, { host_invoked: true, state_written: null, quota_spent: false, @@ -204,7 +206,9 @@ test("recorded effects distinguish checkpoint, prepared uncertainty and no attem test("pending spend and foreign lineage cannot certify a quota charge", () => { const input = request("in_progress"); input.journal.completed_phases = ["host_execute", "typed_result", "validation", "durable_writeback"]; - input.journal.effect_attempts = {quota_spend: {status: "prepared"}}; + input.journal.effect_attempts = {quota_spend: { + status: "prepared", effect_ref: `${effectId("fixture-goal", "fixture-agent")}#quota_spend`, + }}; assert.deepEqual(interpretTurnJournal(input).recorded_effects, { host_invoked: true, state_written: true, quota_spent: null, scheduler_acknowledged: false, @@ -234,6 +238,32 @@ test("malformed attempt facts never certify absence of an effect", () => { assert.equal(interpretTurnJournal(input).recorded_effects.quota_spent, null); }); +test("unknown prepared steps block recovery without false non-execution facts", () => { + for (const completedPhases of [[], ["host_execute", "typed_result", "validation"]]) { + const input = request("in_progress"); + input.journal.completed_phases = completedPhases; + input.journal.effect_attempts = {unknown_provider_step: { + status: "prepared", + effect_ref: `${effectId("fixture-goal", "fixture-agent")}#durable_writeback`, + }}; + input.journal.scheduler = {acknowledged: false}; + const before = structuredClone(input); + const result = interpretTurnJournal(input); + assert.equal(result.journal_consistent, false); + assert.ok(result.violations.includes("prepared_effect_step_unsupported")); + assert.equal(result.recovery_decision.action, "blocked"); + assert.equal(result.recovery_decision.can_continue, false); + assert.equal(result.recovery_decision.reinvoke_host, false); + assert.equal(result.recovery_decision.resume_from, null); + assert.deepEqual(result.recorded_effects, { + host_invoked: completedPhases.length ? true : null, + state_written: null, quota_spent: null, scheduler_acknowledged: null, + }); + assert.deepEqual(result.effects, []); + assert.deepEqual(input, before); + } +}); + test("scheduler recovery resumes only scheduler apply", () => { const input = request("scheduler_action_required"); input.journal.completed_phases = [ @@ -473,7 +503,8 @@ test("prepared effects are delegated to the existing provider readback step", () const input = request("in_progress"); input.journal.completed_phases = ["host_execute", "typed_result", "validation"]; input.journal.effect_attempts = { - durable_writeback: { status: "prepared", effect_ref: "effect:fixture" }, + durable_writeback: { status: "prepared", + effect_ref: `${effectId("fixture-goal", "fixture-agent")}#durable_writeback` }, }; const result = interpretTurnJournal(input); diff --git a/tests/control_plane_ts/turn_journal_effects.test.ts b/tests/control_plane_ts/turn_journal_effects.test.ts index 9e53c2710a..11d6bc6b07 100644 --- a/tests/control_plane_ts/turn_journal_effects.test.ts +++ b/tests/control_plane_ts/turn_journal_effects.test.ts @@ -5,6 +5,7 @@ import { join } from "node:path"; import test from "node:test"; import { commitTurnJournal } from "../../loopx/control_plane/turn_driver/turn_journal_effects.ts"; +import { interpretTurnJournal } from "../../loopx/control_plane/turn_driver/turn_journal.ts"; const turnKey = `sha256:${"a".repeat(64)}`; const todoId = "todo_fixture0001"; @@ -215,6 +216,62 @@ test("prepared effects must use the settlement identity", async () => { }); }); +test("inspector and writer share prepared-intent rejection and valid controls", async () => { + const prepared = {status: "prepared", effect_ref: `${effectId()}#durable_writeback`}; + const cases = [ + {attempts: {unknown_provider_step: prepared}, code: "prepared_effect_step_unsupported"}, + {attempts: [], code: "prepared_effect_attempts_invalid"}, + {attempts: null, code: "prepared_effect_attempts_invalid"}, + {attempts: {}, code: "prepared_effect_count_invalid"}, + {attempts: {durable_writeback: prepared, quota_spend: prepared}, code: "prepared_effect_count_invalid"}, + {attempts: {durable_writeback: null}, code: "prepared_effect_identity_invalid"}, + {attempts: {durable_writeback: {...prepared, status: "committed"}}, code: "prepared_effect_identity_invalid"}, + {attempts: {durable_writeback: {...prepared, effect_ref: "foreign#durable_writeback"}}, code: "prepared_effect_identity_invalid"}, + {attempts: {quota_spend: {status: "prepared", effect_ref: `${effectId()}#quota_spend`}}, code: "prepared_effect_phase_invalid"}, + ]; + await withJournalPath(async (path) => { + await commit(path, journal()); + await commit(path, journal("in_progress", phases.slice(0, 2))); + await commit(path, journal("in_progress", phases.slice(0, 3))); + const inspect = (snapshot: Record) => interpretTurnJournal({ + schema_version: "loopx_turn_journal_interpretation_request_v0", + journal: snapshot, goal_id: "fixture-goal", agent_id: "fixture-agent", turn_key: turnKey, + }); + const before = await readFile(path, "utf8"); + for (const {attempts, code} of cases) { + const invalid = {...journal("in_progress", phases.slice(0, 3)), effect_attempts: attempts}; + const result = inspect(invalid); + assert.ok(result.violations.includes(code), code); + assert.equal(result.journal_consistent, false, code); + assert.equal(result.recovery_decision.can_continue, false, code); + assert.equal(result.recovery_decision.reinvoke_host, false, code); + assert.equal(result.recorded_effects.host_invoked, true, code); + assert.equal(result.recorded_effects.state_written, null, code); + assert.equal(result.recorded_effects.quota_spent, null, code); + assert.deepEqual(result.effects, []); + await assert.rejects(commit(path, invalid), /Turn journal/, code); + assert.equal(await readFile(path, "utf8"), before, code); + } + assert.equal(inspect(journal("in_progress", phases.slice(0, 3))).journal_consistent, true); + const valid = {...journal("in_progress", phases.slice(0, 3)), effect_attempts: {durable_writeback: prepared}}; + assert.equal(inspect(valid).recovery_decision.reason, "resolve_prepared_effect"); + await commit(path, valid); + const validCloseout = {...journal("in_progress", phases.slice(0, 5)), effect_attempts: { + terminal_closeout: {status: "prepared", effect_ref: `${effectId()}#terminal_closeout`}, + }}; + assert.equal(inspect(validCloseout).journal_consistent, true); + const invalidTerminalStatus = {...validCloseout, status: "committed"}; + assert.ok(inspect(invalidTerminalStatus).violations.includes("prepared_effect_status_invalid")); + await assert.rejects(commit(path, invalidTerminalStatus), /terminal state cannot retain/); + const terminalIntent = {...journal("committed", phases), effect_attempts: { + terminal_closeout: {status: "prepared", effect_ref: `${effectId()}#terminal_closeout`}, + }}; + assert.equal(inspect(terminalIntent).journal_consistent, false); + await assert.rejects(commit(path, terminalIntent), /prepared effect/); + assert.equal(inspect(journal("committed", phases)).recovery_decision.action, "return_existing"); + }); +}); + test("failed snapshots must name the next uncompleted phase", async () => { await withJournalPath(async (path) => { const invalid = journal("failed", phases.slice(0, 4)); diff --git a/tests/test_loopx_turn_journal_inspection.py b/tests/test_loopx_turn_journal_inspection.py index 2b4f5271d2..24b5bcfeb1 100644 --- a/tests/test_loopx_turn_journal_inspection.py +++ b/tests/test_loopx_turn_journal_inspection.py @@ -319,11 +319,20 @@ def unexpected_access(*args: object, **kwargs: object) -> None: ) +@pytest.mark.parametrize("unknown_intent", [False, True]) def test_inspect_journal_cli_branches_before_live_or_write_paths( tmp_path: Path, monkeypatch: pytest.MonkeyPatch, + unknown_intent: bool, ) -> None: - _write_journal(tmp_path, _journal()) + journal = _journal(status="in_progress" if unknown_intent else "committed") + if unknown_intent: + journal["completed_phases"] = [] + journal["effect_attempts"] = {"unknown_provider_step": { + "status": "prepared", "effect_ref": "private-effect-ref-not-returned", + }} + journal_path = _write_journal(tmp_path, journal) + before = journal_path.read_bytes() def unexpected_call(*args: object, **kwargs: object) -> None: raise AssertionError("inspect-journal reached a live or write path") @@ -359,8 +368,17 @@ def unexpected_call(*args: object, **kwargs: object) -> None: payload = json.loads(raw_output) assert exit_code == 0 assert payload["schema_version"] == "loopx_turn_journal_inspection_v1" - assert payload["decision"] == "replay_legal" + assert payload["decision"] == ("replay_blocked" if unknown_intent else "replay_legal") assert payload["effects"] == [] + assert journal_path.read_bytes() == before + if unknown_intent: + assert payload["journal_consistent"] is False + assert "prepared_effect_step_unsupported" in payload["violations"] + assert payload["recovery_decision"]["action"] == "blocked" + assert payload["recovery_decision"]["can_continue"] is False + assert payload["recovery_decision"]["reinvoke_host"] is False + assert all(value is None for value in payload["recorded_effects"].values()) + assert "private-effect-ref-not-returned" not in raw_output def test_inspect_journal_cli_json_and_markdown_share_allowlisted_projection(