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 @@ -2642,10 +2642,17 @@ source paths, authorize monitor writeback, or change provider/promotion holds.

**D1 — qualify permanent projection delivery; may overlap T1/T2.**

The T2 monitor successor route owner is now shared across preflight, the legacy
effect adapter and receipt checks. Its result proves only normalized intent,
not actor authority, provider commit or atomic monitor-plus-successor durability.
Keep the monitor writer fence and promotion hold until that transaction closes.
T2 now commits a lease-free native Monitor observation and its independent
successors in one canonical CAS/receipt; the route planner alone still grants
no authority. The CLI delivers committed state through the existing journal/
outbox renderer: display failure is pending, not rollback or successor recreation.
Quota consumes the same v0 business receipt for its separate settlement.
Validation covers the real CLI with missing display, operation replay after a
renderer failure, and complex-data concurrency/lost-acknowledgement recovery on
File, NoKV and isolated real PostgreSQL. Retained Monitor leases and cross-owner
successor claims remain explicitly unsupported. This slice changes neither the
provider default, writer fence nor promotion approval, and does not replace the
independent legacy three-arm comparison or D2 soak.

- Start from `loopx/control_plane/todos/provider_projection.py`, the existing
Todo-section renderer and canonical journal/outbox. #4097 already recovers
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2095,10 +2095,13 @@ summary,之前消费 legacy summary;真实 CLI 覆盖容量变化和 promote

**D1 — 资格化永久投影交付,可与 T1/T2 重叠推进。**

T2 monitor successor 的路由 owner 已由 preflight、legacy effect adapter 和回执校验
共享。其结果仅证明规范化 intent,不证明 actor authority、provider commit 或
monitor-plus-successor 原子持久化;事务闭合前继续保留 monitor writer fence 与
promotion hold。
T2 的无 lease 原生 Monitor 观察与独立后继现由同一 canonical CAS/receipt 提交;
route planner 本身仍不授予权限。CLI 将已提交回执交给既有 journal/outbox renderer,
展示失败标为 pending,不回滚提交、不重新生成后继;quota 继续消费同一 v0 业务回执
完成独立记账。验证覆盖真实 CLI 的缺失 display、renderer 失败后的 operation 重放,
及 File/NoKV/真实隔离 PostgreSQL 的复杂数据、并发和丢回执恢复。
带 lease Monitor、跨 owner claim 等未闭合能力仍明确拒绝;此切片不改变 provider
默认、writer fence 或 promotion 审批,也不替代三臂 legacy 对照及 D2 soak。

- 从 `loopx/control_plane/todos/provider_projection.py`、既有 Todo-section renderer、
canonical journal/outbox 入手。复用 #4097 已有的缺失 Todo section 恢复及
Expand Down
29 changes: 24 additions & 5 deletions docs/architecture/rfcs/typescript-control-plane-migration-v0.md
Original file line number Diff line number Diff line change
Expand Up @@ -449,11 +449,30 @@ are compared as the same route at readback. The original wire observation still
owns the v0 replay digest; normalization must not silently invalidate pending
receipts. The node-independent repository/bootstrap codec remains separately
characterized, not replaced by a runtime dependency.
This is **not** the T2 atomic transaction: monitor mutation and successor writes
still use existing fenced effects. Cross-effect crash recovery, native writer
closure and whole-Goal promotion remain held; do not infer them from a route plan.

- Inventory `monitor_poll_writeback.py` and its event/Todo/lease callers.
The native `coordination.local_authority.monitor_poll` transaction now commits a
lease-free Monitor observation and its requested independent successors against
one canonical revision, with one CAS and durable operation receipt. It composes
the existing generation, successor-route, User authoring-scope and Todo-create
planners. Public create and Monitor batches share create admission/duplicate
planning; target selection is shared by legacy preflight and native commit.
Python only routes provider intent and drains the existing projection outbox.

Explicit semantic corrections: completed/archived Monitor targets are rejected;
target-key selection ignores finished history but never guesses between live matches;
successor authoring requires an actually advanced material-change generation,
not merely a repeated `material_change=true` assertion for the same evidence.
Retrying the original operation recovers the original successors instead of
creating new work. A fresh observation with no successor remains valid. User
gates use the existing actor-bound scope, never an inferred global gate.

Boundaries still open: any retained Monitor lease fails closed in this native
operation; cross-owner successor claims are not implicitly authorized. Unpromoted
Goals retain their legacy writer. Quota accounting stays in its existing
preflight/writeback/settlement protocol and reuses the v0 receipt shape and raw
observation identity. Canonical commit success is independent of pending Markdown
delivery. This does not finish all T2 commands or authorize whole-Goal promotion.

