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
29 changes: 29 additions & 0 deletions docs/reference/protocols/loopx-turn-v0.md
Original file line number Diff line number Diff line change
Expand Up @@ -619,6 +619,35 @@ Every attempted tick returns one result kind:
| `validation_failed` | Host output exists but task validation failed or is inconclusive. | Preserve failure evidence and route to repair/replan. |
| `writeback_failed` | Validated work could not be durably recorded. | Do not spend; retry idempotent writeback before more delivery. |

### In-flight Turn settlement / 在途 Turn 结算

A Todo can stay open across several bounded Turns. An exact accountable
`outcome_progress` writeback with an accepted `vision_checkpoint_v0`
`in_flight_continuation` boundary discharges that Turn's progress obligation,
not the Todo's terminal acceptance. With its matching durable writeback and
quota-spend receipts, the original Turn replays as `heartbeat_settled_skip`:
no more work and no second debit. Without the spend receipt it remains
`settlement_pending`; a missing writeback receipt, unaccepted checkpoint or
wrong Goal/Agent/Todo/Turn cannot prove settlement. A plain progress claim or
`semantic_closeout` checkpoint is not this exception.

Todo 可以跨多个有界 Turn 保持开放。与原始身份精确绑定的 `outcome_progress`
写回,只有携带已获准的 `vision_checkpoint_v0`、`in_flight_continuation` 边界和
当前 Todo 的 trigger,才履行该 Turn 的进展义务,而非 Todo 的最终验收。有匹配
的写回和扣额回执后,同一 Turn 返回 `heartbeat_settled_skip`,不得再次执行或
重复扣额;缺少扣额回执时仍是 `settlement_pending`。缺少写回回执、未获准
checkpoint 或错配 Goal/Agent/Todo/Turn 均不能证明结算,普通进展声明或
`semantic_closeout` 也不能替代这项凭证。

Waiting conditions and frontier/successor changes do not reopen a settled Turn.
A fresh Turn must recompute admission to continue the open Todo or select an
independent successor. The Todo's completion validator, definition revision,
leases and Goal acceptance remain authoritative and unchanged.

等待条件和 frontier/后继变化不能重新打开已结算 Turn。继续开放 Todo 或选择
独立后继必须用新 Turn 重新准入。Todo 完成验证器、定义版本、租约和 Goal 验收
仍由原权威负责,不因在途结算而放宽或改写。

`validated_completion` is admitted only when the Turn caller supplies an
explicit Todo lifecycle adapter. After independent validation, the adapter must
authorize and complete the selected Todo through the existing Todo lifecycle,
Expand Down
26 changes: 26 additions & 0 deletions loopx/control_plane/quota/settlement_phase.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,32 @@ export function isBoundedBlockedRetry(value: unknown, todoId: string | null): bo
return Number.isFinite(delay) && delay >= 60 && delay <= 30 * 60;
}

/** The committed checkpoint accepts progress for a Turn, not Todo completion.
* The caller must first verify this writeback's exact durable receipt. */
export function isAcceptedInFlightWriteback(
value: unknown,
identity: SettlementIdentity,
): boolean {
const run = jsonObject(value);
const checkpoint = jsonObject(run?.vision_checkpoint);
if (identity.binding_kind !== "todo" || identity.todo_id === null ||
!run || run.goal_id !== identity.goal_id ||
run.agent_id !== identity.agent_id || run.todo_id !== identity.todo_id ||
run.turn_instance_id !== identity.turn_instance_id ||
run.delivery_outcome !== "outcome_progress" ||
!checkpoint || checkpoint.schema_version !== "vision_checkpoint_v0" ||
checkpoint.agent_id !== identity.agent_id ||
checkpoint.delivery_boundary !== "in_flight_continuation" ||
checkpoint.satisfied !== true || !Array.isArray(checkpoint.triggers)) {
return false;
}
return checkpoint.triggers.some((value) => {
const trigger = jsonObject(value);
return trigger?.kind === "in_flight_continuation" &&
trigger.todo_id === identity.todo_id;
});
}

