From 118546a3e82de8f19505a953acb8a36d074a4f92 Mon Sep 17 00:00:00 2001 From: Tartar Date: Thu, 10 Sep 2026 17:24:31 +0800 Subject: [PATCH 1/2] fix(refresh): require acknowledgement to resume paused external delivery Signed-off-by: Tartar --- loopx/cli_commands/project_lifecycle.py | 27 +++-- .../quota/refresh_external_delivery.py | 82 ++++++++++++++ .../quota/refresh_external_delivery.ts | 95 ++++++++++++++++ loopx/control_plane/quota/refresh_recovery.ts | 4 + loopx/control_plane/quota/settlement.py | 5 +- .../quota/settlement_readback.ts | 22 ++-- loopx/rollout_event_log.py | 1 + loopx/state_refresh.py | 42 +++----- .../test_refresh_checkpoint_isolation.py | 33 +++++- .../test_refresh_external_delivery.py | 101 ++++++++++++++++++ .../refresh_external_delivery.test.ts | 92 ++++++++++++++++ 11 files changed, 453 insertions(+), 51 deletions(-) create mode 100644 loopx/control_plane/quota/refresh_external_delivery.py create mode 100644 loopx/control_plane/quota/refresh_external_delivery.ts create mode 100644 tests/control_plane/test_refresh_external_delivery.py create mode 100644 tests/control_plane_ts/refresh_external_delivery.test.ts diff --git a/loopx/cli_commands/project_lifecycle.py b/loopx/cli_commands/project_lifecycle.py index 0cf6bbb5d7..88bfa12870 100644 --- a/loopx/cli_commands/project_lifecycle.py +++ b/loopx/cli_commands/project_lifecycle.py @@ -432,15 +432,21 @@ def register_project_lifecycle_commands( action="store_true", help="Do not refresh the shared global registry after writing the state run.", ) - refresh_state_parser.add_argument( + delivery_flags = refresh_state_parser.add_mutually_exclusive_group() + delivery_flags.add_argument( "--suppress-external-sinks", action="store_true", help=( "Keep enabled local projections active but suppress configured external " - "sink writes for this refresh. Pending sink digests remain retryable." + "sink writes for this refresh. Turn-bound retries require an explicit resume acknowledgement." ), ) + delivery_flags.add_argument( + "--resume-external-sinks", metavar="RESUME_KEY", + help="Acknowledge the current pause returned by a Turn-bound refresh recovery. Does not grant provider permissions.", + ) + read_only_map_parser = subparsers.add_parser( "read-only-map", help="Append a generic read-only project-map run for a connected project.", @@ -632,7 +638,11 @@ def handle_project_lifecycle_command( print_payload(payload, fmt, render_state_refresh_markdown) return 1 try: + if getattr(args, "resume_external_sinks", None) and not getattr(args, "turn_instance_id", None): + raise ValueError("--resume-external-sinks requires the original --turn-instance-id") payload = refresh_state_run( + external_delivery={"suppress": bool(args.suppress_external_sinks), + "resume_key": getattr(args, "resume_external_sinks", None)}, registry_path=registry_path, runtime_root_override=args.runtime_root, goal_id=args.goal_id, @@ -702,8 +712,9 @@ def handle_project_lifecycle_command( ) if projected_capabilities: payload["available_capabilities"] = projected_capabilities - payload["external_sink_delivery_authorized"] = not bool( - args.suppress_external_sinks + payload.setdefault( + "external_sink_delivery_authorized", + not bool(args.suppress_external_sinks or getattr(args, "turn_instance_id", None)), ) material_refresh_ready = bool( payload.get("ok") @@ -848,9 +859,7 @@ def handle_project_lifecycle_command( agent_id=args.agent_id, project=Path(args.project).expanduser() if args.project else None, state_file=Path(args.state_file).expanduser() if args.state_file else None, - external_sink_delivery_authorized=not bool( - args.suppress_external_sinks - ), + external_sink_delivery_authorized=payload["external_sink_delivery_authorized"] is True, syncer=lark_explore_graph_syncer( args.runtime_root, registry_path=registry_path, @@ -875,9 +884,7 @@ def handle_project_lifecycle_command( runtime_root_override=args.runtime_root, goal_id=args.goal_id, agent_id=args.agent_id, - external_sink_delivery_authorized=not bool( - args.suppress_external_sinks - ), + external_sink_delivery_authorized=payload["external_sink_delivery_authorized"] is True, ) except Exception: gate_sync = goal_channel_gate_sync_failure( diff --git a/loopx/control_plane/quota/refresh_external_delivery.py b/loopx/control_plane/quota/refresh_external_delivery.py new file mode 100644 index 0000000000..385813584d --- /dev/null +++ b/loopx/control_plane/quota/refresh_external_delivery.py @@ -0,0 +1,82 @@ +"""Persist the typed resume decision inside the existing refresh serialization lock.""" +from pathlib import Path +from typing import Any + +from ...rollout_event_log import append_rollout_event, build_rollout_event, rollout_event_log_path +from .settlement import QuotaSettlementReadback, settlement_result_payload + + +def finish_external_delivery_refresh( + payload: dict[str, Any], readback: QuotaSettlementReadback | None, + runtime_root: Path, *, dry_run: bool, +) -> dict[str, Any]: + if readback is None: + return payload + plan = readback.external_delivery + if not isinstance(plan, dict) or plan.get("schema_version") != "refresh_external_delivery_v0": + raise RuntimeError("TypeScript refresh external delivery result missing or invalid") + payload["external_delivery"] = {k: v for k, v in plan.items() if k != "transition"} + payload["external_sink_delivery_authorized"] = plan["authorized"] is True + transition = plan.get("transition") + if payload.get("ok") and transition and not dry_run: + identity = readback.identity.value + if identity is None: + raise RuntimeError("external delivery transition has no settlement identity") + append_rollout_event( + rollout_event_log_path(runtime_root, identity.goal_id), + build_rollout_event( + goal_id=identity.goal_id, agent_id=identity.agent_id, + todo_id=identity.todo_id, run_id=identity.turn_instance_id, + event_kind="refresh_external_delivery", status=transition["state"], + summary="Refresh external delivery preference recorded.", details=transition, + ), + ) + plan["transition"] = None # The planned journal effect was committed once. + return payload + + +def refresh_recovery_payload( + readback: QuotaSettlementReadback, *, registry_path: Path, + runtime_root: Path, goal_id: str, dry_run: bool, +) -> dict[str, Any] | None: + recovery = readback.refresh_recovery + identity = readback.identity.value + plan = readback.external_delivery + if not recovery or identity is None or not isinstance(plan, dict): + raise RuntimeError("TypeScript refresh admission result missing") + decision = recovery["decision"] + delivery_error = plan.get("error_code") if decision != "reject" else None + if decision not in {"replay", "repair_receipt", "reject"} and not delivery_error: + # Record a requested pause before a new writeback can commit. If later + # validation fails, retaining the pause is conservative and retryable. + if (plan.get("transition") or {}).get("state") == "paused": + finish_external_delivery_refresh({"ok": True}, readback, runtime_root, dry_run=dry_run) + return None + payload = { + **(readback.writeback_run or {}), "ok": decision != "reject" and not delivery_error, + "dry_run": dry_run, "appended": False, + "idempotent_replay": decision == "replay" and not delivery_error, + "receipt_repair_required": decision == "repair_receipt" and not dry_run, + "registry": str(registry_path), "runtime_root": str(runtime_root), "goal_id": goal_id, + "refresh_recovery": recovery, "settlement_identity": identity.as_dict(), + "settlement_result": settlement_result_payload(readback.delivery), + } + if decision == "reject": + payload["error"] = ( + f"{recovery['reason']}: committed writeback is unchanged; " + "do not begin a new Turn or repeat spend to repair it. " + "Retry the original delivery fields with only the missing vision decision; " + "if a newer vision already exists, inspect current quota instead." + ) + elif delivery_error: + payload["error_code"] = delivery_error + payload["refresh_recovery"] = {**recovery, "decision": "reject", "reason": delivery_error} + key = plan.get("resume_key") + payload["error"] = ( + f"{delivery_error}: external delivery was not resumed. " + "Keep --suppress-external-sinks for local recovery. " + + (f"To deliberately resume this operation, retry the same recovery command with " + f"--resume-external-sinks {key} instead of --suppress-external-sinks. " if key else "") + + "Existing provider permissions still apply; do not repeat business mutations or spend." + ) + return finish_external_delivery_refresh(payload, readback, runtime_root, dry_run=dry_run) diff --git a/loopx/control_plane/quota/refresh_external_delivery.ts b/loopx/control_plane/quota/refresh_external_delivery.ts new file mode 100644 index 0000000000..131e4b354c --- /dev/null +++ b/loopx/control_plane/quota/refresh_external_delivery.ts @@ -0,0 +1,95 @@ +/** Per-settlement pause acknowledgement; never a provider permission grant. */ +import { createHash } from "node:crypto"; +import type { JsonObject, SettlementIdentity } from "../effect_program.ts"; +import { EffectRuntimeRequestError } from "../effect_runtime_errors.ts"; +import { jsonObject, requireJsonObject } from "../runtime_decode.ts"; + +export const EXTERNAL_DELIVERY_SCHEMA = "refresh_external_delivery_v0"; +export const EXTERNAL_DELIVERY_EVENT = "refresh_external_delivery"; +export interface ExternalDeliveryRequest { + suppress: boolean; + resume_key: string | null; +} + +export function decodeExternalDelivery(value: unknown): ExternalDeliveryRequest | null { + if (value == null) return null; // Existing non-delivery callers remain local. + const input = requireJsonObject(value, "external_delivery"); + if (typeof input.suppress !== "boolean" || + (input.resume_key !== null && + (typeof input.resume_key !== "string" || !/^[a-f0-9]{64}$/.test(input.resume_key))) || + (input.suppress && input.resume_key !== null)) { + throw new EffectRuntimeRequestError("external_delivery requires suppress and a mutually exclusive resume key"); + } + return { suppress: input.suppress, resume_key: input.resume_key }; +} + +function pauseKey(identity: SettlementIdentity, precedingEvent: JsonObject | null): string { + return createHash("sha256").update(JSON.stringify([ + EXTERNAL_DELIVERY_SCHEMA, identity.effect_id, precedingEvent?.event_id ?? null, + ])).digest("hex"); +} + +export function refreshExternalDelivery( + request: ExternalDeliveryRequest | null, + identity: SettlementIdentity, + events: readonly JsonObject[], + deliveryStage: boolean, +): JsonObject { + let previous: JsonObject | null = null; + let state: "paused" | "ready" | null = null; + let key: string | null = null; + let invalid = false; + for (const event of events) { + if (event.event_kind !== EXTERNAL_DELIVERY_EVENT || event.goal_id !== identity.goal_id || + event.agent_id !== identity.agent_id || event.run_id !== identity.turn_instance_id || + (event.todo_id ?? null) !== identity.todo_id) continue; + const details = jsonObject(event.details); + if ((details?.replan_obligation_id ?? null) !== identity.replan_obligation_id) continue; + if (!details || details.schema_version !== EXTERNAL_DELIVERY_SCHEMA || + details.settlement_effect_id !== identity.effect_id || typeof event.event_id !== "string" || + (details.state !== "paused" && details.state !== "ready") || + typeof details.resume_key !== "string" || !/^[a-f0-9]{64}$/.test(details.resume_key)) { + invalid = true; + break; + } + if (event.event_id === previous?.event_id) continue; + if (details.state === "paused" + ? state === "paused" || details.resume_key !== pauseKey(identity, previous) + : state !== "paused" || details.resume_key !== key) { + invalid = true; + break; + } + state = details.state; + key = details.resume_key; + previous = event; + } + const result = (authorized: boolean, reason: string, error: string | null = null, + transition: JsonObject | null = null): JsonObject => ({ + schema_version: EXTERNAL_DELIVERY_SCHEMA, authorized, reason, + error_code: error, resume_key: key, transition, + }); + // No send phase: do not reject a local closeout or consume a resume acknowledgement. + if (!request || (!deliveryStage && !request.suppress)) return result(false, "local_recovery"); + if (invalid) return request.suppress ? result(false, "suppressed") + : result(false, "invalid_pause_history", "external_delivery_history_invalid"); + const transition = (next: "paused" | "ready", resumeKey: string): JsonObject => ({ + schema_version: EXTERNAL_DELIVERY_SCHEMA, settlement_effect_id: identity.effect_id, + replan_obligation_id: identity.replan_obligation_id, state: next, resume_key: resumeKey, + }); + if (request.suppress) { + if (state === "paused") return result(false, "suppressed"); + key = pauseKey(identity, previous); + return result(false, "suppressed", null, transition("paused", key)); + } + if (request.resume_key !== null) { + if (key === null || request.resume_key !== key) { + return result(false, "resume_key_mismatch", "external_delivery_resume_mismatch"); + } + return result(true, "resume_acknowledged", null, + state === "paused" ? transition("ready", key) : null); + } + if (state === "paused") { + return result(false, "resume_confirmation_required", "external_delivery_resume_required"); + } + return result(true, state === "ready" ? "resumed" : "legacy_per_call"); +} diff --git a/loopx/control_plane/quota/refresh_recovery.ts b/loopx/control_plane/quota/refresh_recovery.ts index 08e27f5a48..f80e0b750e 100644 --- a/loopx/control_plane/quota/refresh_recovery.ts +++ b/loopx/control_plane/quota/refresh_recovery.ts @@ -5,7 +5,10 @@ import { EffectRuntimeRequestError } from "../effect_runtime_errors.ts"; import { jsonObject, requireJsonObject } from "../runtime_decode.ts"; import { normalizeDeliveryWorkspaceSnapshot } from "../agents/delivery_workspace.ts"; +import { decodeExternalDelivery, type ExternalDeliveryRequest } from "./refresh_external_delivery.ts"; + export interface RefreshRetryRequest { + external_delivery?: ExternalDeliveryRequest | null; vision: JsonObject | null; unchanged_reason: string | null; merge_patch: boolean; @@ -34,6 +37,7 @@ export function decodeRefreshRetry(value: unknown): RefreshRetryRequest | null { return value; }; return { + external_delivery: decodeExternalDelivery(input.external_delivery), vision: input.vision === null ? null : requireJsonObject(input.vision, "refresh_retry.vision"), unchanged_reason: nullableString("unchanged_reason"), merge_patch: input.merge_patch === true, diff --git a/loopx/control_plane/quota/settlement.py b/loopx/control_plane/quota/settlement.py index 03543a814c..7a0d87c452 100644 --- a/loopx/control_plane/quota/settlement.py +++ b/loopx/control_plane/quota/settlement.py @@ -44,7 +44,8 @@ def _checkpoint_instructions(checkpoint: Mapping[str, Any]) -> str: "- Preserve: Keep original values and presence for target, scope, and isolation " "options: `--registry`, `--runtime-root`, `--project`, `--state-file`, " "`--progress-scope`, `--agent-lane`, `--available-capability`, " - "`--no-global-sync`, `--suppress-external-sinks`; identity and delivery options: " + "`--no-global-sync`, `--suppress-external-sinks`, `--resume-external-sinks`; " + "identity and delivery options: " "`--goal-id`, `--agent-id`, `--todo-id`, `--replan-obligation-id`, " "`--turn-instance-id`, `--completion-todo-id`, `--completion-turn-key`, " "`--classification`, `--recommended-action`, `--delivery-batch-scale`, " @@ -150,6 +151,7 @@ class QuotaSettlementReadback: monitor_phase: ReceiptBoundMonitorPhase | None replay_phase: ReceiptBoundReplayPhase | None refresh_recovery: dict[str, Any] | None = None + external_delivery: dict[str, Any] | None = None __all__ = [ @@ -304,6 +306,7 @@ def read_heartbeat_settlement( ), writeback_run=_optional_readback_record(payload.get("writeback_run")), refresh_recovery=_optional_readback_record(payload.get("refresh_recovery")), + external_delivery=_optional_readback_record(payload.get("external_delivery")), spend_run=_optional_readback_record(payload.get("spend_run")), heartbeat_receipt=_optional_readback_record(payload.get("heartbeat_receipt")), writeback_event=_optional_readback_record(payload.get("writeback_event")), diff --git a/loopx/control_plane/quota/settlement_readback.ts b/loopx/control_plane/quota/settlement_readback.ts index 1261dfbebe..d7259bb2c9 100644 --- a/loopx/control_plane/quota/settlement_readback.ts +++ b/loopx/control_plane/quota/settlement_readback.ts @@ -36,6 +36,8 @@ import { type RefreshRetryRequest, } from "./refresh_recovery.ts"; +import { refreshExternalDelivery } from "./refresh_external_delivery.ts"; + export const QUOTA_SETTLEMENT_READBACK_REQUEST_SCHEMA = "loopx_quota_settlement_readback_request_v0"; export const QUOTA_SETTLEMENT_READBACK_RESULT_SCHEMA = @@ -763,6 +765,15 @@ export async function readQuotaSettlement(value: unknown): Promise { normalizeDeliveryWorkspaceCausality(nestedCausality, identity.todo_id) ?? normalizeDeliveryWorkspaceCausality(flatCausality, identity.todo_id); + const recovery = request.refresh_retry === null ? null : refreshRecovery( + request.refresh_retry, writebackRun, writeback.failure === null, + workspaceCausality?.requirement, + writebackRun !== null && runs.slice(runs.indexOf(writebackRun) + 1).some((run) => + run.goal_id === identity.goal_id && run.agent_id === identity.agent_id && + (jsonObject(run.agent_vision) !== null || jsonObject(run.vision_checkpoint)?.required === true) + ), + ); + return { schema_version: QUOTA_SETTLEMENT_READBACK_RESULT_SCHEMA, found: true, @@ -776,13 +787,10 @@ export async function readQuotaSettlement(value: unknown): Promise { workspace_causality: workspaceCausality, semantic_replan_guard: projectSemanticReplanGuard(receiptDetails), writeback_run: writebackRun, - refresh_recovery: request.refresh_retry === null ? null : refreshRecovery( - request.refresh_retry, writebackRun, writeback.failure === null, - workspaceCausality?.requirement, - writebackRun !== null && runs.slice(runs.indexOf(writebackRun) + 1).some((run) => - run.goal_id === identity.goal_id && run.agent_id === identity.agent_id && - (jsonObject(run.agent_vision) !== null || jsonObject(run.vision_checkpoint)?.required === true) - ), + refresh_recovery: recovery, + external_delivery: request.refresh_retry === null ? null : refreshExternalDelivery( + request.refresh_retry.external_delivery ?? null, identity, events, + recovery?.decision !== "reject" && recovery?.decision !== "repair_receipt", ), spend_run: spendRun, heartbeat_receipt: heartbeatReceipt, diff --git a/loopx/rollout_event_log.py b/loopx/rollout_event_log.py index d7908334dc..df5c28e7f0 100644 --- a/loopx/rollout_event_log.py +++ b/loopx/rollout_event_log.py @@ -28,6 +28,7 @@ "quota_spend", "quota_void", "refresh_state", + "refresh_external_delivery", "research_evidence", "research_hypothesis", "todo_add", diff --git a/loopx/state_refresh.py b/loopx/state_refresh.py index ce58b9e068..8283611291 100644 --- a/loopx/state_refresh.py +++ b/loopx/state_refresh.py @@ -22,12 +22,14 @@ from .control_plane.agents.workspace_guard import ( capture_delivery_workspace, ) +from .control_plane.quota.refresh_external_delivery import ( + finish_external_delivery_refresh, refresh_recovery_payload, +) from .control_plane.quota.settlement import ( SettlementIdentity, read_heartbeat_settlement, render_first_refresh_checkpoint_hint, render_refresh_recovery_markdown, - settlement_result_payload, ) from .control_plane.quota.settlement_workspace_causality import resolve_settlement_workspace_requirement from .control_plane.quota.codex_session_usage import ( @@ -792,6 +794,7 @@ def refresh_state_run( usage_codex_session: Path | None = None, dry_run: bool, sync_global: bool = True, + external_delivery: dict[str, Any] | None = None, ) -> dict[str, Any]: safe_goal_id = validate_goal_id_path_segment(goal_id) validate_public_safe_text("classification", classification) @@ -874,6 +877,7 @@ def refresh_state_run( turn_instance_id=turn_instance_id, replan_obligation_id=normalized_replan_obligation_id, refresh_retry={ + "external_delivery": external_delivery, "vision": agent_vision_packet, "unchanged_reason": vision_unchanged_reason, "merge_patch": bool(merge_agent_vision_patch), @@ -905,33 +909,13 @@ def refresh_state_run( refresh_recovery = settlement_readback.refresh_recovery if not refresh_recovery: raise RuntimeError("settlement readback omitted refresh recovery admission") - decision = refresh_recovery["decision"] prior_writeback_run = settlement_readback.writeback_run - if decision in {"replay", "repair_receipt", "reject"}: - payload = { - **(prior_writeback_run or {}), - "ok": decision != "reject", - "dry_run": dry_run, - "appended": False, - "idempotent_replay": decision == "replay", - "receipt_repair_required": decision == "repair_receipt" and not dry_run, - "registry": str(registry_path), - "runtime_root": str(runtime_root), - "goal_id": safe_goal_id, - "refresh_recovery": refresh_recovery, - "settlement_identity": settlement_identity.as_dict(), - "settlement_result": settlement_result_payload( - settlement_readback.delivery - ), - } - if decision == "reject": - payload["error"] = ( - f"{refresh_recovery['reason']}: committed writeback is unchanged; " - "do not begin a new Turn or repeat spend to repair it. " - "Retry the original delivery fields with only the missing vision decision; " - "if a newer vision already exists, inspect current quota instead." - ) - return payload + recovery_payload = refresh_recovery_payload( + settlement_readback, registry_path=registry_path, runtime_root=runtime_root, + goal_id=safe_goal_id, dry_run=dry_run, + ) + if recovery_payload is not None: + return recovery_payload settlement_workspace_requirement = resolve_settlement_workspace_requirement( delivery_workspace_causality, settlement_binding_kind=settlement_identity.binding_kind.value ) @@ -1471,4 +1455,6 @@ def refresh_state_run( "raw_artifacts_copied": False, "recommended_action_copied": False, } - return payload + return finish_external_delivery_refresh( + payload, settlement_readback, runtime_root, dry_run=dry_run, + ) diff --git a/tests/control_plane/test_refresh_checkpoint_isolation.py b/tests/control_plane/test_refresh_checkpoint_isolation.py index 3ae11ba9c9..c8b9d12838 100644 --- a/tests/control_plane/test_refresh_checkpoint_isolation.py +++ b/tests/control_plane/test_refresh_checkpoint_isolation.py @@ -101,8 +101,8 @@ def _configure_sinks(registry, runtime): @pytest.mark.parametrize("hint_source", ["first", "replay"]) @pytest.mark.parametrize("baseline", [False, True]) -@pytest.mark.parametrize("isolated", [True, False], ids=["isolated", "send-control"]) -def test_stdout_driven_recovery_preserves_boundaries( +@pytest.mark.parametrize("isolated", [True, False], ids=["isolated", "confirmed-resume"]) +def test_stdout_recovery_requires_confirmation_to_resume_external_delivery( tmp_path, monkeypatch, capsys, hint_source, baseline, isolated, ): project, runtime, registry = _write_fixture(tmp_path) @@ -113,10 +113,10 @@ def test_stdout_driven_recovery_preserves_boundaries( monkeypatch.chdir(project) prefix = ["--registry", str(registry), "--runtime-root", str(runtime)] - def run(argv, output="json"): + def run(argv, output="json", expected_rc=0): rc = cli.main(["--format", output, *argv]) stdout = capsys.readouterr().out - assert rc == 0, stdout + assert rc == expected_rc, stdout return json.loads(stdout) if output == "json" else stdout if baseline: @@ -193,10 +193,22 @@ def send_notification(**kwargs): assert graph_calls == [] and channel_calls == [] assert _snapshot(shared) == shared_before if not isolated: - # Independent control: explicitly widen this call after following the hint. + # Omission must not open either sink or mutate the original writeback. recovery = [arg for arg in recovery if arg not in { "--no-global-sync", "--suppress-external-sinks", }] + before = index.read_bytes() + denied = run(recovery, expected_rc=1) + assert denied["error_code"] == "external_delivery_resume_required" + assert index.read_bytes() == before + assert _snapshot(shared) == shared_before + assert graph_calls == [] and channel_calls == [] + assert state_path.read_bytes() == original_state + key = denied["external_delivery"]["resume_key"] + assert f"--resume-external-sinks {key}" in denied["error"] + wrong = run([*recovery, "--resume-external-sinks", "0" * 64], expected_rc=1) + assert wrong["error_code"] == "external_delivery_resume_mismatch" + recovery += ["--resume-external-sinks", key] repaired = run(recovery) assert repaired["appended"] is True assert repaired["vision_checkpoint"]["satisfied"] is True @@ -220,6 +232,16 @@ def send_notification(**kwargs): else: assert _snapshot(shared) != shared_before assert graph_calls and channel_calls, (graph_calls, channel_calls) + # Suppression during a pure replay must persist without appending a run. + paused = run([*recovery[:-2], "--suppress-external-sinks"]) + assert paused["idempotent_replay"] is True + assert index.read_bytes() == after + stale = run(recovery, expected_rc=1) + assert stale["error_code"] == "external_delivery_resume_mismatch" + assert stale["external_delivery"]["resume_key"] != key + resumed = run([*recovery[:-1], stale["external_delivery"]["resume_key"]]) + assert resumed["external_sink_delivery_authorized"] is True + assert index.read_bytes() == after @pytest.mark.parametrize("hint_source", ["first", "replay"]) @@ -230,6 +252,7 @@ def test_hint_preserves_lane_and_repeated_options_without_adding_isolation(first "--registry", "registry with spaces.json", "refresh-state", "--goal-id", GOAL_ID, "--agent-id", AGENT_ID, "--progress-scope", "agent_lane", "--agent-lane", "validation", "--available-capability", "filesystem_read", "--available-capability", "shell", + "--resume-external-sinks", "a" * 64, ] vision = ["--vision-last-patch", "Validation evidence checked."] assert _recovery_argv(render_state_refresh_markdown(first_refresh), original, vision) == original + vision diff --git a/tests/control_plane/test_refresh_external_delivery.py b/tests/control_plane/test_refresh_external_delivery.py new file mode 100644 index 0000000000..3639607991 --- /dev/null +++ b/tests/control_plane/test_refresh_external_delivery.py @@ -0,0 +1,101 @@ +"""Exercise legacy, local and failed-persistence paths through the real backend.""" +import json + +import pytest + +from loopx import cli +from loopx.control_plane.quota import refresh_external_delivery as bridge +from loopx.state_refresh import refresh_state_run +from tests.control_plane.test_quota_settlement_cli import ( + AGENT_ID, GOAL_ID, TODO_ID, TURN_ID, _write_fixture, +) + + +@pytest.fixture +def session(tmp_path, monkeypatch, capsys): + project, runtime, registry = _write_fixture(tmp_path) + monkeypatch.setenv("LOOPX_RUNTIME_ROOT", str(tmp_path / "shared")) + monkeypatch.chdir(project) + prefix = ["--registry", str(registry), "--runtime-root", str(runtime)] + binding = ["--goal-id", GOAL_ID, "--agent-id", AGENT_ID, + "--todo-id", TODO_ID, "--turn-instance-id", TURN_ID] + + def run(args, expected=0): + rc = cli.main(["--format", "json", *prefix, *args]) + stdout = capsys.readouterr().out + assert rc == expected, stdout + return json.loads(stdout) + + assert run(["quota", "should-run", "--codex-app", *binding, + "--scan-path", str(project)])["decision"] == "run" + args = ["refresh-state", *binding, "--classification", "validated_change", + "--delivery-outcome", "outcome_progress", "--no-global-sync"] + journal = runtime / "goals" / GOAL_ID / "rollout-event-log.jsonl" + index = runtime / "goals" / GOAL_ID / "runs/index.jsonl" + return project, runtime, registry, args, run, journal, index + + +def test_unstated_direct_caller_and_legacy_receipt_repair_remain_usable(session): + project, runtime, registry, args, run, journal, _ = session + local = refresh_state_run( + registry_path=registry, runtime_root_override=str(runtime), goal_id=GOAL_ID, + project=project, state_file=None, classification="validated_change", + recommended_action=None, delivery_outcome="outcome_progress", todo_id=TODO_ID, + turn_instance_id=TURN_ID, agent_id=AGENT_ID, dry_run=False, sync_global=False, + ) + assert local["ok"] and local["appended"] + assert local["external_sink_delivery_authorized"] is False + assert "refresh_external_delivery" not in journal.read_text(encoding="utf-8") + repaired = run(args) # Real old-style writeback lacking both pause and CLI receipt. + assert repaired["receipt_repaired"] is True + assert repaired["external_sink_delivery_authorized"] is False + replay = run(args) + assert replay["idempotent_replay"] is True + assert replay["external_sink_delivery_authorized"] is True + assert replay["external_delivery"]["reason"] == "legacy_per_call" + + +def test_dry_run_cannot_pause_or_resume_a_committed_operation(session): + _, _, _, args, run, journal, index = session + run(args) + before = journal.read_bytes(), index.read_bytes() + run([*args, "--suppress-external-sinks", "--dry-run"]) + assert (journal.read_bytes(), index.read_bytes()) == before + assert run(args)["external_sink_delivery_authorized"] is True + run([*args, "--suppress-external-sinks"]) + denied = run(args, expected=1) + key = denied["external_delivery"]["resume_key"] + before = journal.read_bytes(), index.read_bytes() + run([*args, "--resume-external-sinks", key, "--dry-run"]) + assert (journal.read_bytes(), index.read_bytes()) == before + assert run(args, expected=1)["error_code"] == "external_delivery_resume_required" + + +def test_pause_persistence_failure_stops_before_business_writeback(session, monkeypatch): + _, _, _, args, run, journal, index = session + before = journal.read_bytes() + + def fail(*_args, **_kwargs): + raise OSError("fixture journal unavailable") + + with monkeypatch.context() as patch: + patch.setattr(bridge, "append_rollout_event", fail) + failed = run([*args, "--suppress-external-sinks"], expected=1) + assert "fixture journal unavailable" in failed["error"] + assert journal.read_bytes() == before + assert not index.exists() + assert run([*args, "--suppress-external-sinks"])["ok"] + assert run(args, expected=1)["error_code"] == "external_delivery_resume_required" + + +def test_receipt_only_repair_preserves_pause_without_requiring_confirmation(session): + _, _, _, args, run, journal, index = session + run([*args, "--suppress-external-sinks"]) + rows = [json.loads(line) for line in journal.read_text(encoding="utf-8").splitlines()] + # Synthetic crash fixture: writeback exists, its CLI receipt is absent. + journal.write_text("".join(json.dumps(row) + "\n" for row in rows + if row["event_kind"] != "refresh_state"), encoding="utf-8") + before = index.read_bytes() + assert run(args)["receipt_repaired"] is True + assert index.read_bytes() == before + assert run(args, expected=1)["error_code"] == "external_delivery_resume_required" diff --git a/tests/control_plane_ts/refresh_external_delivery.test.ts b/tests/control_plane_ts/refresh_external_delivery.test.ts new file mode 100644 index 0000000000..c0bd534843 --- /dev/null +++ b/tests/control_plane_ts/refresh_external_delivery.test.ts @@ -0,0 +1,92 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import { settlementIdentity, type JsonObject } from "../../loopx/control_plane/effect_program.ts"; +import { decodeExternalDelivery, EXTERNAL_DELIVERY_EVENT, refreshExternalDelivery } from "../../loopx/control_plane/quota/refresh_external_delivery.ts"; + +const identity = settlementIdentity({ goal_id: "resume-fixture", agent_id: "builder", + turn_instance_id: "turn-a", todo_id: "todo_resume", replan_obligation_id: null }); +const normal = { suppress: false, resume_key: null }; +const suppress = { suppress: true, resume_key: null }; +const plan = (events: JsonObject[], request = normal, stage = true) => + refreshExternalDelivery(request, identity, events, stage); +function event(result: JsonObject, id: string): JsonObject { + assert.ok(result.transition); + return { event_id: id, event_kind: EXTERNAL_DELIVERY_EVENT, + goal_id: identity.goal_id, agent_id: identity.agent_id, run_id: identity.turn_instance_id, + todo_id: identity.todo_id, details: result.transition }; +} + +test("legacy requests and non-delivery callers keep their behavior", () => { + assert.equal(plan([]).authorized, true); + assert.equal(plan([]).reason, "legacy_per_call"); + const local = refreshExternalDelivery(null, identity, [], true); + assert.equal(local.authorized, false); + assert.equal(local.transition, null); + assert.equal(decodeExternalDelivery(undefined), null); + for (const bad of [true, {}, { suppress: "false", resume_key: null }, + { suppress: true, resume_key: "a".repeat(64) }, { ...normal, resume_key: "bad" }]) { + assert.throws(() => decodeExternalDelivery(bad)); + } +}); + +test("pause, omission, explicit resume, and repeated acknowledgement", () => { + const first = plan([], suppress); + const events = [event(first, "pause-a")]; + assert.equal(first.authorized, false); + assert.equal(plan(events).error_code, "external_delivery_resume_required"); + assert.equal(plan(events, suppress).transition, null); + assert.equal(plan(events, suppress).error_code, null); + const resume = { ...normal, resume_key: String(first.resume_key) }; + const approved = plan(events, resume); + assert.equal(approved.authorized, true); + events.push(event(approved, "resume-a")); + assert.equal(plan(events).authorized, true); + assert.equal(plan(events, resume).transition, null); + assert.equal(plan(events, resume).authorized, true); + const pausedAgain = plan(events, suppress); + assert.notEqual(pausedAgain.resume_key, first.resume_key); + events.push(event(pausedAgain, "pause-b")); + assert.equal(plan(events, resume).error_code, "external_delivery_resume_mismatch"); + assert.equal(plan(events).error_code, "external_delivery_resume_required"); +}); + +test("acknowledgements cannot cross operations or silently initialize a resume", () => { + const first = plan([], suppress); + const foreign = settlementIdentity({ ...identity, turn_instance_id: "turn-b" }); + const other = refreshExternalDelivery(suppress, foreign, [], true); + assert.notEqual(other.resume_key, first.resume_key); + assert.equal(plan([], { ...normal, resume_key: String(first.resume_key) }).error_code, + "external_delivery_resume_mismatch"); + const foreignEvent = { ...event(other, "other-pause"), run_id: foreign.turn_instance_id }; + assert.equal(plan([foreignEvent]).authorized, true); + const events = [event(first, "pause-a")]; + assert.equal(plan(events, { ...normal, resume_key: String(other.resume_key) }).authorized, false); +}); + +test("receipt-only repair and unspecified callers do not consume acknowledgement", () => { + const first = plan([], suppress); + const events = [event(first, "pause-a")]; + const receipt = plan(events, normal, false); + assert.equal(receipt.authorized, false); + assert.equal(receipt.error_code, null); + assert.equal(receipt.transition, null); + assert.equal(refreshExternalDelivery(null, identity, events, true).error_code, null); + assert.equal(plan(events, { ...normal, resume_key: String(first.resume_key) }, false).transition, null); + assert.ok(plan([], suppress, false).transition); // Explicit suppression is still remembered. +}); + +test("older writers cannot erase a pause; corrupt relevant history never permits sending", () => { + const first = plan([], suppress); + const pause = event(first, "pause-a"); + const unrelated = { ...pause, event_kind: "refresh_state", details: {} }; + assert.equal(plan([pause, unrelated]).error_code, "external_delivery_resume_required"); + for (const details of [null, { ...(pause.details as JsonObject), schema_version: "future" }, + { ...(pause.details as JsonObject), state: "ready" }, + { ...(pause.details as JsonObject), settlement_effect_id: "wrong" }, + { ...(pause.details as JsonObject), resume_key: "0".repeat(64) }]) { + const corrupted = [{ ...pause, details }]; + assert.equal(plan(corrupted).error_code, "external_delivery_history_invalid"); + assert.equal(plan(corrupted, suppress).authorized, false); + assert.equal(plan(corrupted, suppress).error_code, null); + } +}); From af1d27094d69cf2b9a0a47497b8683c3f1455dd8 Mon Sep 17 00:00:00 2001 From: Tartar Date: Thu, 10 Sep 2026 17:24:35 +0800 Subject: [PATCH 2/2] docs(refresh): explain external delivery pause and resume compatibility Signed-off-by: Tartar --- .../rfcs/goal-channel-collaboration-v0.md | 8 +++- .../goal-channel-collaboration-v0.zh-CN.md | 5 ++- docs/state-interaction-model.md | 37 ++++++++++++++++++- 3 files changed, 45 insertions(+), 5 deletions(-) diff --git a/docs/architecture/rfcs/goal-channel-collaboration-v0.md b/docs/architecture/rfcs/goal-channel-collaboration-v0.md index 3b669bdc26..56ba825dd5 100644 --- a/docs/architecture/rfcs/goal-channel-collaboration-v0.md +++ b/docs/architecture/rfcs/goal-channel-collaboration-v0.md @@ -289,8 +289,12 @@ reuses the same bot verification, semantic idempotency, cooldown, provider idempotency key, and message readback as `notify-gate`. Use `loopx refresh-state ... --suppress-external-sinks` to suppress delivery for -one refresh without disabling the binding. Disable automatic delivery -persistently with: +one refresh without disabling the binding. Turn-bound recovery remembers that +pause and requires `--resume-external-sinks ` to resume delivery using +the returned key. This acknowledges the current operation, not new permissions; +historical operations without pause evidence retain their previous behavior. +See the [recovery handshake](../../state-interaction-model.md). Disable automatic +delivery persistently with: ```bash loopx goal-channel configure --goal-id --no-auto-notify-human-gates --execute diff --git a/docs/architecture/rfcs/goal-channel-collaboration-v0.zh-CN.md b/docs/architecture/rfcs/goal-channel-collaboration-v0.zh-CN.md index 1247a85086..a13c1c575d 100644 --- a/docs/architecture/rfcs/goal-channel-collaboration-v0.zh-CN.md +++ b/docs/architecture/rfcs/goal-channel-collaboration-v0.zh-CN.md @@ -262,7 +262,10 @@ loopx goal-channel configure --goal-id --auto-notify-human-gates --exe 已有的 bot 身份校验、语义幂等、冷却、provider idempotency key 和消息回读。 单次 refresh 可使用 `loopx refresh-state ... --suppress-external-sinks` -临时抑制投递,而无需禁用 binding。持久关闭自动投递: +临时抑制投递,而无需禁用 binding。对带 Turn 绑定的恢复,工具会记住这次暂停; +再次允许外发需使用返回的 `resume_key`,提交 `--resume-external-sinks `。 +这只确认恢复当前操作的投递,不授予新权限;无暂停记录的旧操作保留原行为。 +详见 [恢复握手](../../state-interaction-model.md)。持久关闭自动投递: ```bash loopx goal-channel configure --goal-id --no-auto-notify-human-gates --execute diff --git a/docs/state-interaction-model.md b/docs/state-interaction-model.md index 50e343cffc..edcb08fd64 100644 --- a/docs/state-interaction-model.md +++ b/docs/state-interaction-model.md @@ -770,8 +770,41 @@ isolation (`--no-global-sync`, `--suppress-external-sinks`) options, with their original values and presence. Both first-writeback and replay Markdown list the preserve/remove/add rules; do not add isolation flags absent from the original command or silently change the delivery target. These instructions -preserve the original boundary when followed; they do not persist an authority -ceiling that rejects callers who manually remove flags. +preserve target, scope, and global-sync choices. Turn-bound external delivery +also has a tool-enforced pause/resume handshake: + +- `--suppress-external-sinks` records a pause for the exact settlement operation. + A later send-capable refresh that omits this flag returns + `external_delivery_resume_required` before new business writes or sends. +- To keep recovering locally, retain `--suppress-external-sinks`. To deliberately + resume external delivery, retry the same recovery command with + `--resume-external-sinks ` using the key returned in + `external_delivery.resume_key`. The two flags are mutually exclusive. Preserve + the original identity and omit already executed mutation and spend arguments. +- A matching acknowledgement lifts only that operation's current pause. A new + pause invalidates the old key. Repeating an accepted acknowledgement is safe; + it does not append another business run. Existing provider permissions, + configuration, and delivery de-duplication still apply. This is explicit intent, + not a new authorization grant or an exactly-once transport guarantee. +- Receipt-only repair remains local and does not require or consume a resume + acknowledgement; the next send-capable replay still checks the pause. Direct + Python callers that request no delivery remain local without creating a pause. +- Historical operations without pause evidence keep the previous per-call + behavior. Once a new explicit suppression is recorded, later recovery must + acknowledge it. Non-Turn refresh keeps the existing single-call suppression. + No permanent permission ceiling or new recovery mode is introduced. + +Pause/resume observations are scoped by settlement identity in the existing +local rollout journal and processed by the TypeScript quota owner. A pause is +persisted before a new business writeback; if later validation fails, the pause +remains and can be explicitly resumed. Resume is persisted only after an +accepted recovery, before CLI delivery. Journal failures prevent delivery; +dry-run never persists either transition. Corrupt or unsupported relevant pause +history cannot enable sends, but explicit suppression remains available for +local recovery. Existing original-writeback, one-spend, and sink retry contracts +remain unchanged. Global-sync flags, target migration, general hooks, and +cross-service atomicity are outside this handshake. Older binaries do not enforce +the handshake: rolling back loses this protection and must not erase its journal. Remove previously executed state-mutation options, even when their values are unchanged: `--next-action`, `--autonomous-replan-recorded`, `--repair-delta-kind`, `--usage-json`, and