- Finish the retained lease and event callers of `monitor_poll_writeback.py`.
Reuse existing monitor generation, independent-successor and settlement
owners. Compose one transaction rather than adding a second monitor engine.
- Preserve unchanged polling/reschedule behavior, generation fences,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -346,10 +346,25 @@ guard/resolver 及 TS 回执端独立的默认值/capability 解释。非法 c
action/claim/capability 别名和 Git transport 在回执核对时指向同一路由。v0 replay
digest 仍绑定原始 wire observation,不能因规范化而悄悄使 pending receipt 失效。
无需 Node 的 repository/bootstrap codec 暂留并做跨运行时对照,不引入启动依赖。
这**不是** T2 原子事务:monitor mutation 和 successor 写入仍通过既有 fenced effect
执行;跨 effect crash 恢复、native writer 闭合及整 Goal promotion 仍未放行。

- 盘点 `monitor_poll_writeback.py` 及 event/Todo/lease caller,复用 monitor
原生 `coordination.local_authority.monitor_poll` 现将无 lease Monitor 的观察及请求的
独立后继,绑定同一个 canonical revision,以一次 CAS 和持久 operation receipt
提交。它组合已有 generation、successor route、User authoring scope 和 Todo create
planner;单项 create 与 Monitor 批次共用创建准入/语义去重,legacy preflight 与
native commit 共用目标选择。Python 只路由意图并交付既有 projection outbox。

明确的语义修正:拒绝已完成/归档的 Monitor;target-key 选择排除结束的历史项,
但多个活跃匹配仍要求显式 id;创建后继必须实际推进 material-change
generation,不能对相同证据重复声明 `material_change=true` 就继续生成任务。
原 operation 重试恢复原后继,不创建新工作;不附带后继的新 observation 仍可接受。
User gate 复用既有 actor-bound scope,不推导全局 gate。

尚未闭合:原生操作遇到任何保留的 Monitor lease 仍 fail closed,不隐式授权跨 owner
的 successor claim;未 promotion Goal 仍走旧 writer。Quota 记账继续使用现有
preflight/writeback/settlement 协议,沿用 v0 回执形状及原始 observation identity。
Canonical 提交成功独立于 Markdown delivery pending。这不代表全部 T2 命令或整 Goal
promotion 已完成。

