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
Original file line number Diff line number Diff line change
Expand Up @@ -37,8 +37,13 @@ advancement work remains active.
- An executed turn-scoped poll returns `turn_continuation`. An exact match
between `settlement_todo_id` and the observed `todo_id` closes the no-spend
monitor Turn and requires a fresh `--turn-instance-id`. A different admitted
monitor Todo is auxiliary: it records the observation but leaves the original
advancement Turn open for its durable writeback and single spend. Without an
monitor Todo is auxiliary: it records the observation independently of the
original advancement's closeout. If closeout is pending, the Turn remains open
for its durable writeback and single spend. If it is already settled, even a
**first** due observation is allowed, and continuation reports
`current_turn_settled=true`, `next_turn_required=true`; independent advancement
still needs a fresh Turn. Replays read current verified closeout rather than
reopening it from the receipt's historical continuation. Without an
exact or typed auxiliary binding the response fails closed from claiming the
Turn settled.

Expand All @@ -49,6 +54,13 @@ and replay idempotently in one settlement Turn without spending quota, while a
guard replay continues to select the original advancement Todo. Existing
wrong-Todo tests for receipt-bound monitor Turns must remain passing.

This corrects the previous default rejection of first auxiliary observations
after advancement settlement. Both call orders must retain one primary debit,
one receipt per observation identity and current due/actor/lease admission. It
does not add a configuration switch, relax Monitor authority or retry external
observations. Existing CLI/managed-Turn receipts provide the readback; no new
frontend or Lark command/settings owner is introduced.

### Canonical leased observations and recovery

For a promoted Goal, an existing Monitor execution can supply its current lease
Expand Down Expand Up @@ -207,8 +219,11 @@ using a complete read-only snapshot with disposable File/SQLite/PostgreSQL arms.
- 执行成功的 turn-scoped poll 会返回 `turn_continuation`。仅当
`settlement_todo_id` 与被观察的 `todo_id` 精确一致时,才完成该 monitor Turn 的
不计费结算,并要求后续使用新的 `--turn-instance-id`。不同但已准入的 monitor
Todo 属于辅助观察:只写观察回执,原 advancement Turn 仍需完成 durable
writeback 与唯一一次 spend。既非精确匹配、也无 typed auxiliary binding 时,响应
Todo 属于辅助观察:观察回执与原 advancement 的结算独立。若尚未结算,仍需完成
durable writeback 与唯一一次 spend;若已经结算,**首次**到期观察同样可以登记,
continuation 返回 `current_turn_settled=true`、`next_turn_required=true`,新的独立
推进仍须新 Turn。重放按当前已核验结算读回,不用历史 continuation 重开旧 Turn。
既非精确匹配、也无 typed auxiliary binding 时,响应
必须失败关闭,不能宣称 Turn 已结算。

### 验收
Expand All @@ -217,6 +232,11 @@ CLI 端到端测试必须证明:多个到期 monitor 能在同一结算 Turn
幂等重放、全程不消耗配额;随后重放 guard 仍选择原 advancement Todo。同时,
receipt-bound monitor Turn 的错误 Todo 替换测试必须继续通过。

这是对原默认行为的修正:不再拒绝 advancement 结算后的首次辅助观察。两种调用
顺序均须保留一次主任务扣额、每个观察身份一个回执,以及当前到期/actor/租约
准入。不新增配置开关、不放宽 Monitor 权限、不重新执行外部观察。CLI/managed
Turn 复用现有回执读回;不另建前端或 Lark 命令/配置权威。

### Canonical 带租约观察与恢复

已晋升 Goal 的 Monitor 若已有执行租约,可用上述命令读取租约,并传入当前 key 和
Expand Down
54 changes: 40 additions & 14 deletions loopx/control_plane/quota/monitor_poll_commit.ts
Original file line number Diff line number Diff line change
Expand Up @@ -541,6 +541,18 @@ interface Admission {
external: boolean;
}

