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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 6 additions & 2 deletions docs/architecture/rfcs/goal-channel-collaboration-v0.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 <resume_key>` 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 <goal-id> --no-auto-notify-human-gates --execute
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -262,7 +262,10 @@ loopx goal-channel configure --goal-id <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 <resume_key>`。
这只确认恢复当前操作的投递,不授予新权限;无暂停记录的旧操作保留原行为。
详见 [恢复握手](../../state-interaction-model.md)。持久关闭自动投递:

```bash
loopx goal-channel configure --goal-id <goal-id> --no-auto-notify-human-gates --execute
Expand Down
37 changes: 35 additions & 2 deletions docs/state-interaction-model.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 <resume_key>` 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
Expand Down
27 changes: 17 additions & 10 deletions loopx/cli_commands/project_lifecycle.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.",
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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")
Expand Down Expand Up @@ -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,
Expand All @@ -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(
Expand Down
82 changes: 82 additions & 0 deletions loopx/control_plane/quota/refresh_external_delivery.py
Original file line number Diff line number Diff line change
@@ -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)
95 changes: 95 additions & 0 deletions loopx/control_plane/quota/refresh_external_delivery.ts
Original file line number Diff line number Diff line change
@@ -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");
}
4 changes: 4 additions & 0 deletions loopx/control_plane/quota/refresh_recovery.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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,
Expand Down
5 changes: 4 additions & 1 deletion loopx/control_plane/quota/settlement.py
Original file line number Diff line number Diff line change
Expand Up @@ -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`, "
Expand Down Expand Up @@ -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__ = [
Expand Down Expand Up @@ -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")),
Expand Down
Loading