/** Both same-Turn readback and prior-Turn recovery accept the shipped effect identities. */
export function isCommittedMonitorPollEffect(
effectId: unknown,
Expand Down
5 changes: 4 additions & 1 deletion loopx/control_plane/quota/settlement_readback.ts
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@ import {
import {
isBoundedBlockedRetry,
isCommittedMonitorPollEffect,
isAcceptedInFlightWriteback,
receiptBoundMonitorPhase,
receiptBoundReplayPhase,
} from "./settlement_phase.ts";
Expand Down Expand Up @@ -1006,6 +1007,8 @@ function readQuotaSettlementFromRequest(
const todoBoundReplan = identity.binding_kind === "todo" &&
semanticReplanGuard.scope === "turn_guard" &&
semanticReplanGuard.selected_obligation_id !== null;
const inFlightWriteback = writeback.failure === null &&
isAcceptedInFlightWriteback(writebackRun, identity);

const recovery = request.refresh_retry === null ? null : refreshRecovery(
request.refresh_retry, writebackRun, writeback.failure === null,
Expand Down Expand Up @@ -1051,7 +1054,7 @@ function readQuotaSettlementFromRequest(
}),
replay_phase: receiptBoundReplayPhase({
binding_kind: identity.binding_kind,
writeback_completes_binding: todoBoundReplan || blockedNoSpend,
writeback_completes_binding: todoBoundReplan || blockedNoSpend || inFlightWriteback,
completion_receipt_present: completionEvent !== null,
durable_writeback_present: writeback.failure === null,
quota_spend_present: spend.failure === null,
Expand Down
117 changes: 117 additions & 0 deletions tests/control_plane/test_in_flight_turn_replay_cli.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,117 @@
"""Real CLI characterization: Turn settlement is not terminal Todo acceptance."""
from __future__ import annotations

from pathlib import Path
from typing import Any

import pytest

from canonical_authority_fixture import (
initialize_canonical_authority,
isolate_sqlite_runtime,
)
from loopx.control_plane.coordination.runtime_shadow import (
build_todo_runtime_shadow_projection,
)
from loopx.control_plane.quota.turn_envelope import build_turn_envelope
from test_quota_settlement_cli import (
AGENT_ID,
ALTERNATIVE_TODO_ID,
GOAL_ID,
TODO_ID,
_configure_completion_validation_todo,
_configure_selectable_alternative,
_run_cli,
_spend_run_count,
_write_fixture,
)


@pytest.mark.parametrize("provider", ["legacy", "file", "sqlite"])
@pytest.mark.parametrize("waiting", [False, True])
def test_in_flight_turn_replay_preserves_open_todo_and_fresh_turn_admission(
tmp_path: Path,
monkeypatch: pytest.MonkeyPatch,
provider: str,
waiting: bool,
) -> None:
if provider == "sqlite":
# The shared SQLite fixture predates typed test helpers.
isolate_sqlite_runtime(tmp_path, monkeypatch) # type: ignore[no-untyped-call]
project, runtime, registry = _write_fixture(tmp_path)
state = _configure_completion_validation_todo(project)
_configure_selectable_alternative(project)

def cli(*args: str) -> dict[str, Any]:
rc, result = _run_cli(registry, runtime, *args, cwd=project)
assert rc == 0, result
return result

def todo() -> dict[str, Any]:
listed = cli("todo", "list", "--goal-id", GOAL_ID)
return next(item for item in listed["todos"] if item["todo_id"] == TODO_ID)

if provider != "legacy":
listed = cli("todo", "list", "--goal-id", GOAL_ID)
projection = build_todo_runtime_shadow_projection(
goal_id=GOAL_ID, todos=listed["todos"], handoff_mode="soft_claim",
)
initialize_canonical_authority(
runtime, GOAL_ID, projection, state_path=state, provider=provider,
)
before = todo()
turn_id = f"inflight-{provider}-{waiting}"
binding = ("--agent-id", AGENT_ID, "--todo-id", TODO_ID, "--turn-instance-id", turn_id)
guard = ("quota", "should-run", "--codex-app", "--goal-id", GOAL_ID, "--scan-path", str(project))
first = cli(*guard, *binding)
assert first["should_run"] is True
refresh = cli(
"refresh-state", "--goal-id", GOAL_ID, *binding,
"--classification", "validated_intermediate_progress",
"--delivery-batch-scale", "implementation",
"--delivery-outcome", "outcome_progress",
"--delivery-boundary", "in_flight_continuation",
"--no-global-sync", "--suppress-external-sinks",
)
assert refresh["vision_checkpoint"]["satisfied"] is True
spend_args = ("quota", "spend-slot", "--goal-id", GOAL_ID, *binding,
"--slots", "1", "--source", "heartbeat", "--execute")
spend = cli(*spend_args)
assert spend["appended"] is True
assert spend["settlement_progress"]["state"] == "settled"
if waiting:
cli(
"todo", "update", "--goal-id", GOAL_ID, "--agent-id", AGENT_ID,
"--todo-id", TODO_ID, "--status", "open",
"--resume-when", "resume_at:2099-01-01T00:00:00Z",
"--successor-todo-id", ALTERNATIVE_TODO_ID,
"--reason", "Wait for independent evidence; continue the runnable successor.",
)

# Argument-less and explicit reentry must both respect the original receipt,
# even if the live frontier now recommends a different Todo.
for selection in ((), ("--todo-id", TODO_ID)):
replay = cli(*guard, "--agent-id", AGENT_ID, "--turn-instance-id", turn_id, *selection)
assert replay["should_run"] is False, replay
assert replay["effective_action"] == "heartbeat_settled_skip"
assert replay["execution_obligation"]["kind"] == "heartbeat_settled_skip"
assert replay["execution_obligation"]["must_attempt_work"] is False
assert replay["heartbeat_receipt"]["settlement_identity"]["todo_id"] == TODO_ID
assert replay["interaction_contract"]["agent_channel"]["must_attempt"] is False
assert replay["interaction_contract"]["cli_channel"]["spend_after_validation"] is False
envelope = build_turn_envelope(replay)
assert envelope["contract_capsule"]["execution_obligation"]["must_attempt_work"] is False
assert envelope["writeback"]["spend_after_validation"] is False
spend_replay = cli(*spend_args)
assert spend_replay["appended"] is False
assert _spend_run_count(runtime) == 1
after = todo()
assert after["status"] == "open"
assert after["completion_validation_required"] is True
assert after["completion_validation_sha256"] == before["completion_validation_sha256"]

fresh = cli(*guard, "--agent-id", AGENT_ID, "--turn-instance-id", f"{turn_id}-next")
assert fresh["should_run"] is True, fresh
assert fresh["effective_action"] != "unsettled_host_turn_recovery"
assert fresh["selected_todo"]["todo_id"] == (ALTERNATIVE_TODO_ID if waiting else TODO_ID)
assert _spend_run_count(runtime) == 1
14 changes: 14 additions & 0 deletions tests/control_plane/test_quota_settlement_cli.py
Original file line number Diff line number Diff line change
Expand Up @@ -1570,6 +1570,20 @@ def test_in_flight_progress_settles_while_completion_validation_todo_is_open(
)
assert current_todo["status"] == "open", current_todo

replay_rc, replay = _run_cli(
registry_path, runtime, "quota", "should-run", "--codex-app",
"--goal-id", GOAL_ID, "--agent-id", AGENT_ID,
"--turn-instance-id", turn_id, "--scan-path", str(project),
)
assert replay_rc == 0, replay
assert replay["should_run"] is False, replay
assert replay["effective_action"] == "heartbeat_settled_skip"
assert replay["execution_obligation"]["kind"] == "heartbeat_settled_skip"
assert replay["execution_obligation"]["must_attempt_work"] is False
assert replay["interaction_contract"]["agent_channel"]["must_attempt"] is False
assert replay["interaction_contract"]["cli_channel"]["spend_after_validation"] is False
assert _spend_run_count(runtime) == 1


def test_open_completion_todo_accepts_only_matching_in_flight_writeback(
tmp_path: Path,
Expand Down
65 changes: 65 additions & 0 deletions tests/control_plane_ts/quota_settlement_readback.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -79,6 +79,7 @@ async function fixture(options: {
writebackOutcome?: string;
progressObservation?: Record<string, unknown>;
blockedRetry?: boolean;
visionCheckpoint?: Record<string, unknown>;
} = {}) {
const runtimeRoot = await mkdtemp(join(tmpdir(), "loopx-settlement-readback-"));
const goalRoot = join(runtimeRoot, "goals", goalId);
Expand Down Expand Up @@ -151,6 +152,7 @@ async function fixture(options: {
todo_id: todoId,
turn_instance_id: turnId,
settlement_identity: identity,
...(options.visionCheckpoint ? {vision_checkpoint: options.visionCheckpoint} : {}),
...(options.blockedRetry ? {blocked_retry: {
schema_version: "quota_blocked_retry_v0",
source: "todo",
Expand Down Expand Up @@ -282,6 +284,69 @@ test("settlement progress preserves typed source and rejects malformed sources",
} finally { await rm(root, {recursive: true, force: true}); }
});

test("accepted in-flight writeback closes only the exact Turn, not its Todo", async t => {
const checkpoint = {
schema_version: "vision_checkpoint_v0", agent_id: agentId, satisfied: true,
delivery_boundary: "in_flight_continuation",
triggers: [{kind: "in_flight_continuation", todo_id: todoId}],
};
const cases = [
{name: "both effects", spend: true, expected: "settled"},
{name: "spend still required", spend: false, expected: "settlement_pending"},
{name: "writeback receipt missing", spend: true, remove: "refresh_state", expected: "open"},
{name: "spend receipt missing", spend: true, remove: "quota_spend", expected: "settlement_pending"},
];
for (const entry of cases) await t.test(entry.name, async () => {
const root = await fixture({writeback: true, spend: entry.spend, visionCheckpoint: checkpoint});
try {
if (entry.remove) {
const path = join(root, "goals", goalId, "rollout-event-log.jsonl");
const events = (await readFile(path, "utf8")).trim().split("\n").map(line => JSON.parse(line));
await writeFile(path, events.filter(event => event.event_kind !== entry.remove).map(event => JSON.stringify(event)).join("\n") + "\n");
}
const result = await readQuotaSettlement(request(root));
assert.equal(result.replay_phase, entry.expected);
assert.equal(result.completion_event, null);
assert.equal((result.terminal_closeout as any).payload.ok, false);
} finally { await rm(root, {recursive: true, force: true}); }
});
});

test("in-flight replay rejects unaccepted checkpoints and unrelated identities", async t => {
const checkpoint = {
schema_version: "vision_checkpoint_v0", agent_id: agentId, satisfied: true,
delivery_boundary: "in_flight_continuation",
triggers: [{kind: "in_flight_continuation", todo_id: todoId}],
};
const patches: [string, Record<string, unknown>][] = [
["checkpoint absent", {vision_checkpoint: null}],
["checkpoint malformed", {vision_checkpoint: []}],
["not accepted", {vision_checkpoint: {...checkpoint, satisfied: false}}],
["truthy is not acceptance", {vision_checkpoint: {...checkpoint, satisfied: "true"}}],
["wrong schema", {vision_checkpoint: {...checkpoint, schema_version: "other"}}],
["semantic closeout", {vision_checkpoint: {...checkpoint, delivery_boundary: "semantic_closeout"}}],
["checkpoint other agent", {vision_checkpoint: {...checkpoint, agent_id: "peer"}}],
["no trigger", {vision_checkpoint: {...checkpoint, triggers: []}}],
["trigger other Todo", {vision_checkpoint: {...checkpoint, triggers: [{kind: "in_flight_continuation", todo_id: "todo_other"}]}}],
["trigger wrong kind", {vision_checkpoint: {...checkpoint, triggers: [{kind: "vision_unchanged", todo_id: todoId}]}}],
...["surface_only", "outcome_gap", "primary_goal_outcome"].map(outcome => [outcome, {delivery_outcome: outcome}] as [string, Record<string, unknown>]),
...["goal_id", "agent_id", "todo_id", "turn_instance_id"].map(field => [field, {[field]: "other"}] as [string, Record<string, unknown>]),
["effect mismatch", {settlement_identity: {...identity, effect_id: "other"}}],
];
for (const [name, patch] of patches) await t.test(name, async () => {
const root = await fixture({writeback: true, spend: true, visionCheckpoint: checkpoint});
try {
const path = join(root, "goals", goalId, "runs", "index.jsonl");
const runs = (await readFile(path, "utf8")).trim().split("\n").map(line => JSON.parse(line));
runs[0] = {...runs[0], ...patch};
await writeFile(path, runs.map(run => JSON.stringify(run)).join("\n") + "\n");
const result = await readQuotaSettlement(request(root));
assert.equal(result.replay_phase, "open");
assert.equal(result.completion_event, null);
} finally { await rm(root, {recursive: true, force: true}); }
});
});

test("monitor closeout requires the exact committed effect, not a matching observation row", async t => {
const effect = `quota-monitor-poll:${goalId}:${agentId}:${turnId}`;
const cases: [string, Record<string, unknown>, string][] = [
Expand Down
Loading