From f876989cfbd0afa1ebf5e7f07904edd1bd6f6d2a Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 12:16:22 +0800 Subject: [PATCH] fix(quota): admit due auxiliary monitors after primary settlement Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../quota-monitor-observation-receipt-v0.md | 28 +++++++-- .../quota/monitor_poll_commit.ts | 54 +++++++++++++----- .../test_monitor_observation_admission.py | 48 ++++++++++++---- .../auxiliary_monitor_settlement.test.ts | 57 +++++++++++++------ 4 files changed, 143 insertions(+), 44 deletions(-) diff --git a/docs/reference/protocols/quota-monitor-observation-receipt-v0.md b/docs/reference/protocols/quota-monitor-observation-receipt-v0.md index 4d5ee77eb0..62f8a0c662 100644 --- a/docs/reference/protocols/quota-monitor-observation-receipt-v0.md +++ b/docs/reference/protocols/quota-monitor-observation-receipt-v0.md @@ -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. @@ -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 @@ -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 已结算。 ### 验收 @@ -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 和 diff --git a/loopx/control_plane/quota/monitor_poll_commit.ts b/loopx/control_plane/quota/monitor_poll_commit.ts index 8eb4e7010a..32fbd42e25 100644 --- a/loopx/control_plane/quota/monitor_poll_commit.ts +++ b/loopx/control_plane/quota/monitor_poll_commit.ts @@ -541,6 +541,18 @@ interface Admission { external: boolean; } +async function readAuxiliarySettlement(request: MonitorRequest): Promise { + 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 { @@ -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; @@ -1488,14 +1496,30 @@ function successorReceipts(providerReceipt: JsonObject | null): JsonObject[] { : []; } -function payloadFor( +async function projectAuxiliaryContinuation(request: MonitorRequest, payload: JsonObject): Promise { + 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 { const event = requiredObject(record.monitor_event, "record.monitor_event"); const receipts = successorReceipts(request.provider_receipt); const payload: JsonObject = { @@ -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; } @@ -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, @@ -2045,7 +2071,7 @@ export async function evaluateQuotaMonitorPollCommit( fingerprint, "preview", record, - payloadFor( + await payloadFor( effectiveRequest, record, paths.jsonPath, @@ -2295,7 +2321,7 @@ export async function evaluateQuotaMonitorPollCommit( effectiveRequest.generated_at, request.effect_id, ); - const payload = payloadFor( + const payload = await payloadFor( effectiveRequest, record, jsonPath, diff --git a/tests/control_plane/test_monitor_observation_admission.py b/tests/control_plane/test_monitor_observation_admission.py index 62787d0433..6940220b70 100644 --- a/tests/control_plane/test_monitor_observation_admission.py +++ b/tests/control_plane/test_monitor_observation_admission.py @@ -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 @@ -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, @@ -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 @@ -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, @@ -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( @@ -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 diff --git a/tests/control_plane_ts/auxiliary_monitor_settlement.test.ts b/tests/control_plane_ts/auxiliary_monitor_settlement.test.ts index fb9bdc83e6..07b039d9b6 100644 --- a/tests/control_plane_ts/auxiliary_monitor_settlement.test.ts +++ b/tests/control_plane_ts/auxiliary_monitor_settlement.test.ts @@ -67,7 +67,7 @@ async function settlePrimary(runtime: string): Promise { 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"], @@ -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); @@ -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);