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
79 changes: 79 additions & 0 deletions docs/reference/protocols/quota-blocked-causal-closeout-v0.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,79 @@
# Blocked causal Turn closeout / 因果等待的阻塞 Turn 结算

## English

An admitted advancement Turn can discover a real dependency and register a
`monitor_changed:<todo_id>` or `todo_done:<todo_id>` 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_id>` 或 `todo_done:<todo_id>` 登记等待,并保留独立可执行
的 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 资格已通过。
111 changes: 20 additions & 91 deletions loopx/control_plane/quota/blocked_retry.py
Original file line number Diff line number Diff line change
@@ -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
Expand All @@ -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:<timezone-aware-time> 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(
Expand Down
133 changes: 133 additions & 0 deletions loopx/control_plane/quota/blocked_wait.ts
Original file line number Diff line number Diff line change
@@ -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:<timezone-aware-time> 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 };
}
4 changes: 3 additions & 1 deletion loopx/control_plane/quota/settlement_phase.ts
Original file line number Diff line number Diff line change
@@ -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") ||
Expand Down
4 changes: 4 additions & 0 deletions loopx/control_plane/quota/settlement_readback.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -1071,6 +1072,9 @@ export function readQuotaSettlementFromSnapshot(
}

export async function readQuotaSettlement(value: unknown): Promise<JsonObject> {
if (jsonObject(value)?.schema_version === BLOCKED_WAIT_REQUEST_SCHEMA) {
return prepareBlockedWait(value);
}
const request = decodeRequest(value);
return readQuotaSettlementFromRequest(
request,
Expand Down
Loading
Loading