From ac46be75db82aa7d641cc796d55cba87dea91cc3 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Mon, 28 Sep 2026 05:58:56 +0800 Subject: [PATCH] fix(quota): close registered causal blocked Turns without spending Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../quota-blocked-causal-closeout-v0.md | 79 +++++++++++ loopx/control_plane/quota/blocked_retry.py | 111 +++------------ loopx/control_plane/quota/blocked_wait.ts | 133 ++++++++++++++++++ loopx/control_plane/quota/settlement_phase.ts | 4 +- .../quota/settlement_readback.ts | 4 + .../test_causal_blocked_closeout_cli.py | 121 ++++++++++++++++ .../causal_blocked_wait.test.ts | 66 +++++++++ .../quota_settlement_readback.test.ts | 73 +++++++++- 8 files changed, 497 insertions(+), 94 deletions(-) create mode 100644 docs/reference/protocols/quota-blocked-causal-closeout-v0.md create mode 100644 loopx/control_plane/quota/blocked_wait.ts create mode 100644 tests/control_plane/test_causal_blocked_closeout_cli.py create mode 100644 tests/control_plane_ts/causal_blocked_wait.test.ts diff --git a/docs/reference/protocols/quota-blocked-causal-closeout-v0.md b/docs/reference/protocols/quota-blocked-causal-closeout-v0.md new file mode 100644 index 0000000000..f8cf9998c4 --- /dev/null +++ b/docs/reference/protocols/quota-blocked-causal-closeout-v0.md @@ -0,0 +1,79 @@ +# Blocked causal Turn closeout / 因果等待的阻塞 Turn 结算 + +## English + +An admitted advancement Turn can discover a real dependency and register a +`monitor_changed:` or `todo_done:` wait with an independently +runnable successor. Previously, blocked no-spend closeout accepted only a +1–30-minute `resume_at` retry. A valid causal wait could therefore leave the +original Turn unsettled and prevent independent work. + +The existing TypeScript quota settlement owner now accepts either the bounded +retry or a `quota_blocked_causal_wait_v0` proof. Python transports current Todo +facts; it does not implement a second wait decision. The existing +`quota.settlement.read` method distinguishes the preflight request schema +`loopx_quota_blocked_wait_request_v0` from ordinary durable readback. + +- Preflight recomputes the existing Todo resume condition from the current + waiting Todo and its unique registered dependency. The waiting Agent Todo + must remain open, active and pending. A monitor must remain open, have a + captured non-negative generation, and still match that baseline exactly. + A completed, archived, missing, self-referential, malformed or stale target + is not proof; the monitor's current generation must be explicit. +- The writeback freezes only those dependency facts and an observation clock. + Exact Goal/Agent/Todo/Turn identity, the admitted guard, a typed blocked + observation and the durable writeback receipt remain mandatory. Historical + readback uses the frozen facts, not today's dependency state. +- Closeout returns `typed_blocked_writeback_no_spend`, with validation and + durable-writeback receipts, no debit and no delivery credit. Exact refresh + retry replays the same result. A later spend request is a no-op for this + already-closed identity; an existing debit is never erased. +- The Todo remains open with its original completion validator and canonical + wait. A fresh Turn can select independent work; existing Todo resume semantics + still decide when this Todo is ready. Do not replace a causal dependency with + a short timer, force an early monitor poll, or treat closeout as completion. +- Legacy `quota_blocked_retry_v0` and the promoted Turn-owned five-minute retry + retain their previous behavior. Old runtimes cannot recognize causal proofs; + finish/reconcile them with a compatible runtime before rollback. Do not remove + a writer fence or rewrite historical receipts. + +The affected user path is CLI/managed-Turn blocked writeback and settlement +readback, not a new configuration. Dashboard already reads canonical +`resume_when`, `resume_ready` and resume receipts; Chat/Lark Todo actions use +the same Todo update owner. These projections do not change, so this slice +adds no frontend control, Lark command or separate UI authority. File and SQLite +CLI/provider acceptance checks both dependency kinds, replay, no debit, +original validator preservation and independent next-Turn selection. It does +not claim live research adoption or PostgreSQL qualification. + +## 中文 + +已准入的 advancement Turn 可以发现真实依赖,以 +`monitor_changed:` 或 `todo_done:` 登记等待,并保留独立可执行 +的 successor。过去无扣额阻塞结算只支持 1–30 分钟 `resume_at`,合法因果等待 +反而会卡住旧 Turn 和独立工作。 + +现有 TS quota settlement owner 新增 `quota_blocked_causal_wait_v0` 核验,Python +只传当前 Todo 事实,不另建判断源。`quota.settlement.read` 依据 +`loopx_quota_blocked_wait_request_v0` 区分预检与原持久结算读回。 + +- 复用 Todo resume owner,以当前开放、active、pending 的 Agent Todo 与唯一注册 + 依赖重算条件。Monitor 须仍开放,非负 generation 与登记基线精确相等;目标 + 已完成、归档、缺失、自引用、格式错误、代际推进/倒退或陈旧投影均不算等待 + 证明;Monitor 当前 generation 必须显式存在。 +- 写回冻结必要依赖事实与观察时间;仍要求精确 Goal/Agent/Todo/Turn、准入 guard、 + typed blocked observation 及持久回执。历史重放不按今天的依赖状态重开旧 Turn。 +- `typed_blocked_writeback_no_spend` 仅含 validation 与 durable-writeback 回执, + 不扣额、不计交付进展。精确刷新幂等重放;已关闭身份的 spend 请求不再追加, + 已有真实扣额不会被抹去。 +- Todo 保持开放、原验收器和 canonical 等待不变。新 Turn 可选独立工作;何时恢复 + 仍由原 Todo resume 语义判断。不得用短定时器替换依赖、强迫提前 poll,或将 + Turn 结算当成 Todo 完成。 +- 保留旧 v0 有界等待与 promoted Turn 自有五分钟重试。降级前须用兼容运行时 + 完成或核对因果回执;不删除 writer fence,不改写历史。 + +产品入口变化是 CLI/managed Turn 的阻塞写回和结算读回,不是新增配置。 +Dashboard 已消费 canonical 等待与回执,Chat/Lark 仍复用 Todo update owner; +不新增前端控件、Lark 命令或独立 UI 权威。File、SQLite 的真实 CLI/provider +验收覆盖两种依赖、重放、零扣额、原验收器保留与下一 Turn 独立选择;不据此 +宣称投研真实采用或 PostgreSQL 资格已通过。 diff --git a/loopx/control_plane/quota/blocked_retry.py b/loopx/control_plane/quota/blocked_retry.py index ea304be39b..3775262417 100644 --- a/loopx/control_plane/quota/blocked_retry.py +++ b/loopx/control_plane/quota/blocked_retry.py @@ -1,11 +1,12 @@ -"""Bound a typed blocked Turn's no-spend closeout to a durable retry.""" +"""Bind typed blocked Turn closeout to a durable retry or causal wait.""" from __future__ import annotations -from datetime import datetime, timedelta +from datetime import datetime from typing import Any -from ..todos.contract import normalize_todo_id, normalize_todo_resume_when +from ..todos.contract import normalize_todo_id +from ..effect_runtime import EffectRuntimeRejected, effect_runtime_result BLOCKED_RETRY_SCHEMA_VERSION = "quota_blocked_retry_v0" MIN_RETRY_SECONDS = 60 @@ -20,97 +21,25 @@ def require_blocked_retry_wait( observed_at: str, allow_turn_settlement_retry: bool = False, ) -> dict[str, Any]: - """Require a bounded Todo wait or mint a canonical Turn-owned retry. - - A peer lease can prevent the blocked agent from updating the Todo. In that - case the committed Turn owns a five-minute wait and selection projects it - without changing the canonical Todo or its validator. - """ - - summary = (todo_fields or {}).get("agent_todos") - items = summary.get("items") if isinstance(summary, dict) else None - todo = ( - next( - ( - item - for item in items - if isinstance(item, dict) - and normalize_todo_id(item.get("todo_id")) == todo_id - ), - None, - ) - if isinstance(items, list) - else None - ) - resume = normalize_todo_resume_when(todo.get("resume_when")) if todo else None - condition = todo.get("resume_condition") if isinstance(todo, dict) else None - if ( - not isinstance(todo, dict) - or todo.get("status") not in {"open", "deferred"} - or todo.get("task_class") != "advancement_task" - ): - raise ValueError( - "typed blocked no-spend closeout requires the same unfinished " - "advancement Todo" - ) + """Transport current Todo facts; TS owns wait qualification and receipts.""" + items = [] + for section in ("agent_todos", "user_todos"): + summary = (todo_fields or {}).get(section) + if isinstance(summary, dict) and isinstance(summary.get("items"), list): + items.extend(summary["items"]) try: - observed = datetime.fromisoformat(observed_at.replace("Z", "+00:00")) - if observed.tzinfo is None: - raise ValueError("observation timestamp must be timezone-aware") - except (TypeError, ValueError) as exc: - raise ValueError("typed blocked retry wait has an invalid timestamp") from exc - if ( - not resume - and not todo.get("resume_when") - and allow_turn_settlement_retry - and todo.get("status") == "open" - ): - due_at = ( - (observed + timedelta(seconds=TURN_SETTLEMENT_RETRY_SECONDS)) - .isoformat() - .replace("+00:00", "Z") - ) - return { - "schema_version": BLOCKED_RETRY_SCHEMA_VERSION, - "source": "turn_settlement", + result = effect_runtime_result("quota.settlement.read", { + "schema_version": "loopx_quota_blocked_wait_request_v0", + "todos": items, "todo_id": todo_id, - "resume_when": f"resume_at:{due_at}", "observed_at": observed_at, - "due_at": due_at, - } - if ( - not resume - or not resume.startswith("resume_at:") - or todo.get("resume_ready") is not False - or not isinstance(condition, dict) - or condition.get("kind") != "resume_at" - or condition.get("resume_when") != resume - or condition.get("satisfied") is not False - ): - raise ValueError( - "typed blocked no-spend closeout requires the same unfinished Todo to " - "have a pending resume_when=resume_at: wait; " - "schedule it with todo update, read it back, then retry this Turn" - ) - try: - due_at = resume.partition(":")[2] - due = datetime.fromisoformat(due_at.replace("Z", "+00:00")) - delay = (due - observed).total_seconds() - except (TypeError, ValueError) as exc: - raise ValueError("typed blocked retry wait has an invalid timestamp") from exc - if not MIN_RETRY_SECONDS <= delay <= MAX_RETRY_SECONDS: - raise ValueError( - "typed blocked retry wait must be due in 1–30 minutes; update " - "the Todo resume_at and retry this same Turn" - ) - return { - "schema_version": BLOCKED_RETRY_SCHEMA_VERSION, - "source": "todo", - "todo_id": todo_id, - "resume_when": resume, - "observed_at": observed_at, - "due_at": due_at, - } + "allow_turn_settlement_retry": allow_turn_settlement_retry, + }) + except EffectRuntimeRejected as exc: + raise ValueError(str(exc)) from exc + if not isinstance(result, dict): + raise RuntimeError("TypeScript blocked wait result must be an object") + return result def active_turn_retry_for_run( diff --git a/loopx/control_plane/quota/blocked_wait.ts b/loopx/control_plane/quota/blocked_wait.ts new file mode 100644 index 0000000000..30a69a501f --- /dev/null +++ b/loopx/control_plane/quota/blocked_wait.ts @@ -0,0 +1,133 @@ +import type { JsonObject } from "../effect_program.ts"; +import { EffectRuntimeRequestError } from "../effect_runtime_errors.ts"; +import { jsonObject, requireJsonObject } from "../runtime_decode.ts"; +import { + evaluateTodoResumeConditions, + normalizeTodoResumeWhen, + resumeConditionHasKnownPendingTarget, + TODO_RESUME_EVALUATION_REQUEST_SCHEMA_VERSION, + TODO_RESUME_NORMALIZE_REQUEST_SCHEMA_VERSION, +} from "../todos/resume_condition.ts"; + +export const BLOCKED_WAIT_REQUEST_SCHEMA = "loopx_quota_blocked_wait_request_v0"; +const CAUSAL_WAIT_SCHEMA = "quota_blocked_causal_wait_v0"; + +function reject(message: string): never { + throw new EffectRuntimeRequestError(`typed blocked no-spend closeout ${message}`); +} + +function causalCondition(waiting: JsonObject, target: JsonObject): JsonObject | null { + if (waiting.status !== "open" || waiting.role !== "agent" || + (waiting.archive_state != null && waiting.archive_state !== "active") || + waiting.task_class !== "advancement_task" || + waiting.resume_ready !== false || waiting.todo_id === target.todo_id || + typeof waiting.resume_when !== "string") return null; + const kind = waiting.resume_when.split(":", 1)[0]; + if (kind !== "monitor_changed" && kind !== "todo_done") return null; + if (target.archive_state != null && target.archive_state !== "active") return null; + if (kind === "monitor_changed" && + (typeof target.material_change_generation !== "number" || + !Number.isSafeInteger(target.material_change_generation) || + target.material_change_generation < 0)) return null; + const evaluated = evaluateTodoResumeConditions({ + schema_version: TODO_RESUME_EVALUATION_REQUEST_SCHEMA_VERSION, + items: [waiting], source_items: [target], + }); + const rows = evaluated.conditions as JsonObject[]; + const condition = jsonObject(rows[0]?.condition); + if (kind === "monitor_changed" && + condition?.material_change_generation !== waiting.resume_monitor_generation) return null; + return condition && resumeConditionHasKnownPendingTarget(condition, waiting) + ? condition : null; +} + +/** Frozen canonical dependency facts, not a caller-authored wait string. The + * durable writeback/guard receipt binds these facts to the exact Turn. Readback + * must not re-evaluate a historical wait against today's dependency state. */ +export function isCausalBlockedWait(value: unknown, todoId: string | null): boolean { + const wait = jsonObject(value); + const waiting = jsonObject(wait?.waiting_todo); + const target = jsonObject(wait?.target_todo); + if (!wait || wait.schema_version !== CAUSAL_WAIT_SCHEMA || wait.source !== "todo" || + !todoId || wait.todo_id !== todoId || waiting?.todo_id !== todoId || + !target || wait.resume_when !== waiting.resume_when || + timestamp(wait.observed_at) === null) return false; + try { + return causalCondition(waiting, target) !== null; + } catch { + return false; + } +} + +function timestamp(value: unknown): number | null { + if (typeof value !== "string") return null; + try { + const resume = normalizeTodoResumeWhen({ + schema_version: TODO_RESUME_NORMALIZE_REQUEST_SCHEMA_VERSION, + resume_when: `resume_at:${value}`, + }); + return resume ? Date.parse(resume.slice("resume_at:".length)) : null; + } catch { + return null; + } +} + +function retainedTodo(todo: JsonObject): JsonObject { + const fields = ["todo_id", "role", "status", "task_class", "archive_state", + "resume_when", "resume_ready", "resume_monitor_generation", "material_change_generation"]; + return Object.fromEntries(fields.filter((field) => todo[field] !== undefined) + .map((field) => [field, todo[field]])); +} + +/** Preflight belongs to the same TS settlement owner as durable readback. + * Python only transports the complete current Todo facts and observation clock. */ +export function prepareBlockedWait(value: unknown): JsonObject { + const request = requireJsonObject(value, "blocked wait request"); + if (request.schema_version !== BLOCKED_WAIT_REQUEST_SCHEMA || !Array.isArray(request.todos)) { + reject("requires current Todo facts"); + } + const todos = request.todos.map((value) => requireJsonObject(value, "blocked wait Todo")); + const matches = todos.filter((todo) => todo.todo_id === request.todo_id); + const todo = matches[0]; + if (matches.length !== 1 || !todo || !["open", "deferred"].includes(String(todo.status)) || + todo.task_class !== "advancement_task") reject("requires the same unfinished advancement Todo"); + const observed = timestamp(request.observed_at); + if (observed === null) reject("has an invalid timestamp"); + const resume = todo.resume_when; + const condition = jsonObject(todo.resume_condition); + if (typeof resume === "string" && /^(?:monitor_changed|todo_done):/.test(resume)) { + const targets = todos.filter((target) => target.todo_id === resume.slice(resume.indexOf(":") + 1)); + const target = targets[0]; + const current = targets.length === 1 && target ? causalCondition(todo, target) : null; + if (!target || !condition || !current || + !resumeConditionHasKnownPendingTarget(condition, todo)) { + reject("requires a registered pending causal target with its captured generation"); + } + // A projection must describe exactly the same generation as its source. + if (condition.kind === "monitor_changed" && + condition.material_change_generation !== current.material_change_generation) { + reject("has a stale causal target generation"); + } + return { schema_version: CAUSAL_WAIT_SCHEMA, source: "todo", todo_id: request.todo_id, + resume_when: resume, observed_at: request.observed_at, + waiting_todo: retainedTodo(todo), target_todo: retainedTodo(target) }; + } + if (!resume && request.allow_turn_settlement_retry === true && todo.status === "open") { + const due = new Date(observed + 300_000).toISOString().replace(".000Z", "Z"); + return { schema_version: "quota_blocked_retry_v0", source: "turn_settlement", + todo_id: request.todo_id, resume_when: `resume_at:${due}`, + observed_at: request.observed_at, due_at: due }; + } + if (typeof resume !== "string" || !resume.startsWith("resume_at:") || + todo.resume_ready !== false || !condition || condition.kind !== "resume_at" || + condition.resume_when !== resume || condition.satisfied !== false) { + reject("requires the same unfinished Todo to have a pending resume_when=resume_at: wait or registered causal wait; schedule it with todo update, read it back, then retry this Turn"); + } + const dueAt = resume.slice("resume_at:".length); + const due = timestamp(dueAt); + if (due === null) reject("has an invalid timestamp"); + const delay = (due - observed) / 1000; + if (delay < 60 || delay > 1800) reject("retry wait must be due in 1–30 minutes; update the Todo resume_at and retry this same Turn"); + return { schema_version: "quota_blocked_retry_v0", source: "todo", todo_id: request.todo_id, + resume_when: resume, observed_at: request.observed_at, due_at: dueAt }; +} diff --git a/loopx/control_plane/quota/settlement_phase.ts b/loopx/control_plane/quota/settlement_phase.ts index 17625e978a..67fc2b2878 100644 --- a/loopx/control_plane/quota/settlement_phase.ts +++ b/loopx/control_plane/quota/settlement_phase.ts @@ -1,8 +1,10 @@ import type { SettlementIdentity } from "../effect_program.ts"; import { jsonObject } from "../runtime_decode.ts"; +import { isCausalBlockedWait } from "./blocked_wait.ts"; -/** A typed blocked Turn may close without spend only with a bounded retry. */ +/** A blocked Turn needs a bounded retry or a verified canonical causal wait. */ export function isBoundedBlockedRetry(value: unknown, todoId: string | null): boolean { + if (isCausalBlockedWait(value, todoId)) return true; const retry = jsonObject(value); if (!retry || retry.schema_version !== "quota_blocked_retry_v0" || (retry.source !== "todo" && retry.source !== "turn_settlement") || diff --git a/loopx/control_plane/quota/settlement_readback.ts b/loopx/control_plane/quota/settlement_readback.ts index ca27d892ae..ff41ea9e90 100644 --- a/loopx/control_plane/quota/settlement_readback.ts +++ b/loopx/control_plane/quota/settlement_readback.ts @@ -54,6 +54,7 @@ import { } from "./heartbeat_receipt_identity.ts"; import { refreshExternalDelivery } from "./refresh_external_delivery.ts"; +import { BLOCKED_WAIT_REQUEST_SCHEMA, prepareBlockedWait } from "./blocked_wait.ts"; export const QUOTA_SETTLEMENT_READBACK_REQUEST_SCHEMA = "loopx_quota_settlement_readback_request_v0"; @@ -1071,6 +1072,9 @@ export function readQuotaSettlementFromSnapshot( } export async function readQuotaSettlement(value: unknown): Promise { + if (jsonObject(value)?.schema_version === BLOCKED_WAIT_REQUEST_SCHEMA) { + return prepareBlockedWait(value); + } const request = decodeRequest(value); return readQuotaSettlementFromRequest( request, diff --git a/tests/control_plane/test_causal_blocked_closeout_cli.py b/tests/control_plane/test_causal_blocked_closeout_cli.py new file mode 100644 index 0000000000..f69bc19a9e --- /dev/null +++ b/tests/control_plane/test_causal_blocked_closeout_cli.py @@ -0,0 +1,121 @@ +"""Real CLI/provider causal wait closeout, not a completion or delivery credit.""" +from __future__ import annotations + +import json +import subprocess +from pathlib import Path + +import pytest + +from canonical_authority_fixture import initialize_canonical_authority, isolate_sqlite_runtime +from loopx.cli import main as cli_main +from test_quota_settlement_cli import ( + AGENT_ID, ALTERNATIVE_TODO_ID, GOAL_ID, TODO_ID, + _classification_count, _configure_completion_validation_todo, + _configure_selectable_alternative, _initialize_git_checkout, _spend_run_count, _write_fixture, +) +from loopx.control_plane.coordination.runtime_shadow import build_todo_runtime_shadow_projection + +MONITOR_ID = "todo_causal_monitor" + + +@pytest.mark.parametrize("provider", ["file", "sqlite"]) +@pytest.mark.parametrize("kind", ["monitor_changed", "todo_done"]) +def test_pending_causal_wait_settles_once_and_releases_independent_work( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, capsys: pytest.CaptureFixture[str], + provider: str, kind: str, +) -> None: + if provider == "sqlite": + isolate_sqlite_runtime(tmp_path, monkeypatch) + project, runtime, registry = _write_fixture(tmp_path) + + def cli(*args: str, cwd: Path = project) -> tuple[int, dict]: + # Exercise the production CLI parser/dispatch with real TS processes + # and provider files, without paying Python cold-start on every read. + with monkeypatch.context() as context: + context.chdir(cwd) + rc = cli_main(["--registry", str(registry), "--runtime-root", str(runtime), + "--format", "json", *args]) + return rc, json.loads(capsys.readouterr().out) + + state = _configure_completion_validation_todo(project) + _configure_selectable_alternative(project) + dependency_owner = AGENT_ID if kind == "monitor_changed" else "codex-dependency-peer" + if kind == "todo_done": + configuration = json.loads(registry.read_text()) + configuration["goals"][0]["coordination"]["registered_agents"].append(dependency_owner) + registry.write_text(json.dumps(configuration)) + _initialize_git_checkout(project) + subprocess.run(["git", "-c", "user.name=LoopX Test", "-c", + "user.email=loopx-test@example.invalid", "commit", "--quiet", + "--allow-empty", "-s", "-m", "causal fixture"], cwd=project, check=True) + workspace = tmp_path / "linked-worktree" + subprocess.run(["git", "worktree", "add", "--quiet", "--detach", str(workspace)], + cwd=project, check=True) + else: + workspace = project + dependency_metadata = ( + "task_class=continuous_monitor target_key=causal-test cadence=6h " + "next_due_at=2099-01-01T00:00:00Z expires_at=2099-01-08T00:00:00Z " + "material_change_generation=0" + if kind == "monitor_changed" else "task_class=advancement_task priority=2" + ) + state.write_text(state.read_text() + ( + "\n- [ ] [P2] Observe a dependency change.\n" + f" \n" + )) + rc, listed = cli("todo", "list", "--goal-id", GOAL_ID) + assert rc == 0, listed + initialize_canonical_authority(runtime, GOAL_ID, build_todo_runtime_shadow_projection( + goal_id=GOAL_ID, todos=listed["todos"], handoff_mode="soft_claim", leases=[], + ), state_path=state, provider=provider) + binding = ("--agent-id", AGENT_ID, "--todo-id", TODO_ID, + "--turn-instance-id", f"causal-blocked-{kind}-{provider}") + rc, guard = cli("quota", "should-run", "--codex-app", + "--goal-id", GOAL_ID, *binding, "--scan-path", str(project), cwd=project) + assert rc == 0 and guard["should_run"] is True, guard + rc, original = cli("todo", "list", "--goal-id", GOAL_ID, "--todo-id", TODO_ID) + assert rc == 0, original + digest = original["todos"][0]["completion_validation_sha256"] + rc, wait = cli("todo", "update", "--goal-id", GOAL_ID, + "--todo-id", TODO_ID, "--agent-id", AGENT_ID, + "--resume-when", f"{kind}:{MONITOR_ID}", + "--successor-todo-id", ALTERNATIVE_TODO_ID) + assert rc == 0, wait + refresh_args = ("refresh-state", "--goal-id", GOAL_ID, + "--classification", "causal_wait_writeback", "--delivery-batch-scale", "single_surface", + "--delivery-outcome", "outcome_gap", *binding, + "--progress-result-class", "blocked", "--progress-blocker-id", MONITOR_ID, + "--progress-evidence-id", "evidence:registered-wait", + "--delivery-workspace-path", str(workspace), + "--no-global-sync", "--suppress-external-sinks") + rc, refresh = cli(*refresh_args, cwd=project) + assert rc == 0, json.dumps(refresh, indent=2) + assert refresh["blocked_retry"]["schema_version"] == "quota_blocked_causal_wait_v0" + assert refresh["settlement_progress"]["state"] == "settled" + assert refresh["settlement_progress"]["closeout_kind"] == "typed_blocked_writeback_no_spend" + assert [r["step_kind"] for r in refresh["settlement_result"]["receipts"]] == ["validation", "durable_writeback"] + rc, replay = cli(*refresh_args, cwd=project) + assert rc == 0 and replay["idempotent_replay"] is True, replay + assert replay["blocked_retry"] == refresh["blocked_retry"] + assert _classification_count(runtime, "causal_wait_writeback") == 1 + rc, spend = cli("quota", "spend-slot", "--goal-id", GOAL_ID, + "--slots", "1", "--source", "heartbeat", "--execute", *binding, + "--scan-path", str(project), cwd=project) + assert rc == 0 and spend["appended"] is False, spend + assert _spend_run_count(runtime) == 0 + rc, after = cli("todo", "list", "--goal-id", GOAL_ID, "--todo-id", TODO_ID) + assert rc == 0, after + todo = after["todos"][0] + assert todo["status"] == "open" and todo["resume_ready"] is False + assert todo["resume_when"] == f"{kind}:{MONITOR_ID}" + assert todo["completion_validation_sha256"] == digest + rc, next_turn = cli("quota", "should-run", "--codex-app", "--goal-id", GOAL_ID, + "--agent-id", AGENT_ID, "--turn-instance-id", f"after-causal-{provider}", + "--scan-path", str(project), cwd=project) + assert rc == 0, next_turn + assert next_turn["effective_action"] != "unsettled_host_turn_recovery" + assert next_turn["selected_todo"]["todo_id"] == ALTERNATIVE_TODO_ID + # No continuous monitor poll was needed or forced by the advancement closeout. + assert _classification_count(runtime, "quota_monitor_poll") == 0 diff --git a/tests/control_plane_ts/causal_blocked_wait.test.ts b/tests/control_plane_ts/causal_blocked_wait.test.ts new file mode 100644 index 0000000000..a660e8eb9a --- /dev/null +++ b/tests/control_plane_ts/causal_blocked_wait.test.ts @@ -0,0 +1,66 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import { BLOCKED_WAIT_REQUEST_SCHEMA, prepareBlockedWait } from "../../loopx/control_plane/quota/blocked_wait.ts"; +import { isBoundedBlockedRetry } from "../../loopx/control_plane/quota/settlement_phase.ts"; +import { evaluateTodoResumeConditions, TODO_RESUME_EVALUATION_REQUEST_SCHEMA_VERSION } from "../../loopx/control_plane/todos/resume_condition.ts"; + +function fixture(kind = "monitor_changed") { + const target = { todo_id: "todo_dependency", status: "open", role: "agent", + task_class: kind === "monitor_changed" ? "continuous_monitor" : "advancement_task", + material_change_generation: 2 }; + const waiting = { todo_id: "todo_waiting", status: "open", role: "agent", + task_class: "advancement_task", resume_when: `${kind}:${target.todo_id}`, + resume_ready: false, resume_monitor_generation: 2 }; + const result = evaluateTodoResumeConditions({ + schema_version: TODO_RESUME_EVALUATION_REQUEST_SCHEMA_VERSION, + items: [waiting], source_items: [target], + }); + const condition = (result.conditions as { condition: Record }[])[0].condition; + return { schema_version: BLOCKED_WAIT_REQUEST_SCHEMA, todo_id: waiting.todo_id, + observed_at: "2026-01-01T00:00:00Z", todos: [{ ...waiting, resume_condition: condition }, target] }; +} + +for (const kind of ["monitor_changed", "todo_done"]) { + test(`${kind} closes only the exact pending registered dependency`, () => { + const request = fixture(kind); + const wait = prepareBlockedWait(request); + assert.equal(wait.schema_version, "quota_blocked_causal_wait_v0"); + assert.equal(isBoundedBlockedRetry(wait, "todo_waiting"), true); + assert.equal(isBoundedBlockedRetry(wait, "todo_other"), false); + assert.equal(isBoundedBlockedRetry({ ...wait, resume_when: `${kind}:todo_other` }, "todo_waiting"), false); + assert.equal(isBoundedBlockedRetry({ ...wait, target_todo: { ...request.todos[1], status: "done" } }, "todo_waiting"), false); + assert.throws(() => prepareBlockedWait({ ...request, todos: [request.todos[0]] }), /registered pending/); + assert.throws(() => prepareBlockedWait({ ...request, todos: [...request.todos, request.todos[1]] }), /registered pending/); + assert.throws(() => prepareBlockedWait({ ...request, todos: [{ ...request.todos[0], status: "done" }, request.todos[1]] }), /unfinished/); + assert.throws(() => prepareBlockedWait({ ...request, todos: [{ ...request.todos[0], resume_ready: true }, request.todos[1]] }), /registered pending/); + assert.throws(() => prepareBlockedWait({ ...request, todos: [request.todos[0], { ...request.todos[1], status: "done" }] }), /registered pending/); + assert.throws(() => prepareBlockedWait({ ...request, todos: [request.todos[0], { ...request.todos[1], archive_state: "archived" }] }), /registered pending/); + }); +} + +test("generation advance, missing baseline and stale projection are not pending proofs", () => { + const request = fixture(); + assert.throws(() => prepareBlockedWait({ ...request, todos: [request.todos[0], { ...request.todos[1], material_change_generation: 3 }] }), /registered pending/); + assert.throws(() => prepareBlockedWait({ ...request, todos: [{ ...request.todos[0], resume_monitor_generation: undefined }, request.todos[1]] }), /registered pending/); + for (const generation of [undefined, -1, 1.5, "2"]) { + assert.throws(() => prepareBlockedWait({ ...request, todos: [request.todos[0], { ...request.todos[1], material_change_generation: generation }] }), /registered pending/); + } + assert.throws(() => prepareBlockedWait({ ...request, todos: [request.todos[0], { ...request.todos[1], material_change_generation: 1 }] }), /registered pending/); + assert.throws(() => prepareBlockedWait({ ...request, todos: [{ ...request.todos[0], resume_condition: { ...request.todos[0].resume_condition, material_change_generation: 1 } }, request.todos[1]] }), /stale causal/); + const wait = prepareBlockedWait(request); + assert.equal(isBoundedBlockedRetry({ ...wait, target_todo: { ...request.todos[1], material_change_generation: 3 } }, "todo_waiting"), false); + assert.equal(isBoundedBlockedRetry({ ...wait, target_todo: { ...request.todos[1], material_change_generation: 1 } }, "todo_waiting"), false); + assert.equal(isBoundedBlockedRetry({ ...wait, observed_at: "2026-01-01T00:00:00" }, "todo_waiting"), false); +}); + +test("unregistered, fabricated and self-referential wait strings never suffice", () => { + const request = fixture(); + assert.throws(() => prepareBlockedWait({ ...request, todos: [{ ...request.todos[0], resume_condition: { satisfied: false } }, request.todos[1]] }), /registered pending/); + assert.throws(() => prepareBlockedWait({ ...request, todos: [{ ...request.todos[0], resume_when: "monitor_changed:todo_waiting" }, request.todos[1]] }), /registered pending/); + assert.throws(() => prepareBlockedWait({ ...request, todos: [request.todos[0], { ...request.todos[1], task_class: "advancement_task" }] }), /registered pending/); + assert.throws(() => prepareBlockedWait({ ...request, observed_at: "unknown" }), /timestamp/); + assert.throws(() => prepareBlockedWait({ ...request, observed_at: "2026-02-30T00:00:00Z" }), /timestamp/); + assert.equal(isBoundedBlockedRetry(prepareBlockedWait({ ...request, observed_at: "1970-01-01T00:00:00Z" }), "todo_waiting"), true); + assert.throws(() => prepareBlockedWait({ ...request, todos: [{ ...request.todos[0], role: "user" }, request.todos[1]] }), /registered pending/); + assert.throws(() => prepareBlockedWait({ ...request, todos: [{ ...request.todos[0], archive_state: "archived" }, request.todos[1]] }), /registered pending/); +}); diff --git a/tests/control_plane_ts/quota_settlement_readback.test.ts b/tests/control_plane_ts/quota_settlement_readback.test.ts index 5ae4f0c1c6..d1a102df08 100644 --- a/tests/control_plane_ts/quota_settlement_readback.test.ts +++ b/tests/control_plane_ts/quota_settlement_readback.test.ts @@ -12,6 +12,8 @@ import { join } from "node:path"; import test from "node:test"; import { settlementIdentity } from "../../loopx/control_plane/effect_program.ts"; +import { BLOCKED_WAIT_REQUEST_SCHEMA, prepareBlockedWait } from "../../loopx/control_plane/quota/blocked_wait.ts"; +import { evaluateTodoResumeConditions, TODO_RESUME_EVALUATION_REQUEST_SCHEMA_VERSION } from "../../loopx/control_plane/todos/resume_condition.ts"; import { projectSemanticReplanGuard, QUOTA_SETTLEMENT_READBACK_REQUEST_SCHEMA, @@ -78,7 +80,7 @@ async function fixture(options: { monitor?: boolean; writebackOutcome?: string; progressObservation?: Record; - blockedRetry?: boolean; + blockedRetry?: boolean | Record; visionCheckpoint?: Record; } = {}) { const runtimeRoot = await mkdtemp(join(tmpdir(), "loopx-settlement-readback-")); @@ -153,7 +155,7 @@ async function fixture(options: { turn_instance_id: turnId, settlement_identity: identity, ...(options.visionCheckpoint ? {vision_checkpoint: options.visionCheckpoint} : {}), - ...(options.blockedRetry ? {blocked_retry: { + ...(options.blockedRetry ? {blocked_retry: typeof options.blockedRetry === "object" ? options.blockedRetry : { schema_version: "quota_blocked_retry_v0", source: "todo", todo_id: todoId, @@ -754,6 +756,73 @@ test("does not pair a spend row with malformed persisted settlement identity", a } }); +test("causal no-spend closeout retains exact receipt identity and historical debits", async () => { + const target = { todo_id: "todo_dependency", role: "agent", status: "open", + task_class: "continuous_monitor", material_change_generation: 2 }; + const waiting = { todo_id: todoId, role: "agent", status: "open", + task_class: "advancement_task", resume_when: `monitor_changed:${target.todo_id}`, + resume_ready: false, resume_monitor_generation: 2 }; + const evaluated = evaluateTodoResumeConditions({ + schema_version: TODO_RESUME_EVALUATION_REQUEST_SCHEMA_VERSION, + items: [waiting], source_items: [target], + }); + const condition = (evaluated.conditions as Record[])[0].condition; + const proof = prepareBlockedWait({ schema_version: BLOCKED_WAIT_REQUEST_SCHEMA, + todo_id: todoId, observed_at: "2026-09-24T10:00:00Z", + todos: [{ ...waiting, resume_condition: condition }, target] }); + const options = { writeback: true, writebackOutcome: "outcome_gap", blockedRetry: proof, + progressObservation: { schema_version: "typed_progress_observation_v0", + result_class: "blocked", work_item_id: todoId, + blocker_id: target.todo_id, evidence_ids: ["evidence:canonical-wait"] } }; + const runtime = await fixture(options); + try { + const result = await readQuotaSettlement(request(runtime)); + assert.equal((result.progress as any).state, "settled"); + assert.equal((result.progress as any).closeout_kind, "typed_blocked_writeback_no_spend"); + assert.equal(result.replay_phase, "settled"); + // Frozen historical facts stay closed after a later dependency observation. + await appendFile(join(runtime, "goals", goalId, "runs", "index.jsonl"), + `${JSON.stringify({ classification: "quota_monitor_poll", goal_id: goalId, + agent_id: agentId, todo_id: target.todo_id, turn_instance_id: "later-turn", + material_change_generation: 3 })}\n`); + assert.equal((await readQuotaSettlement(request(runtime))).replay_phase, "settled"); + } finally { + await rm(runtime, { recursive: true, force: true }); + } + for (const field of ["goal_id", "agent_id", "todo_id", "turn_instance_id"]) { + const mismatch = await fixture(options); + try { + const index = join(mismatch, "goals", goalId, "runs", "index.jsonl"); + const run = JSON.parse((await readFile(index, "utf8")).trim()); + await writeFile(index, `${JSON.stringify({ ...run, [field]: "another-identity" })}\n`); + const result = await readQuotaSettlement(request(mismatch)); + assert.notEqual((result.progress as any).state, "settled", field); + assert.notEqual(result.replay_phase, "settled", field); + } finally { + await rm(mismatch, { recursive: true, force: true }); + } + } + const missing = await fixture({ ...options, guard: false }); + const missingWriteback = await fixture(options); + const spent = await fixture({ ...options, spend: true }); + const malformed = await fixture({ ...options, blockedRetry: { ...proof, + waiting_todo: { ...waiting, resume_monitor_generation: 3 } } }); + try { + assert.notEqual((await readQuotaSettlement(request(missing))).replay_phase, "settled"); + const log = join(missingWriteback, "goals", goalId, "rollout-event-log.jsonl"); + const events = (await readFile(log, "utf8")).trim().split("\n").map(line => JSON.parse(line)); + await writeFile(log, `${events.filter(event => event.event_kind !== "refresh_state").map(event => JSON.stringify(event)).join("\n")}\n`); + assert.notEqual((await readQuotaSettlement(request(missingWriteback))).replay_phase, "settled"); + const debited = await readQuotaSettlement(request(spent)); + assert.equal((debited.spend as any).payload.ok, true); + assert.equal((debited.progress as any).closeout_kind, undefined); + assert.equal((await readQuotaSettlement(request(malformed))).replay_phase, "open"); + } finally { + await Promise.all([missing, missingWriteback, spent, malformed].map(path => + rm(path, { recursive: true, force: true }))); + } +}); + test("accepts only an attributable typed blocker as an outcome-gap writeback", async () => { const qualifiedRuntime = await fixture({ writeback: true,