diff --git a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md index 329ab40ccf..0738892d95 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md @@ -3016,6 +3016,16 @@ independent legacy three-arm comparison or D2 soak. move. Exit with deterministic freshness/readback and an actionable repair path; a successful render once is insufficient. +Continuation and closure readback now uses the same TS relation evidence for +completed-work gaps and handoff states before filtering/capping. Capture v1 +retains reachable archived successors and original deferred status; derived +summary evaluations never enter provider records. Completion retries also +recover the matching receipt when a peer commits between receipt and head reads, +without accepting state-only replay with fresh validation evidence. Real CLI and complete-graph +provider conformance cover the consumer family. See [operation and semantic +changes](../../reference/todo-continuation-readback.md). This closes a bounded +L5/L7 gap; permanent projection delivery/recovery, D2 and D3 are still open. + **D2 — qualify exactly one local profile; independent of PostgreSQL deployment.** - Reconcile the SQLite candidate #4121 with Section 7.2 before adding code. @@ -3073,8 +3083,9 @@ The reconciled baseline includes #4286 (command receipts/archive), #4289 #4316 (Goal Channel observation), #4317 (provider opening), #4348 (renew), #4328 (first SQLite D2 batch) and #4334 (PostgreSQL service admission): all are merged at the 2026-09-20 checkpoint. Their existence does not qualify the full -cards. #4732 remains the open Monitor observation/reactivation slice; #4224 -retains contributor ownership of SQLite D2. Re-read actual heads before work. +cards. Monitor observation/reactivation #4732 and linked User completion #4754 +are also merged. #4224 retains contributor ownership of SQLite D2. Re-read +actual heads before work. The identifiers below are **planned PR packages**, not reserved GitHub numbers. A package may split at a real effect/compatibility boundary; changing languages @@ -3096,12 +3107,11 @@ or moving a helper is not by itself a package exit. packages as complete operations while L6/L7 progress independently. B integrates those contracts into complete user flows; C has one reproducible qualification checkpoint; D changes the default in its own reviewable PR. After the linked -User completion slice, the 2026-09-20 planning estimate is **6–9 further cohesive +User completion slice, the 2026-09-20 planning estimate is **5–8 further cohesive PRs**, conditional on the caller audit finding no additional missing effects: | Remaining work package | Estimated PRs | Exit | | --- | --- | --- | -| Monitor observation/reactivation | 1, existing #4732 | Real caller and complete graph acceptance; avoid a duplicate implementation. | | Remaining L2/L3 caller and executor-effect fences | 1–2 | Actual CLI/Turn/Chat command inventory and external-effect boundary closure. | | L5 / D1 consumer and projection closure | 1 | Full consumer parity, lag/recovery and packaged client readback. | | L6 / SQLite D2 | 1–2, contributor-owned #4224 | Capacity, crash/restore and separately authorized elapsed-soak evidence on one profile. | @@ -3113,8 +3123,8 @@ Scope may split only where a real effect/compatibility boundary warrants it. Small Python business-rule deletions can ship with each TS owner; rendering, private command execution and import/export keep their active adapters. -截至 2026-09-20,关联 User 完成链路补齐后,按以上六类完整交付边界估算还需 **6–9 个 PR**。 -Monitor 复用 #4732,SQLite D2 仍归 #4224 contributor;其余顺序是调用方/执行围栏、 +截至 2026-09-20,关联 User 完成链路补齐后,按以上五类完整交付边界估算还需 **5–8 个 PR**。 +Monitor #4732 已合并,SQLite D2 仍归 #4224 contributor;其余顺序是调用方/执行围栏、 消费与投影、capture 与整 Goal 演练,最后独立切换默认值。该估算以未发现更多缺失 effect 为前提,不是合并数承诺,也不要求先删完 Python。TS owner 每收敛一块即可 删除对应旧规则;仍有真实调用方的渲染、私有命令执行和导入导出适配器继续保留。 diff --git a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md index d4b45d98da..b3f60c2e3d 100644 --- a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md +++ b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md @@ -1017,6 +1017,19 @@ history before that checkpoint. Successful schemas, File/NoKV persisted bytes, request identity and revision algorithms remain compatible. This supports T3/D1 readers but does not finish Todo writers, retention/compaction or promotion. +Continuation readback now shares one typed succession resolver, handoff state +machine and summary closure decision. The legacy adapter no longer owns those +rules. Full-source evaluations survive display selection; nonexistent/self +successors cannot certify closure and archived continuation evidence survives +capture. The existing archive-capture request advances to v1 so older runtimes +cannot silently omit the expanded graph. Query subsets do not emit whole-source +closure proofs, and bounded handoff views preserve their state and exclusions. +See [continuation readback](../../reference/todo-continuation-readback.md). +This closes that T3/L5 consumer family and its bounded L7 dependency, not D1–D3 +or every T3 consumer. Python retains codecs, IO and the documented legacy route +prose hint until its remaining writers emit explicit replan flags; no new +capability/provider or parallel business authority is introduced. + **T4 — collect full-writer retirement after durability cutover.** - The 2026-09-19 command audit retires two already-typed but unconsumed diff --git a/docs/reference/todo-continuation-readback.md b/docs/reference/todo-continuation-readback.md new file mode 100644 index 0000000000..91f0bd9fc1 --- /dev/null +++ b/docs/reference/todo-continuation-readback.md @@ -0,0 +1,138 @@ +# Todo continuation and closure readback + +Todo list, status and quota distinguish a completed record from a closed work +slice. A completed tracked advancement Todo still needs an existing successor +or an explicit `no_followup=true`. This read policy lives in +`control_plane/todos/succession.ts`; Python normalizes legacy input and renders +its decisions. It belongs to the existing Todo control plane, uses the selected +AuthorityStore, and adds no capability or extension provider. + +## Relationship evidence + +The policy evaluates the complete available Todo graph before role, status, +Agent, ID or display-limit selection. It recognizes explicit +`successor_todo_ids`, `superseded_by`, and advancement records pointing back +through `unblocks_todo_id` or `resume_when=todo_done:`. + +- A declared successor must exist and differ from the source. A dangling or + self reference does not close work. A retained archived record remains + relationship evidence; it does not become active work. +- Explicit links retain their existing role-neutral meaning. Inferred + successors require an advancement task. `monitor_changed` is a resume + condition, not an inferred work successor. +- An existing successor records continuation lineage. It does not prove that + the successor has executed, been accepted or acquired a lease. This is not a + transitive Goal acceptance proof or a cycle-freedom certificate. +- Basic historical checkboxes without structured execution context keep their + compatibility behavior. Explicit no-follow-up remains an independent closeout + choice. Deferred work is never classified as a completed advancement gap. + +Both completed-work warnings and handoff gates use that same graph. Previously +handoff ignored explicit successor lists, while completed-work warnings accepted +nonexistent/self links. Filtering or archiving a valid inferred successor could +also manufacture a warning that was absent on the full source. + +| Handoff facts, in precedence order | State | +| --- | --- | +| Existing, non-self supersession target | `superseded` | +| Deferred source | `deferred` | +| Source has not completed | `blocking` | +| Completed with explicit no-follow-up | `cleared_no_followup` | +| Completed with a resolved successor | `cleared_with_successor` | +| Completed without either | `cleared_without_successor` | + +Only active dependency-linked executor exclusions are handoff gates. These +states describe the gate; none changes claims, grants, leases or stored Todos. +The existing legacy stale-closeout prose hint remains a compatibility adapter +until route-closeout writers supply the explicit replan flag. Its substring +matching can overmatch narrative and is not used for successor resolution, +handoff state or permission. An explicit boolean replan flag takes precedence. + +## Selection, proofs and transport + +A fresh full-source evaluation accompanies each internal summary row as an +ephemeral Python attribute, outside dictionary fields and JSON serialization. Its fact +digest prevents reuse after relevant item edits; it is a consistency check, +not authentication. Fresh parsing/canonical reads always recompute it rather +than trusting stored evaluations. Shadow capture discards this derived field; +canonical records and durable source digests do not gain a second authority. +Public parser rows keep their existing dictionary schema. Final list/status +responses copy plain dictionaries, retaining decision fields without exposing +the internal evaluation or expanding the hot-path payload. + +Legacy archive/recreate can retain one archived and one active record with the +same logical Todo ID. The active record owns that ID's inferred edges regardless +of source order; archived metadata cannot supply stale edges for the replacement. +Two active or two archived records with the same ID remain ambiguous and reject. +This read precedence does not relax canonical capture's unique-identity contract. + +A status/ID/Agent-filtered list describes that selection but emits no Goal-source +or terminal-closure proof. A display limit alone does not change the source: +counts, warning decisions and proof eligibility are computed first. Handoff +state, successor count and executor exclusions survive the bounded list view. + +Terminal closure additionally requires no deferred/convergent work, unresolved +handoff, successor gap or route-replan obligation. Watch-only monitors retain +the existing convergent-work exception. It remains separate from Goal acceptance. + +Field-presence sets are interned inside a succession RPC request so long archive +histories do not repeat identical metadata shapes. The full graph is retained; +no record sampling, per-page rule evaluation or RPC limit increase is used. + +## Migration and operation + +The shared archive-capture owner retains the reachable continuation graph as +well as resume dependencies and standing decisions. It preserves real record +status, including deferred history: capturing a deferred record does **not** +satisfy `todo_done`. Duplicate identities, invalid archive state and incompatible +role/authority combinations still reject capture. Unrelated archive records +remain outside the bounded canonical capture. + +The internal request is `todo_archive_dependency_capture_request_v1`. Python +and the bundled TS runtime must be upgraded together; an older runtime rejects +the new request instead of silently omitting continuation edges. Existing +historical capture/promotion receipts are not rewritten or upgraded in place. +Requalify capture on this runtime before a future promotion. + +Use existing read commands; no activation or new option is needed: + +```bash +loopx --registry registry.json todo list --goal-id example-goal --format json +loopx --registry registry.json todo list --goal-id example-goal --todo-id todo_source --limit 1 --format json +``` + +Completion retries also close a receipt/head read race: if the first receipt +lookup misses a peer commit but the head already shows completion, recheck the +matching operation receipt before interpreting a supplied validation receipt. +This returns the committed result without repeating effects; no matching +receipt still follows the existing validation and identity guards. + +Reads do not repair Markdown, mutate Todo/lease state or replay a business +operation. Missing promoted Markdown is acceptable; an unavailable provider is +not an empty Goal. No frontend configuration changes are needed: CLI, manager +Chat details and existing status/quota consumers retain their current entry +points. To reverse a business decision, use its ordinary mutation, not a read +model or restored Markdown. Code rollback retains provider state and fences; +old read policies can again misclassify these cases. + +## 中文 + +“这个 Todo 已完成”与“这一段工作已闭环”不同。结构化推进任务完成后,要有真实 +存在的后继,或明确声明 `no_followup=true`。TS 现在统一解析显式后继、替代关系和 +反向交接关系;Python 保留旧输入规范化与展示。不存在的 ID、自指 ID 不再遮住 +未闭环工作,归档与筛选也不再凭空制造后继缺口。 + +关系在完整可用源上判定,然后才筛选、分页。按状态、ID 或 Agent 筛出的列表不能 +为整个源出具闭环证明;仅限制显示条数不会改变完整源上的计数与判断。handoff 的 +状态与排除执行者信息不会在压缩展示时丢失。派生判断带相关事实摘要以防陈旧复用, +但不是授权凭据,也不写回 provider。旧 prose replan 提示仍仅用于兼容;显式 +布尔标记优先,不能靠标题里的几个词推导后继存在或授予权限。 + +归档捕获现在保留与当前工作有关的后继图和原有依赖、standing decision。延后历史 +可以被保留,但状态仍是 deferred,绝不会因此满足 `todo_done`。新的内部 v1 请求 +要求 Python 与 TS 配套升级;旧回执不被重新解释。长历史重复字段集合采用无损共享, +没有放宽 RPC 上限或丢弃历史节点。 + +本阶段关闭一组 T3/L5 读语义及其 L7 捕获依赖,不代表 D1 投影投递、D2 耐久性、 +D3 整 Goal 切换完成,也不修改默认 provider。PostgreSQL 使用相同规则,服务部署与 +资格仍独立。复杂 fixture 和只读快照演练不是长期 soak 或生产晋升许可。 diff --git a/examples/control_plane/hot-path-interface-budget-smoke.py b/examples/control_plane/hot-path-interface-budget-smoke.py index 5172693732..1910a7183a 100644 --- a/examples/control_plane/hot-path-interface-budget-smoke.py +++ b/examples/control_plane/hot-path-interface-budget-smoke.py @@ -388,6 +388,8 @@ def main() -> int: assert "presentation_surfaces" not in status_payload, status_payload status_items = status_payload["attention_queue"]["items"] assert status_items, status_payload + # Internal graph evaluations must not expand public status payloads. + assert "succession_evaluation" not in json.dumps(status_items) assert "task_graph_projection" not in status_items[0], status_items[0] route_health = status_payload["runtime_projection_routes"] assert route_health["healthy"] is True, route_health diff --git a/examples/control_plane/quota-cleared-blocker-successor-gate-smoke.py b/examples/control_plane/quota-cleared-blocker-successor-gate-smoke.py index b8039cfbc3..f3deec9222 100644 --- a/examples/control_plane/quota-cleared-blocker-successor-gate-smoke.py +++ b/examples/control_plane/quota-cleared-blocker-successor-gate-smoke.py @@ -456,12 +456,14 @@ def assert_completed_successor_keeps_old_gate_cleared() -> None: assert summary["current_agent_cleared_without_successor_handoff_count"] == 0, payload -def assert_superseded_completed_blocker_does_not_wake_agent() -> None: +def assert_resolved_supersession_does_not_wake_agent() -> None: payload = build_quota_should_run( status_payload( [ primary_owned_todo(), handoff_review(status="done", superseded_by="todo_value_successor"), + todo_item(todo_id="todo_value_successor", text="Verified replacement", + status="done", no_followup=True), ], recommended_action="Wait for main-control after superseded handoff.", ), @@ -503,7 +505,7 @@ def main() -> int: assert_archived_completed_blocker_does_not_wake_agent() assert_existing_successor_runs_normally() assert_completed_successor_keeps_old_gate_cleared() - assert_superseded_completed_blocker_does_not_wake_agent() + assert_resolved_supersession_does_not_wake_agent() assert_no_followup_completed_blocker_does_not_wake_agent() print("quota-cleared-blocker-successor-gate-smoke ok") return 0 diff --git a/loopx/control_plane/coordination/runtime_shadow.py b/loopx/control_plane/coordination/runtime_shadow.py index 647d4c0f0a..3536472b46 100644 --- a/loopx/control_plane/coordination/runtime_shadow.py +++ b/loopx/control_plane/coordination/runtime_shadow.py @@ -159,10 +159,11 @@ def capture_todo_archive_dependencies(todos: list[dict[str, Any]], state_text: s _, archived, _ = parse_todo_source(state_text) # No prose or wide diagnostics cross the selection transport budget. capture_fields = ("todo_id", "role", "task_class", "status", "done", "archive_state", "resume_when", - "decision_scope", "decision_outcome", "global_gate", "blocks_agent", "bound_agent", "goal_bound") + "decision_scope", "decision_outcome", "global_gate", "blocks_agent", "bound_agent", "goal_bound", + "successor_todo_ids", "superseded_by", "unblocks_todo_id") capture = effect_runtime_result("todo.archive.capture_dependencies", { - "schema_version": "todo_archive_dependency_capture_request_v0", - "active": [{"todo_id": item["todo_id"], "resume_when": item.get("resume_when")} for item in todos], + "schema_version": "todo_archive_dependency_capture_request_v1", + "active": [{key: item[key] for key in capture_fields if key in item} for item in todos], "archived": [{key: item[key] for key in capture_fields if key in item} for item in archived], }) if not isinstance(capture, dict) or capture.get("schema_version") != "todo_archive_dependency_capture_result_v0": diff --git a/loopx/control_plane/coordination/todo_terminal_lifecycle.ts b/loopx/control_plane/coordination/todo_terminal_lifecycle.ts index 96c18d290f..fee12daf9f 100644 --- a/loopx/control_plane/coordination/todo_terminal_lifecycle.ts +++ b/loopx/control_plane/coordination/todo_terminal_lifecycle.ts @@ -805,6 +805,13 @@ export async function executeCoordinationTodoTerminalLifecycle( todo_id: input.todo_id, }, "decision_rejection"); } + // The first receipt read can precede a peer commit while this head already + // observes it. Recover only the matching operation receipt; never discard a + // validation receipt to manufacture a terminal replay from Todo state alone. + if (input.command === "complete" && todo.status === "done" && input.validation_receipt !== null) { + const committedReplay = await terminalReceipt(input, requestSha).read(store); + if (committedReplay !== null) return committedReplay; + } if (input.expected_role !== null && todo.role !== input.expected_role) { return terminalFailure( "todo_role_mismatch", diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index afbcc3b881..6938a9ff13 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -1,4 +1,5 @@ import {evaluateUserCompletion} from "./todos/user_completion.ts"; +import {projectTodoSuccession, projectTodoClosure} from "./todos/succession.ts"; import {projectTodoSummaryLanes, projectLegacyTodoWorkCounts} from "./todos/summary_lanes.ts"; import {recordDelegationAdoption, delegationInventoryItem, delegationInventoryQuery, delegationPreflight, delegationTurnPlanDecision, selectDelegationBinding, transitionDelegationObservation} from "./collaboration/delegation.ts"; import {planChatMode} from "./collaboration/chat_mode.ts"; @@ -412,6 +413,8 @@ export function createEffectRuntimeHandlers( ["todo.public_update.plan", planPublicTodoUpdate], ["todo.standing_decision.project", evaluateStandingDecisionProjection], ["todo.summary_lanes.project", projectTodoSummaryLanes], + ["todo.succession.project", projectTodoSuccession], + ["todo.succession.closure", projectTodoClosure], ["todo.work_counts.project", projectLegacyTodoWorkCounts], ["todo.decision_scope.evaluate", evaluateDecisionScope], ["todo.user_completion.plan", evaluateUserCompletion], diff --git a/loopx/control_plane/todos/active_state_todos.py b/loopx/control_plane/todos/active_state_todos.py index a4b00bec65..26c800bd5c 100644 --- a/loopx/control_plane/todos/active_state_todos.py +++ b/loopx/control_plane/todos/active_state_todos.py @@ -9,6 +9,8 @@ read_canonical_todos_if_promoted, ) +from .succession_warning import public_todo_summary + MONITOR_WRITEBACK_CONTRACT_SCHEMA_VERSION = "monitor_writeback_contract_v0" @@ -46,7 +48,7 @@ def _redacted_status_todo_fields(fields: dict[str, Any]) -> dict[str, Any]: group = redacted.get(key) if not isinstance(group, dict): continue - group_copy = dict(group) + group_copy = public_todo_summary(group) items: list[Any] = [] for item in group_copy.get("items") or []: if not isinstance(item, dict): diff --git a/loopx/control_plane/todos/archive_capture.ts b/loopx/control_plane/todos/archive_capture.ts index 460a6910b5..8f81c4dd01 100644 --- a/loopx/control_plane/todos/archive_capture.ts +++ b/loopx/control_plane/todos/archive_capture.ts @@ -1,3 +1,4 @@ +import {indexInferredSuccessors} from "./succession.ts"; /** Capture actual archived dependency records, never cached resume conclusions. * Missing legacy roles can be reconstructed only from an explicit agent-only * task class. User decision authority always requires a recorded user role. */ @@ -8,7 +9,7 @@ import {normalizeTodoResumeWhen, TODO_RESUME_NORMALIZE_REQUEST_SCHEMA_VERSION} f import {AGENT_TODO_TASK_CLASSES as AGENT_CLASSES, USER_TODO_TASK_CLASSES as USER_CLASSES} from "./authoring_scope.ts"; import {isStandingDecisionReceipt} from "./standing_decision.ts"; -export const ARCHIVE_CAPTURE_REQUEST_SCHEMA = "todo_archive_dependency_capture_request_v0"; +export const ARCHIVE_CAPTURE_REQUEST_SCHEMA = "todo_archive_dependency_capture_request_v1"; const fail = (reason: string): never => {throw new EffectRuntimeRequestError(`archive dependency capture: ${reason}`);}; export function captureArchivedTodoDependencies(value: unknown): JsonObject { @@ -23,6 +24,15 @@ export function captureArchivedTodoDependencies(value: unknown): JsonObject { const records = byId.get(item.todo_id) ?? []; records.push({item, index}); byId.set(item.todo_id, records); } + const normalizedResume = (item: JsonObject) => normalizeTodoResumeWhen({ + schema_version: TODO_RESUME_NORMALIZE_REQUEST_SCHEMA_VERSION, resume_when: item.resume_when ?? null}); + const inferred = indexInferredSuccessors([...active, ...archived].map(item => { + const resume = normalizedResume(item); + return {id: typeof item.todo_id === "string" ? item.todo_id : null, + advancement: item.task_class === "advancement_task", + unblocks: typeof item.unblocks_todo_id === "string" ? item.unblocks_todo_id : null, + resumes: resume?.startsWith("todo_done:") ? resume.slice("todo_done:".length) : null}; + })); const activeIds = new Set(active.map(item => item.todo_id)); const selected = new Map(); const queue = [...active]; @@ -40,27 +50,35 @@ export function captureArchivedTodoDependencies(value: unknown): JsonObject { queue.push(item); } for (let cursor = 0; cursor < queue.length; cursor++) { - const token = normalizeTodoResumeWhen({schema_version: TODO_RESUME_NORMALIZE_REQUEST_SCHEMA_VERSION, - resume_when: queue[cursor]!.resume_when ?? null}); - if (!token || !(token.startsWith("todo_done:") || token.startsWith("monitor_changed:"))) continue; - const id = token.slice(token.indexOf(":") + 1), candidates = byId.get(id); - if (!candidates) continue; // A genuinely absent target remains an unsatisfied condition. - if (candidates.length !== 1 || activeIds.has(id)) fail("duplicate dependency identity"); - if (selected.has(id)) continue; - const {item, index} = candidates[0]!; - if (item.archive_state !== "archive" || item.status !== "done" || item.done !== true) { - fail("dependency is not an archived completion"); + const current = queue[cursor]!; + const token = normalizedResume(current); + const dependency = token && (token.startsWith("todo_done:") || token.startsWith("monitor_changed:")) + ? token.slice(token.indexOf(":") + 1) : null; + const links = [...new Set([dependency, current.superseded_by, + ...(Array.isArray(current.successor_todo_ids) ? current.successor_todo_ids : []), + ...(inferred.get(String(current.todo_id)) ?? [])].filter((value): value is string => typeof value === "string"))]; + for (const id of links) { + const candidates = byId.get(id); + if (!candidates) continue; // A genuinely absent target remains an unsatisfied condition. + if (candidates.length !== 1 || activeIds.has(id)) fail("duplicate dependency identity"); + const {item, index} = candidates[0]!; + // Capture records, not satisfaction. Deferred history is valid evidence + // of an unmet todo_done condition; its status must remain deferred. + if (item.archive_state !== "archive" || !["done", "deferred"].includes(String(item.status)) || item.done !== true) { + fail("dependency is not an archived terminal record"); + } + if (selected.has(id)) continue; + const taskClass = String(item.task_class ?? ""); + const role = item.role ?? (AGENT_CLASSES.has(taskClass) ? "agent" : null); + if (!((role === "agent" && AGENT_CLASSES.has(taskClass)) || + (role === "user" && USER_CLASSES.has(taskClass)))) { + fail("dependency requires a recorded role and compatible explicit task_class"); + } + // Contradictory user authority on an agent record is not repaired by inference. + if (role === "agent" && ["decision_scope", "decision_outcome", "global_gate", "blocks_agent", "bound_agent", "goal_bound"] + .some(field => item[field] != null && item[field] !== false)) fail("agent dependency carries user authority"); + selected.set(id, {index, role}); queue.push(item); } - const taskClass = String(item.task_class ?? ""); - const role = item.role ?? (AGENT_CLASSES.has(taskClass) ? "agent" : null); - if (!((role === "agent" && AGENT_CLASSES.has(taskClass)) || - (role === "user" && USER_CLASSES.has(taskClass)))) { - fail("dependency requires a recorded role and compatible explicit task_class"); - } - // Contradictory user authority on an agent record is not repaired by inference. - if (role === "agent" && ["decision_scope", "decision_outcome", "global_gate", "blocks_agent", "bound_agent", "goal_bound"] - .some(field => item[field] != null && item[field] !== false)) fail("agent dependency carries user authority"); - selected.set(id, {index, role}); queue.push(item); } return {schema_version: "todo_archive_dependency_capture_result_v0", records: [...selected.values()]}; } diff --git a/loopx/control_plane/todos/goal_todo_projection.py b/loopx/control_plane/todos/goal_todo_projection.py index fada523911..bb86b7c643 100644 --- a/loopx/control_plane/todos/goal_todo_projection.py +++ b/loopx/control_plane/todos/goal_todo_projection.py @@ -15,6 +15,7 @@ from .active_state_editing import TODO_SECTION_HEADINGS from .active_state_todo_parser import parse_active_state_todos from .list_projection import compact_explicit_limit_todo_summary +from .succession_warning import public_todo_summary from .contract import ( build_todo_id, normalize_todo_blocks_agent, @@ -97,6 +98,7 @@ def filtered_todo_summary( source_section=source_section, role=role, item_limit=item_limit, + full_selection=not (normalized_status or normalized_todo_id or normalized_agent_id), ) or empty_todo_summary(role=role) ) @@ -341,6 +343,7 @@ def todo_summaries_from_fields( role=item_role, item_limit=limit, ) + summary = public_todo_summary(summary) summaries[key] = summary todos.extend(summary.get("items") or []) uncapped_todo_count += int(summary.get("total_count") or 0) diff --git a/loopx/control_plane/todos/handoff_gate.py b/loopx/control_plane/todos/handoff_gate.py index 2d10511e8d..c7c24bcb53 100644 --- a/loopx/control_plane/todos/handoff_gate.py +++ b/loopx/control_plane/todos/handoff_gate.py @@ -5,23 +5,16 @@ from typing import Any from .contract import ( - TODO_RESUME_KIND_TODO_DONE, - TODO_STATUS_DEFERRED, TODO_STATUS_DONE, TODO_STATUS_OPEN, - TODO_TASK_CLASS_ADVANCEMENT, normalize_todo_claimed_by, normalize_todo_excluded_agents, normalize_todo_id, - normalize_todo_no_followup, - normalize_todo_resume_when, normalize_todo_status, - normalize_todo_task_class, todo_done_for_status, ) TODO_HANDOFF_GATE_SCHEMA_VERSION = "todo_handoff_gate_v1" -TODO_ARCHIVE_STATE_ACTIVE = "active" class HandoffGateState(str, Enum): @@ -48,51 +41,12 @@ def _todo_text(item: dict[str, Any]) -> str: return str(item.get("text") or "").strip() -def _todo_archive_state(item: dict[str, Any]) -> str: - value = str(item.get("archive_state") or TODO_ARCHIVE_STATE_ACTIVE).strip() - return value or TODO_ARCHIVE_STATE_ACTIVE - - -def _successor_todo_ids( - gate: dict[str, Any], - *, - items: list[dict[str, Any]], -) -> list[str]: - gate_id = normalize_todo_id(gate.get("todo_id")) - superseded_by = normalize_todo_id(gate.get("superseded_by")) - successor_ids: list[str] = [] - if superseded_by: - successor_ids.append(superseded_by) - if not gate_id: - return successor_ids - - for item in items: - if normalize_todo_task_class( - item.get("task_class"), - text=_todo_text(item), - action_kind=item.get("action_kind"), - ) != TODO_TASK_CLASS_ADVANCEMENT: - continue - candidate_id = normalize_todo_id(item.get("todo_id")) - if not candidate_id or candidate_id == gate_id: - continue - resume_when = normalize_todo_resume_when(item.get("resume_when")) or "" - resume_kind, separator, resume_target = resume_when.partition(":") - candidate_unblocks = normalize_todo_id(item.get("unblocks_todo_id")) - if ( - candidate_id == superseded_by - or candidate_unblocks == gate_id - or ( - separator - and resume_kind == TODO_RESUME_KIND_TODO_DONE - and normalize_todo_id(resume_target) == gate_id - ) - ) and candidate_id not in successor_ids: - successor_ids.append(candidate_id) - return successor_ids - - def _stale_handoff_closeout_replan_required(gate: dict[str, Any]) -> bool: + # Legacy compatibility only: old authors did not write the typed route flag. + # Do not use this prose hint for successor existence, gate state or permission. + # Retire it after the remaining route-closeout writers emit the typed flag. + if isinstance(gate.get("route_continuation_replan_required"), bool): + return gate["route_continuation_replan_required"] is True if _todo_done(gate): return False label = " ".join( @@ -103,25 +57,6 @@ def _stale_handoff_closeout_replan_required(gate: dict[str, Any]) -> bool: return "stale" in label and "handoff" in label and "closeout" in label -def _handoff_gate_state( - gate: dict[str, Any], - *, - successor_ids: list[str], -) -> HandoffGateState: - status = _todo_status(gate) - if normalize_todo_id(gate.get("superseded_by")): - return HandoffGateState.SUPERSEDED - if status == TODO_STATUS_DEFERRED: - return HandoffGateState.DEFERRED - if not _todo_done(gate): - return HandoffGateState.BLOCKING - if normalize_todo_no_followup(gate.get("no_followup")) is True: - return HandoffGateState.CLEARED_NO_FOLLOWUP - if successor_ids: - return HandoffGateState.CLEARED_WITH_SUCCESSOR - return HandoffGateState.CLEARED_WITHOUT_SUCCESSOR - - def _compact_handoff_gate( gate: dict[str, Any], *, @@ -180,37 +115,21 @@ def _compact_handoff_gate( return {key: value for key, value in payload.items() if value not in (None, "")} -def build_todo_handoff_gate_states(items: Iterable[Any]) -> list[dict[str, Any]]: - """Project dependency-linked executor exclusions into a handoff state machine.""" +def build_todo_handoff_gate_states( + items: Iterable[Any], *, evaluations: list[dict[str, Any]] | None = None, +) -> list[dict[str, Any]]: + """Render the typed full-source handoff decision; retain presentation only.""" + from .succession_warning import project_succession todo_items = [item for item in items if isinstance(item, dict)] - gates: list[dict[str, Any]] = [] - seen: set[tuple[str, str]] = set() - for item in todo_items: - if _todo_archive_state(item) != TODO_ARCHIVE_STATE_ACTIVE: - continue - excluded_agents = normalize_todo_excluded_agents(item.get("excluded_agents")) - if not excluded_agents or not normalize_todo_id(item.get("unblocks_todo_id")): - continue - identity = (str(item.get("todo_id") or ""), _todo_text(item)) - if identity in seen: - continue - seen.add(identity) - successor_ids = _successor_todo_ids(item, items=todo_items) - gates.append( - _compact_handoff_gate( - item, - state=_handoff_gate_state(item, successor_ids=successor_ids), - successor_ids=successor_ids, - ) - ) - return sorted( - gates, - key=lambda item: ( - int(item.get("index") or 999999), - str(item.get("todo_id") or ""), - ), - ) + decisions = evaluations if evaluations is not None else project_succession(todo_items) + gates = [ + _compact_handoff_gate(item, state=HandoffGateState(decision["handoff_state"]), + successor_ids=decision["successor_todo_ids"]) + for item, decision in zip(todo_items, decisions, strict=True) + if decision["handoff_state"] is not None + ] + return sorted(gates, key=lambda item: (int(item.get("index") or 999999), str(item.get("todo_id") or ""))) def todo_summary_handoff_gates(value: dict[str, Any]) -> list[dict[str, Any]]: diff --git a/loopx/control_plane/todos/list_projection.py b/loopx/control_plane/todos/list_projection.py index f630dda4bc..0c3c98fb63 100644 --- a/loopx/control_plane/todos/list_projection.py +++ b/loopx/control_plane/todos/list_projection.py @@ -48,6 +48,9 @@ } _ITEM_FIELDS = ( "schema_version", + "gate_state", + "successor_count", + "excluded_agents", "index", "todo_id", "role", diff --git a/loopx/control_plane/todos/succession.ts b/loopx/control_plane/todos/succession.ts new file mode 100644 index 0000000000..6d29b1b63a --- /dev/null +++ b/loopx/control_plane/todos/succession.ts @@ -0,0 +1,188 @@ +/** Read-only continuation evidence over the complete Todo graph. + * These evaluations are derived display state, never completion or execution authority. */ +import {canonicalAuthoritySha256} from "../coordination/authority_store_codec.ts"; +import {EffectRuntimeRequestError} from "../effect_runtime_errors.ts"; +import type {JsonObject} from "../effect_program.ts"; +import {requireBoolean, requireJsonObject, requireStringLiteral} from "../runtime_decode.ts"; + +export const HANDOFF_STATES = ["blocking", "cleared_without_successor", "cleared_with_successor", + "cleared_no_followup", "superseded", "deferred"] as const; +export type HandoffState = typeof HANDOFF_STATES[number]; +const CONTEXT_FIELDS = ["action_kind", "task_repository", "continuation_policy", "claimed_by", + "completed_at", "updated_at", "required_write_scopes", "required_capabilities", "target_capabilities", + "explore_result_node_refs", "decision_scope", "required_decision_scopes", "unblocks_todo_id", + "resume_when", "blocks_agent", "excluded_agents", "global_gate"] as const; +const EVALUATION_SCHEMA = "todo_succession_evaluation_v0"; +interface Row { + facts: JsonObject; id: string | null; status: "open" | "blocked" | "done" | "deferred"; + active: boolean; advancement: boolean; noFollowup: boolean; tracked: boolean; + successors: string[]; supersededBy: string | null; unblocks: string | null; + resumes: string | null; handoff: boolean; +} +function id(value: unknown): string | null { + if (value === null) return null; + if (typeof value !== "string" || !/^todo_[a-z0-9_-]{3,64}$/.test(value)) { + throw new EffectRuntimeRequestError("succession requires normalized Todo identities"); + } + return value; +} +function decode(value: unknown): Row { + const raw = requireJsonObject(value, "succession facts"); + const facts: JsonObject = {...raw, context_fields: Array.isArray(raw.context_fields) + ? CONTEXT_FIELDS.filter(field => (raw.context_fields as unknown[]).includes(field)) : raw.context_fields}; + const status = requireStringLiteral(facts.status, ["open", "blocked", "done", "deferred"], "status"); + const active = requireBoolean(facts.active, "active"); + const advancement = requireBoolean(facts.advancement, "advancement"); + if (!Array.isArray(facts.successors) || !Array.isArray(facts.context_fields) || + facts.context_fields.some(field => typeof field !== "string")) { + throw new EffectRuntimeRequestError("succession lists must be arrays"); + } + return {facts, id: id(facts.todo_id), status, active, advancement, + noFollowup: requireBoolean(facts.no_followup, "no_followup"), + tracked: active && status === "done" && advancement && + CONTEXT_FIELDS.some(field => (facts.context_fields as unknown[]).includes(field)), + successors: facts.successors.map(value => { + const target = id(value); + if (!target) throw new EffectRuntimeRequestError("successor identity cannot be null"); + return target; + }), supersededBy: id(facts.superseded_by), unblocks: id(facts.unblocks), + resumes: id(facts.resumes), handoff: requireBoolean(facts.handoff, "handoff")}; +} +function handoffState(row: Row, successors: readonly string[]): HandoffState | null { + if (!row.active || !row.handoff) return null; + if (row.supersededBy && successors.includes(row.supersededBy)) return "superseded"; + if (row.status === "deferred") return "deferred"; + if (row.status !== "done") return "blocking"; + if (row.noFollowup) return "cleared_no_followup"; + return successors.length ? "cleared_with_successor" : "cleared_without_successor"; +} + +/** The same edge index drives live readback and bounded archive capture. */ +export function indexInferredSuccessors(rows: readonly Pick[]): Map { + const inferred = new Map(); + for (const row of rows) { + if (!row.id || !row.advancement) continue; + for (const source of new Set([row.unblocks, row.resumes])) { + if (source && source !== row.id) { + const targets = inferred.get(source) ?? []; + targets.push(row.id); + inferred.set(source, targets); + } + } + } + return inferred; +} + +/** Index inferred edges once; both completion and handoff use the same resolver. */ +export function evaluateTodoSuccession(values: readonly unknown[]): JsonObject[] { + const rows = values.map(decode), byId = new Map(); + const generations = new Set(); + for (const row of rows) { + if (!row.id) continue; + const generation = `${row.id}:${row.active ? "active" : "archive"}`; + if (generations.has(generation)) throw new EffectRuntimeRequestError(`duplicate succession identity: ${row.id}`); + generations.add(generation); + // Legacy archive/recreate retains an older record with the same logical id. + // Only the current active record contributes edges for that identity. + if (row.active || !byId.has(row.id)) byId.set(row.id, row); + } + const inferred = indexInferredSuccessors([...byId.values()]); + return rows.map(row => { + const declared = [...new Set([...row.successors, ...(row.supersededBy ? [row.supersededBy] : [])])]; + // A retained archived target is still evidence. A missing/self target is not. + const resolved = declared.filter(target => target !== row.id && byId.has(target)); + const successors = [...new Set([...resolved, ...(row.id ? inferred.get(row.id) ?? [] : [])])]; + return {schema_version: EVALUATION_SCHEMA, item_sha256: canonicalAuthoritySha256(row.facts), + successor_todo_ids: successors, unresolved_successor_ids: declared.filter(target => !resolved.includes(target)), + tracked_completion: row.tracked, successor_gap: row.tracked && !row.noFollowup && successors.length === 0, + handoff_state: handoffState(row, successors)}; + }); +} + +/** Filtering may reuse a fresh full-source result, but may not change its item facts. */ +export function validateTodoSuccession(facts: unknown, value: unknown): JsonObject { + const row = decode(facts), evaluation = requireJsonObject(value, "succession evaluation"); + if (evaluation.schema_version !== EVALUATION_SCHEMA || evaluation.item_sha256 !== canonicalAuthoritySha256(row.facts) || + !Array.isArray(evaluation.successor_todo_ids) || !Array.isArray(evaluation.unresolved_successor_ids)) { + throw new EffectRuntimeRequestError("Todo display requires a matching full-source succession evaluation"); + } + for (const target of [...evaluation.successor_todo_ids, ...evaluation.unresolved_successor_ids]) { + if (id(target) === null) throw new EffectRuntimeRequestError("successor identity cannot be null"); + } + requireBoolean(evaluation.tracked_completion, "tracked_completion"); + requireBoolean(evaluation.successor_gap, "successor_gap"); + if (evaluation.handoff_state !== null) requireStringLiteral(evaluation.handoff_state, HANDOFF_STATES, "handoff_state"); + const successors = evaluation.successor_todo_ids as string[]; + if (evaluation.tracked_completion !== row.tracked || + evaluation.successor_gap !== (row.tracked && !row.noFollowup && successors.length === 0) || + evaluation.handoff_state !== handoffState(row, successors) || + new Set(successors).size !== successors.length || successors.includes(row.id ?? "")) { + throw new EffectRuntimeRequestError("inconsistent Todo succession evaluation"); + } + return evaluation; +} + +export function projectTodoSuccession(value: unknown): JsonObject { + const request = requireJsonObject(value, "Todo succession request"); + if (request.schema_version !== "todo_succession_request_v0" || !Array.isArray(request.rows)) { + throw new EffectRuntimeRequestError("Todo succession request schema mismatch"); + } + const contexts = request.context_field_sets; + if (contexts !== undefined && (!Array.isArray(contexts) || contexts.some(value => + !Array.isArray(value) || value.some(field => typeof field !== "string")))) { + throw new EffectRuntimeRequestError("invalid succession context field sets"); + } + const rows = request.rows.map(value => { + const row = requireJsonObject(value, "succession row"); + if (!Array.isArray(contexts)) return row; + const index = row.context_fields; + if (typeof index !== "number" || !Number.isSafeInteger(index) || index < 0 || index >= contexts.length) { + throw new EffectRuntimeRequestError("invalid succession context field index"); + } + return {...row, context_fields: contexts[index]}; + }); + const evaluations = request.evaluations; + if (evaluations !== undefined && (!Array.isArray(evaluations) || evaluations.length !== request.rows.length)) { + throw new EffectRuntimeRequestError("succession evaluation cardinality mismatch"); + } + return {schema_version: "todo_succession_result_v0", evaluations: Array.isArray(evaluations) + ? rows.map((row, index) => validateTodoSuccession(row, evaluations[index])) + : evaluateTodoSuccession(rows)}; +} + +/** Summary proofs are derived from every selected row before display caps. + * A query subset can describe items but cannot certify closure of the source. */ +export function projectTodoClosure(value: unknown): JsonObject { + const request = requireJsonObject(value, "Todo closure request"); + if (request.schema_version !== "todo_closure_request_v0" || !Array.isArray(request.rows)) { + throw new EffectRuntimeRequestError("Todo closure request schema mismatch"); + } + const source = typeof request.source_section === "string" ? request.source_section : ""; + const role = request.role; + const valid = (role === "user" || role === "agent") && source.trim() !== "" && + requireBoolean(request.full_selection, "full_selection"); + const rows = request.rows.map(value => { + const row = requireJsonObject(value, "closure row"); + return {status: requireStringLiteral(row.status, ["open", "blocked", "done", "deferred"], "status"), + watch: requireBoolean(row.watch_only, "watch_only"), noFollowup: requireBoolean(row.no_followup, "no_followup"), + gap: requireBoolean(row.successor_gap, "successor_gap"), replan: requireBoolean(row.replan, "replan"), + handoff: row.handoff_state === null ? null : requireStringLiteral(row.handoff_state, HANDOFF_STATES, "handoff_state")}; + }); + const noFollowup = rows.filter(row => (row.status === "done" || row.status === "deferred") && row.noFollowup).length; + const watches = rows.filter(row => row.status !== "done" && row.status !== "deferred" && row.watch).length; + const convergent = rows.filter(row => row.status !== "done" && row.status !== "deferred" && !row.watch).length; + const result: JsonObject = {}; + if (valid && convergent === 0 && rows.every(row => row.status !== "deferred")) { + result.source_proof = {schema_version: "todo_source_proof_v0", role, item_count: rows.length, derived: true}; + } + if (valid && rows.every(row => (row.status === "done" || row.watch) && row.status !== "deferred" && + !row.gap && !row.replan && row.handoff !== "cleared_without_successor" && row.handoff !== "blocking")) { + result.terminal_closure_proof = {schema_version: "todo_terminal_closure_proof_v0", role, + source_section: source, item_count: rows.length, all_todos_done: watches === 0, + monitor_open_count: watches, successor_gap_count: 0, route_replan_count: 0, + no_followup_count: noFollowup, derived: true, + ...(watches ? {all_convergent_todos_done: true, watch_only_monitor_count: watches} : {})}; + } + if (noFollowup) result.closure_intent = {schema_version: "todo_closure_intent_v0", kind: "no_followup", derived: true, count: noFollowup}; + return result; +} diff --git a/loopx/control_plane/todos/succession_warning.py b/loopx/control_plane/todos/succession_warning.py index 3c76d62cb0..cc0dd14380 100644 --- a/loopx/control_plane/todos/succession_warning.py +++ b/loopx/control_plane/todos/succession_warning.py @@ -1,13 +1,20 @@ +"""Succession read-policy adapter and existing warning presentation.""" from __future__ import annotations from typing import Any +from ..effect_runtime import effect_runtime_result from ..agents.agent_scope import agent_scope_item_claimed_by_agent_or_unclaimed from .compact_projection import compact_todo_projection_item from .contract import ( TODO_STATUS_OPEN, normalize_todo_id, normalize_todo_id_list, + normalize_todo_status, + normalize_todo_task_class, + normalize_todo_no_followup, + normalize_todo_resume_when, + normalize_todo_excluded_agents, ) @@ -130,3 +137,88 @@ def todo_succession_gap_items( for item in items if agent_scope_item_claimed_by_agent_or_unclaimed(item, agent_id=agent_id) ] + + +class _EvaluatedTodo(dict[str, Any]): + """Ephemeral graph evidence stays outside the public row/JSON contract.""" + + def __init__(self, item: dict[str, Any], evaluation: dict[str, Any]) -> None: + super().__init__(item) + self.succession_evaluation = evaluation + + +def succession_facts(item: dict[str, Any]) -> dict[str, Any]: + resume = normalize_todo_resume_when(item.get("resume_when")) or "" + return { + "todo_id": normalize_todo_id(item.get("todo_id")), + "status": normalize_todo_status(item.get("status")) or ("done" if item.get("done") else "open"), + "active": (item.get("archive_state") or "active") == "active", + "advancement": normalize_todo_task_class(item.get("task_class"), text=str(item.get("text") or ""), + action_kind=item.get("action_kind")) == "advancement_task", + "no_followup": normalize_todo_no_followup(item.get("no_followup")) is True, + "successors": normalize_todo_id_list(item.get("successor_todo_ids")), + "superseded_by": normalize_todo_id(item.get("superseded_by")), + "unblocks": normalize_todo_id(item.get("unblocks_todo_id")), + "resumes": normalize_todo_id(resume.partition(":")[2]) if resume.startswith("todo_done:") else None, + "handoff": bool(normalize_todo_excluded_agents(item.get("excluded_agents"))) + and bool(normalize_todo_id(item.get("unblocks_todo_id"))), + "context_fields": sorted(key for key, value in item.items() + if value is not None and key != "succession_evaluation"), + } + + +def project_succession(items: list[dict[str, Any]], *, reuse: bool = False) -> list[dict[str, Any]]: + rows = [succession_facts(item) for item in items] + # Archived histories repeat field-presence shapes thousands of times. Intern + # those shapes, preserving the complete graph inside the existing RPC budget. + contexts: list[list[str]] = [] + context_ids: dict[tuple[str, ...], int] = {} + for row in rows: + shape = tuple(row["context_fields"]) + if shape not in context_ids: + context_ids[shape] = len(contexts) + contexts.append(list(shape)) + row["context_fields"] = context_ids[shape] + request: dict[str, Any] = {"schema_version": "todo_succession_request_v0", + "rows": rows, "context_field_sets": contexts} + if reuse: + request["evaluations"] = [item.succession_evaluation if isinstance(item, _EvaluatedTodo) else None for item in items] + result = effect_runtime_result("todo.succession.project", request) + if not isinstance(result, dict) or result.get("schema_version") != "todo_succession_result_v0": + raise ValueError("invalid typed Todo succession result") + evaluations = result.get("evaluations") + if not isinstance(evaluations, list) or len(evaluations) != len(items) or any(not isinstance(row, dict) for row in evaluations): + raise ValueError("invalid typed Todo succession cardinality") + return evaluations + + +def evaluate_succession(items: list[dict[str, Any]], lineage: list[dict[str, Any]] | None = None) -> list[dict[str, Any]]: + # Active evaluated rows override their unevaluated source versions; retained + # history and other roles remain graph evidence but never become active rows. + if not items: + return [] + # Replace one matching source row, not every occurrence of its identity. + # Retained generations and conflicting authorities stay visible to TS. + selected = {(normalize_todo_id(item.get("todo_id")), item.get("role"), item.get("archive_state") or "active") + for item in items} + source = [] + for item in lineage or []: + key = (normalize_todo_id(item.get("todo_id")), item.get("role"), item.get("archive_state") or "active") + if key in selected: + selected.remove(key) + else: + source.append(item) + source.extend(items) + evaluations = project_succession(source) + selected_evaluations = evaluations[len(source) - len(items):] + items[:] = [_EvaluatedTodo(item, evaluation) + for item, evaluation in zip(items, selected_evaluations, strict=True)] + return selected_evaluations + + +def public_todo_summary(summary: dict[str, Any]) -> dict[str, Any]: + """Drop the internal full-graph handoff once a consumer has selected rows.""" + return {**summary, "items": [ + dict(item) if isinstance(item, dict) else item + for item in summary.get("items") or [] + ]} diff --git a/loopx/control_plane/todos/todo_summary.py b/loopx/control_plane/todos/todo_summary.py index 168a7c3e20..bb3ef318f4 100644 --- a/loopx/control_plane/todos/todo_summary.py +++ b/loopx/control_plane/todos/todo_summary.py @@ -7,7 +7,6 @@ from ..goals.goal_vision_wait_projection import attach_active_vision_waits from .contract import ( - TODO_RESUME_KIND_TODO_DONE, TODO_STATUS_DONE, TODO_STATUS_OPEN, TODO_TASK_CLASS_ADVANCEMENT, @@ -32,7 +31,6 @@ normalize_todo_generation, normalize_todo_goal_bound, normalize_todo_id, - normalize_todo_id_list, normalize_todo_no_followup, normalize_removed_todo_continuation_policy, normalize_todo_required_decision_scopes, @@ -88,9 +86,6 @@ MAX_COMPLETED_SUCCESSION_WARNING_ITEMS = 5 MAX_RECENT_COMPLETED_ADVANCEMENT_ITEMS = MAX_TODO_VISIBILITY_LANE_ITEMS -TODO_SOURCE_PROOF_SCHEMA_VERSION = "todo_source_proof_v0" -TODO_CLOSURE_INTENT_SCHEMA_VERSION = "todo_closure_intent_v0" -TODO_TERMINAL_CLOSURE_PROOF_SCHEMA_VERSION = "todo_terminal_closure_proof_v0" TASK_ORCHESTRATION_AUTHORITY_SCHEMA_VERSION = "task_orchestration_authority_v0" TODO_ARCHIVE_STATE_ACTIVE = "active" AttentionItemBuilder = Callable[..., dict[str, Any]] @@ -469,8 +464,10 @@ def canonical_todo_read_record( ) -> dict[str, Any]: """Copy one already-normalized Todo consumer record without re-deriving it.""" + # Read-policy evaluations are recomputed from a complete source. They must + # never enter shadow capture, provider records or the durable source digest. record = canonical_record_fields( - item, + {key: value for key, value in item.items() if key != "succession_evaluation"}, fields=TODO_CANONICAL_READ_RECORD_FIELDS, required_fields=TODO_CANONICAL_REQUIRED_READ_FIELDS, label="canonical Todo read record", @@ -809,77 +806,19 @@ def active_next_action_todo_ids(value: Any) -> set[str]: return todo_ids -def _normalized_todo_id_list(value: Any) -> list[str]: - return normalize_todo_id_list(value) +def todo_successor_todo_ids(item: dict[str, Any], *, items: list[dict[str, Any]]) -> list[str]: + """Compatibility call site for repair-delta; the typed graph owns links.""" + from .succession_warning import evaluate_succession - -def todo_successor_todo_ids( - item: dict[str, Any], - *, - items: list[dict[str, Any]], -) -> list[str]: - successor_ids = _normalized_todo_id_list(item.get("successor_todo_ids")) - superseded_by = normalize_todo_id(item.get("superseded_by")) - if superseded_by and superseded_by not in successor_ids: - successor_ids.append(superseded_by) - - source_todo_id = normalize_todo_id(item.get("todo_id")) - if not source_todo_id: - return successor_ids - - for candidate in items: - if not isinstance(candidate, dict): - continue - candidate_id = normalize_todo_id(candidate.get("todo_id")) - if not candidate_id or candidate_id == source_todo_id: - continue - if todo_item_task_class(candidate) != TODO_TASK_CLASS_ADVANCEMENT: - continue - resume_when = normalize_todo_resume_when(candidate.get("resume_when")) or "" - resume_kind, separator, resume_target = resume_when.partition(":") - candidate_unblocks = normalize_todo_id(candidate.get("unblocks_todo_id")) - if candidate_unblocks != source_todo_id and not ( - separator - and resume_kind == TODO_RESUME_KIND_TODO_DONE - and normalize_todo_id(resume_target) == source_todo_id - ): - continue - if candidate_id not in successor_ids: - successor_ids.append(candidate_id) - return successor_ids + selected = dict(item) + evaluation = evaluate_succession([selected], items)[0] + return list(evaluation["successor_todo_ids"]) def todo_item_is_succession_tracked_completion(item: dict[str, Any]) -> bool: - if todo_archive_state(item) != TODO_ARCHIVE_STATE_ACTIVE: - return False - if not item.get("done"): - return False - if todo_item_is_deferred(item): - return False - if todo_item_task_class(item) != TODO_TASK_CLASS_ADVANCEMENT: - return False - return any( - item.get(key) is not None - for key in ( - "action_kind", - "task_repository", - "continuation_policy", - "claimed_by", - "completed_at", - "updated_at", - "required_write_scopes", - "required_capabilities", - "target_capabilities", - "explore_result_node_refs", - "decision_scope", - "required_decision_scopes", - "unblocks_todo_id", - "resume_when", - "blocks_agent", - "excluded_agents", - "global_gate", - ) - ) + from .succession_warning import project_succession + + return project_succession([item])[0]["tracked_completion"] is True def _completed_succession_sort_key(item: dict[str, Any]) -> tuple[str, int]: @@ -893,25 +832,17 @@ def _completed_succession_sort_key(item: dict[str, Any]) -> tuple[str, int]: def completed_without_successor_items( - done_items: list[dict[str, Any]], - *, - all_items: list[dict[str, Any]], + items: list[dict[str, Any]], *, evaluations: list[dict[str, Any]], ) -> list[dict[str, Any]]: - gap_items: list[dict[str, Any]] = [] - for item in done_items: - if not todo_item_is_succession_tracked_completion(item): - continue - if normalize_todo_no_followup(item.get("no_followup")) is True: - continue - if todo_successor_todo_ids(item, items=all_items): + gap_items = [] + for item, evaluation in zip(items, evaluations, strict=True): + if not evaluation["successor_gap"]: continue compact = compact_todo_item(item) for key in ("note", "evidence", "reason"): compact.pop(key, None) compact["succession_tracked"] = True - compact["recommended_action"] = ( - "record no_followup=true or add/link a successor todo" - ) + compact["recommended_action"] = "record no_followup=true or add/link a successor todo" gap_items.append(compact) return sorted(gap_items, key=_completed_succession_sort_key, reverse=True) @@ -1034,6 +965,9 @@ def compact_todo_group( available_capabilities=available_capabilities, evaluated_at=evaluated_at, ) + from .succession_warning import evaluate_succession + + evaluate_succession(items, resume_source_items) return compact_evaluated_todo_group( items, source_section=source_section, role=role, include_empty_source=include_empty_source, preferred_todo_ids=preferred_todo_ids, @@ -1054,6 +988,7 @@ def compact_evaluated_todo_group( include_task_orchestration_authority: bool = False, vision_runs: list[dict[str, Any]] | None = None, lineage_items: list[dict[str, Any]] | None = None, + full_selection: bool = True, ) -> dict[str, Any] | None: """Filter/display an already evaluated snapshot, never re-evaluate topology. @@ -1064,17 +999,10 @@ def compact_evaluated_todo_group( return None projected = _project_summary_lanes(items, preferred_todo_ids) lanes = _TodoGroupLanes(**projected["lanes"]) - source_valid = role in {"user", "agent"} and bool(str(source_section or "").strip()) - no_followup_items = [ - item - for item in items - if todo_done_for_status(item.get("status")) - and normalize_todo_no_followup(item.get("no_followup")) is True - ] - successor_gap_items = completed_without_successor_items( - lanes.done_items, - all_items=items, - ) + from .succession_warning import project_succession + + succession = project_succession(items, reuse=True) + successor_gap_items = completed_without_successor_items(items, evaluations=succession) recent_completed_advancement_items = [ compact_todo_item(item) for item in sorted( @@ -1091,11 +1019,7 @@ def compact_evaluated_todo_group( for item in recent_completed_advancement_items: for key in ("note", "evidence", "reason"): item.pop(key, None) - handoff_gates = build_todo_handoff_gate_states(items) - route_replan_required = any( - item.get("route_continuation_replan_required") is True - for item in [*items, *handoff_gates] - ) + handoff_gates = build_todo_handoff_gate_states(items, evaluations=succession) watch_only_monitor_items = [ item for item in lanes.monitor_items @@ -1208,46 +1132,23 @@ def compact_evaluated_todo_group( summary["blocker_items"] = [ compact_todo_item(item) for item in lanes.blocker_items ] - if not convergent_open_items and not lanes.deferred_items: - summary["source_proof"] = { - "schema_version": TODO_SOURCE_PROOF_SCHEMA_VERSION, - "role": role, - "item_count": len(items), - "derived": source_valid, - } - if ( - source_valid - and len(lanes.done_items) + len(watch_only_monitor_items) == len(items) - and not convergent_open_items - and not successor_gap_items - and not route_replan_required - ): - summary["terminal_closure_proof"] = { - "schema_version": TODO_TERMINAL_CLOSURE_PROOF_SCHEMA_VERSION, - "role": role, - "source_section": source_section, - "item_count": len(items), - "all_todos_done": not watch_only_monitor_items, - "monitor_open_count": len(watch_only_monitor_items), - "successor_gap_count": 0, - "route_replan_count": 0, - "no_followup_count": len(no_followup_items), - "derived": True, - } - if watch_only_monitor_items: - summary["terminal_closure_proof"].update( - { - "all_convergent_todos_done": True, - "watch_only_monitor_count": len(watch_only_monitor_items), - } - ) - if no_followup_items: - summary["closure_intent"] = { - "schema_version": TODO_CLOSURE_INTENT_SCHEMA_VERSION, - "kind": "no_followup", - "derived": True, - "count": len(no_followup_items), - } + from ..effect_runtime import effect_runtime_result + + replan_gates = {gate.get("todo_id") for gate in handoff_gates + if gate.get("route_continuation_replan_required") is True} + closure = effect_runtime_result("todo.succession.closure", { + "schema_version": "todo_closure_request_v0", "role": role, + "source_section": source_section, "full_selection": full_selection, + "rows": [{"status": item.get("status") or ("done" if item.get("done") else "open"), + "watch_only": projection_todo_item_is_watch_only_monitor(item), + "no_followup": normalize_todo_no_followup(item.get("no_followup")) is True, + "successor_gap": evaluation["successor_gap"], "handoff_state": evaluation["handoff_state"], + "replan": item.get("route_continuation_replan_required") is True or item.get("todo_id") in replan_gates} + for item, evaluation in zip(items, succession, strict=True)], + }) + if not isinstance(closure, dict): + raise ValueError("invalid typed Todo closure projection") + summary.update(closure) if lanes.resume_blocked_items: summary["resume_blocked_count"] = len(lanes.resume_blocked_items) summary["resume_blocked_items"] = [ diff --git a/tests/control_plane/test_succession_provider_readback.py b/tests/control_plane/test_succession_provider_readback.py new file mode 100644 index 0000000000..4636f25050 --- /dev/null +++ b/tests/control_plane/test_succession_provider_readback.py @@ -0,0 +1,71 @@ +"""Public consumers use full relationship evidence without writing the provider.""" +from __future__ import annotations + +import json +import subprocess +import sys +from pathlib import Path + +import pytest + +from canonical_authority_fixture import initialize_canonical_authority, isolate_sqlite_runtime +from loopx.chat_manager_details import read_manager_goal_details +from loopx.control_plane.coordination.local_authority import read_canonical_todos_if_promoted +from loopx.control_plane.todos.machine_section_projection import render_canonical_todo_sections +from loopx.control_plane.coordination.runtime_shadow import build_todo_runtime_shadow_projection +from loopx.todos import list_goal_todos + + +def fixture_projection(): + script = """import {productionScaleSuccessionFixture,PRODUCTION_SCALE_VALIDATION_DECLARATION} from './tests/control_plane_ts/production_scale_coordination_fixture.ts'; +process.stdout.write(JSON.stringify({...productionScaleSuccessionFixture('goal-a','legacy'), validation:PRODUCTION_SCALE_VALIDATION_DECLARATION}));""" + process = subprocess.run(['node', '--no-warnings', '--experimental-strip-types', '--input-type=module', '-e', script], + capture_output=True, text=True, check=True, timeout=30) + return json.loads(process.stdout) + + +def cli(registry: Path, *args: str): + child = subprocess.run([sys.executable, '-m', 'loopx.cli', '--registry', str(registry), '--format', 'json', + 'todo', 'list', '--goal-id', 'goal-a', *args], capture_output=True, text=True, timeout=60) + assert child.returncode == 0, child.stdout or child.stderr + return json.loads(child.stdout) + + +@pytest.mark.parametrize('provider', ['legacy', 'file', 'sqlite']) +def test_real_cli_graph_selection_history_and_read_only_manager(tmp_path, monkeypatch, provider): + isolate_sqlite_runtime(tmp_path, monkeypatch) + fixture = fixture_projection() + projection, cases = fixture['projection'], fixture['cases'] + runtime = tmp_path / 'runtime' + state = tmp_path / 'state.md' + state.write_text(render_canonical_todo_sections('# Goal\n\nIndependent narrative.\n\n## Agent Todo\n', + projection['todos'], provider_revision='fixture-source', + private_validation_declarations={row['todo_id']: fixture['validation'] for row in projection['todos'] + if row.get('completion_validation_required') is True}).markdown) + registry = tmp_path / 'registry.json' + registry.write_text(json.dumps({'common_runtime_root': str(runtime), 'goals': [{ + 'id': 'goal-a', 'repo': str(tmp_path), 'state_file': state.name, 'status': 'active', + 'coordination': {'registered_agents': ['agent-a', 'agent-b'], 'handoff_mode': 'soft_claim'}, + }]})) + if provider != 'legacy': + initialize_canonical_authority(runtime, 'goal-a', projection, state_path=state, provider=provider) + state.unlink() # Promoted readback must not rely on or regenerate Markdown. + state_before = state.read_bytes() if state.exists() else None + before = read_canonical_todos_if_promoted(runtime_root=runtime, goal_id='goal-a') + for name, gap in [('inferred_source', 0), ('missing_source', 1), ('self_source', 1), ('handoff_source', 0)]: + result = cli(registry, '--todo-id', cases[name], '--limit', '1') + summary = result['agent_todos'] + assert summary.get('completed_without_successor_count', 0) == gap + assert 'terminal_closure_proof' not in summary + if name == 'handoff_source': + assert summary['handoff_gates'][0]['gate_state'] == 'cleared_with_successor' + details = read_manager_goal_details(registry, runtime, 'goal-a', owner_scope=True, limit=3) + assert details['status'] == 'read' + assert details['coverage']['active'] > details['coverage']['included'] + # A read-model evaluation must not become another persisted fact on capture. + rows = list_goal_todos(registry_path=registry, goal_id='goal-a')['todos'] + assert all("succession_evaluation" not in row for row in rows) + captured = build_todo_runtime_shadow_projection(goal_id='goal-a', todos=rows, handoff_mode='soft_claim') + assert all('succession_evaluation' not in row for row in captured['todos']) + assert read_canonical_todos_if_promoted(runtime_root=runtime, goal_id='goal-a') == before + assert (state.read_bytes() if state.exists() else None) == state_before diff --git a/tests/control_plane/test_todo_succession_read_model.py b/tests/control_plane/test_todo_succession_read_model.py new file mode 100644 index 0000000000..a7d32a5a83 --- /dev/null +++ b/tests/control_plane/test_todo_succession_read_model.py @@ -0,0 +1,128 @@ +"""Closure is a property of the source graph, not the selected display rows.""" +from __future__ import annotations + +import pytest + +from loopx.control_plane.todos.todo_summary import compact_todo_group +from loopx.control_plane.todos.goal_todo_projection import filtered_todo_summary +from loopx.control_plane.todos.handoff_gate import build_todo_handoff_gate_states + + +def work(todo_id, **fields): + return {"todo_id": todo_id, "text": "Deliver a verified result", "role": "agent", + "status": "done", "done": True, "archive_state": "active", + "task_class": "advancement_task", "claimed_by": "agent-a", **fields} + + +def summary(items, **options): + return compact_todo_group(items, source_section="Agent Todo", role="agent", item_limit=None, **options) + + +@pytest.mark.parametrize("successors", [["todo_missing"], ["todo_source"]]) +def test_unresolved_or_self_reference_does_not_certify_closure(successors): + result = summary([work("todo_source", successor_todo_ids=successors)]) + assert result.get("completed_without_successor_count") == 1 + assert "terminal_closure_proof" not in result + + +def test_filtered_done_source_keeps_its_inferred_open_successor(): + result = summary([work("todo_source"), work("todo_next", status="open", done=False, + resume_when="todo_done:todo_source")]) + assert not result.get("completed_without_successor_count") + selected = filtered_todo_summary(result, role="agent", todo_id="todo_source") + assert not selected.get("completed_without_successor_count") + + +def test_archived_successor_remains_relationship_evidence(): + source = work("todo_source") + archived = work("todo_next", archive_state="archive", resume_when="todo_done:todo_source", no_followup=True) + result = summary([source], resume_source_items=[source, archived]) + assert not result.get("completed_without_successor_count") + + +def test_explicit_successor_is_shared_by_handoff_and_completion(): + gate = work("todo_gate", excluded_agents=["agent-b"], unblocks_todo_id="todo_work", + successor_todo_ids=["todo_next"]) + next_item = work("todo_next", status="open", done=False) + result = build_todo_handoff_gate_states([gate, next_item]) + assert result[0]["gate_state"] == "cleared_with_successor" + assert result[0]["successor_todo_ids"] == ["todo_next"] + + +def test_missing_supersession_does_not_clear_blocking_handoff(): + gate = work("todo_gate", status="open", done=False, excluded_agents=["agent-b"], + unblocks_todo_id="todo_work", superseded_by="todo_missing") + assert build_todo_handoff_gate_states([gate])[0]["gate_state"] == "blocking" + + +def test_query_subset_does_not_certify_goal_closure(): + result = summary([work("todo_closed", no_followup=True), work("todo_open", status="open", done=False)]) + selected = filtered_todo_summary(result, role="agent", status="done") + assert selected["total_count"] == 1 + assert "terminal_closure_proof" not in selected + assert "source_proof" not in selected + + +def test_selection_after_warning_display_cap_keeps_exact_gap(): + records = [work(f"todo_work_{index:03}", index=index) for index in range(80)] + result = summary(records) + assert result["completed_without_successor_count"] == 80 + assert len(result["completed_without_successor_items"]) < 80 + selected = filtered_todo_summary(result, role="agent", todo_id="todo_work_000") + assert selected["completed_without_successor_count"] == 1 + + +def test_changed_item_cannot_reuse_old_graph_evaluation(): + result = summary([work("todo_closed", no_followup=True)]) + result["items"][0]["no_followup"] = False + with pytest.raises(Exception, match="matching full-source"): + filtered_todo_summary(result, role="agent", todo_id="todo_closed") + + +def test_explicit_route_flag_is_not_overridden_by_legacy_prose_hint(): + gate = work("todo_gate", status="open", done=False, text="stale handoff closeout", + excluded_agents=["agent-b"], unblocks_todo_id="todo_work", route_continuation_replan_required=False) + result = build_todo_handoff_gate_states([gate])[0] + assert result["route_continuation_replan_required"] is False + assert result["gate_state"] == "blocking" + + +def test_archived_identity_can_be_recreated_but_duplicate_active_authority_rejects(): + source = work("todo_source", no_followup=True) + result = summary([source], resume_source_items=[source, {**source, "archive_state": "archive"}]) + assert result["terminal_closure_proof"]["all_todos_done"] is True + with pytest.raises(Exception, match="duplicate succession identity"): + summary([source], resume_source_items=[source, dict(source)]) + + +def test_public_summary_drops_internal_evaluation_without_mutating_source(): + from loopx.control_plane.todos.succession_warning import public_todo_summary + + source = summary([work("todo_source", no_followup=True)]) + public = public_todo_summary(source) + assert "succession_evaluation" not in public["items"][0] + assert "succession_evaluation" not in source["items"][0] + assert filtered_todo_summary(source, role="agent", todo_id="todo_source")["total_count"] == 1 + assert public["terminal_closure_proof"] == source["terminal_closure_proof"] + + +def test_public_parser_rows_do_not_carry_internal_evaluations(): + from loopx.control_plane.todos.active_state_todo_parser import parse_active_state_todos + + state = "## Agent Todo\n\n- [ ] Work\n \n" + fields = parse_active_state_todos(state, item_limit=None) + assert fields["agent_todos"]["items"] + assert all("succession_evaluation" not in row for row in fields["agent_todos"]["items"]) + + +def test_evaluation_does_not_mutate_or_expand_input_rows(): + import json + from loopx.control_plane.todos.succession_warning import evaluate_succession + + original = work("todo_source", no_followup=True) + rows = [original] + before = json.dumps(rows, sort_keys=True) + evaluate_succession(rows) + assert rows[0] is not original + assert json.dumps(rows, sort_keys=True) == before + assert "succession_evaluation" not in original diff --git a/tests/control_plane_ts/archive_capture.test.ts b/tests/control_plane_ts/archive_capture.test.ts index de5bdd72ce..d0e19f1d3e 100644 --- a/tests/control_plane_ts/archive_capture.test.ts +++ b/tests/control_plane_ts/archive_capture.test.ts @@ -44,3 +44,37 @@ test("duplicate identities cannot be selected by storage order; archived cycles assert.deepEqual(capture(request([{...dependency, resume_when: "todo_done:todo_history"}])).records, [{index: 0, role: "agent"}]); }); + +test("capture retains explicit and inferred archived continuations through the same edge index", () => { + const active = [{todo_id: "todo_active", successor_todo_ids: ["todo_explicit"], superseded_by: "todo_replaced"}]; + const archived = [{...dependency, todo_id: "todo_explicit"}, {...dependency, todo_id: "todo_replaced"}, + {...dependency, todo_id: "todo_inferred", unblocks_todo_id: "todo_active"}, + {...dependency, todo_id: "todo_transitive", resume_when: "todo_done:todo_inferred"}, + {...dependency, todo_id: "todo_unrelated"}]; + assert.deepEqual(capture(request(archived, active)).records, + [{index: 1, role: "agent"}, {index: 0, role: "agent"}, {index: 2, role: "agent"}, {index: 3, role: "agent"}]); +}); + +test("deferred successor lineage does not become a completed resume prerequisite", () => { + const deferred = {...dependency, status: "deferred", done: true}; + assert.deepEqual(capture(request([deferred], [{todo_id: "todo_active", successor_todo_ids: [dependency.todo_id]}])).records, + [{index: 0, role: "agent"}]); + assert.deepEqual(capture(request([deferred])).records, [{index: 0, role: "agent"}]); + assert.equal(deferred.status, "deferred"); +}); + +test("retained deferred dependencies remain unsatisfied after capture", async () => { + const {evaluateTodoResumeConditions, TODO_RESUME_EVALUATION_REQUEST_SCHEMA_VERSION} = + await import("../../loopx/control_plane/todos/resume_condition.ts"); + const historical = {...dependency, role: "agent", status: "deferred", done: true}; + const selection = capture(request([historical])); + assert.deepEqual(selection.records, [{index: 0, role: "agent"}]); + const result = evaluateTodoResumeConditions({schema_version: TODO_RESUME_EVALUATION_REQUEST_SCHEMA_VERSION, + items: [{todo_id: "todo_active", role: "agent", status: "open", task_class: "advancement_task", resume_when: "todo_done:todo_history"}], + source_items: [historical], rollout_events: [], evaluated_at: "2026-09-20T00:00:00Z"}); + assert.equal(((result.conditions as Record[])[0].condition as Record).satisfied, false); +}); + +test("old capture request cannot silently omit the expanded relationship contract", () => { + assert.throws(() => capture({...request([]), schema_version: "todo_archive_dependency_capture_request_v0"}), /invalid request/); +}); diff --git a/tests/control_plane_ts/authority_store_conformance.ts b/tests/control_plane_ts/authority_store_conformance.ts index 3726b4b72f..a70be10f21 100644 --- a/tests/control_plane_ts/authority_store_conformance.ts +++ b/tests/control_plane_ts/authority_store_conformance.ts @@ -1,4 +1,5 @@ import {registerUserCompletionFollowthroughConformance} from "./user_completion_followthrough_conformance.ts"; +import {registerSuccessionReadConformance} from "./succession_read_conformance.ts"; import {registerUserCompletionUpdateConformance} from "./user_completion_update_conformance.ts"; import {registerLeaseAcquisitionConformance} from "./lease_acquisition_conformance.ts"; import {registerClaimTransferConformance} from "./claim_transfer_conformance.ts"; @@ -221,6 +222,7 @@ export function registerAuthorityStoreConformance( registerLeaseAcquisitionConformance(providerName, factory); registerAuthorityScanConformance(providerName, factory); registerOwnershipObservationConformance(providerName, factory); + registerSuccessionReadConformance(providerName, factory); registerNativePlanningUpdateConformance(providerName, factory); registerUserCompletionUpdateConformance(providerName, factory); registerUserCompletionFollowthroughConformance(providerName, factory); @@ -599,6 +601,26 @@ export function registerAuthorityStoreConformance( ["applied", "recovered", "replayed", "conflict"].includes(status)), JSON.stringify([first, second]), ); + // Pin a receipt lookup before a peer commit and a head read after it. + let receiptMiss = true; + const crossedRead = new Proxy(contender, {get(target, property) { + if (property === "readReceipt") return async (operationId: string) => { + if (receiptMiss) {receiptMiss = false; return {status: "missing"};} + return target.readReceipt(operationId); + }; + const value = Reflect.get(target, property, target); + return typeof value === "function" ? value.bind(target) : value; + }}); + const committedHead = await store.loadAuthority(); + const crossed = await executeCoordinationTodoTerminalLifecycle(crossedRead, commitRequest); + assert.equal(crossed.status, "replayed", JSON.stringify(crossed)); + assert.equal(crossed.changed, false); + const noReceipt = await executeCoordinationTodoTerminalLifecycle(contender, + {...commitRequest, operation_id: "unknown-terminal-operation"}); + assert.equal(noReceipt.status, "failed"); + assert.equal(noReceipt.reason_code, "invalid_todo_completion_transaction"); + assert.deepEqual(await store.loadAuthority(), committedHead); + const committedRevisions = [first, second] .filter((item) => item.status !== "conflict") .map((item) => item.provider_revision); diff --git a/tests/control_plane_ts/production_scale_coordination_fixture.ts b/tests/control_plane_ts/production_scale_coordination_fixture.ts index 830ac9e17c..08dee79258 100644 --- a/tests/control_plane_ts/production_scale_coordination_fixture.ts +++ b/tests/control_plane_ts/production_scale_coordination_fixture.ts @@ -577,3 +577,36 @@ export function productionScaleUserCompletionFixture(goalId: string, schema, {handoff_mode: "hard_lease"}), source: String(source.todo_id), target: String(target.todo_id), scope, registered_agents: fixture.registered_agents}; } + +/** Complete graph evidence must survive archive, selection and display limits. */ +export function productionScaleSuccessionFixture(goalId: string, schema: AuthorityProjectionSchema = "native") { + const fixture = productionScaleCoordinationFixture(goalId, schema); + const projection = structuredClone(fixture.projection); + const todos = projection.todos as Record[]; + const cases = envelope.semantic_cases.succession as Record; + const base = todos.find(todo => todo.role === "agent")!; + const source = (key: string, fields: Record = {}): Record => ({...base, + todo_id: cases[key], text: "Verify the continuation relationship", index: todos.length + Object.keys(cases).indexOf(key) + 1, + status: "done", done: true, archive_state: "active", claimed_by: "agent-a", task_class: "advancement_task", + successor_todo_ids: [], superseded_by: null, resume_when: null, unblocks_todo_id: null, + no_followup: false, excluded_agents: [], ...fields}); + const added = [source("inferred_source"), source("archived_target", {archive_state: "archive", + resume_when: `todo_done:${cases.inferred_source}`, no_followup: true}), + source("missing_source", {successor_todo_ids: ["todo_missing_continuation"]}), + source("self_source", {successor_todo_ids: [cases.self_source]}), + source("handoff_source", {excluded_agents: ["agent-b"], unblocks_todo_id: cases.inferred_source, + successor_todo_ids: [cases.explicit_target]}), + source("explicit_target", {status: "open", done: false}), + source("closed_source", {no_followup: true})]; + for (const record of added) { + Reflect.deleteProperty(record, "material_change_generation"); + if (schema === "native") {Reflect.deleteProperty(record, "index"); Reflect.deleteProperty(record, "source_section");} + else if (record.archive_state === "archive") record.source_section = "Completed Work Archive"; + } + todos.push(...added); + todos.sort((left, right) => authorityUnicodeCompare(String(left.todo_id), String(right.todo_id))); + const readModel = projection.todo_read_model as Record; + readModel.todo_count = todos.length; + readModel.records_sha256 = canonicalAuthoritySha256(todos); + return {projection, cases}; +} diff --git a/tests/control_plane_ts/succession_read_conformance.ts b/tests/control_plane_ts/succession_read_conformance.ts new file mode 100644 index 0000000000..0daf18b45e --- /dev/null +++ b/tests/control_plane_ts/succession_read_conformance.ts @@ -0,0 +1,43 @@ +import assert from "node:assert/strict"; +import {spawnSync} from "node:child_process"; +import test from "node:test"; +import type {AuthorityStoreConformanceFactory} from "./authority_store_conformance.ts"; +import {productionScaleSuccessionFixture} from "./production_scale_coordination_fixture.ts"; +import {validateCoordinationTodoReadModel} from "../../loopx/control_plane/coordination/coordination_projection.ts"; + +// Exercise the shipped Python consumer → typed policy, not a second test reducer. +const CONSUMER = ` +import json, sys +from loopx.control_plane.coordination.local_authority import canonical_todo_summary_fields +from loopx.control_plane.todos.goal_todo_projection import filtered_todo_summary +source=json.load(sys.stdin) +summary=canonical_todo_summary_fields(source['todos'])['agent_todos'] +result={} +for key, todo_id in source['cases'].items(): + selected=filtered_todo_summary(summary,role='agent',todo_id=todo_id) + result[key]={'gap':selected.get('completed_without_successor_count',0), + 'gates':[g['gate_state'] for g in selected.get('handoff_gates',[])], + 'closure':'terminal_closure_proof' in selected} +print(json.dumps(result)) +`; +export function registerSuccessionReadConformance(name: string, factory: AuthorityStoreConformanceFactory): void { + for (const schema of ["native", "legacy"] as const) test(`${name}: full-graph succession readback (${schema})`, async context => { + const {store} = await factory(context), goal = "succession-goal"; + const {projection, cases} = productionScaleSuccessionFixture(goal, schema); + validateCoordinationTodoReadModel(projection, goal); + assert.equal((await store.commitAuthority({operation_id: "succession-source", expected_provider_revision: null, + next_projection: projection, events: [], receipts: []})).status, "applied"); + const before = await store.loadAuthority(); + assert.equal(before.status, "loaded"); + const child = spawnSync("python3", ["-c", CONSUMER], {encoding: "utf8", timeout: 90_000, + input: JSON.stringify({todos: before.head.todos, cases})}); + assert.equal(child.status, 0, child.stderr); + const actual = JSON.parse(child.stdout); + assert.deepEqual(actual.inferred_source, {gap: 0, gates: [], closure: false}); + assert.deepEqual(actual.missing_source, {gap: 1, gates: [], closure: false}); + assert.deepEqual(actual.self_source, {gap: 1, gates: [], closure: false}); + assert.deepEqual(actual.handoff_source, {gap: 0, gates: ["cleared_with_successor"], closure: false}); + assert.deepEqual(actual.closed_source, {gap: 0, gates: [], closure: false}); + assert.deepEqual(await store.loadAuthority(), before); + }); +} diff --git a/tests/control_plane_ts/todo_succession.test.ts b/tests/control_plane_ts/todo_succession.test.ts new file mode 100644 index 0000000000..6a2423d1e0 --- /dev/null +++ b/tests/control_plane_ts/todo_succession.test.ts @@ -0,0 +1,97 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import {evaluateTodoSuccession, projectTodoSuccession, projectTodoClosure} from "../../loopx/control_plane/todos/succession.ts"; + +const row = (todo_id: string, overrides = {}) => ({todo_id, status: "done", active: true, + advancement: true, no_followup: false, successors: [], superseded_by: null, + unblocks: null, resumes: null, handoff: false, context_fields: ["claimed_by"], ...overrides}); + +test("one resolver recognizes explicit, supersession and inferred links in source order", () => { + const rows = [row("todo_source", {successors: ["todo_explicit", "todo_source", "todo_missing"], superseded_by: "todo_replaced", handoff: true}), + row("todo_explicit", {advancement: false}), row("todo_replaced"), + row("todo_inferred", {resumes: "todo_source", unblocks: "todo_source", active: false})]; + const result = evaluateTodoSuccession(rows)[0]; + assert.deepEqual(result.successor_todo_ids, ["todo_explicit", "todo_replaced", "todo_inferred"]); + assert.deepEqual(result.unresolved_successor_ids, ["todo_source", "todo_missing"]); + assert.equal(result.handoff_state, "superseded"); + assert.equal(result.successor_gap, false); +}); + +for (const [status, extra, expected] of [ + ["open", {}, "blocking"], ["blocked", {}, "blocking"], ["deferred", {}, "deferred"], + ["done", {}, "cleared_without_successor"], ["done", {no_followup: true}, "cleared_no_followup"], + ["done", {successors: ["todo_next"]}, "cleared_with_successor"], + ["open", {superseded_by: "todo_next"}, "superseded"], + ["open", {superseded_by: "todo_missing"}, "blocking"], +] as const) test(`handoff ${status} ${JSON.stringify(extra)} → ${expected}`, () => { + assert.equal(evaluateTodoSuccession([row("todo_gate", {status, handoff: true, ...extra}), row("todo_next")])[0].handoff_state, expected); +}); + +test("archived and deferred rows are evidence but never unfinished completed work", () => { + const results = evaluateTodoSuccession([row("todo_archive", {active: false}), row("todo_deferred", {status: "deferred"}), + row("todo_plain", {context_fields: []}), row("todo_closed", {no_followup: true})]); + assert.deepEqual(results.map(value => value.successor_gap), [false, false, false, false]); +}); + +test("filtering preserves full-source evidence; editing relevant facts invalidates it", () => { + const source = row("todo_source"); + const evaluations = evaluateTodoSuccession([source, row("todo_next", {resumes: "todo_source"})]); + const request = {schema_version: "todo_succession_request_v0", rows: [source], evaluations: [evaluations[0]]}; + assert.equal((projectTodoSuccession(request).evaluations as Record[])[0].successor_gap, false); + assert.throws(() => projectTodoSuccession({...request, rows: [{...source, no_followup: true}]}), /matching full-source/); + assert.throws(() => projectTodoSuccession({...request, evaluations: []}), /cardinality/); + assert.throws(() => evaluateTodoSuccession([source, source]), /duplicate succession identity/); +}); + +const closed = {status: "done", watch_only: false, no_followup: true, successor_gap: false, replan: false, handoff_state: null}; +const closure = (rows: unknown[], full_selection = true) => projectTodoClosure({schema_version: "todo_closure_request_v0", + role: "agent", source_section: "Agent Todo", full_selection, rows}); +test("terminal proofs require full selection and absence of every unresolved obligation", () => { + assert.ok(closure([closed]).terminal_closure_proof); + assert.ok(closure([]).terminal_closure_proof); + for (const override of [{successor_gap: true}, {replan: true}, {status: "deferred"}, + {status: "open"}, {handoff_state: "cleared_without_successor"}]) { + assert.equal(closure([{...closed, ...override}]).terminal_closure_proof, undefined); + } + assert.equal(closure([closed], false).terminal_closure_proof, undefined); + assert.equal(closure([closed], false).source_proof, undefined); + const watch = closure([{...closed, status: "open", watch_only: true}]); + assert.equal((watch.terminal_closure_proof as Record).all_convergent_todos_done, true); + assert.equal((watch.terminal_closure_proof as Record).all_todos_done, false); +}); + +test("interned field sets are lossless and reject invalid references", () => { + const rows = Array.from({length: 4000}, (_, index) => row(`todo_history_${index}`, { + active: false, context_fields: index % 2 ? ["claimed_by", "completed_at"] : [], + })); + const request = {schema_version: "todo_succession_request_v0", context_field_sets: [[], ["claimed_by", "completed_at"]], + rows: rows.map((value, index) => ({...value, context_fields: index % 2}))}; + assert.deepEqual(projectTodoSuccession(request).evaluations, evaluateTodoSuccession(rows)); + for (const index of [-1, 2, 0.5, "0"]) assert.throws(() => projectTodoSuccession({...request, + rows: [{...rows[0], context_fields: index}]}), /context field index/); +}); + +test("a cached result cannot contradict its matched item facts", () => { + const source = row("todo_source"), evaluation = evaluateTodoSuccession([source])[0]; + for (const mutation of [{successor_gap: false}, {tracked_completion: false}, {handoff_state: "superseded"}, + {successor_todo_ids: [null]}, {successor_todo_ids: [source.todo_id]}]) { + assert.throws(() => projectTodoSuccession({schema_version: "todo_succession_request_v0", + rows: [source], evaluations: [{...evaluation, ...mutation}]}), /succession evaluation|successor identity/); + } +}); + + +test("active identity replaces archived edges regardless of source order", () => { + const source = row("todo_source"); + const archived = row("todo_reused", {active: false, resumes: "todo_source"}); + const active = row("todo_reused", {status: "open"}); + for (const generations of [[archived, active], [active, archived]]) { + const result = evaluateTodoSuccession([source, ...generations]); + assert.equal(result[0].successor_gap, true); + assert.deepEqual(result[0].successor_todo_ids, []); + const explicit = evaluateTodoSuccession([{...source, successors: ["todo_reused"]}, ...generations]); + assert.deepEqual(explicit[0].successor_todo_ids, ["todo_reused"]); + } + assert.throws(() => evaluateTodoSuccession([active, active]), /duplicate succession identity/); + assert.throws(() => evaluateTodoSuccession([archived, archived, active]), /duplicate succession identity/); +}); diff --git a/tests/fixtures/control_plane/coordination_production_scale_v0.json b/tests/fixtures/control_plane/coordination_production_scale_v0.json index 0364a76577..677ce0d6d3 100644 --- a/tests/fixtures/control_plane/coordination_production_scale_v0.json +++ b/tests/fixtures/control_plane/coordination_production_scale_v0.json @@ -55,6 +55,15 @@ "scope_key": "publish_artifact", "other_blocker_id": "todo_zz_remaining_user_decision" }, + "succession": { + "inferred_source": "todo_succession_inferred_source", + "archived_target": "todo_succession_archived_target", + "missing_source": "todo_succession_missing_source", + "self_source": "todo_succession_self_source", + "handoff_source": "todo_succession_handoff_source", + "explicit_target": "todo_succession_explicit_target", + "closed_source": "todo_succession_closed_source" + }, "title_only_monitor": { "todo_id": "todo_fixture_title_monitor", "role": "agent",