- 继续闭合 `monitor_poll_writeback.py` 保留的 lease 与 event caller,复用 monitor
generation、独立 successor 和 settlement owner,组成一笔事务,不建第二套引擎。
- 保持 unchanged poll/reschedule、generation fence、material-change successor
去重和可归属 settlement。Monitor 不是 delivery 执行任务;独立 advancement Todo
Expand Down
3 changes: 2 additions & 1 deletion loopx/cli_commands/quota.py
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,7 @@
from ..control_plane.coordination.legacy_writer_fence import (
LegacyCoordinationWriterFenced,
)
from ..control_plane.coordination.local_authority import LocalCoordinationAuthorityUnavailable
from ..control_plane.effect_runtime import EffectRuntimeRejected
from ..control_plane.scheduler.execution_context import (
GUIDED_START_TURN_RUNTIME_PROFILES,
Expand Down Expand Up @@ -241,7 +242,7 @@ def _quota_failure_payload(
)
if error.agent_id is not None:
payload["agent_id"] = error.agent_id
elif isinstance(error, LegacyCoordinationWriterFenced):
elif isinstance(error, (LegacyCoordinationWriterFenced, LocalCoordinationAuthorityUnavailable)):
payload.update(
{
"error_code": error.code,
Expand Down
11 changes: 5 additions & 6 deletions loopx/cli_commands/quota_context.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,8 +5,8 @@
from dataclasses import dataclass
from pathlib import Path

from ..control_plane.coordination.legacy_writer_fence import (
require_legacy_coordination_write_allowed,
from ..control_plane.scheduler.provider_monitor_poll import (
require_monitor_poll_source_available,
)
from ..control_plane.quota.error_codes import QuotaCommandValidationError
from ..control_plane.runtime.status_projection_cache import (
Expand Down Expand Up @@ -225,10 +225,9 @@ def prepare_quota_command_context(
runtime_root_override=runtime_root_arg,
)
if command == "monitor-poll" and args.execute and (args.todo_id or args.target_key):
# This command still uses the legacy Todo writer. Preserve its typed
# rejection before collecting a promoted read model (which may itself
# be unavailable). The writer repeats the check under its mutation lock.
require_legacy_coordination_write_allowed(
# Canonical availability precedes unrelated status/quota preparation.
# The eventual transaction repeats its fence check under the writer lock.
require_monitor_poll_source_available(
runtime_root=runtime_root, goal_id=args.goal_id,
)
status_goal_id = args.goal_id if command not in {"status", "plan"} else None
Expand Down
30 changes: 30 additions & 0 deletions loopx/control_plane/coordination/local_authority_runtime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,8 @@ import { ShadowManagementError, requireShadowPrimaryWriteAllowed, shadowMaintena
import { isAbsolute, join } from "node:path";

import type { JsonObject } from "../effect_program.ts";
import {executeCoordinationMonitorPoll, COORDINATION_MONITOR_POLL_REQUEST_SCHEMA,
COORDINATION_MONITOR_POLL_RESULT_SCHEMA} from "./todo_monitor_poll.ts";
import { requireJsonObject } from "../runtime_decode.ts";
import {
LOCAL_COORDINATION_MUTATION_REQUEST_SCHEMA,
Expand Down Expand Up @@ -101,6 +103,34 @@ async function withCanonicalWriter<T>(root: string, goalId: string, dryRun: bool
});
}

/** Monitor observation and successors share the existing writer/fence lifetime. */
export async function pollLocalCoordinationMonitor(value: unknown,
dependencies: LocalAuthorityRuntimeDependencies = {}): Promise<JsonObject> {
const evidence = {source_authority: "file_v0", decision_read_from_provider: true, legacy_fallback_used: false};
try {
const input = requireJsonObject(value, "local Monitor poll request");
if (input.schema_version !== COORDINATION_MONITOR_POLL_REQUEST_SCHEMA) throw new TypeError("Monitor poll schema mismatch");
const root = runtimeRoot(input.runtime_root);
const goalId = requireAuthorityStoreId(input.goal_id, "goal id");
if (!Array.isArray(input.registered_agents)) throw new TypeError("registered_agents must be an array");
const registered = input.registered_agents.map(agent => claimAgentValue(agent, "registered agent"));
return await withCanonicalWriter(root, goalId, input.dry_run === true, async () => ({
...await executeCoordinationMonitorPoll(dependencies.createStore?.(authorityDirectory(root), goalId) ??
new FileAuthorityStore(authorityDirectory(root), goalId), {
goal_id: goalId, operation_id: requireAuthorityStoreId(input.operation_id, "operation id"),
actor_agent_id: input.actor_agent_id == null ? null : claimAgentValue(input.actor_agent_id, "actor_agent_id"),
registered_agents: registered, dry_run: input.dry_run as boolean,
observation: requireJsonObject(input.observation, "Monitor observation"),
intent: requireJsonObject(input.intent, "Monitor successor intent"),
}), ...evidence,
}));
} catch (error) {
return {schema_version: COORDINATION_MONITOR_POLL_RESULT_SCHEMA, status: "failed", changed: false,
reason_code: error instanceof ShadowManagementError ? error.reason_code : "invalid_local_monitor_poll_request",
reason: error instanceof Error ? error.message : String(error), ...evidence};
}
}

interface LocalAuthorityRuntimeDependencies {
createStore?: (directory: string, goalId: string) => AuthorityStore;
createShadowStore?: (directory: string, goalId: string) => AuthorityStore;
Expand Down
34 changes: 23 additions & 11 deletions loopx/control_plane/coordination/todo_create.ts
Original file line number Diff line number Diff line change
Expand Up @@ -166,6 +166,26 @@ function createCandidate(
};
}

/** In-process create planning for a caller-owned canonical transaction. Never
* commits or authorizes the enclosing operation; single create and Monitor
* batches share record admission, attribution and semantic duplicate rules. */
export function planCoordinationTodoCreate(
rawInput: CoordinationTodoCreateInput,
todos: ReadonlyMap<string, JsonObject>,
readModelSchema: unknown,
): CoordinationTodoCreateResult {
const input = normalizeCreateInput(rawInput);
const duplicate = [...todos.values()].find((todo) =>
todo.role === input.todo.role && todo.archive_state === "active" &&
todo.status !== "done" && todo.status !== "deferred" && todo.text === input.todo.text);
if (duplicate !== undefined) return semanticDuplicateResult(input.todo, duplicate, null, null);
if (todos.has(String(input.todo.todo_id))) {
return failure("todo_already_exists", "Todo id already exists in canonical authority", {todo_id: input.todo.todo_id});
}
return {schema_version: COORDINATION_TODO_CREATE_RESULT_SCHEMA, status: "planned",
changed: true, todo_id: input.todo.todo_id, todo: createCandidate(input, readModelSchema)};
}

async function commitCreate(
store: AuthorityStore,
input: CoordinationTodoCreateInput,
Expand Down Expand Up @@ -241,18 +261,10 @@ export async function executeCoordinationTodoCreate(
);
}
const todoId = requireAuthorityStoreId(input.todo.todo_id, "todo id");
const duplicate = [...projection.todos.values()].find((todo) =>
todo.role === input.todo.role && todo.archive_state === "active" &&
todo.status !== "done" && todo.status !== "deferred" && todo.text === input.todo.text
);
if (duplicate !== undefined) {
return semanticDuplicateResult(input.todo, duplicate, head.provider_revision, head.cursor);
}
if (projection.todos.has(todoId)) {
return failure("todo_already_exists", "Todo id already exists in canonical authority", {todo_id: todoId});
}
const readModel = canonicalAuthorityObject(head.head.todo_read_model, "Todo read model");
const created = createCandidate(input, readModel.schema_version);
const plan = planCoordinationTodoCreate(input, projection.todos, readModel.schema_version);
if (plan.status !== "planned") return {...plan, provider_revision: head.provider_revision, cursor: head.cursor};
const created = canonicalAuthorityObject(plan.todo, "created Todo");
if (input.dry_run) {
return {
schema_version: COORDINATION_TODO_CREATE_RESULT_SCHEMA,
Expand Down
Loading