async function readAuxiliarySettlement(request: MonitorRequest): Promise<JsonObject | null> {
if (!request.runtime_root || !request.turn_instance_id ||
!request.observation.actor_agent_id || !request.observation.settlement_todo_id) return null;
return await readQuotaSettlement({
schema_version: QUOTA_SETTLEMENT_READBACK_REQUEST_SCHEMA,
runtime_root: request.runtime_root, goal_id: request.goal_id,
agent_id: request.observation.actor_agent_id, todo_id: request.observation.settlement_todo_id,
turn_instance_id: request.turn_instance_id, replan_obligation_id: null,
infer_turn_instance_id: false, allow_unbound_binding: false,
});
}

async function auxiliaryMonitorAllowed(
request: MonitorRequest, historicalAdmission: boolean,
): Promise<boolean | null> {
Expand All @@ -566,17 +578,13 @@ async function auxiliaryMonitorAllowed(
todo.excluded_agents.some(value => typeof value !== "string") ||
todo.excluded_agents.includes(decision.agent_id)))) throw conflict();
if (!historicalAdmission) {
const settlement = await readQuotaSettlement({
schema_version: QUOTA_SETTLEMENT_READBACK_REQUEST_SCHEMA,
runtime_root: request.runtime_root, goal_id: request.goal_id,
agent_id: decision.agent_id, todo_id: settlementTodo,
turn_instance_id: request.turn_instance_id, replan_obligation_id: null,
infer_turn_instance_id: false, allow_unbound_binding: false,
});
const progress = jsonObject(settlement.progress);
if (settlement.found !== true || jsonObject(jsonObject(settlement.identity)?.result)?.failure !== null ||
const settlement = await readAuxiliarySettlement(request);
const progress = jsonObject(settlement?.progress);
// The advancement's single debit is not an observation-admission fence.
// A settled exact binding still permits independently admitted due Monitors.
if (settlement?.found !== true || jsonObject(jsonObject(settlement.identity)?.result)?.failure !== null ||
progress?.schema_version !== "quota_settlement_progress_v0" ||
progress.state === "identity_required" || progress.state === "settled") throw conflict();
progress.state === "identity_required") throw conflict();
}
const monitor = decision.registry_due_monitor;
// Retain ordinary quota/due-work admission and capability/gate projections;
Expand Down Expand Up @@ -1488,14 +1496,30 @@ function successorReceipts(providerReceipt: JsonObject | null): JsonObject[] {
: [];
}

