diff --git a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md index b01f4d79af..86ad1f13e8 100644 --- a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md +++ b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md @@ -16,6 +16,22 @@ ## Current implementation checkpoint +The projection-delivery stage now closes the cross-language boundary: typed +TypeScript mutation results and the Python compatibility provider share the +same four-state contract (`pending`, `delivered`, `current`, `not_required`). +Provider readback is validated before acknowledgement decisions, and the +end-to-end causal chain is covered by a shared composition fixture. This is a +completed delivery stage, not a promotion of Markdown or a claim that the +remaining lifecycle writers have migrated. + +The same stage also removes duplicated Python read policy around that boundary. +Task-class resolution, title-aware actionability, dependency readiness, agent +eligibility, priority ordering, and canonical Todo read records now have one +Python semantic owner while TypeScript remains the transaction owner. The old +projection module is an import-only compatibility facade. This keeps the +replacement-first rule intact: compatibility remains available, but it cannot +silently become a second semantic implementation. + Native update now composes `todos/public_update.ts` for a bounded nonterminal planning intent (status, evidence/reason, resume/clear and successor links), against the same complete canonical head used for authority checks and CAS. diff --git a/docs/architecture/rfcs/typescript-control-plane-migration-v0.zh-CN.md b/docs/architecture/rfcs/typescript-control-plane-migration-v0.zh-CN.md index 9d536804c1..52439147c5 100644 --- a/docs/architecture/rfcs/typescript-control-plane-migration-v0.zh-CN.md +++ b/docs/architecture/rfcs/typescript-control-plane-migration-v0.zh-CN.md @@ -15,6 +15,18 @@ ## 当前实现检查点 +投影交付阶段现已闭合跨语言边界:typed TypeScript mutation 结果与 Python +兼容 provider 共用四态契约(`pending`、`delivered`、`current`、`not_required`)。 +Provider readback 在 acknowledgement 决策前进行校验,端到端因果链由共享组合 +fixture 覆盖。这是一个已完成的交付阶段,不代表 Markdown 晋升,也不声称其余 +lifecycle writer 已全部迁移。 + +同一阶段也删除了该边界周围重复的 Python read policy。task-class 解析、识别 title +的 actionable 判断、依赖就绪、Agent eligibility、priority 排序和 canonical Todo +read record 现在只有一个 Python 语义 owner,而 TypeScript 仍是事务 owner。旧 +projection 模块只保留 import-only 兼容 facade。这样继续遵守 replacement-first: +兼容路径仍可用,但不能静默形成第二份语义实现。 + Native update 现通过 `todos/public_update.ts` 组合有界的非终态 planning intent (status、evidence/reason、resume/clear、successor links),使用权限检查与 CAS 同一份完整 canonical head。独立 intent 命名空间不扩大原 text/note patch allowlist, diff --git a/docs/reference/contracts/interface-budget-contract.md b/docs/reference/contracts/interface-budget-contract.md index ee307582b9..bcb3c83f5f 100644 --- a/docs/reference/contracts/interface-budget-contract.md +++ b/docs/reference/contracts/interface-budget-contract.md @@ -11,7 +11,7 @@ and size/count budgets. | `heartbeat_prompt_json` | heartbeat automation | wake and route one bounded turn | `quota should-run`, `status`, or `review-packet --handoff-only` | `json_chars <= 4800` plus `interface_budget.within_budget=true` | `nested_keys <= 40` | `top_level_keys <= 30` | | `review_packet_handoff_only_json` | project-agent handoff | forward the smallest sufficient task packet | full `review-packet` or run-history artifact | `json_chars <= 3000` plus `handoff_interface_budget.within_budget=true` | `nested_keys <= 40` | `top_level_keys <= 18` | | `quota_should_run_json` | quota guard | decide whether the selected goal may spend compute | `status`, `history`, or active state | `json_chars <= 13000` | `nested_keys <= 330` | `top_level_keys <= 52` | -| `dashboard_status_json` | operator dashboard | render first-screen operator state | `history`, run artifacts, or project-local adapter output | `json_chars <= 18500` | `nested_keys <= 260` | `top_level_keys <= 25` | +| `dashboard_status_json` | operator dashboard | render first-screen operator state | `history`, run artifacts, or project-local adapter output | `json_chars <= 19500` | `nested_keys <= 260` | `top_level_keys <= 25` | These four budgets measure compact machine payloads. For `heartbeat_prompt_json`, the measured payload is the actual diff --git a/examples/control_plane/hot-path-interface-budget-smoke.py b/examples/control_plane/hot-path-interface-budget-smoke.py index 05f39b87a4..5eb564a7fb 100644 --- a/examples/control_plane/hot-path-interface-budget-smoke.py +++ b/examples/control_plane/hot-path-interface-budget-smoke.py @@ -76,7 +76,7 @@ "owner": "operator dashboard", "consumer": "render first-screen operator state", "cold_path": "history, run artifacts, or project-local adapter output", - "max_json_chars": 18_500, + "max_json_chars": 19_500, "max_nested_keys": 260, "max_top_level_keys": 25, }, diff --git a/loopx/capabilities/explore/composition_frontier.py b/loopx/capabilities/explore/composition_frontier.py index 5752e1d1a4..cdfe6c3917 100644 --- a/loopx/capabilities/explore/composition_frontier.py +++ b/loopx/capabilities/explore/composition_frontier.py @@ -12,7 +12,7 @@ normalize_explore_result_node_refs, normalize_todo_claimed_by, ) -from ...control_plane.todos.projection import ( +from ...control_plane.todos.todo_semantics import ( todo_item_is_actionable_open, todo_item_task_class, ) diff --git a/loopx/capabilities/explore/todo_branch_plan.py b/loopx/capabilities/explore/todo_branch_plan.py index 881f765f71..b1269d9616 100644 --- a/loopx/capabilities/explore/todo_branch_plan.py +++ b/loopx/capabilities/explore/todo_branch_plan.py @@ -13,7 +13,7 @@ normalize_todo_id, normalize_todo_status, ) -from ...control_plane.todos.projection import ( +from ...control_plane.todos.todo_semantics import ( todo_item_is_actionable_open, todo_item_task_class, todo_priority_rank, diff --git a/loopx/capabilities/explore/worker_branch_plan.py b/loopx/capabilities/explore/worker_branch_plan.py index 52ced39055..244947c19a 100644 --- a/loopx/capabilities/explore/worker_branch_plan.py +++ b/loopx/capabilities/explore/worker_branch_plan.py @@ -46,7 +46,7 @@ _shared_dependency_capabilities, _scopes_overlap, ) -from ...control_plane.todos.projection import todo_item_task_class, todo_projection_sort_key +from ...control_plane.todos.todo_semantics import todo_item_task_class, todo_projection_sort_key WORKER_BRANCH_PLAN_SCHEMA_VERSION = "loopx_explore_worker_branch_plan_v0" diff --git a/loopx/capabilities/issue_fix/pr_gate_reconcile.py b/loopx/capabilities/issue_fix/pr_gate_reconcile.py index 20757a6aa1..9e4942f959 100644 --- a/loopx/capabilities/issue_fix/pr_gate_reconcile.py +++ b/loopx/capabilities/issue_fix/pr_gate_reconcile.py @@ -8,7 +8,7 @@ from ...control_plane.todos.contract import ( normalize_todo_decision_scope, ) -from ...control_plane.todos.projection import todo_item_task_class +from ...control_plane.todos.todo_semantics import todo_item_task_class from ...todos import complete_goal_todo, list_goal_todos from .pr_lifecycle import build_issue_fix_pr_lifecycle_monitor_packet from .pr_lifecycle_rollout import append_pr_merge_rollout_event diff --git a/loopx/capabilities/issue_fix/pr_review_ack.py b/loopx/capabilities/issue_fix/pr_review_ack.py index b1352e01a5..549891ae0a 100644 --- a/loopx/capabilities/issue_fix/pr_review_ack.py +++ b/loopx/capabilities/issue_fix/pr_review_ack.py @@ -10,7 +10,7 @@ TODO_TASK_CLASS_USER_ACTION, normalize_todo_bound_agent, ) -from ...control_plane.todos.projection import todo_item_task_class +from ...control_plane.todos.todo_semantics import todo_item_task_class from ...history import load_registry from ...paths import resolve_runtime_root from ...rollout_event_log import ( diff --git a/loopx/capabilities/periodic_report/project_progress_snapshot.py b/loopx/capabilities/periodic_report/project_progress_snapshot.py index f6b41a0df0..7c0bfc3d49 100644 --- a/loopx/capabilities/periodic_report/project_progress_snapshot.py +++ b/loopx/capabilities/periodic_report/project_progress_snapshot.py @@ -6,7 +6,7 @@ from typing import Any from ...control_plane.todos.active_state_todo_parser import parse_active_state_todos -from ...control_plane.todos.projection import todo_item_is_actionable_open +from ...control_plane.todos.todo_semantics import todo_item_is_actionable_open from ...registry import find_registry_goal, read_json, resolve_state_file from .incremental import select_incremental_project_progress diff --git a/loopx/control_plane/agents/agent_lane_recommendation.py b/loopx/control_plane/agents/agent_lane_recommendation.py index 4e5c85ccfe..7f83a349f6 100644 --- a/loopx/control_plane/agents/agent_lane_recommendation.py +++ b/loopx/control_plane/agents/agent_lane_recommendation.py @@ -11,7 +11,7 @@ normalize_todo_claimed_by, normalize_todo_id, ) -from ..todos.projection import todo_item_is_due_monitor +from ..todos.todo_semantics import todo_item_is_due_monitor from ..todos.summary_item import compact_todo_summary_item from ..work_items.primary_action import protocol_action_text from ..work_items.work_lane import ( diff --git a/loopx/control_plane/agents/agent_scope.py b/loopx/control_plane/agents/agent_scope.py index dd3a48f265..af8931bf2b 100644 --- a/loopx/control_plane/agents/agent_scope.py +++ b/loopx/control_plane/agents/agent_scope.py @@ -27,7 +27,7 @@ ) from ..todos.handoff_gate import HandoffGateState from ..todos.resume_planning import project_todo_resume_planning -from ..todos.projection import ( +from ..todos.todo_semantics import ( todo_item_claimed_by_agent_or_unclaimed, todo_item_excludes_agent, todo_item_is_actionable_open, diff --git a/loopx/control_plane/agents/capability_gate.py b/loopx/control_plane/agents/capability_gate.py index 150935d68d..383639e316 100644 --- a/loopx/control_plane/agents/capability_gate.py +++ b/loopx/control_plane/agents/capability_gate.py @@ -9,7 +9,7 @@ normalize_target_capabilities, normalize_todo_claimed_by, ) -from ..todos.projection import ( +from ..todos.todo_semantics import ( todo_index_rank, todo_item_is_actionable_open, todo_item_task_class, diff --git a/loopx/control_plane/coordination/todo_claim.ts b/loopx/control_plane/coordination/todo_claim.ts index ed348cf03e..48e2ec7a9f 100644 --- a/loopx/control_plane/coordination/todo_claim.ts +++ b/loopx/control_plane/coordination/todo_claim.ts @@ -8,6 +8,7 @@ import { requireAuthorityStoreId, } from "./authority_store_codec.ts"; import {validateContinuationNote, computeContinuationTodoFacts} from "./continuation_note.ts"; +import {projectionDelivery} from "../todos/projection_delivery.ts"; import {normalizeRegisteredTodoAgents, normalizeTodoAgent} from "./todo_agents.ts"; import { prepareCoordinationProjectionCommit, @@ -450,7 +451,7 @@ export async function executeCoordinationTodoClaim( provider_revision: receipt.provider_revision, cursor: receipt.cursor, original_receipt: original, - projection_delivery: result.changed === false ? "not_required" : "pending", + projection_delivery: projectionDelivery(result.changed !== false), projection_source: "committed_authority_journal", }; }; diff --git a/loopx/control_plane/coordination/todo_compatibility_edit.ts b/loopx/control_plane/coordination/todo_compatibility_edit.ts index e47006b26e..94cde3e4c2 100644 --- a/loopx/control_plane/coordination/todo_compatibility_edit.ts +++ b/loopx/control_plane/coordination/todo_compatibility_edit.ts @@ -2,6 +2,7 @@ import type { JsonObject } from "../effect_program.ts"; import type { AuthorityStore } from "./authority_store.ts"; import { canonicalAuthorityBytes, canonicalAuthorityObject, canonicalAuthoritySha256, requireAuthorityStoreId } from "./authority_store_codec.ts"; import { prepareCoordinationProjectionCommit, indexCoordinationProjection, validateCoordinationTodoReadModel } from "./coordination_projection.ts"; +import { projectionDelivery } from "../todos/projection_delivery.ts"; export const TODO_COMPATIBILITY_EDIT_SCHEMA = "loopx_todo_compatibility_edit_request_v0"; export const TODO_COMPATIBILITY_EDIT_RESULT_SCHEMA = "loopx_todo_compatibility_edit_result_v0"; @@ -75,7 +76,7 @@ export async function editCoordinationTodo( status: status === "applied" && !original.changed ? "no_change" : status, changed: status !== "replayed" && original.changed, provider_revision: receipt.provider_revision, cursor: receipt.cursor, - projection_delivery: original.changed ? "pending" : "not_required", + projection_delivery: projectionDelivery(original.changed), projection_source: "committed_authority_journal", }; }; diff --git a/loopx/control_plane/coordination/todo_create.ts b/loopx/control_plane/coordination/todo_create.ts index eafb21d76e..93fcd38891 100644 --- a/loopx/control_plane/coordination/todo_create.ts +++ b/loopx/control_plane/coordination/todo_create.ts @@ -1,4 +1,5 @@ import type { JsonObject } from "../effect_program.ts"; +import {projectionDelivery} from "../todos/projection_delivery.ts"; import type { AuthorityStore, AuthorityStoreReceiptResult } from "./authority_store.ts"; import { AuthorityStoreProtocolError, @@ -77,7 +78,7 @@ function replayCreate( provider_revision: receipt.provider_revision, cursor: receipt.cursor, original_receipt: original, - projection_delivery: "pending", + projection_delivery: projectionDelivery(true), projection_source: "committed_authority_journal", }; } diff --git a/loopx/control_plane/coordination/todo_monitor_poll.ts b/loopx/control_plane/coordination/todo_monitor_poll.ts index 264ee05390..b5129d7c3a 100644 --- a/loopx/control_plane/coordination/todo_monitor_poll.ts +++ b/loopx/control_plane/coordination/todo_monitor_poll.ts @@ -11,6 +11,7 @@ import {planMonitorMetadata, TODO_MONITOR_METADATA_REQUEST_SCHEMA} from "../todo import {planMonitorSuccessor, selectMonitorTodo, MONITOR_SUCCESSOR_REQUEST_SCHEMA} from "../scheduler/monitor_successor.ts"; import {optionalNonEmptyString, requireBoolean} from "../runtime_decode.ts"; import {planTodoAuthoringScope, TODO_AUTHORING_SCOPE_REQUEST_SCHEMA} from "../todos/authoring_scope.ts"; +import {projectionDelivery} from "../todos/projection_delivery.ts"; export const COORDINATION_MONITOR_POLL_REQUEST_SCHEMA = "loopx_coordination_monitor_poll_request_v0"; export const COORDINATION_MONITOR_POLL_RESULT_SCHEMA = "loopx_coordination_monitor_poll_result_v0"; @@ -42,7 +43,7 @@ function replay(receipt: AuthorityStoreReceiptResult, input: CoordinationMonitor return {schema_version: COORDINATION_MONITOR_POLL_RESULT_SCHEMA, status, changed: status !== "replayed", provider_revision: receipt.provider_revision, cursor: receipt.cursor, writeback: {...canonicalAuthorityObject(original.writeback, "Monitor writeback"), provider_replayed: status === "replayed"}, - projection_delivery: "pending", projection_source: "committed_authority_journal"}; + projection_delivery: projectionDelivery(true), projection_source: "committed_authority_journal"}; } function normalize(raw: CoordinationMonitorPollInput): CoordinationMonitorPollInput { diff --git a/loopx/control_plane/coordination/todo_terminal_lifecycle.ts b/loopx/control_plane/coordination/todo_terminal_lifecycle.ts index 2e2bc2e691..c2ed4cf84e 100644 --- a/loopx/control_plane/coordination/todo_terminal_lifecycle.ts +++ b/loopx/control_plane/coordination/todo_terminal_lifecycle.ts @@ -1,4 +1,5 @@ import { createHash } from "node:crypto"; +import {projectionDelivery} from "../todos/projection_delivery.ts"; import type { JsonObject } from "../effect_program.ts"; import type { @@ -451,7 +452,7 @@ function replayTerminal( provider_revision: receipt.provider_revision, cursor: receipt.cursor, original_receipt: original, - projection_delivery: result.changed === true ? "pending" : "not_required", + projection_delivery: projectionDelivery(result.changed === true), projection_source: "committed_authority_journal", }; } @@ -1134,7 +1135,7 @@ function replayArchive( provider_revision: receipt.provider_revision, cursor: receipt.cursor, original_receipt: original, - projection_delivery: result.changed === true ? "pending" : "not_required", + projection_delivery: projectionDelivery(result.changed === true), projection_source: "committed_authority_journal", }; } diff --git a/loopx/control_plane/coordination/todo_update.ts b/loopx/control_plane/coordination/todo_update.ts index 11756d7202..68d6a4c50e 100644 --- a/loopx/control_plane/coordination/todo_update.ts +++ b/loopx/control_plane/coordination/todo_update.ts @@ -27,6 +27,7 @@ import { evaluateCoordinationTerminalFence, COORDINATION_TERMINAL_FENCE_REQUEST_ import { leaseEpoch } from "../work_items/task_lease_acquire.ts"; import { parseIsoTimestamp } from "../runtime_timestamp.ts"; import { normalizeNativePlanningIntent, planNativeTodoUpdate } from "../todos/native_update_plan.ts"; +import { projectionDelivery } from "../todos/projection_delivery.ts"; export const COORDINATION_TODO_UPDATE_REQUEST_SCHEMA = "loopx_local_coordination_todo_update_request_v0"; @@ -139,7 +140,7 @@ function replayUpdate( changed: status !== "replayed" && original.changed, todo_id: input.todo_id, provider_revision: receipt.provider_revision, cursor: receipt.cursor, original_receipt: original, - projection_delivery: original.changed ? "pending" : "not_required", + projection_delivery: projectionDelivery(original.changed), projection_source: "committed_authority_journal"}; } diff --git a/loopx/control_plane/goals/goal_frontier/__init__.py b/loopx/control_plane/goals/goal_frontier/__init__.py index 88bdbe311c..83637cf3de 100644 --- a/loopx/control_plane/goals/goal_frontier/__init__.py +++ b/loopx/control_plane/goals/goal_frontier/__init__.py @@ -12,7 +12,7 @@ from ...agents.runtime_model import peer_work_key, select_peer_for_work from ...runtime.time import parse_timestamp from ...todos.contract import normalize_todo_replan_obligation_id -from ...todos.projection import ( +from ...todos.todo_semantics import ( todo_advancement_frontier_counts, todo_item_is_watch_only_monitor, ) diff --git a/loopx/control_plane/goals/goal_frontier/ack_policy.py b/loopx/control_plane/goals/goal_frontier/ack_policy.py index c50f95805b..3447ebb995 100644 --- a/loopx/control_plane/goals/goal_frontier/ack_policy.py +++ b/loopx/control_plane/goals/goal_frontier/ack_policy.py @@ -10,7 +10,7 @@ normalize_todo_replan_obligation_id, replan_successor_semantic_binding, ) -from ...todos.projection import todo_item_is_actionable_open, todo_item_task_class +from ...todos.todo_semantics import todo_item_is_actionable_open, todo_item_task_class from ...work_items.progress_observation import ( replan_obligation_trigger_checkpoints, replan_obligation_trigger_kinds, diff --git a/loopx/control_plane/goals/goal_frontier/fallback_disposition.py b/loopx/control_plane/goals/goal_frontier/fallback_disposition.py index 630a8bd278..a99b9ba531 100644 --- a/loopx/control_plane/goals/goal_frontier/fallback_disposition.py +++ b/loopx/control_plane/goals/goal_frontier/fallback_disposition.py @@ -7,7 +7,7 @@ normalize_todo_id, ) from ...todos.resume_planning import project_todo_resume_planning -from ...todos.projection import ( +from ...todos.todo_semantics import ( agent_scoped_selectable_advancement_todo_ids, ) from ..goal_vision_read_model import ( diff --git a/loopx/control_plane/goals/goal_vision_wait.py b/loopx/control_plane/goals/goal_vision_wait.py index 317d903423..f453c5a6ec 100644 --- a/loopx/control_plane/goals/goal_vision_wait.py +++ b/loopx/control_plane/goals/goal_vision_wait.py @@ -5,7 +5,7 @@ from typing import Any from ..effect_runtime import effect_runtime_result -from ..todos.projection import todo_item_excludes_agent +from ..todos.todo_semantics import todo_item_excludes_agent from ..todos.contract import ( TODO_TASK_CLASS_BLOCKER, normalize_todo_claimed_by, diff --git a/loopx/control_plane/goals/goal_vision_wait_projection.py b/loopx/control_plane/goals/goal_vision_wait_projection.py index 7a3b62dd2c..383365869e 100644 --- a/loopx/control_plane/goals/goal_vision_wait_projection.py +++ b/loopx/control_plane/goals/goal_vision_wait_projection.py @@ -2,7 +2,7 @@ from typing import Any -from ..todos.projection import agent_scoped_selectable_advancement_todo_ids +from ..todos.todo_semantics import agent_scoped_selectable_advancement_todo_ids from .goal_vision_read_model import ( acceptance_gaps_from_agent_vision, latest_agent_vision_from_runs, diff --git a/loopx/control_plane/goals/start_goal_todo_delta.py b/loopx/control_plane/goals/start_goal_todo_delta.py index 2bbb0d5206..8152323c03 100644 --- a/loopx/control_plane/goals/start_goal_todo_delta.py +++ b/loopx/control_plane/goals/start_goal_todo_delta.py @@ -21,7 +21,7 @@ from ...control_plane.todos.contract import ( TODO_TASK_CLASS_ADVANCEMENT, ) -from ...control_plane.todos.projection import todo_item_is_actionable_open +from ...control_plane.todos.todo_semantics import todo_item_is_actionable_open from ...project_prompt import render_cli_command_prefix, shell_arg from ...registry import registry_goals, resolve_state_file diff --git a/loopx/control_plane/quota/monitor_poll.py b/loopx/control_plane/quota/monitor_poll.py index 4bf5d4b5d4..6532a329d1 100644 --- a/loopx/control_plane/quota/monitor_poll.py +++ b/loopx/control_plane/quota/monitor_poll.py @@ -25,7 +25,7 @@ from ..todos.external_wait_contract import ( build_monitor_advancement_authoring_contract, ) -from ..todos.projection import todo_item_task_class +from ..todos.todo_semantics import todo_item_task_class from .decision_summary import compact_quota_decision, quota_decision_agent_id from .spend_sources import DEFAULT_SLOT_SPEND_SOURCE diff --git a/loopx/control_plane/quota/monitor_poll_commit.ts b/loopx/control_plane/quota/monitor_poll_commit.ts index c120f4baa3..625d140170 100644 --- a/loopx/control_plane/quota/monitor_poll_commit.ts +++ b/loopx/control_plane/quota/monitor_poll_commit.ts @@ -3,6 +3,7 @@ import { access, readFile, rm } from "node:fs/promises"; import { basename, dirname, join, resolve } from "node:path"; import { monitorSuccessorIntent, monitorSuccessorRoute } from "../scheduler/monitor_successor.ts"; import { normalizeTodoCapabilities } from "../todos/work_requirements.ts"; +import { parseProjectionDelivery } from "../todos/projection_delivery.ts"; import type { JsonObject } from "../effect_program.ts"; import { EffectRuntimeRequestError } from "../effect_runtime_errors.ts"; @@ -642,9 +643,13 @@ function compactProviderWriteback(receipt: JsonObject): JsonObject { * Omitted on the legacy path to retain its exact v0 response shape. */ function monitorProjectionDelivery(receipt: JsonObject): JsonObject { if (receipt.projection_delivery == null) return {}; - const status = optionalString(receipt.projection_delivery, "projection_delivery"); - if (!["delivered", "pending", "not_required"].includes(String(status))) { - throw new EffectRuntimeRequestError("invalid Monitor projection delivery status"); + let status: ReturnType; + try { + status = parseProjectionDelivery(receipt.projection_delivery); + } catch (error) { + throw new EffectRuntimeRequestError( + error instanceof Error ? error.message : "invalid Monitor projection delivery status", + ); } const outbox = requiredObject(receipt.projection_outbox, "projection_outbox"); const diagnostic: JsonObject = {}; diff --git a/loopx/control_plane/quota/projection_repair.py b/loopx/control_plane/quota/projection_repair.py index 7f8b225858..df427715aa 100644 --- a/loopx/control_plane/quota/projection_repair.py +++ b/loopx/control_plane/quota/projection_repair.py @@ -8,7 +8,7 @@ TODO_TASK_CLASS_ADVANCEMENT, normalize_required_write_scopes, ) -from ..todos.projection import ( +from ..todos.todo_semantics import ( todo_item_is_actionable_open, todo_item_task_class, ) diff --git a/loopx/control_plane/quota/should_run_packet.py b/loopx/control_plane/quota/should_run_packet.py index 1e6c9ea217..5b460b427d 100644 --- a/loopx/control_plane/quota/should_run_packet.py +++ b/loopx/control_plane/quota/should_run_packet.py @@ -95,7 +95,7 @@ from ..todos.contract import ( normalize_todo_claimed_by, ) -from ..todos.projection import ( +from ..todos.todo_semantics import ( todo_item_is_actionable_open as projection_todo_item_is_actionable_open, ) from ..todos.quota_summary import ( diff --git a/loopx/control_plane/quota/should_run_prepare.py b/loopx/control_plane/quota/should_run_prepare.py index 632d7053cf..2a66abe948 100644 --- a/loopx/control_plane/quota/should_run_prepare.py +++ b/loopx/control_plane/quota/should_run_prepare.py @@ -79,19 +79,11 @@ normalize_todo_resume_when, normalize_todo_status, ) -from ..todos.projection import ( +from ..todos.todo_semantics import ( todo_item_is_actionable_open as projection_todo_item_is_actionable_open, -) -from ..todos.projection import ( todo_item_is_due_monitor as projection_todo_item_is_due_monitor, -) -from ..todos.projection import ( todo_item_is_expired_monitor as projection_todo_item_is_expired_monitor, -) -from ..todos.projection import ( todo_item_next_due_at as projection_todo_item_next_due_at, -) -from ..todos.projection import ( todo_item_task_class as projection_todo_item_task_class, ) from ..todos.quota_summary import ( diff --git a/loopx/control_plane/quota/task_orchestration_admission.py b/loopx/control_plane/quota/task_orchestration_admission.py index 5809372c2e..dd2c40cb26 100644 --- a/loopx/control_plane/quota/task_orchestration_admission.py +++ b/loopx/control_plane/quota/task_orchestration_admission.py @@ -14,7 +14,7 @@ normalize_todo_task_domain, normalize_todo_task_repository, ) -from ..todos.projection import todo_item_is_actionable_open +from ..todos.todo_semantics import todo_item_is_actionable_open from ..work_items.primary_action import protocol_action_text from .projection_repair import write_scope_allowed diff --git a/loopx/control_plane/scheduler/external_evidence_observation.py b/loopx/control_plane/scheduler/external_evidence_observation.py index 05a25c0fb9..e8be659e77 100644 --- a/loopx/control_plane/scheduler/external_evidence_observation.py +++ b/loopx/control_plane/scheduler/external_evidence_observation.py @@ -4,7 +4,7 @@ from typing import Any from ..todos.contract import normalize_todo_claimed_by, normalize_todo_id -from ..todos.projection import ( +from ..todos.todo_semantics import ( todo_summary_claim_scope_agent_id, todo_summary_has_only_future_scoped_monitor_work, todo_summary_monitor_due_count, diff --git a/loopx/control_plane/todos/projection.py b/loopx/control_plane/todos/projection.py index 048c0fa04e..982623f84b 100644 --- a/loopx/control_plane/todos/projection.py +++ b/loopx/control_plane/todos/projection.py @@ -1,813 +1,8 @@ -from __future__ import annotations +"""Compatibility exports for the canonical Todo semantic kernel. -from datetime import datetime -import re -from typing import Any +New production code should import :mod:`todo_semantics` directly. This module +remains a stable import path for extensions and older integrations while the +Python/TypeScript control-plane migration removes duplicate decision rules. +""" -from ..scheduler.monitor_todo import ( - monitor_todo_expires_at, - monitor_todo_has_schedule, - monitor_todo_is_actionable_open, - monitor_todo_is_due, - monitor_todo_is_expired, - monitor_todo_missing_schedule, - monitor_todo_next_due_at, - monitor_todo_task_class, -) -from .contract import ( - TODO_STATUS_DEFERRED, - TODO_TASK_CLASS_ADVANCEMENT, - TODO_TASK_CLASS_MONITOR, - normalize_todo_claimed_by, - normalize_todo_excluded_agents, - normalize_removed_todo_continuation_policy, - normalize_todo_id, - normalize_todo_status, - normalize_todo_watch_only, -) - - -TODO_MISSING_PRIORITY_RANK = 50 -TODO_MISSING_INDEX = 999999 -TODO_PRIORITY_PREFIX_PATTERN = re.compile( - r"^\s*\[(P[0-4][^\]]*)\]\s*(.+)$", - re.IGNORECASE, -) -TODO_PRIORITY_LABEL_PATTERN = re.compile(r"\bP([0-4])\b", re.IGNORECASE) - - -def todo_item_is_watch_only_monitor(item: dict[str, Any]) -> bool: - return bool( - todo_item_task_class(item) == TODO_TASK_CLASS_MONITOR - and normalize_todo_watch_only(item.get("watch_only")) is True - ) - - -def todo_priority_parts(text: str) -> tuple[str | None, str]: - match = TODO_PRIORITY_PREFIX_PATTERN.match(text) - if not match: - return None, text - return match.group(1).strip().upper(), match.group(2).strip() - - -def todo_priority_label( - item: dict[str, Any], - *, - text_mode: str = "label", -) -> str | None: - priority = item.get("priority") - if isinstance(priority, str) and priority.strip(): - return priority.strip().upper() - text = " ".join( - str(value or "") - for value in (item.get("title"), item.get("text")) - if str(value or "").strip() - ) - if text_mode == "prefix": - priority, _ = todo_priority_parts(text) - return priority - match = TODO_PRIORITY_LABEL_PATTERN.search(text.upper()) - if not match: - return None - return f"P{match.group(1)}" - - -def todo_priority_rank(value: Any, *, text_mode: str = "label") -> int: - if isinstance(value, dict): - priority = todo_priority_label(value, text_mode=text_mode) - elif isinstance(value, str): - priority = value.strip().upper() - else: - priority = None - if not priority: - return TODO_MISSING_PRIORITY_RANK - match = re.match(r"P([0-4])", priority) - if not match: - return TODO_MISSING_PRIORITY_RANK - return int(match.group(1)) - - -def todo_index_rank(item: dict[str, Any]) -> int: - raw_index = item.get("index") - try: - return int(raw_index) if raw_index is not None else TODO_MISSING_INDEX - except (TypeError, ValueError): - return TODO_MISSING_INDEX - - -def todo_projection_sort_key( - item: dict[str, Any], - *, - text_mode: str = "label", -) -> tuple[int, int]: - return (todo_priority_rank(item, text_mode=text_mode), todo_index_rank(item)) - - -def todo_claimed_visibility_items( - items: list[dict[str, Any]], - *, - limit: int, -) -> list[dict[str, Any]]: - if limit <= 0 or len(items) <= limit: - return items[:limit] - claim_order: list[str] = [] - buckets: dict[str, list[dict[str, Any]]] = {} - for item in items: - claimed_by = normalize_todo_claimed_by(item.get("claimed_by")) - if not claimed_by: - continue - if claimed_by not in buckets: - buckets[claimed_by] = [] - claim_order.append(claimed_by) - buckets[claimed_by].append(item) - if not buckets: - return items[:limit] - - original_index = {id(item): index for index, item in enumerate(items)} - per_claimant_cap = max(1, limit // len(buckets)) - selected: list[dict[str, Any]] = [] - selected_ids: set[int] = set() - for claimed_by in claim_order: - taken = 0 - for item in buckets[claimed_by]: - if taken >= per_claimant_cap: - break - if len(selected) >= limit: - break - selected.append(item) - selected_ids.add(id(item)) - taken += 1 - if len(selected) >= limit: - break - - if len(selected) < limit: - for item in items: - if id(item) in selected_ids: - continue - selected.append(item) - selected_ids.add(id(item)) - if len(selected) >= limit: - break - - return sorted( - selected, key=lambda item: original_index.get(id(item), TODO_MISSING_INDEX) - )[:limit] - - -def todo_item_task_text( - item: dict[str, Any], - *, - keys: tuple[str, ...] = ("title", "text"), -) -> str: - return " ".join( - str(item.get(key) or "") for key in keys if str(item.get(key) or "").strip() - ) - - -def todo_item_task_class( - item: dict[str, Any], - *, - task_text_keys: tuple[str, ...] = ("title", "text"), -) -> str: - return monitor_todo_task_class( - item, - task_text=todo_item_task_text(item, keys=task_text_keys), - ) - - -def todo_item_is_actionable_open(item: dict[str, Any]) -> bool: - return monitor_todo_is_actionable_open(item) - - -def todo_item_is_deferred(item: dict[str, Any]) -> bool: - return (normalize_todo_status(item.get("status")) or "") == TODO_STATUS_DEFERRED - - -def todo_item_next_due_at(item: dict[str, Any]) -> datetime | None: - return monitor_todo_next_due_at(item) - - -def todo_item_has_monitor_schedule(item: dict[str, Any]) -> bool: - return monitor_todo_has_schedule(item) - - -def todo_item_expires_at(item: dict[str, Any]) -> datetime | None: - return monitor_todo_expires_at(item) - - -def todo_item_is_expired_monitor( - item: dict[str, Any], *, now: datetime | None = None -) -> bool: - return monitor_todo_is_expired(item, now=now) - - -def todo_item_is_due_monitor( - item: dict[str, Any], - *, - now: datetime | None = None, - task_text_keys: tuple[str, ...] = ("title", "text"), -) -> bool: - return monitor_todo_is_due( - item, - now=now, - task_text=todo_item_task_text(item, keys=task_text_keys), - ) - - -def todo_item_missing_monitor_schedule( - item: dict[str, Any], - *, - now: datetime | None = None, - task_text_keys: tuple[str, ...] = ("title", "text"), -) -> bool: - return monitor_todo_missing_schedule( - item, - now=now, - task_text=todo_item_task_text(item, keys=task_text_keys), - ) - - -def todo_item_claimed_by_agent_or_unclaimed( - item: dict[str, Any], - *, - agent_id: str | None, -) -> bool: - if todo_item_has_removed_continuation_policy(item): - return False - normalized_agent_id = normalize_todo_claimed_by(agent_id) - if not normalized_agent_id: - return True - if normalized_agent_id in normalize_todo_excluded_agents( - item.get("excluded_agents") - ): - return False - claimed_by = normalize_todo_claimed_by(item.get("claimed_by")) - return not claimed_by or claimed_by == normalized_agent_id - - -def todo_advancement_frontier_items( - summary: dict[str, Any] | None, - *, - agent_id: str | None, -) -> dict[str, list[dict[str, Any]]]: - """Return the authoritative advancement frontier items grouped by claim ownership. - - Preserves the slot precedence of executable backlog first, falling back to - unclaimed priority and claimed advancement open items when the executable backlog - is omitted. Peer-claimed items are tracked separately and excluded from the current - agent's selectable advancement frontier. - """ - - empty: dict[str, list[dict[str, Any]]] = { - "current_agent_claimed_items": [], - "unclaimed_items": [], - "other_agent_claimed_items": [], - } - if not isinstance(summary, dict): - return empty - - normalized_agent_id = normalize_todo_claimed_by(agent_id) - executable_items = summary.get("executable_backlog_items") - if isinstance(executable_items, list): - current_items: list[dict[str, Any]] = [] - unclaimed_items: list[dict[str, Any]] = [] - other_items: list[dict[str, Any]] = [] - for value in executable_items: - if not isinstance(value, dict): - continue - if not todo_item_is_actionable_open(value): - continue - if todo_item_task_class(value) != TODO_TASK_CLASS_ADVANCEMENT: - continue - claimed_by = normalize_todo_claimed_by(value.get("claimed_by")) - if claimed_by: - if normalized_agent_id and claimed_by == normalized_agent_id: - if not todo_item_excludes_agent( - value, agent_id=normalized_agent_id - ): - current_items.append(value) - elif normalized_agent_id: - other_items.append(value) - else: - current_items.append(value) - continue - if not todo_item_excludes_agent(value, agent_id=normalized_agent_id): - unclaimed_items.append(value) - return { - "current_agent_claimed_items": current_items, - "unclaimed_items": unclaimed_items, - "other_agent_claimed_items": other_items, - } - - unclaimed_items = [ - value - for value in summary.get("unclaimed_priority_open_items") or [] - if isinstance(value, dict) - and todo_item_is_actionable_open(value) - and todo_item_task_class(value) == TODO_TASK_CLASS_ADVANCEMENT - and not todo_item_excludes_agent(value, agent_id=normalized_agent_id) - ] - current_items = [ - value - for value in summary.get("claimed_advancement_open_items") or [] - if isinstance(value, dict) - and todo_item_is_actionable_open(value) - and todo_item_task_class(value) == TODO_TASK_CLASS_ADVANCEMENT - and ( - not normalized_agent_id - or normalize_todo_claimed_by(value.get("claimed_by")) == normalized_agent_id - ) - and not todo_item_excludes_agent(value, agent_id=normalized_agent_id) - ] - other_items = [ - value - for value in summary.get("claimed_advancement_open_items") or [] - if isinstance(value, dict) - and todo_item_is_actionable_open(value) - and todo_item_task_class(value) == TODO_TASK_CLASS_ADVANCEMENT - and normalized_agent_id - and normalize_todo_claimed_by(value.get("claimed_by")) - and normalize_todo_claimed_by(value.get("claimed_by")) != normalized_agent_id - ] - return { - "current_agent_claimed_items": current_items, - "unclaimed_items": unclaimed_items, - "other_agent_claimed_items": other_items, - } - - -def agent_scoped_selectable_advancement_todo_ids( - agent_todo_summary: dict[str, Any] | None, - *, - agent_id: str | None, -) -> set[str]: - """Return the ids the agent-scoped selectable advancement frontier holds. - - Derived directly from the authoritative ``todo_advancement_frontier_items`` - helper so that slot precedence and claim ownership predicates never diverge - from the frontier counter. - """ - - frontier_items = todo_advancement_frontier_items( - agent_todo_summary, - agent_id=agent_id, - ) - selectable: set[str] = set() - for item in ( - frontier_items["current_agent_claimed_items"] - + frontier_items["unclaimed_items"] - ): - if todo_id := normalize_todo_id(item.get("todo_id")): - selectable.add(todo_id) - return selectable - - -def todo_advancement_frontier_counts( - summary: dict[str, Any] | None, - *, - agent_id: str | None, -) -> dict[str, int]: - """Classify the durable advancement frontier by exact claim ownership.""" - - if not isinstance(summary, dict): - return { - "current_agent_claimed_advancement_count": 0, - "unclaimed_advancement_count": 0, - "other_agent_claimed_advancement_count": 0, - } - frontier_items = todo_advancement_frontier_items(summary, agent_id=agent_id) - claim_scope = summary.get("claim_scope") - other_items = ( - claim_scope.get("other_agent_claimed_items") - if isinstance(claim_scope, dict) - else [] - ) - diagnostic_other_count = sum( - 1 - for value in other_items or [] - if isinstance(value, dict) - and todo_item_is_actionable_open(value) - and todo_item_task_class(value) == TODO_TASK_CLASS_ADVANCEMENT - ) - return { - "current_agent_claimed_advancement_count": max( - len(frontier_items["current_agent_claimed_items"]), - _positive_int(summary.get("current_agent_claimed_advancement_count")), - ), - "unclaimed_advancement_count": len(frontier_items["unclaimed_items"]), - "other_agent_claimed_advancement_count": max( - len(frontier_items["other_agent_claimed_items"]), - diagnostic_other_count, - ), - } - - -def todo_item_has_removed_continuation_policy(item: dict[str, Any]) -> bool: - return bool( - normalize_removed_todo_continuation_policy( - item.get("removed_continuation_policy") - ) - ) - - -def todo_item_excludes_agent( - item: dict[str, Any], - *, - agent_id: str | None, -) -> bool: - normalized_agent_id = normalize_todo_claimed_by(agent_id) - return bool( - normalized_agent_id - and normalized_agent_id - in normalize_todo_excluded_agents(item.get("excluded_agents")) - ) - - -def todo_summary_claim_scope_agent_id(summary: dict[str, Any] | None) -> str | None: - if not isinstance(summary, dict): - return None - claim_scope = summary.get("claim_scope") - if not isinstance(claim_scope, dict): - return None - return normalize_todo_claimed_by(claim_scope.get("agent_id")) - - -def todo_summary_monitor_writeback_contract( - summary: dict[str, Any] | None, -) -> dict[str, Any] | None: - if not isinstance(summary, dict): - return None - contract = summary.get("monitor_writeback") - if not isinstance(contract, dict): - return None - if contract.get("supported") is not False: - return None - compact: dict[str, Any] = {"supported": False} - source = str(contract.get("source") or "").strip() - if source: - compact["source"] = source - return compact - - -def todo_summary_monitor_writeback_supported(summary: dict[str, Any] | None) -> bool: - contract = todo_summary_monitor_writeback_contract(summary) - if not contract: - return True - return contract.get("supported") is not False - - -def todo_summary_monitor_items(summary: dict[str, Any] | None) -> list[dict[str, Any]]: - if not isinstance(summary, dict): - return [] - items: list[dict[str, Any]] = [] - seen: set[tuple[str, int]] = set() - for key in ( - "monitor_due_items", - "current_agent_claimed_monitor_items", - "monitor_open_items", - "claimed_monitor_open_items", - "first_open_items", - ): - values = summary.get(key) - if not isinstance(values, list): - continue - for value in values: - if not isinstance(value, dict): - continue - if not todo_item_is_actionable_open(value): - continue - if todo_item_task_class(value) != TODO_TASK_CLASS_MONITOR: - continue - identity = (normalize_todo_id(value.get("todo_id")) or "", id(value)) - if identity in seen: - continue - seen.add(identity) - items.append(value) - return items - - -def _summary_monitor_items( - summary: dict[str, Any] | None, - *, - projected_key: str, - predicate: Any, - task_text_keys: tuple[str, ...], - text_mode: str, -) -> list[dict[str, Any]]: - if not isinstance(summary, dict): - return [] - if not todo_summary_monitor_writeback_supported(summary): - return [] - projected_items = summary.get(projected_key) - if isinstance(projected_items, list): - items = [ - item - for item in projected_items - if isinstance(item, dict) - if todo_item_is_actionable_open(item) - if todo_item_task_class(item, task_text_keys=task_text_keys) - == TODO_TASK_CLASS_MONITOR - if predicate(item) - ] - else: - raw_items = summary.get("monitor_open_items") - items = [ - item - for item in (raw_items if isinstance(raw_items, list) else []) - if isinstance(item, dict) - if predicate(item) - ] - agent_id = todo_summary_claim_scope_agent_id(summary) - if agent_id: - items = [ - item - for item in items - if todo_item_claimed_by_agent_or_unclaimed(item, agent_id=agent_id) - ] - return sorted( - items, - key=lambda item: todo_projection_sort_key(item, text_mode=text_mode), - ) - - -def todo_summary_monitor_due_items( - summary: dict[str, Any] | None, - *, - task_text_keys: tuple[str, ...] = ("title", "text"), - text_mode: str = "label", -) -> list[dict[str, Any]]: - return _summary_monitor_items( - summary, - projected_key="monitor_due_items", - predicate=lambda item: todo_item_is_due_monitor( - item, - task_text_keys=task_text_keys, - ), - task_text_keys=task_text_keys, - text_mode=text_mode, - ) - - -def todo_summary_monitor_due_count( - summary: dict[str, Any] | None, - *, - due_items: list[dict[str, Any]] | None = None, - task_text_keys: tuple[str, ...] = ("title", "text"), - text_mode: str = "label", -) -> int: - if not isinstance(summary, dict): - return 0 - if not todo_summary_monitor_writeback_supported(summary): - return 0 - projected_count = summary.get("monitor_due_count") - if isinstance(projected_count, int): - return max(0, projected_count) - agent_id = todo_summary_claim_scope_agent_id(summary) - if agent_id: - raw_items = summary.get("monitor_open_items") - if isinstance(raw_items, list): - return len( - [ - item - for item in raw_items - if isinstance(item, dict) - if todo_item_is_due_monitor(item, task_text_keys=task_text_keys) - if todo_item_claimed_by_agent_or_unclaimed(item, agent_id=agent_id) - ] - ) - return len( - due_items - if due_items is not None - else todo_summary_monitor_due_items( - summary, - task_text_keys=task_text_keys, - text_mode=text_mode, - ) - ) - return len( - due_items - if due_items is not None - else todo_summary_monitor_due_items( - summary, - task_text_keys=task_text_keys, - text_mode=text_mode, - ) - ) - - -def todo_summary_monitor_schedule_gap_items( - summary: dict[str, Any] | None, - *, - task_text_keys: tuple[str, ...] = ("title", "text"), - text_mode: str = "label", -) -> list[dict[str, Any]]: - return _summary_monitor_items( - summary, - projected_key="monitor_schedule_gap_items", - predicate=lambda item: todo_item_missing_monitor_schedule( - item, - task_text_keys=task_text_keys, - ), - task_text_keys=task_text_keys, - text_mode=text_mode, - ) - - -def todo_summary_monitor_schedule_gap_count( - summary: dict[str, Any] | None, - *, - gap_items: list[dict[str, Any]] | None = None, - task_text_keys: tuple[str, ...] = ("title", "text"), - text_mode: str = "label", -) -> int: - if not isinstance(summary, dict): - return 0 - if not todo_summary_monitor_writeback_supported(summary): - return 0 - agent_id = todo_summary_claim_scope_agent_id(summary) - if agent_id: - raw_items = summary.get("monitor_open_items") - if isinstance(raw_items, list): - return len( - [ - item - for item in raw_items - if isinstance(item, dict) - if todo_item_missing_monitor_schedule( - item, - task_text_keys=task_text_keys, - ) - if todo_item_claimed_by_agent_or_unclaimed(item, agent_id=agent_id) - ] - ) - return len( - gap_items - if gap_items is not None - else todo_summary_monitor_schedule_gap_items( - summary, - task_text_keys=task_text_keys, - text_mode=text_mode, - ) - ) - projected_count = summary.get("monitor_schedule_gap_count") - if isinstance(projected_count, int): - return max(0, projected_count) - return len( - gap_items - if gap_items is not None - else todo_summary_monitor_schedule_gap_items( - summary, - task_text_keys=task_text_keys, - text_mode=text_mode, - ) - ) - - -def todo_summary_open_count(summary: dict[str, Any] | None) -> int: - if not isinstance(summary, dict): - return 0 - try: - return max(0, int(summary.get("open_count") or 0)) - except (TypeError, ValueError): - return 0 - - -def todo_summary_open_task_counts(summary: dict[str, Any] | None) -> dict[str, int]: - open_count = todo_summary_open_count(summary) - classified_items: list[dict[str, Any]] = [] - seen: set[tuple[Any, str]] = set() - executable_backlog_items: list[dict[str, Any]] | None = None - monitor_open_items: list[dict[str, Any]] | None = None - if isinstance(summary, dict): - raw_executable_backlog = summary.get("executable_backlog_items") - if isinstance(raw_executable_backlog, list): - executable_backlog_items = [ - item - for item in raw_executable_backlog - if isinstance(item, dict) - if todo_item_is_actionable_open(item) - if todo_item_task_class(item) == TODO_TASK_CLASS_ADVANCEMENT - ] - raw_monitor_open = summary.get("monitor_open_items") - if isinstance(raw_monitor_open, list): - monitor_open_items = [ - item - for item in raw_monitor_open - if isinstance(item, dict) - if todo_item_is_actionable_open(item) - if todo_item_task_class(item) == TODO_TASK_CLASS_MONITOR - ] - for key in ( - "first_executable_items", - "first_open_items", - "monitor_open_items", - ): - source_items = summary.get(key) - if not isinstance(source_items, list): - continue - for item in source_items: - if not isinstance(item, dict): - continue - text = str(item.get("text") or "").strip() - if not text: - continue - identity = (item.get("index"), text) - if identity in seen: - continue - seen.add(identity) - classified_items.append(item) - if executable_backlog_items is not None: - advancement_count = len(executable_backlog_items) - else: - visible_open = min(open_count, len(classified_items)) - advancement_visible_count = sum( - 1 - for item in classified_items[:visible_open] - if todo_item_is_actionable_open(item) - and todo_item_task_class(item) == TODO_TASK_CLASS_ADVANCEMENT - ) - hidden_count = max(0, open_count - visible_open) - advancement_count = advancement_visible_count + hidden_count - if monitor_open_items is not None: - monitor_visible_count = len(monitor_open_items) - else: - visible_open = min(open_count, len(classified_items)) - monitor_visible_count = sum( - 1 - for item in classified_items[:visible_open] - if todo_item_is_actionable_open(item) - and todo_item_task_class(item) == TODO_TASK_CLASS_MONITOR - ) - hidden_count = max(0, open_count - len(classified_items)) - return { - "open": open_count, - "advancement": advancement_count, - "monitor": monitor_visible_count, - "monitor_due": todo_summary_monitor_due_count(summary), - "monitor_schedule_gap": todo_summary_monitor_schedule_gap_count(summary), - "hidden": hidden_count, - } - - -def todo_summary_has_only_future_scoped_monitor_work( - summary: dict[str, Any] | None, -) -> bool: - """Return true when the scoped agent has only non-due monitor work left.""" - - agent_id = todo_summary_claim_scope_agent_id(summary) - if not agent_id or not isinstance(summary, dict): - return False - if not todo_summary_monitor_items(summary): - return False - if todo_summary_monitor_due_count(summary) > 0: - return False - if todo_summary_monitor_schedule_gap_count(summary) > 0: - return False - if _positive_int(summary.get("current_agent_claimed_advancement_count")) > 0: - return False - - for key in ( - "current_agent_claimed_advancement_items", - "unclaimed_priority_open_items", - "first_executable_items", - "executable_backlog_items", - ): - values = summary.get(key) - if not isinstance(values, list): - continue - for item in values: - if not isinstance(item, dict): - continue - if not todo_item_is_actionable_open(item): - continue - if todo_item_task_class(item) != TODO_TASK_CLASS_ADVANCEMENT: - continue - if todo_item_claimed_by_agent_or_unclaimed(item, agent_id=agent_id): - return False - return True - - -def _positive_int(value: Any) -> int: - try: - parsed = int(value) - except (TypeError, ValueError): - return 0 - return max(0, parsed) - - -def todo_summary_first_executable_item( - summary: dict[str, Any] | None, -) -> dict[str, Any] | None: - if not isinstance(summary, dict): - return None - raw_items = summary.get("first_executable_items") - items = raw_items if isinstance(raw_items, list) else [] - for item in items: - if not isinstance(item, dict): - continue - if not todo_item_is_actionable_open(item): - continue - if todo_item_task_class(item) != TODO_TASK_CLASS_ADVANCEMENT: - continue - return item - return None +from .todo_semantics import * # noqa: F401,F403 diff --git a/loopx/control_plane/todos/projection_delivery.ts b/loopx/control_plane/todos/projection_delivery.ts new file mode 100644 index 0000000000..eea6577409 --- /dev/null +++ b/loopx/control_plane/todos/projection_delivery.ts @@ -0,0 +1,22 @@ +/** Canonical projection-delivery state returned by Todo mutations. */ +export type TodoProjectionDelivery = "pending" | "delivered" | "current" | "not_required"; +const PROJECTION_DELIVERY_VALUES = new Set([ + "pending", "delivered", "current", "not_required", +]); + +/** Keep mutation results consistent and make the no-op meaning explicit. */ +export function projectionDelivery(changed: boolean): TodoProjectionDelivery { + return changed ? "pending" : "not_required"; +} + +/** Decode provider readback without letting ad-hoc strings cross the boundary. */ +export function parseProjectionDelivery(value: unknown): TodoProjectionDelivery { + if (typeof value === "string" && PROJECTION_DELIVERY_VALUES.has(value as TodoProjectionDelivery)) { + return value as TodoProjectionDelivery; + } + throw new Error(`projection_delivery is unsupported: ${String(value)}`); +} + +export function isProjectionDelivery(value: unknown): value is TodoProjectionDelivery { + return typeof value === "string" && PROJECTION_DELIVERY_VALUES.has(value as TodoProjectionDelivery); +} diff --git a/loopx/control_plane/todos/provider_projection.py b/loopx/control_plane/todos/provider_projection.py index 987fc67773..bdf88ace6d 100644 --- a/loopx/control_plane/todos/provider_projection.py +++ b/loopx/control_plane/todos/provider_projection.py @@ -14,6 +14,7 @@ import json import stat import tempfile +from enum import StrEnum from collections.abc import Mapping from pathlib import Path from typing import Any @@ -35,6 +36,37 @@ TODO_PROJECTION_DELIVERY_SCHEMA = "loopx_todo_projection_delivery_v0" +class ProjectionDeliveryStatus(StrEnum): + """Stable cross-language states for canonical Todo display delivery.""" + PENDING = "pending" + DELIVERED = "delivered" + CURRENT = "current" + NOT_REQUIRED = "not_required" + + +def parse_projection_delivery(value: object) -> ProjectionDeliveryStatus: + try: + return ProjectionDeliveryStatus(value) # type: ignore[arg-type] + except (TypeError, ValueError) as error: + raise ValueError(f"unsupported projection_delivery: {value!r}") from error + + +def projection_delivery_requires_ack(value: object) -> bool: + return parse_projection_delivery(value) in { + ProjectionDeliveryStatus.DELIVERED, + ProjectionDeliveryStatus.CURRENT, + } + + +def projection_delivery_for_mutation(changed: bool) -> ProjectionDeliveryStatus: + """Map a committed mutation to its display outbox state.""" + return ( + ProjectionDeliveryStatus.PENDING + if changed + else ProjectionDeliveryStatus.NOT_REQUIRED + ) + + def _read_text_exact(path: Path) -> str: with path.open("r", encoding="utf-8", newline="") as handle: return handle.read() @@ -202,10 +234,10 @@ def settle_canonical_todo_projection( """Drain the committed provider head, preserving a successful mutation.""" if payload.get("dry_run") is True or payload.get("status") == "planned": - payload["projection_delivery"] = "not_required" + payload["projection_delivery"] = projection_delivery_for_mutation(False).value payload["projection_outbox"] = { "schema_version": TODO_PROJECTION_DELIVERY_SCHEMA, - "status": "not_required", + "status": ProjectionDeliveryStatus.NOT_REQUIRED.value, "source": "committed_authority_journal", } return payload @@ -246,13 +278,17 @@ def settle_canonical_todo_projection( delivery["trigger_provider_revision"] = trigger_revision if isinstance(trigger_cursor, str) and trigger_cursor: delivery["trigger_cursor"] = trigger_cursor - payload["projection_delivery"] = delivery["status"] + payload["projection_delivery"] = parse_projection_delivery(delivery["status"]).value payload["projection_outbox"] = delivery return payload __all__ = [ "TODO_PROJECTION_DELIVERY_SCHEMA", + "ProjectionDeliveryStatus", + "parse_projection_delivery", + "projection_delivery_requires_ack", + "projection_delivery_for_mutation", "project_current_canonical_todos", "settle_canonical_todo_projection", ] diff --git a/loopx/control_plane/todos/provider_terminal_lifecycle.py b/loopx/control_plane/todos/provider_terminal_lifecycle.py index c911eabf44..f6b9556cb8 100644 --- a/loopx/control_plane/todos/provider_terminal_lifecycle.py +++ b/loopx/control_plane/todos/provider_terminal_lifecycle.py @@ -36,7 +36,7 @@ from .contract import resolve_next_user_task_class from .mutation_authority import normalize_todo_lifecycle_authority from .path_resolution import resolve_todo_state_path -from .provider_projection import settle_canonical_todo_projection +from .provider_projection import projection_delivery_requires_ack, settle_canonical_todo_projection from .successor_derivation import build_successor_intents _TERMINAL_REQUEST_SCHEMA = "loopx_local_coordination_todo_terminal_lifecycle_request_v0" @@ -564,7 +564,7 @@ def archive_canonical_todos_if_promoted( if ( not dry_run and response.get("moved_count", 0) > 0 - and response.get("projection_delivery") in {"delivered", "current"} + and projection_delivery_requires_ack(response.get("projection_delivery")) ): # The native owner retains the attempt until its external projection # provider succeeds. An ACK failure must preserve the committed result diff --git a/loopx/control_plane/todos/todo_semantics.py b/loopx/control_plane/todos/todo_semantics.py new file mode 100644 index 0000000000..048c0fa04e --- /dev/null +++ b/loopx/control_plane/todos/todo_semantics.py @@ -0,0 +1,813 @@ +from __future__ import annotations + +from datetime import datetime +import re +from typing import Any + +from ..scheduler.monitor_todo import ( + monitor_todo_expires_at, + monitor_todo_has_schedule, + monitor_todo_is_actionable_open, + monitor_todo_is_due, + monitor_todo_is_expired, + monitor_todo_missing_schedule, + monitor_todo_next_due_at, + monitor_todo_task_class, +) +from .contract import ( + TODO_STATUS_DEFERRED, + TODO_TASK_CLASS_ADVANCEMENT, + TODO_TASK_CLASS_MONITOR, + normalize_todo_claimed_by, + normalize_todo_excluded_agents, + normalize_removed_todo_continuation_policy, + normalize_todo_id, + normalize_todo_status, + normalize_todo_watch_only, +) + + +TODO_MISSING_PRIORITY_RANK = 50 +TODO_MISSING_INDEX = 999999 +TODO_PRIORITY_PREFIX_PATTERN = re.compile( + r"^\s*\[(P[0-4][^\]]*)\]\s*(.+)$", + re.IGNORECASE, +) +TODO_PRIORITY_LABEL_PATTERN = re.compile(r"\bP([0-4])\b", re.IGNORECASE) + + +def todo_item_is_watch_only_monitor(item: dict[str, Any]) -> bool: + return bool( + todo_item_task_class(item) == TODO_TASK_CLASS_MONITOR + and normalize_todo_watch_only(item.get("watch_only")) is True + ) + + +def todo_priority_parts(text: str) -> tuple[str | None, str]: + match = TODO_PRIORITY_PREFIX_PATTERN.match(text) + if not match: + return None, text + return match.group(1).strip().upper(), match.group(2).strip() + + +def todo_priority_label( + item: dict[str, Any], + *, + text_mode: str = "label", +) -> str | None: + priority = item.get("priority") + if isinstance(priority, str) and priority.strip(): + return priority.strip().upper() + text = " ".join( + str(value or "") + for value in (item.get("title"), item.get("text")) + if str(value or "").strip() + ) + if text_mode == "prefix": + priority, _ = todo_priority_parts(text) + return priority + match = TODO_PRIORITY_LABEL_PATTERN.search(text.upper()) + if not match: + return None + return f"P{match.group(1)}" + + +def todo_priority_rank(value: Any, *, text_mode: str = "label") -> int: + if isinstance(value, dict): + priority = todo_priority_label(value, text_mode=text_mode) + elif isinstance(value, str): + priority = value.strip().upper() + else: + priority = None + if not priority: + return TODO_MISSING_PRIORITY_RANK + match = re.match(r"P([0-4])", priority) + if not match: + return TODO_MISSING_PRIORITY_RANK + return int(match.group(1)) + + +def todo_index_rank(item: dict[str, Any]) -> int: + raw_index = item.get("index") + try: + return int(raw_index) if raw_index is not None else TODO_MISSING_INDEX + except (TypeError, ValueError): + return TODO_MISSING_INDEX + + +def todo_projection_sort_key( + item: dict[str, Any], + *, + text_mode: str = "label", +) -> tuple[int, int]: + return (todo_priority_rank(item, text_mode=text_mode), todo_index_rank(item)) + + +def todo_claimed_visibility_items( + items: list[dict[str, Any]], + *, + limit: int, +) -> list[dict[str, Any]]: + if limit <= 0 or len(items) <= limit: + return items[:limit] + claim_order: list[str] = [] + buckets: dict[str, list[dict[str, Any]]] = {} + for item in items: + claimed_by = normalize_todo_claimed_by(item.get("claimed_by")) + if not claimed_by: + continue + if claimed_by not in buckets: + buckets[claimed_by] = [] + claim_order.append(claimed_by) + buckets[claimed_by].append(item) + if not buckets: + return items[:limit] + + original_index = {id(item): index for index, item in enumerate(items)} + per_claimant_cap = max(1, limit // len(buckets)) + selected: list[dict[str, Any]] = [] + selected_ids: set[int] = set() + for claimed_by in claim_order: + taken = 0 + for item in buckets[claimed_by]: + if taken >= per_claimant_cap: + break + if len(selected) >= limit: + break + selected.append(item) + selected_ids.add(id(item)) + taken += 1 + if len(selected) >= limit: + break + + if len(selected) < limit: + for item in items: + if id(item) in selected_ids: + continue + selected.append(item) + selected_ids.add(id(item)) + if len(selected) >= limit: + break + + return sorted( + selected, key=lambda item: original_index.get(id(item), TODO_MISSING_INDEX) + )[:limit] + + +def todo_item_task_text( + item: dict[str, Any], + *, + keys: tuple[str, ...] = ("title", "text"), +) -> str: + return " ".join( + str(item.get(key) or "") for key in keys if str(item.get(key) or "").strip() + ) + + +def todo_item_task_class( + item: dict[str, Any], + *, + task_text_keys: tuple[str, ...] = ("title", "text"), +) -> str: + return monitor_todo_task_class( + item, + task_text=todo_item_task_text(item, keys=task_text_keys), + ) + + +def todo_item_is_actionable_open(item: dict[str, Any]) -> bool: + return monitor_todo_is_actionable_open(item) + + +def todo_item_is_deferred(item: dict[str, Any]) -> bool: + return (normalize_todo_status(item.get("status")) or "") == TODO_STATUS_DEFERRED + + +def todo_item_next_due_at(item: dict[str, Any]) -> datetime | None: + return monitor_todo_next_due_at(item) + + +def todo_item_has_monitor_schedule(item: dict[str, Any]) -> bool: + return monitor_todo_has_schedule(item) + + +def todo_item_expires_at(item: dict[str, Any]) -> datetime | None: + return monitor_todo_expires_at(item) + + +def todo_item_is_expired_monitor( + item: dict[str, Any], *, now: datetime | None = None +) -> bool: + return monitor_todo_is_expired(item, now=now) + + +def todo_item_is_due_monitor( + item: dict[str, Any], + *, + now: datetime | None = None, + task_text_keys: tuple[str, ...] = ("title", "text"), +) -> bool: + return monitor_todo_is_due( + item, + now=now, + task_text=todo_item_task_text(item, keys=task_text_keys), + ) + + +def todo_item_missing_monitor_schedule( + item: dict[str, Any], + *, + now: datetime | None = None, + task_text_keys: tuple[str, ...] = ("title", "text"), +) -> bool: + return monitor_todo_missing_schedule( + item, + now=now, + task_text=todo_item_task_text(item, keys=task_text_keys), + ) + + +def todo_item_claimed_by_agent_or_unclaimed( + item: dict[str, Any], + *, + agent_id: str | None, +) -> bool: + if todo_item_has_removed_continuation_policy(item): + return False + normalized_agent_id = normalize_todo_claimed_by(agent_id) + if not normalized_agent_id: + return True + if normalized_agent_id in normalize_todo_excluded_agents( + item.get("excluded_agents") + ): + return False + claimed_by = normalize_todo_claimed_by(item.get("claimed_by")) + return not claimed_by or claimed_by == normalized_agent_id + + +def todo_advancement_frontier_items( + summary: dict[str, Any] | None, + *, + agent_id: str | None, +) -> dict[str, list[dict[str, Any]]]: + """Return the authoritative advancement frontier items grouped by claim ownership. + + Preserves the slot precedence of executable backlog first, falling back to + unclaimed priority and claimed advancement open items when the executable backlog + is omitted. Peer-claimed items are tracked separately and excluded from the current + agent's selectable advancement frontier. + """ + + empty: dict[str, list[dict[str, Any]]] = { + "current_agent_claimed_items": [], + "unclaimed_items": [], + "other_agent_claimed_items": [], + } + if not isinstance(summary, dict): + return empty + + normalized_agent_id = normalize_todo_claimed_by(agent_id) + executable_items = summary.get("executable_backlog_items") + if isinstance(executable_items, list): + current_items: list[dict[str, Any]] = [] + unclaimed_items: list[dict[str, Any]] = [] + other_items: list[dict[str, Any]] = [] + for value in executable_items: + if not isinstance(value, dict): + continue + if not todo_item_is_actionable_open(value): + continue + if todo_item_task_class(value) != TODO_TASK_CLASS_ADVANCEMENT: + continue + claimed_by = normalize_todo_claimed_by(value.get("claimed_by")) + if claimed_by: + if normalized_agent_id and claimed_by == normalized_agent_id: + if not todo_item_excludes_agent( + value, agent_id=normalized_agent_id + ): + current_items.append(value) + elif normalized_agent_id: + other_items.append(value) + else: + current_items.append(value) + continue + if not todo_item_excludes_agent(value, agent_id=normalized_agent_id): + unclaimed_items.append(value) + return { + "current_agent_claimed_items": current_items, + "unclaimed_items": unclaimed_items, + "other_agent_claimed_items": other_items, + } + + unclaimed_items = [ + value + for value in summary.get("unclaimed_priority_open_items") or [] + if isinstance(value, dict) + and todo_item_is_actionable_open(value) + and todo_item_task_class(value) == TODO_TASK_CLASS_ADVANCEMENT + and not todo_item_excludes_agent(value, agent_id=normalized_agent_id) + ] + current_items = [ + value + for value in summary.get("claimed_advancement_open_items") or [] + if isinstance(value, dict) + and todo_item_is_actionable_open(value) + and todo_item_task_class(value) == TODO_TASK_CLASS_ADVANCEMENT + and ( + not normalized_agent_id + or normalize_todo_claimed_by(value.get("claimed_by")) == normalized_agent_id + ) + and not todo_item_excludes_agent(value, agent_id=normalized_agent_id) + ] + other_items = [ + value + for value in summary.get("claimed_advancement_open_items") or [] + if isinstance(value, dict) + and todo_item_is_actionable_open(value) + and todo_item_task_class(value) == TODO_TASK_CLASS_ADVANCEMENT + and normalized_agent_id + and normalize_todo_claimed_by(value.get("claimed_by")) + and normalize_todo_claimed_by(value.get("claimed_by")) != normalized_agent_id + ] + return { + "current_agent_claimed_items": current_items, + "unclaimed_items": unclaimed_items, + "other_agent_claimed_items": other_items, + } + + +def agent_scoped_selectable_advancement_todo_ids( + agent_todo_summary: dict[str, Any] | None, + *, + agent_id: str | None, +) -> set[str]: + """Return the ids the agent-scoped selectable advancement frontier holds. + + Derived directly from the authoritative ``todo_advancement_frontier_items`` + helper so that slot precedence and claim ownership predicates never diverge + from the frontier counter. + """ + + frontier_items = todo_advancement_frontier_items( + agent_todo_summary, + agent_id=agent_id, + ) + selectable: set[str] = set() + for item in ( + frontier_items["current_agent_claimed_items"] + + frontier_items["unclaimed_items"] + ): + if todo_id := normalize_todo_id(item.get("todo_id")): + selectable.add(todo_id) + return selectable + + +def todo_advancement_frontier_counts( + summary: dict[str, Any] | None, + *, + agent_id: str | None, +) -> dict[str, int]: + """Classify the durable advancement frontier by exact claim ownership.""" + + if not isinstance(summary, dict): + return { + "current_agent_claimed_advancement_count": 0, + "unclaimed_advancement_count": 0, + "other_agent_claimed_advancement_count": 0, + } + frontier_items = todo_advancement_frontier_items(summary, agent_id=agent_id) + claim_scope = summary.get("claim_scope") + other_items = ( + claim_scope.get("other_agent_claimed_items") + if isinstance(claim_scope, dict) + else [] + ) + diagnostic_other_count = sum( + 1 + for value in other_items or [] + if isinstance(value, dict) + and todo_item_is_actionable_open(value) + and todo_item_task_class(value) == TODO_TASK_CLASS_ADVANCEMENT + ) + return { + "current_agent_claimed_advancement_count": max( + len(frontier_items["current_agent_claimed_items"]), + _positive_int(summary.get("current_agent_claimed_advancement_count")), + ), + "unclaimed_advancement_count": len(frontier_items["unclaimed_items"]), + "other_agent_claimed_advancement_count": max( + len(frontier_items["other_agent_claimed_items"]), + diagnostic_other_count, + ), + } + + +def todo_item_has_removed_continuation_policy(item: dict[str, Any]) -> bool: + return bool( + normalize_removed_todo_continuation_policy( + item.get("removed_continuation_policy") + ) + ) + + +def todo_item_excludes_agent( + item: dict[str, Any], + *, + agent_id: str | None, +) -> bool: + normalized_agent_id = normalize_todo_claimed_by(agent_id) + return bool( + normalized_agent_id + and normalized_agent_id + in normalize_todo_excluded_agents(item.get("excluded_agents")) + ) + + +def todo_summary_claim_scope_agent_id(summary: dict[str, Any] | None) -> str | None: + if not isinstance(summary, dict): + return None + claim_scope = summary.get("claim_scope") + if not isinstance(claim_scope, dict): + return None + return normalize_todo_claimed_by(claim_scope.get("agent_id")) + + +def todo_summary_monitor_writeback_contract( + summary: dict[str, Any] | None, +) -> dict[str, Any] | None: + if not isinstance(summary, dict): + return None + contract = summary.get("monitor_writeback") + if not isinstance(contract, dict): + return None + if contract.get("supported") is not False: + return None + compact: dict[str, Any] = {"supported": False} + source = str(contract.get("source") or "").strip() + if source: + compact["source"] = source + return compact + + +def todo_summary_monitor_writeback_supported(summary: dict[str, Any] | None) -> bool: + contract = todo_summary_monitor_writeback_contract(summary) + if not contract: + return True + return contract.get("supported") is not False + + +def todo_summary_monitor_items(summary: dict[str, Any] | None) -> list[dict[str, Any]]: + if not isinstance(summary, dict): + return [] + items: list[dict[str, Any]] = [] + seen: set[tuple[str, int]] = set() + for key in ( + "monitor_due_items", + "current_agent_claimed_monitor_items", + "monitor_open_items", + "claimed_monitor_open_items", + "first_open_items", + ): + values = summary.get(key) + if not isinstance(values, list): + continue + for value in values: + if not isinstance(value, dict): + continue + if not todo_item_is_actionable_open(value): + continue + if todo_item_task_class(value) != TODO_TASK_CLASS_MONITOR: + continue + identity = (normalize_todo_id(value.get("todo_id")) or "", id(value)) + if identity in seen: + continue + seen.add(identity) + items.append(value) + return items + + +def _summary_monitor_items( + summary: dict[str, Any] | None, + *, + projected_key: str, + predicate: Any, + task_text_keys: tuple[str, ...], + text_mode: str, +) -> list[dict[str, Any]]: + if not isinstance(summary, dict): + return [] + if not todo_summary_monitor_writeback_supported(summary): + return [] + projected_items = summary.get(projected_key) + if isinstance(projected_items, list): + items = [ + item + for item in projected_items + if isinstance(item, dict) + if todo_item_is_actionable_open(item) + if todo_item_task_class(item, task_text_keys=task_text_keys) + == TODO_TASK_CLASS_MONITOR + if predicate(item) + ] + else: + raw_items = summary.get("monitor_open_items") + items = [ + item + for item in (raw_items if isinstance(raw_items, list) else []) + if isinstance(item, dict) + if predicate(item) + ] + agent_id = todo_summary_claim_scope_agent_id(summary) + if agent_id: + items = [ + item + for item in items + if todo_item_claimed_by_agent_or_unclaimed(item, agent_id=agent_id) + ] + return sorted( + items, + key=lambda item: todo_projection_sort_key(item, text_mode=text_mode), + ) + + +def todo_summary_monitor_due_items( + summary: dict[str, Any] | None, + *, + task_text_keys: tuple[str, ...] = ("title", "text"), + text_mode: str = "label", +) -> list[dict[str, Any]]: + return _summary_monitor_items( + summary, + projected_key="monitor_due_items", + predicate=lambda item: todo_item_is_due_monitor( + item, + task_text_keys=task_text_keys, + ), + task_text_keys=task_text_keys, + text_mode=text_mode, + ) + + +def todo_summary_monitor_due_count( + summary: dict[str, Any] | None, + *, + due_items: list[dict[str, Any]] | None = None, + task_text_keys: tuple[str, ...] = ("title", "text"), + text_mode: str = "label", +) -> int: + if not isinstance(summary, dict): + return 0 + if not todo_summary_monitor_writeback_supported(summary): + return 0 + projected_count = summary.get("monitor_due_count") + if isinstance(projected_count, int): + return max(0, projected_count) + agent_id = todo_summary_claim_scope_agent_id(summary) + if agent_id: + raw_items = summary.get("monitor_open_items") + if isinstance(raw_items, list): + return len( + [ + item + for item in raw_items + if isinstance(item, dict) + if todo_item_is_due_monitor(item, task_text_keys=task_text_keys) + if todo_item_claimed_by_agent_or_unclaimed(item, agent_id=agent_id) + ] + ) + return len( + due_items + if due_items is not None + else todo_summary_monitor_due_items( + summary, + task_text_keys=task_text_keys, + text_mode=text_mode, + ) + ) + return len( + due_items + if due_items is not None + else todo_summary_monitor_due_items( + summary, + task_text_keys=task_text_keys, + text_mode=text_mode, + ) + ) + + +def todo_summary_monitor_schedule_gap_items( + summary: dict[str, Any] | None, + *, + task_text_keys: tuple[str, ...] = ("title", "text"), + text_mode: str = "label", +) -> list[dict[str, Any]]: + return _summary_monitor_items( + summary, + projected_key="monitor_schedule_gap_items", + predicate=lambda item: todo_item_missing_monitor_schedule( + item, + task_text_keys=task_text_keys, + ), + task_text_keys=task_text_keys, + text_mode=text_mode, + ) + + +def todo_summary_monitor_schedule_gap_count( + summary: dict[str, Any] | None, + *, + gap_items: list[dict[str, Any]] | None = None, + task_text_keys: tuple[str, ...] = ("title", "text"), + text_mode: str = "label", +) -> int: + if not isinstance(summary, dict): + return 0 + if not todo_summary_monitor_writeback_supported(summary): + return 0 + agent_id = todo_summary_claim_scope_agent_id(summary) + if agent_id: + raw_items = summary.get("monitor_open_items") + if isinstance(raw_items, list): + return len( + [ + item + for item in raw_items + if isinstance(item, dict) + if todo_item_missing_monitor_schedule( + item, + task_text_keys=task_text_keys, + ) + if todo_item_claimed_by_agent_or_unclaimed(item, agent_id=agent_id) + ] + ) + return len( + gap_items + if gap_items is not None + else todo_summary_monitor_schedule_gap_items( + summary, + task_text_keys=task_text_keys, + text_mode=text_mode, + ) + ) + projected_count = summary.get("monitor_schedule_gap_count") + if isinstance(projected_count, int): + return max(0, projected_count) + return len( + gap_items + if gap_items is not None + else todo_summary_monitor_schedule_gap_items( + summary, + task_text_keys=task_text_keys, + text_mode=text_mode, + ) + ) + + +def todo_summary_open_count(summary: dict[str, Any] | None) -> int: + if not isinstance(summary, dict): + return 0 + try: + return max(0, int(summary.get("open_count") or 0)) + except (TypeError, ValueError): + return 0 + + +def todo_summary_open_task_counts(summary: dict[str, Any] | None) -> dict[str, int]: + open_count = todo_summary_open_count(summary) + classified_items: list[dict[str, Any]] = [] + seen: set[tuple[Any, str]] = set() + executable_backlog_items: list[dict[str, Any]] | None = None + monitor_open_items: list[dict[str, Any]] | None = None + if isinstance(summary, dict): + raw_executable_backlog = summary.get("executable_backlog_items") + if isinstance(raw_executable_backlog, list): + executable_backlog_items = [ + item + for item in raw_executable_backlog + if isinstance(item, dict) + if todo_item_is_actionable_open(item) + if todo_item_task_class(item) == TODO_TASK_CLASS_ADVANCEMENT + ] + raw_monitor_open = summary.get("monitor_open_items") + if isinstance(raw_monitor_open, list): + monitor_open_items = [ + item + for item in raw_monitor_open + if isinstance(item, dict) + if todo_item_is_actionable_open(item) + if todo_item_task_class(item) == TODO_TASK_CLASS_MONITOR + ] + for key in ( + "first_executable_items", + "first_open_items", + "monitor_open_items", + ): + source_items = summary.get(key) + if not isinstance(source_items, list): + continue + for item in source_items: + if not isinstance(item, dict): + continue + text = str(item.get("text") or "").strip() + if not text: + continue + identity = (item.get("index"), text) + if identity in seen: + continue + seen.add(identity) + classified_items.append(item) + if executable_backlog_items is not None: + advancement_count = len(executable_backlog_items) + else: + visible_open = min(open_count, len(classified_items)) + advancement_visible_count = sum( + 1 + for item in classified_items[:visible_open] + if todo_item_is_actionable_open(item) + and todo_item_task_class(item) == TODO_TASK_CLASS_ADVANCEMENT + ) + hidden_count = max(0, open_count - visible_open) + advancement_count = advancement_visible_count + hidden_count + if monitor_open_items is not None: + monitor_visible_count = len(monitor_open_items) + else: + visible_open = min(open_count, len(classified_items)) + monitor_visible_count = sum( + 1 + for item in classified_items[:visible_open] + if todo_item_is_actionable_open(item) + and todo_item_task_class(item) == TODO_TASK_CLASS_MONITOR + ) + hidden_count = max(0, open_count - len(classified_items)) + return { + "open": open_count, + "advancement": advancement_count, + "monitor": monitor_visible_count, + "monitor_due": todo_summary_monitor_due_count(summary), + "monitor_schedule_gap": todo_summary_monitor_schedule_gap_count(summary), + "hidden": hidden_count, + } + + +def todo_summary_has_only_future_scoped_monitor_work( + summary: dict[str, Any] | None, +) -> bool: + """Return true when the scoped agent has only non-due monitor work left.""" + + agent_id = todo_summary_claim_scope_agent_id(summary) + if not agent_id or not isinstance(summary, dict): + return False + if not todo_summary_monitor_items(summary): + return False + if todo_summary_monitor_due_count(summary) > 0: + return False + if todo_summary_monitor_schedule_gap_count(summary) > 0: + return False + if _positive_int(summary.get("current_agent_claimed_advancement_count")) > 0: + return False + + for key in ( + "current_agent_claimed_advancement_items", + "unclaimed_priority_open_items", + "first_executable_items", + "executable_backlog_items", + ): + values = summary.get(key) + if not isinstance(values, list): + continue + for item in values: + if not isinstance(item, dict): + continue + if not todo_item_is_actionable_open(item): + continue + if todo_item_task_class(item) != TODO_TASK_CLASS_ADVANCEMENT: + continue + if todo_item_claimed_by_agent_or_unclaimed(item, agent_id=agent_id): + return False + return True + + +def _positive_int(value: Any) -> int: + try: + parsed = int(value) + except (TypeError, ValueError): + return 0 + return max(0, parsed) + + +def todo_summary_first_executable_item( + summary: dict[str, Any] | None, +) -> dict[str, Any] | None: + if not isinstance(summary, dict): + return None + raw_items = summary.get("first_executable_items") + items = raw_items if isinstance(raw_items, list) else [] + for item in items: + if not isinstance(item, dict): + continue + if not todo_item_is_actionable_open(item): + continue + if todo_item_task_class(item) != TODO_TASK_CLASS_ADVANCEMENT: + continue + return item + return None diff --git a/loopx/control_plane/todos/todo_summary.py b/loopx/control_plane/todos/todo_summary.py index 8558b0e867..6b5524785d 100644 --- a/loopx/control_plane/todos/todo_summary.py +++ b/loopx/control_plane/todos/todo_summary.py @@ -547,7 +547,7 @@ def compact_active_next_action_todo_item(item: dict[str, Any]) -> dict[str, Any] def todo_item_task_class(item: dict[str, Any]) -> str: - return projection_todo_item_task_class(item, task_text_keys=("text",)) + return projection_todo_item_task_class(item) def count_advancement_todos(items: list[dict[str, Any]]) -> int: diff --git a/loopx/control_plane/turn_driver/delivery_continuity.py b/loopx/control_plane/turn_driver/delivery_continuity.py index 4a6643dffc..52dc693874 100644 --- a/loopx/control_plane/turn_driver/delivery_continuity.py +++ b/loopx/control_plane/turn_driver/delivery_continuity.py @@ -16,7 +16,7 @@ normalize_todo_id, normalize_todo_status, ) -from ..todos.projection import todo_item_task_class +from ..todos.todo_semantics import todo_item_task_class DELIVERY_BOUNDARY_IN_FLIGHT = "in_flight_continuation" DELIVERY_BOUNDARY_SEMANTIC_CLOSEOUT = "semantic_closeout" diff --git a/loopx/control_plane/work_items/capability_monitor_fallback.py b/loopx/control_plane/work_items/capability_monitor_fallback.py index c830ef450c..801fb620c0 100644 --- a/loopx/control_plane/work_items/capability_monitor_fallback.py +++ b/loopx/control_plane/work_items/capability_monitor_fallback.py @@ -4,7 +4,7 @@ from ..agents.capability_gate import build_capability_gate from ..todos.contract import TODO_TASK_CLASS_ADVANCEMENT, TODO_TASK_CLASS_MONITOR -from ..todos.projection import todo_item_task_class +from ..todos.todo_semantics import todo_item_task_class from ..todos.summary_item import compact_todo_summary_item diff --git a/loopx/control_plane/work_items/delivery_history.py b/loopx/control_plane/work_items/delivery_history.py index 5de7cf52fb..e0cc34b590 100644 --- a/loopx/control_plane/work_items/delivery_history.py +++ b/loopx/control_plane/work_items/delivery_history.py @@ -88,7 +88,7 @@ def project_delivery_response( ) -> dict[str, Any]: """Select a canonical source row; TS alone decides its supervision meaning.""" from ..todos.summary_item import todo_planning_source_items - from ..todos.projection import todo_summary_claim_scope_agent_id + from ..todos.todo_semantics import todo_summary_claim_scope_agent_id source = next((item for item in todo_planning_source_items(summary, include_terminal=True) if item.get("todo_id") == run.get("todo_id")), None) if summary else None diff --git a/loopx/control_plane/work_items/interaction_contract.py b/loopx/control_plane/work_items/interaction_contract.py index 79d38377ff..6fcf2fd5e4 100644 --- a/loopx/control_plane/work_items/interaction_contract.py +++ b/loopx/control_plane/work_items/interaction_contract.py @@ -30,7 +30,7 @@ normalize_todo_id, normalize_todo_replan_obligation_id, ) -from ..todos.projection import todo_item_task_class +from ..todos.todo_semantics import todo_item_task_class from ..todos.user_gate import open_todo_count from ..todos.write_hint import build_capability_resolution_writeback_actions from .autonomous_replan_obligation import ( diff --git a/loopx/control_plane/work_items/planning_inventory.py b/loopx/control_plane/work_items/planning_inventory.py index 5d115d5e41..12283fd730 100644 --- a/loopx/control_plane/work_items/planning_inventory.py +++ b/loopx/control_plane/work_items/planning_inventory.py @@ -9,7 +9,7 @@ normalize_todo_claimed_by, normalize_todo_id, ) -from ..todos.projection import ( +from ..todos.todo_semantics import ( todo_item_is_actionable_open, todo_item_task_class, ) diff --git a/loopx/control_plane/work_items/primary_action.py b/loopx/control_plane/work_items/primary_action.py index d49a54ea74..a4c53af5f7 100644 --- a/loopx/control_plane/work_items/primary_action.py +++ b/loopx/control_plane/work_items/primary_action.py @@ -7,7 +7,7 @@ agent_scope_frontier_action as _agent_scope_frontier_action, ) from ..todos.contract import TODO_TASK_CLASS_ADVANCEMENT -from ..todos.projection import todo_item_is_actionable_open, todo_item_task_class +from ..todos.todo_semantics import todo_item_is_actionable_open, todo_item_task_class from .autonomous_replan_obligation import todo_lifecycle_settlement_obligation diff --git a/loopx/control_plane/work_items/repair_delta.py b/loopx/control_plane/work_items/repair_delta.py index 5a64b26c51..008e6a7037 100644 --- a/loopx/control_plane/work_items/repair_delta.py +++ b/loopx/control_plane/work_items/repair_delta.py @@ -19,7 +19,7 @@ normalize_todo_resume_when, normalize_todo_status, ) -from ..todos.projection import ( +from ..todos.todo_semantics import ( todo_item_claimed_by_agent_or_unclaimed, todo_item_expires_at, todo_item_is_actionable_open, diff --git a/loopx/control_plane/work_items/user_action_frontier.py b/loopx/control_plane/work_items/user_action_frontier.py index 3c02e1eefd..07cdaa0612 100644 --- a/loopx/control_plane/work_items/user_action_frontier.py +++ b/loopx/control_plane/work_items/user_action_frontier.py @@ -3,7 +3,7 @@ from typing import Any from ..todos.contract import TODO_TASK_CLASS_USER_ACTION -from ..todos.projection import todo_item_task_class +from ..todos.todo_semantics import todo_item_task_class def user_action_owns_empty_agent_lane(payload: dict[str, Any]) -> bool: diff --git a/loopx/control_plane/work_items/work_lane.py b/loopx/control_plane/work_items/work_lane.py index 4525333377..ea8ba4d2f4 100644 --- a/loopx/control_plane/work_items/work_lane.py +++ b/loopx/control_plane/work_items/work_lane.py @@ -4,7 +4,7 @@ from ..effect_program import ReceiptBoundMonitorPhase from ..todos.contract import TODO_TASK_CLASS_MONITOR, normalize_todo_id -from ..todos.projection import todo_priority_label, todo_priority_rank +from ..todos.todo_semantics import todo_priority_label, todo_priority_rank WORK_LANE_CONTRACT_SCHEMA_VERSION = "work_lane_contract_v1" WORK_LANE_RECEIPT_BOUND_MONITOR_SETTLEMENT_OBLIGATION = ( diff --git a/loopx/control_plane/work_items/work_lane_context.py b/loopx/control_plane/work_items/work_lane_context.py index 94aabf3bd0..248eb2c017 100644 --- a/loopx/control_plane/work_items/work_lane_context.py +++ b/loopx/control_plane/work_items/work_lane_context.py @@ -5,7 +5,7 @@ from ..agents.agent_scope import _agent_scope_monitor_blocked_resume_candidates from ..scheduler.external_evidence_observation import build_external_evidence_poll_signal from ..todos.contract import next_action_requires_advancement_text -from ..todos.projection import ( +from ..todos.todo_semantics import ( todo_summary_claim_scope_agent_id, todo_summary_first_executable_item, todo_summary_monitor_due_count, diff --git a/loopx/extensions/lark/presentation/kanban.py b/loopx/extensions/lark/presentation/kanban.py index 62296ffcfb..166f22c6f7 100644 --- a/loopx/extensions/lark/presentation/kanban.py +++ b/loopx/extensions/lark/presentation/kanban.py @@ -2196,7 +2196,7 @@ def sync_loopx_todos_to_lark_kanban( from ....capabilities.issue_fix.outcome_projection import ( build_issue_fix_outcome_collection_from_domain_state, ) - from ....control_plane.todos.projection import todo_priority_label + from ....control_plane.todos.todo_semantics import todo_priority_label from ....todos import resolve_todo_state_path, section_bounds, todo_blocks resolved_project, resolved_state_file = resolve_todo_state_path( diff --git a/loopx/quota.py b/loopx/quota.py index e84daed8d8..410cec4a87 100644 --- a/loopx/quota.py +++ b/loopx/quota.py @@ -98,7 +98,7 @@ normalize_todo_claimed_by, normalize_todo_id, ) -from .control_plane.todos.projection import ( +from .control_plane.todos.todo_semantics import ( todo_index_rank as projection_todo_index_rank, todo_item_expires_at as projection_todo_item_expires_at, todo_item_is_due_monitor as projection_todo_item_is_due_monitor, diff --git a/loopx/status.py b/loopx/status.py index 5f605fcecf..fb021fdcd8 100644 --- a/loopx/status.py +++ b/loopx/status.py @@ -196,7 +196,7 @@ normalize_todo_task_class as normalize_todo_task_class, todo_done_for_status, ) -from .control_plane.todos.projection import ( +from .control_plane.todos.todo_semantics import ( todo_item_is_expired_monitor as todo_item_is_expired_monitor, ) @@ -216,7 +216,7 @@ "todo_item_next_due_at": "loopx.control_plane.todos.todo_summary", "todo_projection_sort_key": "loopx.control_plane.todos.todo_summary", "normalize_todo_task_class": "loopx.control_plane.todos.contract", - "todo_item_is_expired_monitor": "loopx.control_plane.todos.projection", + "todo_item_is_expired_monitor": "loopx.control_plane.todos.todo_semantics", } diff --git a/tests/control_plane/test_todo_provider_projection.py b/tests/control_plane/test_todo_provider_projection.py index c3aaa05d27..8548646b64 100644 --- a/tests/control_plane/test_todo_provider_projection.py +++ b/tests/control_plane/test_todo_provider_projection.py @@ -181,3 +181,39 @@ def test_explicit_projection_fences_requested_revision( ) assert state_file.read_text(encoding="utf-8") == SOURCE + + +def test_projection_delivery_status_contract_is_strict(): + assert provider_projection.parse_projection_delivery("delivered") is provider_projection.ProjectionDeliveryStatus.DELIVERED + assert provider_projection.projection_delivery_requires_ack("current") is True + assert provider_projection.projection_delivery_requires_ack("pending") is False + with pytest.raises(ValueError, match="unsupported projection_delivery"): + provider_projection.parse_projection_delivery("completed") + + +def test_projection_delivery_composition_fixture_matches_provider_semantics(): + fixture_path = Path(__file__).parents[1] / "fixtures" / "control_plane" / "projection_delivery_composition_v0.json" + fixture = json.loads(fixture_path.read_text(encoding="utf-8")) + for case in fixture["cases"]: + status = ( + provider_projection.parse_projection_delivery(case["readback"]).value + if "readback" in case + else provider_projection.projection_delivery_for_mutation(case.get("changed", False)).value + ) + assert status == case["expected"], case["name"] + assert provider_projection.projection_delivery_requires_ack(status) is case["requires_ack"], case["name"] + + +def test_projection_delivery_e2e_fixture_preserves_causal_states(): + fixture_path = Path(__file__).parents[1] / "fixtures" / "control_plane" / "projection_delivery_e2e_v1.json" + fixture = json.loads(fixture_path.read_text(encoding="utf-8")) + observed = [] + for transition in fixture["transitions"]: + if "readback" in transition: + status = provider_projection.parse_projection_delivery(transition["readback"]).value + else: + status = provider_projection.projection_delivery_for_mutation(transition["changed"]).value + observed.append(status) + assert status == transition["delivery"], transition["step"] + assert provider_projection.projection_delivery_requires_ack(status) is transition["ack"], transition["step"] + assert observed == ["pending", "delivered", "current", "not_required", "pending"] diff --git a/tests/control_plane/test_todo_semantic_kernel.py b/tests/control_plane/test_todo_semantic_kernel.py new file mode 100644 index 0000000000..45235b12ce --- /dev/null +++ b/tests/control_plane/test_todo_semantic_kernel.py @@ -0,0 +1,40 @@ +from __future__ import annotations + +import json +from pathlib import Path + +from loopx.control_plane.todos.todo_semantics import ( + TODO_TASK_CLASS_ADVANCEMENT, + TODO_TASK_CLASS_MONITOR, + todo_item_claimed_by_agent_or_unclaimed, + todo_item_is_due_monitor, + todo_item_task_class, +) +from loopx.control_plane.todos.todo_summary import todo_item_task_class as summary_task_class + + +FIXTURE = Path(__file__).parents[1] / "fixtures/control_plane/coordination_production_scale_v0.json" + + +def test_summary_and_projection_share_title_aware_task_classification() -> None: + item = {"title": "Observe dependency health", "text": "", "status": "open"} + assert todo_item_task_class(item) == TODO_TASK_CLASS_MONITOR + assert summary_task_class(item) == TODO_TASK_CLASS_MONITOR + assert todo_item_is_due_monitor( + {**item, "next_due_at": "2025-01-01T00:00:00Z"}, + now=__import__("datetime").datetime.fromisoformat("2025-01-01T00:00:00+00:00"), + ) + + +def test_exclusion_is_part_of_selectability_not_claim_ownership() -> None: + item = {"task_class": TODO_TASK_CLASS_ADVANCEMENT, "status": "open", "excluded_agents": ["agent-a"]} + assert not todo_item_claimed_by_agent_or_unclaimed(item, agent_id="agent-a") + assert todo_item_claimed_by_agent_or_unclaimed(item, agent_id="agent-b") + + +def test_complex_fixture_declares_cross_rfc_semantic_edges() -> None: + cases = json.loads(FIXTURE.read_text())["semantic_cases"] + assert cases["global_gate_without_goal_binding"]["global_gate"] is True + assert cases["global_gate_without_goal_binding"]["goal_bound"] is False + assert cases["expired_lease"]["lease_epoch"] == 7 + assert cases["excluded_unclaimed_advancement"]["claimed_by"] is None diff --git a/tests/control_plane_ts/production_scale_coordination_fixture.ts b/tests/control_plane_ts/production_scale_coordination_fixture.ts index 27a3f6c9ab..05cab6a47a 100644 --- a/tests/control_plane_ts/production_scale_coordination_fixture.ts +++ b/tests/control_plane_ts/production_scale_coordination_fixture.ts @@ -23,6 +23,7 @@ const envelope = JSON.parse(readFileSync(new URL( linked_decision_count: number; completion_target_index: number; supersede_target_index: number; + semantic_cases: Record>; }; export const PRODUCTION_SCALE_FIXTURE_SCHEMA = @@ -48,6 +49,7 @@ export interface ProductionScaleCoordinationFixture { readonly expected_agent_archive_count_after_terminals: number; readonly expected_user_archive_count: number; readonly expected_standing_user_decision_count: number; + readonly semantic_cases: Readonly>>; } function statusSeries(counts: Record): string[] { @@ -209,6 +211,7 @@ export function productionScaleCoordinationFixture( expected_agent_archive_count_after_terminals: initialAgentDone + 2 - 5, expected_user_archive_count: (envelope.user_status_counts.done ?? 0) - 5, expected_standing_user_decision_count: envelope.standing_user_decision_count, + semantic_cases: envelope.semantic_cases, }; } diff --git a/tests/control_plane_ts/projection_delivery.test.ts b/tests/control_plane_ts/projection_delivery.test.ts new file mode 100644 index 0000000000..bc5c1a0a1c --- /dev/null +++ b/tests/control_plane_ts/projection_delivery.test.ts @@ -0,0 +1,35 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { readFile } from "node:fs/promises"; +import { isProjectionDelivery, parseProjectionDelivery, projectionDelivery } from "../../loopx/control_plane/todos/projection_delivery.ts"; + +test("projection delivery maps mutation and no-op outcomes", () => { + assert.equal(projectionDelivery(true), "pending"); + assert.equal(projectionDelivery(false), "not_required"); +}); + +test("projection delivery parser accepts provider readback states", () => { + for (const value of ["pending", "delivered", "current", "not_required"]) { + assert.equal(parseProjectionDelivery(value), value); + } + assert.throws(() => parseProjectionDelivery("unknown")); + assert.equal(isProjectionDelivery("delivered"), true); + assert.equal(isProjectionDelivery("DELIVERED"), false); + assert.equal(isProjectionDelivery(null), false); +}); + +test("composition fixture keeps mutation and provider states distinct", async () => { + const fixture = JSON.parse(await readFile("tests/fixtures/control_plane/projection_delivery_composition_v0.json", "utf8")); + for (const item of fixture.cases) { + const actual = item.changed === undefined ? parseProjectionDelivery(item.readback) : projectionDelivery(item.changed); + assert.equal(actual, item.expected, item.name); + assert.equal(item.requires_ack, actual === "delivered" || actual === "current", item.name); + } +}); + +test("end-to-end fixture preserves delivery causal chain", async () => { + const fixture = JSON.parse(await readFile("tests/fixtures/control_plane/projection_delivery_e2e_v1.json", "utf8")); + const observed = fixture.transitions.map((item: { changed?: boolean; readback?: unknown }) => + item.readback === undefined ? projectionDelivery(item.changed === true) : parseProjectionDelivery(item.readback)); + assert.deepEqual(observed, ["pending", "delivered", "current", "not_required", "pending"]); +}); diff --git a/tests/control_plane_ts/todo_semantic_fixture.test.ts b/tests/control_plane_ts/todo_semantic_fixture.test.ts new file mode 100644 index 0000000000..fffaef5fc1 --- /dev/null +++ b/tests/control_plane_ts/todo_semantic_fixture.test.ts @@ -0,0 +1,16 @@ +import test from "node:test"; +import assert from "node:assert/strict"; + +import {productionScaleCoordinationFixture} from "./production_scale_coordination_fixture.ts"; + +test("production-scale fixture carries cross-RFC semantic edge cases", () => { + const fixture = productionScaleCoordinationFixture("fixture-goal"); + const cases = fixture.semantic_cases; + assert.equal(cases.title_only_monitor.title, "Observe dependency health"); + assert.equal(cases.title_only_monitor.text, ""); + assert.equal(cases.excluded_unclaimed_advancement.claimed_by, null); + assert.deepEqual(cases.excluded_unclaimed_advancement.excluded_agents, ["agent-a"]); + assert.equal(cases.global_gate_without_goal_binding.global_gate, true); + assert.equal(cases.global_gate_without_goal_binding.goal_bound, false); + assert.equal(cases.expired_lease.lease_epoch, 7); +}); diff --git a/tests/fixtures/control_plane/coordination_production_scale_v0.json b/tests/fixtures/control_plane/coordination_production_scale_v0.json index 813203601c..5f0489bc45 100644 --- a/tests/fixtures/control_plane/coordination_production_scale_v0.json +++ b/tests/fixtures/control_plane/coordination_production_scale_v0.json @@ -17,5 +17,47 @@ "scoped_without_outcome_count": 12, "linked_decision_count": 12, "completion_target_index": 160, - "supersede_target_index": 161 + "supersede_target_index": 161, + "semantic_cases": { + "title_only_monitor": { + "todo_id": "todo_fixture_title_monitor", + "role": "agent", + "status": "open", + "done": false, + "title": "Observe dependency health", + "text": "", + "task_class": null, + "next_due_at": "2025-01-01T00:00:00Z", + "watch_only": false + }, + "excluded_unclaimed_advancement": { + "todo_id": "todo_fixture_excluded_advancement", + "role": "agent", + "status": "open", + "done": false, + "text": "Implement the isolated fixture path", + "task_class": "advancement_task", + "excluded_agents": [ + "agent-a" + ], + "claimed_by": null + }, + "global_gate_without_goal_binding": { + "todo_id": "todo_fixture_global_gate", + "role": "user", + "status": "open", + "done": false, + "text": "Approve the shared direction", + "task_class": "user_gate", + "global_gate": true, + "goal_bound": false + }, + "expired_lease": { + "todo_id": "todo_fixture_expired_lease", + "owner": "agent-a", + "status": "active", + "expires_at": "2024-01-01T00:00:00Z", + "lease_epoch": 7 + } + } } diff --git a/tests/fixtures/control_plane/projection_delivery_composition_v0.json b/tests/fixtures/control_plane/projection_delivery_composition_v0.json new file mode 100644 index 0000000000..0bc5898123 --- /dev/null +++ b/tests/fixtures/control_plane/projection_delivery_composition_v0.json @@ -0,0 +1,41 @@ +{ + "schema_version": "todo_projection_delivery_composition_v0", + "cases": [ + { + "name": "mutation_changed", + "changed": true, + "expected": "pending", + "requires_ack": false + }, + { + "name": "mutation_noop", + "changed": false, + "expected": "not_required", + "requires_ack": false + }, + { + "name": "provider_ack", + "readback": "delivered", + "expected": "delivered", + "requires_ack": true + }, + { + "name": "provider_already_current", + "readback": "current", + "expected": "current", + "requires_ack": true + }, + { + "name": "provider_pending", + "readback": "pending", + "expected": "pending", + "requires_ack": false + }, + { + "name": "mutation_noop_does_not_ack", + "changed": false, + "expected": "not_required", + "requires_ack": false + } + ] +} diff --git a/tests/fixtures/control_plane/projection_delivery_e2e_v1.json b/tests/fixtures/control_plane/projection_delivery_e2e_v1.json new file mode 100644 index 0000000000..33b750794d --- /dev/null +++ b/tests/fixtures/control_plane/projection_delivery_e2e_v1.json @@ -0,0 +1,11 @@ +{ + "schema_version": "todo_projection_delivery_e2e_v1", + "goal_id": "fixture-goal", + "transitions": [ + {"step": "create", "changed": true, "delivery": "pending", "ack": false}, + {"step": "display-retry", "readback": "delivered", "delivery": "delivered", "ack": true}, + {"step": "idempotent-replay", "readback": "current", "delivery": "current", "ack": true}, + {"step": "terminal-no-op", "changed": false, "delivery": "not_required", "ack": false}, + {"step": "failed-display", "readback": "pending", "delivery": "pending", "ack": false} + ] +} diff --git a/tsconfig.control-plane.json b/tsconfig.control-plane.json index 44551b9cf4..6577ad2148 100644 --- a/tsconfig.control-plane.json +++ b/tsconfig.control-plane.json @@ -54,6 +54,7 @@ "loopx/control_plane/scheduler/heartbeat_followup_cli.ts", "loopx/control_plane/scheduler/state_transition_rules.ts", "loopx/control_plane/todos/completion_fence.ts", + "loopx/control_plane/todos/projection_delivery.ts", "loopx/control_plane/todos/completion_state.ts", "loopx/control_plane/todos/completion_validation_plan.ts", "loopx/control_plane/todos/next_action.ts",