function payloadFor(
async function projectAuxiliaryContinuation(request: MonitorRequest, payload: JsonObject): Promise<void> {
const continuation = jsonObject(payload.turn_continuation);
if (!request.execute || continuation?.settlement_binding_matches_observation !== false) return;
const settlement = await readAuxiliarySettlement(request);
if (jsonObject(settlement?.progress)?.state !== "settled") return;
// This is current readback, not a change to the historical Monitor receipt
// or permission for another advancement. Replays must not reopen a paid Turn.
payload.turn_continuation = {
...continuation,
current_turn_settled: true,
next_turn_required: true,
next_action: "rerun quota should-run with a fresh --turn-instance-id before independent work",
reason: "the original advancement Turn is already settled; the auxiliary observation adds no spend or delivery identity",
};
}

async function payloadFor(
request: MonitorRequest,
record: JsonObject,
jsonPath: string,
markdownPath: string,
indexPath: string,
options: { appended: boolean; replayed: boolean; repaired: boolean },
): JsonObject {
): Promise<JsonObject> {
const event = requiredObject(record.monitor_event, "record.monitor_event");
const receipts = successorReceipts(request.provider_receipt);
const payload: JsonObject = {
Expand Down Expand Up @@ -1577,6 +1601,7 @@ function payloadFor(
if (request.status_reload_warning) {
payload.status_reload_warning = request.status_reload_warning;
}
await projectAuxiliaryContinuation(request, payload);
return payload;
}

Expand Down Expand Up @@ -1882,6 +1907,7 @@ async function replayDurableReceipt(
? "quota monitor-poll commit repaired its durable transaction artifacts"
: "replayed existing monitor poll event for the same effect identity",
};
await projectAuxiliaryContinuation(request, replayPayload);
return result(
request,
fingerprint,
Expand Down Expand Up @@ -2045,7 +2071,7 @@ export async function evaluateQuotaMonitorPollCommit(
fingerprint,
"preview",
record,
payloadFor(
await payloadFor(
effectiveRequest,
record,
paths.jsonPath,
Expand Down Expand Up @@ -2295,7 +2321,7 @@ export async function evaluateQuotaMonitorPollCommit(
effectiveRequest.generated_at,
request.effect_id,
);
const payload = payloadFor(
const payload = await payloadFor(
effectiveRequest,
record,
jsonPath,
Expand Down
48 changes: 38 additions & 10 deletions tests/control_plane/test_monitor_observation_admission.py
Original file line number Diff line number Diff line change
Expand Up @@ -98,9 +98,7 @@ def test_receipt_bound_advancement_allows_one_auxiliary_due_monitor_receipt(
),
}
assert (
poll["after"]["interaction_contract"]["cli_channel"][
"spend_after_validation"
]
poll["after"]["interaction_contract"]["cli_channel"]["spend_after_validation"]
is True
)
assert poll_replay_rc == 0, poll_replay
Expand All @@ -122,11 +120,23 @@ def test_receipt_bound_advancement_allows_one_auxiliary_due_monitor_receipt(


@pytest.mark.parametrize(
("provider", "writeback"),
[("legacy", False), ("legacy", True), ("file", True), ("sqlite", True)],
("provider", "writeback", "spend_first"),
[
("legacy", False, False),
("legacy", True, False),
("file", True, False),
("sqlite", True, False),
("legacy", True, True),
("file", True, True),
("sqlite", True, True),
],
)
def test_completed_advancement_retains_auxiliary_monitor_admission(
tmp_path: Path, monkeypatch: pytest.MonkeyPatch, provider: str, writeback: bool
tmp_path: Path,
monkeypatch: pytest.MonkeyPatch,
provider: str,
writeback: bool,
spend_first: bool,
) -> None:
from canonical_authority_fixture import (
initialize_canonical_authority,
Expand Down Expand Up @@ -251,11 +261,21 @@ def test_completed_advancement_retains_auxiliary_monitor_admission(
"--suppress-external-sinks",
)
assert rc == 0, refresh
if spend_first:
spend_command = refresh["settlement_owed"]["command"]
rc, spent = _run_cli(
registry_path,
runtime,
*_projected_cli_args(spend_command, turn_instance_id=turn),
)
assert rc == 0, spent
assert _spend_run_count(runtime) == 1
rc, poll = _run_cli(registry_path, runtime, *poll_args)
assert rc == 0, (poll.get("error_code"), poll.get("reason"), poll.get("error"))
assert poll["todo_id"] == DUE_MONITOR_TODO_ID
assert poll["settlement_todo_id"] == TODO_ID
assert poll["turn_continuation"]["current_turn_settled"] is False
assert poll["turn_continuation"]["current_turn_settled"] is spend_first
assert poll["turn_continuation"]["next_turn_required"] is spend_first
if provider != "legacy":
canonical = read_canonical_todos_if_promoted(
runtime_root=runtime, goal_id=GOAL_ID, include_leases=True
Expand All @@ -281,13 +301,17 @@ def test_completed_advancement_retains_auxiliary_monitor_admission(
assert after_cli["settlement_resume_ref"] == "$.settlement_resume"
assert after_cli["spend_after_validation"] is False
assert poll["settlement_progress"]["state"] == (
"spend_required" if writeback else "writeback_required"
"settled"
if spend_first
else "spend_required"
if writeback
else "writeback_required"
)
rc, replay = _run_cli(registry_path, runtime, *poll_args)
assert rc == 0, replay
assert replay["replayed"] is True
assert _classification_count(runtime, "quota_monitor_poll") == 1
assert _spend_run_count(runtime) == 0
assert _spend_run_count(runtime) == int(spend_first)
if provider != "legacy":
rc, released = _run_cli(
registry_path,
Expand All @@ -306,7 +330,10 @@ def test_completed_advancement_retains_auxiliary_monitor_admission(
"3",
)
assert rc == 0, released
if not writeback:
if spend_first:
assert resume["next_step"] is None
assert "settlement_owed" not in poll
elif not writeback:
assert resume["next_step"]["kind"] == "durable_writeback"
writeback_command = resume["next_step"]["command_template"]
args = tuple(
Expand Down Expand Up @@ -341,5 +368,6 @@ def test_completed_advancement_retains_auxiliary_monitor_admission(
assert rc == 0, replay
assert replay["replayed"] is True
assert replay["settlement_progress"]["state"] == "settled"
assert replay["turn_continuation"]["current_turn_settled"] is True
assert replay["settlement_resume"]["next_step"] is None
assert _classification_count(runtime, "quota_monitor_poll") == 1
57 changes: 41 additions & 16 deletions tests/control_plane_ts/auxiliary_monitor_settlement.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,7 @@ async function settlePrimary(runtime: string): Promise<void> {
turn_instance_id: turn, settlement_identity: identity })).join("\n") + "\n");
}

test("fresh auxiliary admission requires the exact unsettled advancement and ordinary due-work gates", async t => {
test("fresh auxiliary admission requires exact advancement identity and ordinary due-work gates before and after settlement", async t => {
const cases: Array<[string, (params: JsonObject) => JsonObject, string]> = [
["other Turn", p => ({ ...p, turn_instance_id: "other-turn" }), "heartbeat_receipt_identity_conflict"],
["other binding", p => ({ ...p, observation: { ...p.observation as JsonObject, settlement_todo_id: "todo_other" } }), "heartbeat_receipt_identity_conflict"],
Expand All @@ -86,9 +86,11 @@ test("fresh auxiliary admission requires the exact unsettled advancement and ord
work_lane_contract: { must_attempt_work: false }, should_run: false } }), "monitor_poll_admission_rejected"],
["user action", p => ({ ...p, decision: { ...p.decision as JsonObject, requires_user_action: true } }), "monitor_poll_admission_rejected"],
];
for (const [name, change, code] of cases) {
await t.test(name, async st => {
for (const settled of [false, true]) for (const [name, change, code] of cases) {
await t.test(`${settled ? "settled" : "pending"}: ${name}`, async st => {
const { runtime, params } = await fixture(st);
if (settled) await settlePrimary(runtime);
const initialIndex = await readFile(join(runtime, "goals", goal, "runs", "index.jsonl"), "utf8");
const changed = change(params);
await assert.rejects(evaluateQuotaMonitorPollCommit(changed), error => {
assert.ok(error instanceof EffectRuntimeRequestError);
Expand All @@ -102,31 +104,54 @@ test("fresh auxiliary admission requires the exact unsettled advancement and ord
});
await assert.rejects(readFile(join(runtime, "goals", goal, "runs", ".transactions", "quota-monitor-poll")),
{ code: "ENOENT" });
assert.equal(await readFile(join(runtime, "goals", goal, "runs", "index.jsonl"), "utf8"), "");
assert.equal(await readFile(join(runtime, "goals", goal, "runs", "index.jsonl"), "utf8"), initialIndex);
});
}
});

test("completed primary preserves an admitted pending effect across settlement, but cannot admit a new one", async t => {
const { runtime, params } = await fixture(t);
const preflight = await evaluateQuotaMonitorPollCommit(params);
assert.equal(preflight.status, "provider_required", JSON.stringify(preflight));
await settlePrimary(runtime);
const receipt = { schema_version: "monitor_poll_todo_writeback_v0", monitor_effect_id: params.effect_id,
function providerReceipt(params: JsonObject): JsonObject {
return { schema_version: "monitor_poll_todo_writeback_v0", monitor_effect_id: params.effect_id,
goal_id: goal, todo_id: monitor, target_key: "public-watch", result_hash: "unchanged",
dry_run: false, material_change: false, material_change_generation: 0, consecutive_no_change: 1,
last_checked_at: params.generated_at, next_due_at: "2026-09-01T01:00:00Z", cadence: "1h",
todo_update: { ok: true }, next_todos: [], successor_receipts: [] };
}

test("settled primary admits the first due auxiliary observation without another debit or delivery", async t => {
const { runtime, params } = await fixture(t);
await settlePrimary(runtime);
const index = await readFile(join(runtime, "goals", goal, "runs", "index.jsonl"));
const request = { ...params, expected_index_digest: `sha256:${createHash("sha256").update(index).digest("hex")}` };
const preflight = await evaluateQuotaMonitorPollCommit(request);
assert.equal(preflight.status, "provider_required", JSON.stringify(preflight));
const commit = { ...request, phase: "commit", provider_receipt: providerReceipt(request) };
const written = await evaluateQuotaMonitorPollCommit(commit);
assert.equal(written.status, "written", JSON.stringify(written));
assert.equal((written.payload.turn_continuation as JsonObject).current_turn_settled, true);
assert.equal((written.payload.turn_continuation as JsonObject).next_turn_required, true);
assert.equal((written.payload.turn_continuation as JsonObject).same_turn_independent_settlement_allowed, false);
assert.equal((await evaluateQuotaMonitorPollCommit(commit)).status, "replayed");
const rows = (await readFile(join(runtime, "goals", goal, "runs", "index.jsonl"), "utf8"))
.trim().split("\n").map(line => JSON.parse(line));
assert.equal(rows.filter(row => row.classification === "quota_monitor_poll").length, 1);
assert.equal(rows.filter(row => row.classification === "state_refreshed").length, 1);
assert.equal(rows.filter(row => row.classification === "quota_slot_spent").length, 1);
assert.equal(rows.length, 3);
});

test("completed primary preserves an admitted pending effect across settlement and replay reads current closeout", async t => {
const { runtime, params } = await fixture(t);
const preflight = await evaluateQuotaMonitorPollCommit(params);
assert.equal(preflight.status, "provider_required", JSON.stringify(preflight));
await settlePrimary(runtime);
const receipt = providerReceipt(params);
const postBusiness = { ...params, phase: "commit", provider_receipt: receipt,
decision: { ...params.decision as JsonObject, registry_due_monitor: {}, auxiliary_settlement_todo: null,
due_monitor_candidates: [], work_lane_contract: { must_attempt_work: false }, should_run: false } };
assert.equal((await evaluateQuotaMonitorPollCommit(postBusiness)).status, "written");
assert.equal((await evaluateQuotaMonitorPollCommit(postBusiness)).status, "replayed");
await assert.rejects(evaluateQuotaMonitorPollCommit({ ...params, effect_id: "new-after-settlement" }), error => {
assert.ok(error instanceof EffectRuntimeRequestError);
assert.equal(error.code, "heartbeat_receipt_identity_conflict");
return true;
});
const replay = await evaluateQuotaMonitorPollCommit(postBusiness);
assert.equal(replay.status, "replayed");
assert.equal((replay.payload.turn_continuation as JsonObject).current_turn_settled, true);
const rows = (await readFile(join(runtime, "goals", goal, "runs", "index.jsonl"), "utf8"))
.trim().split("\n").map(line => JSON.parse(line));
assert.equal(rows.filter(row => row.classification === "quota_monitor_poll").length, 1);
Expand Down
Loading