diff --git a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md index db1e310ed0..9545a4e44c 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md @@ -2500,6 +2500,12 @@ The next complete stage packages are: monitor/resume effects close together. Neither an admission result nor a lease-fence result is a commit receipt. Keep provider CAS/replay and existing writer lock lifetimes unchanged while collecting this deletion payoff. + Waiting/resume lane selection is now one TS read-policy owner shared by quota, + vision-wait, agent-scope and replan. The obsolete Python selector module is + deleted; the adapter accepts the same canonical summary after promotion and + legacy summary before it. Real CLI coverage includes capacity changes and + missing promoted display without writing it. This does not close all quota + source paths, authorize monitor writeback, or change provider/promotion holds. 2. **Permanent projection closure.** Reuse `provider_projection.py`, the Todo-section renderer and existing journal/outbox. Preserve non-owned human narrative; render owned sections from a known canonical revision, with diff --git a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md index 02a83e38a3..f77b66c045 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md @@ -1982,6 +1982,11 @@ backend、实时双向同步或按命令拆开的权威;晋升后不支持的 这不是完整 native 字段编辑:在 update 的字段、ownership、validation 和 monitor/resume effect 一起闭合前,保留严格 text/note 事务边界。准入结果和 lease-fence 结果都不是 commit receipt;兑现删除收益时,provider CAS/replay 与既有 writer 持锁生命周期不变。 + 等待/恢复 lane 选择现由 quota、vision-wait、agent-scope、replan 共用一个 TS 读取 + 策略 owner,删除旧 Python selector 模块。适配层在 promotion 后消费同一 canonical + summary,之前消费 legacy summary;真实 CLI 覆盖容量变化和 promoted display + 缺失且不写回的场景。这不代表所有 quota source 路径已闭合,不授予 monitor 写回 + 权限,也不改变 provider 默认与 promotion hold。 2. **永久投影闭合。** 复用 `provider_projection.py`、Todo-section renderer 和既有 journal/outbox。保留非托管的人工叙述,从已知 canonical revision 渲染托管 section, 提供幂等修复与 freshness/readback 证据。投影 pending 独立于业务 commit/replay。 diff --git a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md index d057cb7784..086ce7e83a 100644 --- a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md +++ b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md @@ -218,6 +218,40 @@ this slice reduces semantic owners, not crossing count. Native transactions stay in-process. Fold the remaining crossings into that complete transaction rather than extending these adapters field by field. +The waiting/resume planning slice now uses `todos/resume_planning.ts` for the +complete deferred, resume-blocked, monitor-repair and blocked-successor selection. +Quota composes capacity evaluation with these lanes in one request per source summary, reusing the +existing TS resume evaluator in-process; vision-wait, agent-scope, frontier and +replan consumers use the same projection. The old `deferred_resume.py` rule owner +is removed, not retained behind a second implementation. The Python adapter keeps +the reader compatibility boundary, not claim/exclusion selection or wait routing. +Resume, route-continuation and succession-warning share `compact_projection.py` +for field omission and scope normalization; caller-specific text inference and +succession-only fields remain explicit. Priority rank normalization stays in the +resume adapter. This retires +one read-policy family, not the whole quota reducer or the monitor/lease writers. +Equal public sort keys retain source order; full counts precede display limits; +`monitor_changed` is not the legacy `todo_done:` repair path. This +read-only result grants neither execution authority nor a lifecycle receipt. +The adapter exits when its callers consume typed Todo records in-process. + +Resume condition diagnosis is now shared by the evaluator and planning owner; +agent-scope consumes the selected repair lane rather than reinterpreting target +type/status. Old compact inputs may recover omitted kind/class from typed +`resume_when` and the same snapshot's monitor records, never from narrative. +This refinement includes explicit behavior corrections: self-dependencies and +`todo_done` dependencies on unfinished monitors are `resume_condition_invalid`, +not ordinary pending waits. Completed historical monitor dependencies remain +satisfied; missing completion targets remain pending because absence in a +partial snapshot is not proof of an invalid dependency. Valid generation fences, +claim/exclusion, capacity and PR waits retain their existing semantics. Invalid +conditions cannot become exact blocked-successor waits. Monitor completion +repair stays visible and selectable only in the permitted executor scope. +No automatic conversion to `monitor_changed`, baseline reset, persisted-state +rewrite or new writer admission is implied. General add/update admission and a +generic repair action for every invalid condition remain separate scopes; this +is not a claim of zero behavior change or full Todo writer closure. + 1. **Close the actual command and consumer inventory.** Build on the merged create/claim/update and #4053 terminal/successor/archive transactions; do not recreate them. Inventory remaining field-edit, monitor, lease, and event 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 c4d8055549..817647a549 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 @@ -169,6 +169,30 @@ legacy update writer**。字段 patch、省略/清空、monitor/resume effect 语义 owner,不宣称减少 crossings,native transaction 仍进程内调用。下一步将这些 crossing 一起折叠进完整事务,不能沿着 adapter 逐字段继续加桥。 +等待/恢复规划现由 `todos/resume_planning.ts` 一次完成 deferred、resume-blocked、 +monitor-repair 和 blocked-successor 选择。Quota 为每个 source summary 将容量条件与这些 lane 合为一个请求, +在 TS 进程内复用既有 resume evaluator;vision-wait、agent-scope、frontier、replan +共用此投影。删除旧 `deferred_resume.py` 规则 owner,不保留第二份实现。Python 适配层 +只保留 reader 兼容边界,不再决定 claim/exclusion 选择或等待路由。Resume、 +route-continuation、succession-warning 共用 `compact_projection.py` 的字段省略与 scope +归一化;各 caller 的文本推断差异及 succession 独有字段显式保留,priority rank +归一化仍在 resume adapter。这闭合一个读取策略族,不是整个 quota reducer,也未 +迁移 monitor/lease writer。相同公开排序键保持 source 顺序;完整计数先于展示截断; +`monitor_changed` 不进入旧 `todo_done:` 修复路径。只读结果不授予执行权限, +也不是生命周期 receipt;caller 在进程内消费 typed Todo record 后可删除此适配层。 + +条件 evaluator 与规划 owner 现在共用恢复条件诊断;agent-scope 消费已选好的修复 lane, +不再重新解释 target 类型/状态。旧 compact 输入缺少 kind/class 时,只从 typed +`resume_when` 和同一快照的 Monitor 记录补足,不从叙述猜测。本次 refinement 包含 +明确行为修正:自依赖,以及对未完成 Monitor 的 `todo_done` 依赖,被诊断为 +`resume_condition_invalid`,不再当作普通 pending wait。历史已完成 Monitor 依赖仍可 +满足;完成依赖的目标缺失仍为 pending,因局部快照中的缺失不能证明依赖非法。合法的 +generation fence、claim/exclusion、capacity 和 PR 等待语义保持。非法条件不进入 +精确 blocked-successor 等待;Monitor 完成依赖的修复仍可见,且仅在合法执行者范围内 +可选。此诊断不自动改写为 `monitor_changed`、重置 baseline、重写持久状态或增加写入 +准入。普通 add/update 准入及覆盖全部非法条件的通用修复动作仍是独立范围;不能宣称 +全量零行为变化或全部 Todo writer 已闭合。 + 1. **闭合实际命令与 consumer 清单。** 基于已合入的 create/claim/update 和 #4053 terminal/successor/archive transaction 推进,不重复建设。按真实合同盘点剩余 字段编辑、monitor、lease、event caller,把规则迁入既有 TS owner,并在同一切片 diff --git a/examples/control_plane/todo-deferred-resume-lanes-smoke.py b/examples/control_plane/todo-deferred-resume-lanes-smoke.py index 6f87b60589..361029f1ef 100644 --- a/examples/control_plane/todo-deferred-resume-lanes-smoke.py +++ b/examples/control_plane/todo-deferred-resume-lanes-smoke.py @@ -11,14 +11,7 @@ if str(REPO_ROOT) not in sys.path: sys.path.insert(0, str(REPO_ROOT)) -from loopx.control_plane.todos.deferred_resume import ( # noqa: E402 - TODO_DEFERRED_RESUME_SELECTION_POLICY, - TODO_MONITOR_BLOCKED_RESUME_SELECTION_POLICY, - build_todo_deferred_visibility_lanes, - build_todo_resume_blocked_visibility_lanes, - todo_summary_monitor_blocked_resume_items, - todo_summary_resume_blocked_items, -) +from loopx.control_plane.todos.resume_planning import project_todo_resume_planning # noqa: E402 CURRENT_AGENT = "codex-product-capability" @@ -116,11 +109,11 @@ def assert_deferred_resume_lanes_filter_current_unclaimed_and_other_agents() -> ], } - lanes = build_todo_deferred_visibility_lanes( + lanes = project_todo_resume_planning( summary, - agent_identity={"agent_id": CURRENT_AGENT}, + agent_id=CURRENT_AGENT, item_limit=10, - ) + )["deferred_lanes"] assert lanes["deferred_count"] == 1, lanes assert lanes["deferred_visibility_limit"] == 10, lanes assert lanes["deferred_items"][0]["todo_id"] == "todo_deferred_backlog", lanes @@ -140,9 +133,6 @@ def assert_deferred_resume_lanes_filter_current_unclaimed_and_other_agents() -> current = lanes["current_agent_deferred_resume_candidates"][0] assert current["required_write_scopes"] == ["loopx/**"], current assert current["decision_scope"]["scope_key"] == "resume", current - assert lanes["deferred_resume_selection_policy"] == ( - TODO_DEFERRED_RESUME_SELECTION_POLICY - ), lanes def assert_monitor_blocked_resume_lanes_filter_by_claim_and_monitor_target() -> None: @@ -181,7 +171,9 @@ def assert_monitor_blocked_resume_lanes_filter_by_claim_and_monitor_target() -> "backlog_items": [current], } - resume_blocked = todo_summary_resume_blocked_items(summary) + projection = project_todo_resume_planning(summary, agent_id=CURRENT_AGENT, item_limit=10) + lanes = projection["resume_blocked_lanes"] + resume_blocked = lanes["resume_blocked_items"] assert [item["todo_id"] for item in resume_blocked] == [ "todo_current_blocked", "todo_unclaimed_blocked", @@ -189,7 +181,7 @@ def assert_monitor_blocked_resume_lanes_filter_by_claim_and_monitor_target() -> "todo_non_monitor_blocked", "todo_excluded_blocked", ], resume_blocked - monitor_blocked = todo_summary_monitor_blocked_resume_items(summary) + monitor_blocked = projection["monitor_blocked_items"] assert [item["todo_id"] for item in monitor_blocked] == [ "todo_current_blocked", "todo_unclaimed_blocked", @@ -201,20 +193,12 @@ def assert_monitor_blocked_resume_lanes_filter_by_claim_and_monitor_target() -> for item in monitor_blocked ), monitor_blocked - lanes = build_todo_resume_blocked_visibility_lanes( - summary, - agent_identity={"agent_id": CURRENT_AGENT}, - item_limit=10, - ) assert lanes["resume_blocked_count"] == 5, lanes assert lanes["monitor_blocked_resume_count"] == 4, lanes assert lanes["current_agent_monitor_blocked_resume_count"] == 1, lanes assert lanes["unclaimed_monitor_blocked_resume_count"] == 1, lanes assert lanes["other_agent_monitor_blocked_resume_count"] == 1, lanes assert lanes["executor_excluded_self_monitor_blocked_resume_count"] == 1, lanes - assert lanes["monitor_blocked_resume_selection_policy"] == ( - TODO_MONITOR_BLOCKED_RESUME_SELECTION_POLICY - ), lanes def main() -> int: diff --git a/loopx/control_plane/agents/agent_scope.py b/loopx/control_plane/agents/agent_scope.py index 9802ec2f1a..ccf10cbb86 100644 --- a/loopx/control_plane/agents/agent_scope.py +++ b/loopx/control_plane/agents/agent_scope.py @@ -17,7 +17,6 @@ work_lane_contract_requires_current_agent_attempt, ) from ..todos.contract import ( - TODO_STATUS_OPEN, TODO_TASK_CLASS_ADVANCEMENT, TODO_TASK_CLASS_MONITOR, normalize_todo_blocks_agent, @@ -26,11 +25,9 @@ normalize_todo_excluded_agents, normalize_todo_global_gate, normalize_todo_id, - normalize_todo_status, - normalize_todo_task_class, ) from ..todos.handoff_gate import HandoffGateState -from ..todos.deferred_resume import todo_summary_blocked_successor_items +from ..todos.resume_planning import project_todo_resume_planning from ..todos.projection import ( todo_item_claimed_by_agent_or_unclaimed, todo_item_excludes_agent, @@ -936,28 +933,16 @@ def _agent_scope_monitor_blocked_resume_candidates( continue if item.get("resume_ready") is not False: continue - raw_condition = item.get("resume_condition") - condition = raw_condition if isinstance(raw_condition, dict) else {} - if normalize_todo_status(condition.get("target_status")) != TODO_STATUS_OPEN: - continue - target_todo_id = normalize_todo_id( - item.get("blocking_monitor_todo_id") - or condition.get("target_todo_id") - or condition.get("target") - ) - target_task_class = normalize_todo_task_class( - condition.get("target_task_class"), - text="", - ) - if target_task_class != TODO_TASK_CLASS_MONITOR and not target_todo_id: - continue + # The typed resume-planning owner has already diagnosed and selected + # this repair lane. This consumer keeps executor scope and presentation, + # not a second interpretation of condition target/class/status. identity = str(item.get("todo_id") or item.get("index") or item.get("text") or "") if identity in seen: continue seen.add(identity) compact = compact_todo_summary_item(item, text=str(item.get("text") or "").strip()) - if target_todo_id: - compact["blocking_monitor_todo_id"] = target_todo_id + if item.get("blocking_monitor_todo_id"): + compact["blocking_monitor_todo_id"] = item["blocking_monitor_todo_id"] unique.append(compact) return sorted(unique, key=_todo_projection_sort_key) @@ -1261,10 +1246,10 @@ def _deferred_resume_frontier( def _blocked_successor_wait_frontier( context: _AgentScopeNoCandidateContext, ) -> dict[str, Any] | None: - candidates = todo_summary_blocked_successor_items( + candidates = project_todo_resume_planning( context.summary, agent_id=context.agent_id, - ) + )["blocked_successor_items"] if not candidates: return None first = candidates[0] diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index 8cef4b2c98..c5bf91fa2f 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -71,6 +71,7 @@ import { planTodoExternalWaitTransition, } from "./todos/resume_condition.ts"; import { evaluateSchedulerStateTransition } from "./scheduler/state_transition_rules.ts"; +import { projectTodoResumePlanning } from "./todos/resume_planning.ts"; import { evaluateSchedulerStateOperation, loadSchedulerState, @@ -394,6 +395,7 @@ export function createEffectRuntimeHandlers( ["todo.next_action.transition", transitionTodoNextAction], ["todo.resume_condition.normalize", normalizeTodoResumeWhen], ["todo.resume_condition.evaluate", evaluateTodoResumeConditions], + ["todo.resume_planning.project", projectTodoResumePlanning], ["todo.external_wait.plan", planTodoExternalWaitTransition], ["scheduler.state_transition.evaluate", evaluateSchedulerStateTransition], ["scheduler.state.evaluate", evaluateSchedulerStateOperation], diff --git a/loopx/control_plane/goals/goal_frontier/fallback_disposition.py b/loopx/control_plane/goals/goal_frontier/fallback_disposition.py index 5382e49282..630a8bd278 100644 --- a/loopx/control_plane/goals/goal_frontier/fallback_disposition.py +++ b/loopx/control_plane/goals/goal_frontier/fallback_disposition.py @@ -6,7 +6,7 @@ from ...todos.contract import ( normalize_todo_id, ) -from ...todos.deferred_resume import todo_summary_blocked_successor_items +from ...todos.resume_planning import project_todo_resume_planning from ...todos.projection import ( agent_scoped_selectable_advancement_todo_ids, ) @@ -125,10 +125,10 @@ def _blocked_successor_todo_ids( todo_id for todo_id in ( normalize_todo_id(item.get("todo_id")) - for item in todo_summary_blocked_successor_items( + for item in project_todo_resume_planning( agent_todo_summary, agent_id=agent_id, - ) + )["blocked_successor_items"] if isinstance(item, dict) ) if todo_id @@ -148,10 +148,10 @@ def _blocked_primary_waiting( if isinstance(blocker_items, list) and blocker_items: return True return bool( - todo_summary_blocked_successor_items( + project_todo_resume_planning( agent_todo_summary, agent_id=agent_id, - ) + )["blocked_successor_items"] ) diff --git a/loopx/control_plane/goals/goal_vision_wait.py b/loopx/control_plane/goals/goal_vision_wait.py index ef4f29e4f4..317d903423 100644 --- a/loopx/control_plane/goals/goal_vision_wait.py +++ b/loopx/control_plane/goals/goal_vision_wait.py @@ -12,7 +12,7 @@ normalize_todo_id, normalize_todo_status, ) -from ..todos.deferred_resume import todo_summary_blocked_successor_items +from ..todos.resume_planning import project_todo_resume_planning GOAL_VISION_WAIT_STATE_SCHEMA_VERSION = "goal_vision_wait_state_v0" VISION_ACCEPTANCE_GAP_KIND = "vision_acceptance_gap" @@ -170,9 +170,9 @@ def _covered_wait_items( and normalize_todo_claimed_by(item.get("claimed_by")) == safe_agent_id and not todo_item_excludes_agent(item, agent_id=safe_agent_id) ] - candidates = todo_summary_blocked_successor_items( + candidates = project_todo_resume_planning( agent_todo_summary or {}, agent_id=agent_id - ) + )["blocked_successor_items"] coverage = effect_runtime_result( "goal.vision_wait.coverage", { diff --git a/loopx/control_plane/todos/compact_projection.py b/loopx/control_plane/todos/compact_projection.py new file mode 100644 index 0000000000..2b6cbf320f --- /dev/null +++ b/loopx/control_plane/todos/compact_projection.py @@ -0,0 +1,118 @@ +"""Shared compatibility display codec; no Todo selection or mutation authority. + +Callers own text presentation and explicit extension fields. Omission and +scope normalization stay common; canonical records never round-trip through it. +""" + +from __future__ import annotations + +from typing import Any + +from .contract import ( + normalize_required_write_scopes, + normalize_todo_decision_scope, + normalize_todo_required_decision_scopes, + normalize_todo_task_class, +) + + +def projection_task_class(item: dict[str, Any]) -> str: + text = " ".join( + str(value or "") + for value in (item.get("title"), item.get("text")) + if str(value or "").strip() + ) + return normalize_todo_task_class( + item.get("task_class"), + text=text, + action_kind=item.get("action_kind"), + ) + + +def compact_todo_projection_item( + item: dict[str, Any], + *, + text: Any, + extra_fields: tuple[str, ...] = (), + task_class_text: str | None = None, +) -> dict[str, Any]: + compact: dict[str, Any] = { + "index": item.get("index"), + "text": text, + } + for key in ( + "schema_version", + "todo_id", + "role", + "status", + "priority", + "title", + "archive_state", + "source_section", + "task_class", + "action_kind", + "task_domain", + "task_repository", + "continuation_policy", + "required_write_scopes", + "required_capabilities", + "target_capabilities", + "decision_scope", + "required_decision_scopes", + "claimed_by", + "blocks_agent", + "excluded_agents", + "unblocks_todo_id", + "resume_when", + "resume_monitor_generation", + "resume_condition", + "resume_ready", + "no_followup", + "successor_todo_ids", + "target_key", + "cadence", + "next_due_at", + "expires_at", + "last_checked_at", + "result_hash", + "consecutive_no_change", + "material_change", + "material_change_generation", + "max_no_change_before_replan", + "route_continuation_replan_required", + "route_continuation_reason", + "route_id", + "route_key", + "completed_at", + "updated_at", + "superseded_by", + ) + extra_fields: + if item.get(key) is not None: + compact[key] = item.get(key) + required_write_scopes = normalize_required_write_scopes(compact.get("required_write_scopes")) + if required_write_scopes: + compact["required_write_scopes"] = required_write_scopes + else: + compact.pop("required_write_scopes", None) + decision_scope = normalize_todo_decision_scope(compact.get("decision_scope")) + if decision_scope: + compact["decision_scope"] = decision_scope + else: + compact.pop("decision_scope", None) + required_decision_scopes = normalize_todo_required_decision_scopes( + compact.get("required_decision_scopes") + ) + if required_decision_scopes: + compact["required_decision_scopes"] = required_decision_scopes + else: + compact.pop("required_decision_scopes", None) + compact["task_class"] = ( + projection_task_class(compact) + if task_class_text is None + else normalize_todo_task_class( + compact.get("task_class"), + text=task_class_text, + action_kind=compact.get("action_kind"), + ) + ) + return compact diff --git a/loopx/control_plane/todos/deferred_resume.py b/loopx/control_plane/todos/deferred_resume.py deleted file mode 100644 index 10572a4e90..0000000000 --- a/loopx/control_plane/todos/deferred_resume.py +++ /dev/null @@ -1,542 +0,0 @@ -from __future__ import annotations - -from typing import Any - -from .contract import ( - TODO_STATUS_OPEN, - TODO_TASK_CLASS_ADVANCEMENT, - TODO_TASK_CLASS_MONITOR, - normalize_required_capabilities, - normalize_required_write_scopes, - normalize_todo_claimed_by, - normalize_todo_decision_scope, - normalize_todo_id, - normalize_todo_required_decision_scopes, - normalize_todo_resume_when, - normalize_todo_status, - normalize_todo_task_class, -) -from .resume_condition import evaluate_todo_resume_conditions -from .projection import ( - todo_item_excludes_agent, - todo_item_is_deferred, - todo_projection_sort_key, -) - - -TODO_DEFERRED_RESUME_SELECTION_POLICY = ( - "quota may wake the current peer only for ready deferred todos " - "claimed by that agent or unclaimed; other-agent deferred todos remain " - "diagnostic visibility and executor-excluded todos remain visible but " - "non-selectable" -) -TODO_MONITOR_BLOCKED_RESUME_SELECTION_POLICY = ( - "open advancement todos gated by todo_done: must " - "project as successor replan/state repair instead of quiet monitor wait" -) - - -def resolve_capacity_resume_summary( - value: Any, - *, - available_capabilities: Any, -) -> Any: - """Resolve capacity-backed deferred todos from the quota host capability set.""" - if not isinstance(value, dict): - return value - available = set(normalize_required_capabilities(available_capabilities)) - resolved = dict(value) - deferred_items = todo_summary_deferred_items(value, "deferred_items") - source_items = [ - item - for item in value.get("items", []) - if isinstance(item, dict) - ] - conditions = evaluate_todo_resume_conditions( - deferred_items, - source_items=[*source_items, *deferred_items], - available_capabilities=available, - kinds=["capacity_available"], - ) - for item in deferred_items: - todo_id = normalize_todo_id(item.get("todo_id")) - condition = conditions.get(todo_id or "") - if condition is None: - continue - item["resume_condition"] = condition - item["resume_ready"] = condition.get("satisfied") is True - resolved["deferred_items"] = deferred_items - resolved["deferred_resume_candidates"] = [ - item for item in deferred_items if item.get("resume_ready") is True - ] - return resolved - - -def _todo_task_class(item: dict[str, Any]) -> str: - text = " ".join( - str(value or "") - for value in (item.get("title"), item.get("text")) - if str(value or "").strip() - ) - return normalize_todo_task_class( - item.get("task_class"), - text=text, - action_kind=item.get("action_kind"), - ) - - -def _compact_deferred_resume_item( - item: dict[str, Any], - *, - text: str, -) -> dict[str, Any]: - compact: dict[str, Any] = { - "index": item.get("index"), - "text": text, - } - for key in ( - "schema_version", - "todo_id", - "role", - "status", - "priority", - "title", - "archive_state", - "source_section", - "task_class", - "action_kind", - "task_domain", - "task_repository", - "continuation_policy", - "required_write_scopes", - "required_capabilities", - "target_capabilities", - "decision_scope", - "required_decision_scopes", - "claimed_by", - "blocks_agent", - "excluded_agents", - "unblocks_todo_id", - "resume_when", - "resume_monitor_generation", - "resume_condition", - "resume_ready", - "no_followup", - "successor_todo_ids", - "target_key", - "cadence", - "next_due_at", - "expires_at", - "last_checked_at", - "result_hash", - "consecutive_no_change", - "material_change", - "material_change_generation", - "max_no_change_before_replan", - "route_continuation_replan_required", - "route_continuation_reason", - "route_id", - "route_key", - "completed_at", - "updated_at", - "superseded_by", - ): - if item.get(key) is not None: - compact[key] = item.get(key) - required_write_scopes = normalize_required_write_scopes(compact.get("required_write_scopes")) - if required_write_scopes: - compact["required_write_scopes"] = required_write_scopes - else: - compact.pop("required_write_scopes", None) - decision_scope = normalize_todo_decision_scope(compact.get("decision_scope")) - if decision_scope: - compact["decision_scope"] = decision_scope - else: - compact.pop("decision_scope", None) - required_decision_scopes = normalize_todo_required_decision_scopes( - compact.get("required_decision_scopes") - ) - if required_decision_scopes: - compact["required_decision_scopes"] = required_decision_scopes - else: - compact.pop("required_decision_scopes", None) - compact["task_class"] = _todo_task_class(compact) - return compact - - -def todo_summary_deferred_items( - value: dict[str, Any], - key: str, -) -> list[dict[str, Any]]: - if not isinstance(value, dict): - return [] - raw_source_items = value.get(key) - source_items = raw_source_items if isinstance(raw_source_items, list) else [] - if not source_items and key == "deferred_items": - raw_items = value.get("items") - source_items = [ - item - for item in (raw_items if isinstance(raw_items, list) else []) - if isinstance(item, dict) and todo_item_is_deferred(item) - ] - items: list[dict[str, Any]] = [] - for item in source_items: - if not isinstance(item, dict): - continue - text = str(item.get("text") or "").strip() - if not text: - continue - compact = _compact_deferred_resume_item(item, text=text) - resume_when = normalize_todo_resume_when(item.get("resume_when")) - if resume_when: - compact["resume_when"] = resume_when - if item.get("resume_condition") is not None: - compact["resume_condition"] = item.get("resume_condition") - if item.get("resume_ready") is not None: - compact["resume_ready"] = bool(item.get("resume_ready")) - if todo_item_is_deferred(compact): - items.append(compact) - return sorted(items, key=todo_projection_sort_key) - - -def _dedupe_todo_items(items: list[dict[str, Any]]) -> list[dict[str, Any]]: - unique: list[dict[str, Any]] = [] - seen: set[tuple[str, str, Any]] = set() - for item in items: - todo_id = normalize_todo_id(item.get("todo_id")) or "" - text = str(item.get("text") or "").strip() - identity = (todo_id, text, item.get("index")) - if identity in seen: - continue - seen.add(identity) - unique.append(item) - return unique - - -def todo_summary_resume_blocked_items(value: dict[str, Any]) -> list[dict[str, Any]]: - if not isinstance(value, dict): - return [] - raw_source_items = value.get("resume_blocked_items") - source_items = ( - list(raw_source_items) if isinstance(raw_source_items, list) else [] - ) - if not source_items: - for key in ("items", "backlog_items", "first_open_items"): - raw_value_items = value.get(key) - raw_items = raw_value_items if isinstance(raw_value_items, list) else [] - source_items.extend(item for item in raw_items if isinstance(item, dict)) - items: list[dict[str, Any]] = [] - for item in source_items: - if not isinstance(item, dict): - continue - text = str(item.get("text") or "").strip() - if not text: - continue - if item.get("done") is True: - continue - if not normalize_todo_resume_when(item.get("resume_when")): - continue - if item.get("resume_ready") is not False: - continue - items.append(_compact_deferred_resume_item(item, text=text)) - return sorted(_dedupe_todo_items(items), key=todo_projection_sort_key) - - -def _monitor_target_todo_ids(value: dict[str, Any]) -> set[str]: - ids: set[str] = set() - for key in ( - "monitor_open_items", - "current_agent_claimed_monitor_items", - "claimed_monitor_open_items", - "items", - "backlog_items", - "first_open_items", - ): - raw_items = value.get(key) - source_items = raw_items if isinstance(raw_items, list) else [] - for item in source_items: - if not isinstance(item, dict): - continue - todo_id = normalize_todo_id(item.get("todo_id")) - if not todo_id: - continue - if _todo_task_class(item) == TODO_TASK_CLASS_MONITOR: - ids.add(todo_id) - return ids - - -def todo_summary_monitor_blocked_resume_items( - value: dict[str, Any], -) -> list[dict[str, Any]]: - monitor_ids = _monitor_target_todo_ids(value) - candidates: list[dict[str, Any]] = [] - for item in todo_summary_resume_blocked_items(value): - if _todo_task_class(item) != TODO_TASK_CLASS_ADVANCEMENT: - continue - raw_condition = item.get("resume_condition") - condition = raw_condition if isinstance(raw_condition, dict) else {} - if condition.get("kind") == "monitor_changed": - # This is the intentional typed wait lifecycle. The monitor stays - # independently schedulable and its generation fence owns resume; - # do not reinterpret it as the legacy todo_done: repair gap. - continue - target_todo_id = normalize_todo_id( - condition.get("target_todo_id") or condition.get("target") - ) - target_status = normalize_todo_status(condition.get("target_status")) - target_task_class = normalize_todo_task_class( - condition.get("target_task_class"), - text="", - ) - if target_status != TODO_STATUS_OPEN: - continue - if target_task_class != TODO_TASK_CLASS_MONITOR and target_todo_id not in monitor_ids: - continue - candidate = dict(item) - if target_todo_id: - candidate["blocking_monitor_todo_id"] = target_todo_id - candidates.append(candidate) - return sorted(_dedupe_todo_items(candidates), key=todo_projection_sort_key) - - -def todo_summary_blocked_successor_items( - value: dict[str, Any], - *, - agent_id: str | None, -) -> list[dict[str, Any]]: - """Return exact non-monitor successor waits selectable by one agent lane. - - Open resume-gated todos and deferred todos share the same wait contract. - A standing continuous monitor is deliberately excluded because it has a - dedicated gate-repair route: an endless monitor must not become a hidden - ``todo_done`` prerequisite for ordinary advancement. Current-agent claims - outrank unclaimed waits before ordinary priority ordering, matching the - agent-scoped open-todo selection contract. - """ - - if not isinstance(value, dict) or not agent_id: - return [] - candidates = [ - *todo_summary_resume_blocked_items(value), - *[ - item - for item in todo_summary_deferred_items(value, "deferred_items") - if item.get("resume_ready") is False - ], - ] - selected: list[dict[str, Any]] = [] - for item in _dedupe_todo_items(candidates): - if _todo_task_class(item) != TODO_TASK_CLASS_ADVANCEMENT: - continue - if todo_item_excludes_agent(item, agent_id=agent_id): - continue - claimed_by = normalize_todo_claimed_by(item.get("claimed_by")) - if claimed_by and claimed_by != agent_id: - continue - resume_when = normalize_todo_resume_when(item.get("resume_when")) - raw_condition = item.get("resume_condition") - condition = raw_condition if isinstance(raw_condition, dict) else {} - if not resume_when or condition.get("satisfied") is not False: - continue - target_task_class = normalize_todo_task_class( - condition.get("target_task_class"), - text="", - ) - if target_task_class == TODO_TASK_CLASS_MONITOR: - continue - compact = dict(item) - compact["resume_when"] = resume_when - compact["resume_ready"] = False - selected.append(compact) - return sorted( - _dedupe_todo_items(selected), - key=lambda item: ( - 0 - if normalize_todo_claimed_by(item.get("claimed_by")) == agent_id - else 1, - *todo_projection_sort_key(item), - ), - ) - - -def _agent_claim_filtered_deferred_items( - items: list[dict[str, Any]], - *, - agent_id: str | None, - claim: str, -) -> list[dict[str, Any]]: - selected: list[dict[str, Any]] = [] - for item in items: - claimed_by = normalize_todo_claimed_by(item.get("claimed_by")) - excluded = todo_item_excludes_agent(item, agent_id=agent_id) - if claim == "excluded": - if not excluded: - continue - selected.append(item) - continue - if excluded: - continue - if claim == "current" and claimed_by != agent_id: - continue - if claim == "unclaimed" and claimed_by: - continue - if claim == "other" and (not claimed_by or claimed_by == agent_id): - continue - selected.append(item) - return selected - - -def build_todo_resume_blocked_visibility_lanes( - value: dict[str, Any], - *, - agent_identity: dict[str, Any] | None, - item_limit: int, -) -> dict[str, Any]: - resume_blocked_items = todo_summary_resume_blocked_items(value) - monitor_blocked_items = todo_summary_monitor_blocked_resume_items(value) - if not resume_blocked_items and not monitor_blocked_items: - return {} - lanes: dict[str, Any] = { - "resume_blocked_count": len(resume_blocked_items), - "resume_blocked_items": resume_blocked_items[:item_limit], - } - if monitor_blocked_items: - lanes.update( - { - "monitor_blocked_resume_count": len(monitor_blocked_items), - "monitor_blocked_resume_candidates": monitor_blocked_items[ - :item_limit - ], - } - ) - agent_id = ( - normalize_todo_claimed_by(agent_identity.get("agent_id")) - if isinstance(agent_identity, dict) - else None - ) - if agent_id and monitor_blocked_items: - current_agent_candidates = _agent_claim_filtered_deferred_items( - monitor_blocked_items, - agent_id=agent_id, - claim="current", - ) - unclaimed_candidates = _agent_claim_filtered_deferred_items( - monitor_blocked_items, - agent_id=agent_id, - claim="unclaimed", - ) - other_agent_candidates = _agent_claim_filtered_deferred_items( - monitor_blocked_items, - agent_id=agent_id, - claim="other", - ) - excluded_self_candidates = _agent_claim_filtered_deferred_items( - monitor_blocked_items, - agent_id=agent_id, - claim="excluded", - ) - lanes.update( - { - "current_agent_monitor_blocked_resume_candidates": ( - current_agent_candidates[:item_limit] - ), - "unclaimed_monitor_blocked_resume_candidates": ( - unclaimed_candidates[:item_limit] - ), - "other_agent_monitor_blocked_resume_candidates": ( - other_agent_candidates[:item_limit] - ), - "current_agent_monitor_blocked_resume_count": len( - current_agent_candidates - ), - "unclaimed_monitor_blocked_resume_count": len(unclaimed_candidates), - "other_agent_monitor_blocked_resume_count": len(other_agent_candidates), - "executor_excluded_self_monitor_blocked_resume_candidates": ( - excluded_self_candidates[:item_limit] - ), - "executor_excluded_self_monitor_blocked_resume_count": len( - excluded_self_candidates - ), - "monitor_blocked_resume_selection_policy": ( - TODO_MONITOR_BLOCKED_RESUME_SELECTION_POLICY - ), - } - ) - return lanes - - -def build_todo_deferred_visibility_lanes( - value: dict[str, Any], - *, - agent_identity: dict[str, Any] | None, - item_limit: int, -) -> dict[str, Any]: - if not isinstance(value, dict): - return {} - deferred_items = todo_summary_deferred_items(value, "deferred_items") - deferred_resume_candidates = [ - item - for item in todo_summary_deferred_items(value, "deferred_resume_candidates") - if item.get("resume_ready") is True - ] - if not deferred_items and not deferred_resume_candidates and not value.get("deferred_count"): - return {} - - lanes: dict[str, Any] = { - "deferred_count": value.get("deferred_count", len(deferred_items)), - "deferred_visibility_limit": item_limit, - "deferred_items": deferred_items[:item_limit], - "deferred_resume_candidates": deferred_resume_candidates[:item_limit], - } - agent_id = ( - normalize_todo_claimed_by(agent_identity.get("agent_id")) - if isinstance(agent_identity, dict) - else None - ) - if agent_id: - current_agent_candidates = _agent_claim_filtered_deferred_items( - deferred_resume_candidates, - agent_id=agent_id, - claim="current", - ) - unclaimed_candidates = _agent_claim_filtered_deferred_items( - deferred_resume_candidates, - agent_id=agent_id, - claim="unclaimed", - ) - other_agent_candidates = _agent_claim_filtered_deferred_items( - deferred_resume_candidates, - agent_id=agent_id, - claim="other", - ) - excluded_self_candidates = _agent_claim_filtered_deferred_items( - deferred_resume_candidates, - agent_id=agent_id, - claim="excluded", - ) - lanes.update( - { - "current_agent_deferred_resume_candidates": ( - current_agent_candidates[:item_limit] - ), - "unclaimed_deferred_resume_candidates": ( - unclaimed_candidates[:item_limit] - ), - "other_agent_deferred_resume_candidates": ( - other_agent_candidates[:item_limit] - ), - "current_agent_deferred_resume_count": len(current_agent_candidates), - "unclaimed_deferred_resume_count": len(unclaimed_candidates), - "other_agent_deferred_resume_count": len(other_agent_candidates), - "executor_excluded_self_deferred_resume_candidates": ( - excluded_self_candidates[:item_limit] - ), - "executor_excluded_self_deferred_resume_count": len( - excluded_self_candidates - ), - "deferred_resume_selection_policy": ( - TODO_DEFERRED_RESUME_SELECTION_POLICY - ), - } - ) - return lanes diff --git a/loopx/control_plane/todos/quota_summary.py b/loopx/control_plane/todos/quota_summary.py index 8761fbc3dc..5c81ed9b51 100644 --- a/loopx/control_plane/todos/quota_summary.py +++ b/loopx/control_plane/todos/quota_summary.py @@ -21,11 +21,7 @@ TODO_TASK_CLASS_BLOCKER, TODO_TASK_CLASS_MONITOR, ) -from .deferred_resume import ( - build_todo_deferred_visibility_lanes, - build_todo_resume_blocked_visibility_lanes, - resolve_capacity_resume_summary, -) +from .resume_planning import project_todo_resume_planning from .frontier_deadline import todo_summary_frontier_deadline from .handoff_gate import build_todo_handoff_gate_lanes from .projection import ( @@ -488,10 +484,16 @@ def summarize_user_todos_for_quota( agent_identity: dict[str, Any] | None = None, filter_user_gate_blocks_agent: bool = False, available_capabilities: Any = None, + resume_planning: dict[str, Any] | None = None, ) -> dict[str, Any] | None: if not isinstance(value, dict): return None source_completeness, closure_intent = validate_todo_source_contract(value) + if resume_planning is None: + resume_planning = project_todo_resume_planning( + value, agent_id=(agent_identity or {}).get("agent_id"), + item_limit=TODO_DEFERRED_VISIBILITY_LIMIT, + ) all_open_items = sorted( todo_summary_source_items(value), key=todo_projection_sort_key, @@ -594,20 +596,8 @@ def summarize_user_todos_for_quota( visibility_lane_limit=TODO_VISIBILITY_LANE_LIMIT, ) ) - summary.update( - build_todo_deferred_visibility_lanes( - value, - agent_identity=agent_identity, - item_limit=TODO_DEFERRED_VISIBILITY_LIMIT, - ) - ) - summary.update( - build_todo_resume_blocked_visibility_lanes( - value, - agent_identity=agent_identity, - item_limit=TODO_DEFERRED_VISIBILITY_LIMIT, - ) - ) + summary.update(resume_planning["deferred_lanes"]) + summary.update(resume_planning["resume_blocked_lanes"]) summary.update( build_todo_handoff_gate_lanes( value, @@ -924,6 +914,7 @@ def summarize_project_asset_todos_for_quota( agent_identity: dict[str, Any] | None = None, filter_user_gate_blocks_agent: bool = False, available_capabilities: Any = None, + resume_planning: dict[str, Any] | None = None, ) -> dict[str, Any] | None: if not isinstance(value, dict): return None @@ -938,6 +929,7 @@ def summarize_project_asset_todos_for_quota( agent_identity=agent_identity, filter_user_gate_blocks_agent=filter_user_gate_blocks_agent, available_capabilities=available_capabilities, + resume_planning=resume_planning, ) all_open_items = sorted( @@ -992,13 +984,12 @@ def summarize_project_asset_todos_for_quota( visibility_lane_limit=TODO_VISIBILITY_LANE_LIMIT, ) ) - summary.update( - build_todo_deferred_visibility_lanes( - value, - agent_identity=agent_identity, + if resume_planning is None: + resume_planning = project_todo_resume_planning( + value, agent_id=(agent_identity or {}).get("agent_id"), item_limit=TODO_DEFERRED_VISIBILITY_LIMIT, ) - ) + summary.update(resume_planning["deferred_lanes"]) summary.update( build_todo_handoff_gate_lanes( value, @@ -1073,25 +1064,31 @@ def select_quota_todo_summary( filter_user_gate_blocks_agent: bool = False, available_capabilities: Any = None, ) -> dict[str, Any] | None: - canonical_value = resolve_capacity_resume_summary( - canonical_value, - available_capabilities=available_capabilities, - ) - project_asset_value = resolve_capacity_resume_summary( - project_asset_value, - available_capabilities=available_capabilities, - ) + plans = [] + for source in (canonical_value, project_asset_value): + plans.append(project_todo_resume_planning( + source, agent_id=(agent_identity or {}).get("agent_id"), + item_limit=TODO_DEFERRED_VISIBILITY_LIMIT, + available_capabilities=available_capabilities or [], + ) if isinstance(source, dict) else None) + canonical_plan, project_asset_plan = plans + if canonical_plan is not None: + canonical_value = {**canonical_value, **canonical_plan["capacity_fields"]} + if project_asset_plan is not None: + project_asset_value = {**project_asset_value, **project_asset_plan["capacity_fields"]} canonical_summary = summarize_user_todos_for_quota( canonical_value, agent_identity=agent_identity, filter_user_gate_blocks_agent=filter_user_gate_blocks_agent, available_capabilities=available_capabilities, + resume_planning=canonical_plan, ) project_asset_summary = summarize_project_asset_todos_for_quota( project_asset_value, agent_identity=agent_identity, filter_user_gate_blocks_agent=filter_user_gate_blocks_agent, available_capabilities=available_capabilities, + resume_planning=project_asset_plan, ) if is_canonical_attention_todo_summary(canonical_value): return canonical_summary or project_asset_summary diff --git a/loopx/control_plane/todos/resume_condition.ts b/loopx/control_plane/todos/resume_condition.ts index 24b4284450..72342906fc 100644 --- a/loopx/control_plane/todos/resume_condition.ts +++ b/loopx/control_plane/todos/resume_condition.ts @@ -390,15 +390,44 @@ function conditionFor( return condition; } -function resumeAvailabilityReason(condition: JsonObject): string { - if (condition.satisfied === true) return "resume_condition_satisfied"; - if ( - condition.invalid_target === true || - typeof condition.invalid_state === "string" - ) { - return "resume_condition_invalid"; - } - return "resume_condition_pending"; +export type ResumeConditionDiagnosis = { + kind: TodoResumeKind | null; +} & ( + | { state: "satisfied" | "pending" } + | { state: "invalid"; reason: string } +); + +/** Shared by canonical evaluation and old compact projections. Missing source + * facts are not proof of an invalid dependency. A historical completed monitor + * remains satisfied; a live monitor completion wait requires an explicit replan, + * never an inferred monitor_changed generation baseline. */ +export function diagnoseTodoResumeCondition( + condition: JsonObject, waitingTodoId: string | null = null, +): ResumeConditionDiagnosis { + const parsed = parseResumeWhen(condition.resume_when); + const kind = TODO_RESUME_KINDS.find((value) => value === condition.kind) ?? parsed?.kind ?? null; + const targetId = condition.target_todo_id ?? condition.target ?? parsed?.target; + if ((kind === "todo_done" || kind === "monitor_changed") && waitingTodoId && targetId === waitingTodoId) { + return { kind, state: "invalid", reason: "dependency_self_reference" }; + } + if (typeof condition.invalid_state === "string") { + return { kind, state: "invalid", reason: condition.invalid_state }; + } + if (condition.invalid_target === true) return { kind, state: "invalid", reason: "invalid_target" }; + if (kind === "todo_done" && condition.target_task_class === "continuous_monitor" + && condition.target_status && condition.target_status !== "done") { + return { kind, state: "invalid", reason: "monitor_completion_requires_replan" }; + } + return { kind, state: condition.satisfied === true ? "satisfied" : "pending" }; +} + +function diagnosedCondition(condition: JsonObject, waitingTodoId: string): JsonObject { + const diagnosis = diagnoseTodoResumeCondition(condition, waitingTodoId); + return { + ...condition, + availability_reason: `resume_condition_${diagnosis.state}`, + ...(diagnosis.state === "invalid" ? { satisfied: false, invalid_state: diagnosis.reason } : {}), + }; } export function evaluateTodoResumeConditions(value: unknown): JsonObject { @@ -441,10 +470,9 @@ export function evaluateTodoResumeConditions(value: unknown): JsonObject { rolloutEvents, availableCapabilities, ); - condition.availability_reason = resumeAvailabilityReason(condition); conditions.push({ todo_id: item.todo_id, - condition, + condition: diagnosedCondition(condition, item.todo_id), }); } return { diff --git a/loopx/control_plane/todos/resume_planning.py b/loopx/control_plane/todos/resume_planning.py new file mode 100644 index 0000000000..88d72f586c --- /dev/null +++ b/loopx/control_plane/todos/resume_planning.py @@ -0,0 +1,86 @@ +"""Compatibility codecs for the typed, read-only Todo resume planning owner.""" + +from __future__ import annotations + +from typing import Any + +from ..effect_runtime import EffectRuntimeRejected, effect_runtime_result +from .contract import ( + normalize_required_capabilities, + normalize_todo_claimed_by, normalize_todo_id, + normalize_todo_resume_when, + normalize_todo_status, normalize_todo_task_class, normalize_todo_excluded_agents, +) +from .compact_projection import compact_todo_projection_item +from .projection import todo_projection_sort_key + +_SOURCE_KEYS = ( + "items", "backlog_items", "first_open_items", "deferred_items", + "deferred_resume_candidates", "resume_blocked_items", "monitor_open_items", + "current_agent_claimed_monitor_items", "claimed_monitor_open_items", +) + + +def _planning_item(item: dict[str, Any]) -> dict[str, Any]: + payload = compact_todo_projection_item(item, text=str(item.get("text") or "").strip()) + condition = item.get("resume_condition") + condition = condition if isinstance(condition, dict) else {} + priority, index = todo_projection_sort_key(payload) + ready = item.get("resume_ready") + return { + "payload": payload, "id": normalize_todo_id(item.get("todo_id")), + "status": normalize_todo_status(item.get("status")), + "claim": normalize_todo_claimed_by(item.get("claimed_by")), + "excluded": normalize_todo_excluded_agents(item.get("excluded_agents")), + "resume": normalize_todo_resume_when(item.get("resume_when")), + "done": item.get("done") is True, + "ready": ready if isinstance(ready, bool) else None, + "ready_truthy": bool(ready), "priority": priority, "index": index, + "target_id": normalize_todo_id(condition.get("target_todo_id") or condition.get("target")), + "target_status": normalize_todo_status(condition.get("target_status")), + "target_class": normalize_todo_task_class(condition.get("target_task_class"), text=""), + } + + +def _planning_sources(value: dict[str, Any]) -> dict[str, list[dict[str, Any]]]: + encoded: dict[int, dict[str, Any]] = {} + + def encode(item: Any) -> dict[str, Any]: + # Ignored rows still make an explicit lane nonempty; dropping them + # here could activate the legacy fallback to a different source. + identity = id(item) if isinstance(item, dict) else -1 + if identity not in encoded: + encoded[identity] = _planning_item(item if isinstance(item, dict) else {}) + return encoded[identity] + + sources = {} + for key in _SOURCE_KEYS: + raw = value.get(key) + sources[key] = [encode(item) for item in raw] if isinstance(raw, list) else [] + return sources + + +def project_todo_resume_planning( + value: Any, *, agent_id: str | None = None, item_limit: int = 5, + available_capabilities: Any = None, +) -> dict[str, Any]: + """Project one snapshot; no business write or additional authority is granted.""" + value = value if isinstance(value, dict) else {} + try: + result = effect_runtime_result("todo.resume_planning.project", { + "schema_version": "todo_resume_planning_request_v0", + "sources": _planning_sources(value), + "agent_id": normalize_todo_claimed_by(agent_id), "item_limit": item_limit, + "has_deferred_count": "deferred_count" in value, + "has_visible_deferred_count": bool(value.get("deferred_count")), + "deferred_count": value.get("deferred_count"), + "available_capabilities": ( + normalize_required_capabilities(available_capabilities) + if available_capabilities is not None else None + ), + }) + except EffectRuntimeRejected as exc: + raise ValueError(str(exc)) from None + if not isinstance(result, dict) or result.get("schema_version") != "todo_resume_planning_v0": + raise RuntimeError("TypeScript Todo resume planning shape mismatch") + return result diff --git a/loopx/control_plane/todos/resume_planning.ts b/loopx/control_plane/todos/resume_planning.ts new file mode 100644 index 0000000000..3391de346e --- /dev/null +++ b/loopx/control_plane/todos/resume_planning.ts @@ -0,0 +1,274 @@ +import type { JsonObject } from "../effect_program.ts"; +import { EffectRuntimeRequestError } from "../effect_runtime_errors.ts"; +import { + jsonObject, requireJsonObject, requireBoolean, requireInteger, + requireStringArray, optionalNonEmptyString, +} from "../runtime_decode.ts"; +import { + diagnoseTodoResumeCondition, type ResumeConditionDiagnosis, + evaluateTodoResumeConditions, TODO_RESUME_EVALUATION_REQUEST_SCHEMA_VERSION, +} from "./resume_condition.ts"; + +export const RESUME_PLANNING_REQUEST = "todo_resume_planning_request_v0"; +export const RESUME_PLANNING_RESULT = "todo_resume_planning_v0"; +const DEFERRED_POLICY = "quota may wake the current peer only for ready deferred todos " + + "claimed by that agent or unclaimed; other-agent deferred todos remain " + + "diagnostic visibility and executor-excluded todos remain visible but non-selectable"; +const MONITOR_POLICY = "open advancement todos gated by todo_done: must " + + "project as successor replan/state repair instead of quiet monitor wait"; +const SOURCE_KEYS = ["items", "backlog_items", "first_open_items", "deferred_items", + "deferred_resume_candidates", "resume_blocked_items", "monitor_open_items", + "current_agent_claimed_monitor_items", "claimed_monitor_open_items"] as const; +type SourceKey = typeof SOURCE_KEYS[number]; +type ClaimLane = "current_agent" | "unclaimed" | "other_agent" | "executor_excluded_self"; + +/** Codec facts are separate from policy. Payload fields remain lossless while + * comparisons use the existing reader's normalized values and stable order. */ +interface Item { + payload: JsonObject; + id: string | null; + status: string | null; + claim: string | null; + excluded: readonly string[]; + resume: string | null; + done: boolean; + ready: boolean | null; + readyTruthy: boolean; + priority: number; + index: number; + targetId: string | null; + targetStatus: string | null; + targetClass: string; + diagnosis?: ResumeConditionDiagnosis; +} + +function item(value: unknown): Item { + const raw = requireJsonObject(value, "resume planning item"); + const payload = requireJsonObject(raw.payload, "payload"); + if (typeof payload.text !== "string" || typeof payload.task_class !== "string") { + throw new EffectRuntimeRequestError("resume planning payload requires text and task_class"); + } + const optional = (key: string) => optionalNonEmptyString(raw[key], key); + return { + payload, id: optional("id"), status: optional("status"), claim: optional("claim"), + excluded: requireStringArray(raw.excluded, "excluded"), resume: optional("resume"), + done: requireBoolean(raw.done, "done"), + ready: raw.ready === null ? null : requireBoolean(raw.ready, "ready"), + readyTruthy: requireBoolean(raw.ready_truthy, "ready_truthy"), + priority: requireInteger(raw.priority, "priority"), index: requireInteger(raw.index, "index"), + targetId: optional("target_id"), targetStatus: optional("target_status"), + targetClass: optional("target_class") ?? "advancement_task", + }; +} + +function ordered(items: readonly Item[]): Item[] { + // Stable sort retains persisted/source-lane order for equal public keys. + return [...items].sort((a, b) => a.priority - b.priority || a.index - b.index); +} + +function unique(items: readonly Item[]): Item[] { + const seen = new Set(); + return items.filter((entry) => { + const key = JSON.stringify([entry.id ?? "", entry.payload.text, entry.payload.index]); + if (seen.has(key)) return false; + seen.add(key); + return true; + }); +} + +function deferred(source: readonly Item[]): Item[] { + return ordered(source.filter((entry) => entry.payload.text && entry.status === "deferred") + .map((entry) => ({ ...entry, payload: { + ...entry.payload, + ...(entry.resume ? { resume_when: entry.resume } : {}), + ...(entry.payload.resume_ready === undefined ? {} : { resume_ready: entry.readyTruthy }), + } }))); +} + +function blocked(source: readonly Item[]): Item[] { + return ordered(unique(source.filter((entry) => entry.payload.text && !entry.done && + entry.resume !== null && entry.ready === false))); +} + +function monitorBlocked(source: readonly Item[]): Item[] { + return ordered(unique(source.filter((entry) => entry.payload.task_class === "advancement_task" && + entry.diagnosis?.state === "invalid" && entry.diagnosis.reason === "monitor_completion_requires_replan" + ).map((entry) => ({ ...entry, payload: { + ...entry.payload, ...(entry.targetId ? { blocking_monitor_todo_id: entry.targetId } : {}), + } })))); +} + +function claimLane(entry: Item, agent: string): ClaimLane { + if (entry.excluded.includes(agent)) return "executor_excluded_self"; + if (!entry.claim) return "unclaimed"; + return entry.claim === agent ? "current_agent" : "other_agent"; +} + +function payloads(items: readonly Item[], limit?: number): JsonObject[] { + return (limit === undefined ? items : items.slice(0, limit)).map((entry) => entry.payload); +} + +function claimLanes(items: readonly Item[], agent: string, limit: number, + family: "deferred_resume" | "monitor_blocked_resume"): JsonObject { + const result: JsonObject = {}; + for (const lane of ["current_agent", "unclaimed", "other_agent", "executor_excluded_self"] as const) { + const selected = items.filter((entry) => claimLane(entry, agent) === lane); + result[`${lane}_${family}_candidates`] = payloads(selected, limit); + result[`${lane}_${family}_count`] = selected.length; + } + result[`${family}_selection_policy`] = family === "deferred_resume" ? DEFERRED_POLICY : MONITOR_POLICY; + return result; +} + +function successorWaits(resumeBlocked: readonly Item[], deferredItems: readonly Item[], agent: string | null): Item[] { + if (!agent) return []; + return ordered(unique([...resumeBlocked, ...deferredItems.filter((entry) => entry.payload.resume_ready === false)]) + .filter((entry) => { + const lane = claimLane(entry, agent); + return entry.payload.task_class === "advancement_task" && + (lane === "current_agent" || lane === "unclaimed") && entry.resume !== null && + jsonObject(entry.payload.resume_condition)?.satisfied === false && + entry.diagnosis?.state === "pending" && entry.diagnosis.kind !== "monitor_changed"; + }).map((entry) => ({ ...entry, payload: { + ...entry.payload, resume_when: entry.resume, resume_ready: false, + } }))).sort((a, b) => Number(b.claim === agent) - Number(a.claim === agent)); +} + +function resolveCapacity(items: readonly Item[], source: readonly Item[], capabilities: string[]): Item[] { + const result = evaluateTodoResumeConditions({ + schema_version: TODO_RESUME_EVALUATION_REQUEST_SCHEMA_VERSION, + items: payloads(items).filter((entry) => entry.todo_id), + source_items: payloads([...source, ...items]).filter((entry) => entry.todo_id), + available_capabilities: capabilities, kinds: ["capacity_available"], rollout_events: [], + }); + const conditions = new Map(); + for (const row of result.conditions as JsonObject[]) { + conditions.set(String(row.todo_id), requireJsonObject(row.condition, "resume condition")); + } + return items.map((entry) => { + const condition = conditions.get(entry.id ?? ""); + if (!condition) return entry; + const ready = condition.satisfied === true; + return { ...entry, ready, readyTruthy: ready, payload: { + ...entry.payload, resume_condition: condition, resume_ready: ready, + } }; + }); +} + +function decodeSources(value: unknown): Record { + const rawSources = requireJsonObject(value, "sources"); + const sources = {} as Record; + for (const key of SOURCE_KEYS) { + if (!Array.isArray(rawSources[key])) throw new EffectRuntimeRequestError(`${key} must be an array`); + sources[key] = rawSources[key].map(item); + } + return sources; +} + +function deferredPlan(sources: Record, capabilities: unknown) { + let deferredItems = deferred(sources.deferred_items.length ? sources.deferred_items : sources.items); + let candidates = deferred(sources.deferred_resume_candidates).filter((entry) => entry.payload.resume_ready === true); + let capacityFields: JsonObject | null = null; + if (capabilities !== null) { + deferredItems = resolveCapacity(deferredItems, sources.items, + requireStringArray(capabilities, "available_capabilities")); + candidates = deferredItems.filter((entry) => entry.payload.resume_ready === true); + capacityFields = { deferred_items: payloads(deferredItems), deferred_resume_candidates: payloads(candidates) }; + // Existing summary readers treat an empty explicit deferred lane as absent, + // including after capacity resolution; retain that compatibility fallback. + if (!deferredItems.length) deferredItems = deferred(sources.items); + } + return { deferredItems, candidates, capacityFields }; +} + +function monitorIds(sources: Record): Set { + return new Set(["monitor_open_items", "current_agent_claimed_monitor_items", + "claimed_monitor_open_items", "items", "backlog_items", "first_open_items"] + .flatMap((key) => sources[key as SourceKey]) + .filter((entry) => entry.id && entry.payload.task_class === "continuous_monitor") + .map((entry) => entry.id!)); +} + +function diagnoseSources(sources: Record): void { + const monitors = monitorIds(sources); + // Compatibility belongs here once, not in every agent-scope consumer. Old + // compact conditions may omit kind/class while the same snapshot has them. + for (const key of SOURCE_KEYS) sources[key] = sources[key].map((entry) => { + const condition = jsonObject(entry.payload.resume_condition); + if (!condition || !entry.resume) return entry; + const facts = { + ...condition, resume_when: entry.resume, + target_todo_id: entry.targetId, + target_status: entry.targetStatus, + target_task_class: monitors.has(entry.targetId ?? "") ? "continuous_monitor" : entry.targetClass, + }; + const diagnosis = diagnoseTodoResumeCondition(facts, entry.id); + if (diagnosis.state !== "invalid") return { ...entry, diagnosis }; + return { ...entry, diagnosis, ready: false, readyTruthy: false, payload: { + ...entry.payload, resume_ready: false, resume_condition: { + ...condition, satisfied: false, invalid_state: diagnosis.reason, + availability_reason: "resume_condition_invalid", + }, + } }; + }); +} + +interface DisplayOptions { + agent: string | null; + limit: number; + hasCount: boolean; + hasVisibleCount: boolean; + count: JsonObject[string]; +} + +function deferredVisibility(items: Item[], candidates: Item[], options: DisplayOptions): JsonObject { + const { agent, limit, hasCount, hasVisibleCount, count } = options; + return items.length || candidates.length || hasVisibleCount ? { + deferred_count: hasCount ? count : items.length, + deferred_visibility_limit: limit, deferred_items: payloads(items, limit), + deferred_resume_candidates: payloads(candidates, limit), + ...(agent ? claimLanes(candidates, agent, limit, "deferred_resume") : {}), + } : {}; +} + +function blockedVisibility(resumeBlocked: Item[], monitorItems: Item[], options: DisplayOptions): JsonObject { + const { agent, limit } = options; + return resumeBlocked.length ? { + resume_blocked_count: resumeBlocked.length, resume_blocked_items: payloads(resumeBlocked, limit), + ...(monitorItems.length ? { + monitor_blocked_resume_count: monitorItems.length, + monitor_blocked_resume_candidates: payloads(monitorItems, limit), + ...(agent ? claimLanes(monitorItems, agent, limit, "monitor_blocked_resume") : {}), + } : {}), + } : {}; +} + +/** One read-only decision for quota lanes, replan candidates and exact waits. + * No mutation, claim, lease grant, monitor poll, or recovery effect is emitted. */ +export function projectTodoResumePlanning(value: unknown): JsonObject { + const request = requireJsonObject(value, "resume planning request"); + if (request.schema_version !== RESUME_PLANNING_REQUEST) { + throw new EffectRuntimeRequestError("resume planning schema mismatch"); + } + const sources = decodeSources(request.sources); + diagnoseSources(sources); + const display: DisplayOptions = { + agent: optionalNonEmptyString(request.agent_id, "agent_id"), + limit: requireInteger(request.item_limit, "item_limit"), + hasCount: requireBoolean(request.has_deferred_count, "has_deferred_count"), + hasVisibleCount: requireBoolean(request.has_visible_deferred_count, "has_visible_deferred_count"), + count: request.deferred_count, + }; + const { deferredItems, candidates, capacityFields } = deferredPlan(sources, request.available_capabilities); + const resumeBlocked = blocked(sources.resume_blocked_items.length ? sources.resume_blocked_items + : [...sources.items, ...sources.backlog_items, ...sources.first_open_items]); + const monitorItems = monitorBlocked(resumeBlocked); + return { + schema_version: RESUME_PLANNING_RESULT, + deferred_lanes: deferredVisibility(deferredItems, candidates, display), + resume_blocked_lanes: blockedVisibility(resumeBlocked, monitorItems, display), + capacity_fields: capacityFields, + deferred_items: payloads(deferredItems), monitor_blocked_items: payloads(monitorItems), + blocked_successor_items: payloads(successorWaits(resumeBlocked, deferredItems, display.agent)), + }; +} diff --git a/loopx/control_plane/todos/route_continuation.py b/loopx/control_plane/todos/route_continuation.py index 111e8e42ba..8c5ce2283e 100644 --- a/loopx/control_plane/todos/route_continuation.py +++ b/loopx/control_plane/todos/route_continuation.py @@ -4,13 +4,10 @@ from .contract import ( TODO_TASK_CLASS_ADVANCEMENT, - normalize_required_write_scopes, normalize_todo_claimed_by, - normalize_todo_decision_scope, normalize_todo_excluded_agents, - normalize_todo_required_decision_scopes, - normalize_todo_task_class, ) +from .compact_projection import compact_todo_projection_item, projection_task_class from .handoff_gate import todo_summary_handoff_gates from .projection import todo_projection_sort_key @@ -22,98 +19,6 @@ ) -def _todo_task_class(item: dict[str, Any]) -> str: - text = " ".join( - str(value or "") - for value in (item.get("title"), item.get("text")) - if str(value or "").strip() - ) - return normalize_todo_task_class( - item.get("task_class"), - text=text, - action_kind=item.get("action_kind"), - ) - - -def _compact_route_continuation_item( - item: dict[str, Any], - *, - text: str, -) -> dict[str, Any]: - compact: dict[str, Any] = { - "index": item.get("index"), - "text": text, - } - for key in ( - "schema_version", - "todo_id", - "role", - "status", - "priority", - "title", - "archive_state", - "source_section", - "task_class", - "action_kind", - "task_domain", - "task_repository", - "continuation_policy", - "required_write_scopes", - "required_capabilities", - "target_capabilities", - "decision_scope", - "required_decision_scopes", - "claimed_by", - "blocks_agent", - "excluded_agents", - "unblocks_todo_id", - "resume_when", - "resume_monitor_generation", - "resume_condition", - "resume_ready", - "no_followup", - "successor_todo_ids", - "target_key", - "cadence", - "next_due_at", - "expires_at", - "last_checked_at", - "result_hash", - "consecutive_no_change", - "material_change", - "material_change_generation", - "max_no_change_before_replan", - "route_continuation_replan_required", - "route_continuation_reason", - "route_id", - "route_key", - "completed_at", - "updated_at", - "superseded_by", - ): - if item.get(key) is not None: - compact[key] = item.get(key) - required_write_scopes = normalize_required_write_scopes(compact.get("required_write_scopes")) - if required_write_scopes: - compact["required_write_scopes"] = required_write_scopes - else: - compact.pop("required_write_scopes", None) - decision_scope = normalize_todo_decision_scope(compact.get("decision_scope")) - if decision_scope: - compact["decision_scope"] = decision_scope - else: - compact.pop("decision_scope", None) - required_decision_scopes = normalize_todo_required_decision_scopes( - compact.get("required_decision_scopes") - ) - if required_decision_scopes: - compact["required_decision_scopes"] = required_decision_scopes - else: - compact.pop("required_decision_scopes", None) - compact["task_class"] = _todo_task_class(compact) - return compact - - def todo_summary_route_continuation_candidates( value: dict[str, Any], ) -> list[dict[str, Any]]: @@ -143,7 +48,7 @@ def todo_summary_route_continuation_candidates( if item.get("route_continuation_replan_required") is False: continue task_class = item.get("task_class") - if task_class is not None and _todo_task_class(item) != TODO_TASK_CLASS_ADVANCEMENT: + if task_class is not None and projection_task_class(item) != TODO_TASK_CLASS_ADVANCEMENT: continue text = str( item.get("text") @@ -162,7 +67,7 @@ def todo_summary_route_continuation_candidates( if not identity or identity in seen: continue seen.add(identity) - compact = _compact_route_continuation_item(item, text=text) + compact = compact_todo_projection_item(item, text=text) compact["route_continuation_replan_required"] = True if item.get("route_continuation_reason") is not None: compact["route_continuation_reason"] = item.get("route_continuation_reason") diff --git a/loopx/control_plane/todos/succession_warning.py b/loopx/control_plane/todos/succession_warning.py index 985630c365..3c76d62cb0 100644 --- a/loopx/control_plane/todos/succession_warning.py +++ b/loopx/control_plane/todos/succession_warning.py @@ -3,14 +3,11 @@ from typing import Any from ..agents.agent_scope import agent_scope_item_claimed_by_agent_or_unclaimed +from .compact_projection import compact_todo_projection_item from .contract import ( TODO_STATUS_OPEN, - normalize_required_write_scopes, normalize_todo_id, normalize_todo_id_list, - normalize_todo_decision_scope, - normalize_todo_required_decision_scopes, - normalize_todo_task_class, ) @@ -48,91 +45,6 @@ def build_open_parent_successor_advisory( } -def _compact_succession_warning_item(item: dict[str, Any]) -> dict[str, Any]: - compact: dict[str, Any] = { - "index": item.get("index"), - "text": item.get("text"), - } - for key in ( - "schema_version", - "todo_id", - "role", - "status", - "priority", - "title", - "archive_state", - "source_section", - "task_class", - "action_kind", - "task_domain", - "task_repository", - "continuation_policy", - "required_write_scopes", - "required_capabilities", - "target_capabilities", - "decision_scope", - "required_decision_scopes", - "claimed_by", - "blocks_agent", - "excluded_agents", - "unblocks_todo_id", - "resume_when", - "resume_monitor_generation", - "resume_condition", - "resume_ready", - "no_followup", - "successor_todo_ids", - "completion_continuation", - "completion_recovery", - "target_key", - "cadence", - "next_due_at", - "expires_at", - "last_checked_at", - "result_hash", - "consecutive_no_change", - "material_change", - "material_change_generation", - "max_no_change_before_replan", - "route_continuation_replan_required", - "route_continuation_reason", - "route_id", - "route_key", - "completed_at", - "completion_turn_key", - "updated_at", - "superseded_by", - "done", - "succession_tracked", - "recommended_action", - ): - if item.get(key) is not None: - compact[key] = item.get(key) - required_write_scopes = normalize_required_write_scopes(compact.get("required_write_scopes")) - if required_write_scopes: - compact["required_write_scopes"] = required_write_scopes - else: - compact.pop("required_write_scopes", None) - decision_scope = normalize_todo_decision_scope(compact.get("decision_scope")) - if decision_scope: - compact["decision_scope"] = decision_scope - else: - compact.pop("decision_scope", None) - required_decision_scopes = normalize_todo_required_decision_scopes( - compact.get("required_decision_scopes") - ) - if required_decision_scopes: - compact["required_decision_scopes"] = required_decision_scopes - else: - compact.pop("required_decision_scopes", None) - compact["task_class"] = normalize_todo_task_class( - compact.get("task_class"), - text=str(compact.get("text") or ""), - action_kind=compact.get("action_kind"), - ) - return compact - - def build_todo_succession_warning_lanes( summary: dict[str, Any], *, @@ -146,7 +58,11 @@ def build_todo_succession_warning_lanes( else summary.get("completed_without_successor_items") ) items = [ - _compact_succession_warning_item(item) + compact_todo_projection_item( + item, text=item.get("text"), task_class_text=str(item.get("text") or ""), + extra_fields=("completion_continuation", "completion_recovery", "completion_turn_key", + "done", "succession_tracked", "recommended_action"), + ) for item in (source_items or []) if isinstance(item, dict) ][:item_limit] diff --git a/loopx/control_plane/work_items/autonomous_replan_obligation.py b/loopx/control_plane/work_items/autonomous_replan_obligation.py index d9684708db..bd7a923fa2 100644 --- a/loopx/control_plane/work_items/autonomous_replan_obligation.py +++ b/loopx/control_plane/work_items/autonomous_replan_obligation.py @@ -13,10 +13,7 @@ normalize_todo_id_list, normalize_todo_replan_obligation_id, ) -from ..todos.deferred_resume import ( - todo_summary_deferred_items, - todo_summary_monitor_blocked_resume_items, -) +from ..todos.resume_planning import project_todo_resume_planning from .progress_observation import typed_progress_repeat_trigger from .replan_settlement import ( project_todo_lifecycle_settlement_reentry as project_todo_lifecycle_reentry_effect, @@ -472,9 +469,10 @@ def _future_due_blocking_monitor( monitors_by_id[monitor_todo_id] = monitor blocking_monitor_ids: set[str] = set() + resume_planning = project_todo_resume_planning(agent_todos) blocked_items = [ - *todo_summary_monitor_blocked_resume_items(agent_todos), - *todo_summary_deferred_items(agent_todos, "deferred_items"), + *resume_planning["monitor_blocked_items"], + *resume_planning["deferred_items"], ] for item in blocked_items: claimed_by = normalize_todo_claimed_by(item.get("claimed_by")) diff --git a/tests/control_plane/test_compact_todo_projection.py b/tests/control_plane/test_compact_todo_projection.py new file mode 100644 index 0000000000..8986d40f03 --- /dev/null +++ b/tests/control_plane/test_compact_todo_projection.py @@ -0,0 +1,67 @@ +"""The three public read projections share codecs, not field visibility policy.""" + +from copy import deepcopy + +import pytest + +from loopx.control_plane.todos.resume_planning import project_todo_resume_planning +from loopx.control_plane.todos.route_continuation import todo_summary_route_continuation_candidates +from loopx.control_plane.todos.succession_warning import build_todo_succession_warning_lanes + + +def projected(item: dict) -> list[dict]: + return [ + project_todo_resume_planning({"deferred_items": [item]})["deferred_items"][0], + todo_summary_route_continuation_candidates({"route_continuation_candidates": [item]})[0], + build_todo_succession_warning_lanes( + {"completed_without_successor_items": [item]}, item_limit=5, + )["completed_without_successor_items"][0], + ] + + +@pytest.mark.parametrize("scope", [None, "", [], ["../invalid"]]) +def test_empty_scope_omission_and_domain_extension_visibility(scope) -> None: + item = { + "todo_id": "todo_example", "index": 3, "text": " Build change ", + "status": "deferred", "required_write_scopes": scope, + "decision_scope": None, "required_decision_scopes": [], + "resume_ready": False, "no_followup": False, "claimed_by": "agent-a", + "completion_continuation": {}, "completion_recovery": "", "done": False, + "completion_turn_key": "turn-a", "succession_tracked": False, + "recommended_action": "Review successor", "unregistered_field": "must not leak", + } + before = deepcopy(item) + resume, route, succession = projected(item) + for result in (resume, route, succession): + for omitted in ("required_write_scopes", "decision_scope", "required_decision_scopes", + "unregistered_field"): + assert omitted not in result + assert result["claimed_by"] == "agent-a" + assert result["resume_ready"] is False + assert result["no_followup"] is False + for key in ("completion_continuation", "completion_recovery", "completion_turn_key", + "done", "succession_tracked", "recommended_action"): + assert key not in resume and key not in route + assert succession[key] == item[key] + assert resume["text"] == route["text"] == "Build change" + assert succession["text"] == item["text"] + assert item == before + + +def test_scope_normalization_is_shared_without_normalizing_unowned_fields() -> None: + for result in projected({ + "text": "Build", "status": "deferred", "required_write_scopes": "src/**,src/**;tests/**", + "required_capabilities": "compiler", "note": "not in this display contract", + }): + assert result["required_write_scopes"] == ["src/**", "tests/**"] + assert result["required_capabilities"] == "compiler" + assert "note" not in result + + +def test_legacy_task_class_text_inputs_remain_consumer_owned() -> None: + resume, route, succession = projected({ + "text": "Build change", "title": "Do not execute until approved", + "status": "deferred", + }) + assert resume["task_class"] == route["task_class"] == "continuous_monitor" + assert succession["task_class"] == "advancement_task" diff --git a/tests/control_plane/test_resume_planning_projection.py b/tests/control_plane/test_resume_planning_projection.py new file mode 100644 index 0000000000..502110fe58 --- /dev/null +++ b/tests/control_plane/test_resume_planning_projection.py @@ -0,0 +1,231 @@ +"""Public wait-lane invariants, independent of the projection implementation.""" + +from copy import deepcopy +import json +from pathlib import Path + +import pytest + +from canonical_authority_fixture import initialize_canonical_authority +from loopx.control_plane.testing.canary_harness import write_fixture_registry, run_json_cli_result +from loopx.control_plane.coordination.runtime_shadow import build_todo_runtime_shadow_projection +from loopx.control_plane.todos.active_state_todo_parser import parse_active_state_todos + +from loopx.control_plane.todos.resume_planning import project_todo_resume_planning +from loopx.control_plane.todos.quota_summary import select_quota_todo_summary + + +def waiting(todo_id: str, **fields: object) -> dict: + return { + "todo_id": todo_id, "index": 1, "text": "[P1] Wait for dependency", + "role": "agent", "task_class": "advancement_task", "status": "open", + "resume_when": "todo_done:todo_dependency", "resume_ready": False, + "resume_condition": {"satisfied": False, "kind": "todo_done", + "target_status": "open", "target_task_class": "advancement_task"}, + **fields, + } + + +def ids(items: list[dict]) -> list[str]: + return [item["todo_id"] for item in items] + + +def test_claim_priority_and_exclusion_are_independent_of_display_priority() -> None: + items = [waiting("todo_unclaimed", priority="P0"), + waiting("todo_current", claimed_by="agent-a", priority="P4"), + waiting("todo_other", claimed_by="agent-b"), + waiting("todo_excluded", claimed_by="agent-a", excluded_agents=["agent-a"])] + source = {"items": items} + before = deepcopy(source) + assert ids(project_todo_resume_planning(source, agent_id="agent-a")["blocked_successor_items"]) == [ + "todo_current", "todo_unclaimed", + ] + assert project_todo_resume_planning(source)["blocked_successor_items"] == [] + assert source == before + + +def test_monitor_generation_wait_is_not_the_monitor_completion_repair_lane() -> None: + condition = {"satisfied": False, "target_status": "open", + "target_task_class": "continuous_monitor", "target_todo_id": "todo_monitor"} + source = {"items": [ + waiting("todo_repair", resume_condition={**condition, "kind": "todo_done"}), + waiting("todo_wait", resume_when="monitor_changed:todo_monitor", + resume_condition={**condition, "kind": "monitor_changed"}), + ]} + result = project_todo_resume_planning(source, agent_id="agent-a", item_limit=1) + lanes = result["resume_blocked_lanes"] + assert lanes["resume_blocked_count"] == 2 + assert lanes["monitor_blocked_resume_count"] == 1 + assert ids(lanes["unclaimed_monitor_blocked_resume_candidates"]) == ["todo_repair"] + assert result["blocked_successor_items"] == [] + + +def test_legacy_monitor_lookup_is_diagnosed_once_before_agent_scope_selection() -> None: + from loopx.control_plane.agents.agent_scope import _agent_scope_monitor_blocked_resume_candidates + + rows = [waiting(f"todo_{lane}", resume_when="todo_done:todo_monitor", **scope, + resume_condition={"satisfied": False, "target_status": "open", "target": "todo_monitor"}) + for lane, scope in [("current", {"claimed_by": "agent-a"}), ("unclaimed", {}), + ("other", {"claimed_by": "agent-b"}), ("excluded", {"excluded_agents": ["agent-a"]})]] + result = project_todo_resume_planning({"items": rows, "monitor_open_items": [ + {"todo_id": "todo_monitor", "task_class": "continuous_monitor", "status": "open", "text": "Observe"}, + ]}, agent_id="agent-a") + lanes = result["resume_blocked_lanes"] + selected = _agent_scope_monitor_blocked_resume_candidates(lanes, agent_id="agent-a") + assert ids(selected) == ["todo_current", "todo_unclaimed"] + assert all(row["resume_condition"]["invalid_state"] == "monitor_completion_requires_replan" for row in selected) + assert all(row["blocking_monitor_todo_id"] == "todo_monitor" for row in selected) + assert lanes["other_agent_monitor_blocked_resume_count"] == 1 + assert lanes["executor_excluded_self_monitor_blocked_resume_count"] == 1 + + +def test_invalid_condition_cannot_become_an_exact_pending_successor_wait() -> None: + source = {"items": [waiting("todo_self", resume_when="todo_done:todo_self", + resume_condition={"kind": "todo_done", "target_todo_id": "todo_self", "satisfied": False})]} + result = project_todo_resume_planning(source, agent_id="agent-a") + assert result["blocked_successor_items"] == [] + row = result["resume_blocked_lanes"]["resume_blocked_items"][0] + assert row["resume_condition"]["invalid_state"] == "dependency_self_reference" + + +def test_invalid_deferred_condition_cannot_reuse_a_stale_ready_flag() -> None: + row = waiting("todo_stale", status="deferred", resume_ready=True, + resume_when="todo_done:todo_monitor", resume_condition={ + "kind": "todo_done", "satisfied": True, "target_status": "open", + "target_todo_id": "todo_monitor", "target_task_class": "continuous_monitor", + }) + result = project_todo_resume_planning({"deferred_items": [row], "deferred_resume_candidates": [row]}, agent_id="agent-a") + assert result["deferred_lanes"]["unclaimed_deferred_resume_count"] == 0 + assert result["deferred_items"][0]["resume_ready"] is False + assert row["resume_ready"] is True # Diagnosis is a projection, not an edit. + + +@pytest.mark.parametrize("promoted", [False, True]) +@pytest.mark.parametrize("kind", ["todo_done", "monitor_changed"]) +def test_public_quota_distinguishes_monitor_repair_from_generation_wait(tmp_path: Path, promoted: bool, kind: str) -> None: + runtime, registry, state = tmp_path / "runtime", tmp_path / "registry.json", tmp_path / "state.md" + state.write_text( + "---\nstatus: active\n---\n# Goal\n## Objective\nBuild a checked change.\n\n## Agent Todo\n" + "- [ ] [P1] Continue after observation.\n" + f" \n" + "- [ ] [P2] Observe dependency.\n" + " \n", + encoding="utf-8", + ) + write_fixture_registry(project=tmp_path, runtime_root=runtime, registry_path=registry, + goal_id="goal-a", domain="resume-planning", adapter_kind="generic_project_goal_v0", + state_file=str(state), registered_agents=["agent-a", "agent-b"], quota_allowed_slots=None) + if promoted: + goal = json.loads(registry.read_text())["goals"][0] + fields = parse_active_state_todos(state.read_text(), goal=goal, item_limit=None) + projection = build_todo_runtime_shadow_projection(goal_id="goal-a", todos=fields["agent_todos"]["items"], handoff_mode="soft_claim") + initialize_canonical_authority(runtime, "goal-a", projection, state_path=state) + state.unlink() + before = state.read_bytes() if state.exists() else None + code, packet = run_json_cli_result("quota", "should-run", "--goal-id", "goal-a", "--agent-id", "agent-a", + "--include-detail", "agent-todos", "--scan-path", str(tmp_path), registry_path=registry, runtime_root=runtime) + assert code == 0, packet + summary = packet["agent_todo_summary"] + condition = summary["resume_blocked_items"][0]["resume_condition"] + assert condition["availability_reason"] == ("resume_condition_invalid" if kind == "todo_done" else "resume_condition_pending") + if kind == "todo_done": + assert summary["current_agent_monitor_blocked_resume_count"] == 1 + assert packet["work_lane_contract"]["obligation"] == "repair_resume_gate_or_close_standing_monitor" + assert packet["work_lane_contract"]["selected_todo_id"] == "todo_waiting" + else: + assert not summary.get("monitor_blocked_resume_count") + assert condition["baseline_generation"] == 3 + assert condition["material_change_generation"] == 3 + assert (state.read_bytes() if state.exists() else None) == before + + +def test_ready_deferred_lanes_keep_full_counts_before_truncation() -> None: + items = [waiting(f"todo_ready_{i}", status="deferred", resume_ready=True, + claimed_by="agent-a") for i in range(4)] + items.append(waiting("todo_excluded", status="deferred", resume_ready=True, + excluded_agents=["agent-a"])) + source = {"deferred_items": items, "deferred_resume_candidates": items} + lanes = project_todo_resume_planning(source, agent_id="agent-a", item_limit=1)["deferred_lanes"] + assert lanes["deferred_count"] == 5 + assert lanes["current_agent_deferred_resume_count"] == 4 + assert len(lanes["current_agent_deferred_resume_candidates"]) == 1 + assert lanes["executor_excluded_self_deferred_resume_count"] == 1 + assert lanes["unclaimed_deferred_resume_count"] == 0 + + +def test_ignored_rows_do_not_turn_an_explicit_lane_into_fallback_input() -> None: + source = {"items": [waiting("todo_fallback", status="deferred")], + "deferred_items": [None], "resume_blocked_items": [None]} + result = project_todo_resume_planning(source, agent_id="agent-a") + assert result["deferred_items"] == [] + assert result["resume_blocked_lanes"] == {} + + +def test_quota_composes_capacity_and_visibility_in_one_request_per_source(monkeypatch) -> None: + from loopx.control_plane.todos import resume_planning + + original = resume_planning.effect_runtime_result + calls = [] + + def record(method, params): + calls.append(method) + return original(method, params) + + monkeypatch.setattr(resume_planning, "effect_runtime_result", record) + summary = select_quota_todo_summary( + {"schema_version": "todo_summary_v0", "items": [waiting("todo_capacity", status="deferred", + resume_when="capacity_available:compiler")], "total_count": 1, "deferred_count": 1}, + None, agent_identity={"agent_id": "agent-a"}, available_capabilities=["compiler"], + ) + assert summary["unclaimed_deferred_resume_count"] == 1 + assert calls == ["todo.resume_planning.project"] + + +def test_unavailable_typed_owner_does_not_fall_back_to_python_selection(monkeypatch) -> None: + from loopx.control_plane.todos import resume_planning + + def unavailable(*_args, **_kwargs): + raise RuntimeError("isolated runtime unavailable") + + monkeypatch.setattr(resume_planning, "effect_runtime_result", unavailable) + source = {"items": [waiting("todo_wait")]} + with pytest.raises(RuntimeError, match="isolated runtime unavailable"): + project_todo_resume_planning(source) + + +@pytest.mark.parametrize("promoted", [False, True]) +def test_public_quota_capacity_wait_uses_provider_after_cutover_without_display_write( + tmp_path: Path, promoted: bool, +) -> None: + runtime, registry, state = tmp_path / "runtime", tmp_path / "registry.json", tmp_path / "state.md" + state.write_text( + "---\nstatus: active\n---\n# Goal\n## Objective\nBuild a checked change.\n\n" + "## Agent Todo\n- [ ] [P1] Continue when compiler capacity returns.\n" + " \n", + encoding="utf-8", + ) + write_fixture_registry( + project=tmp_path, runtime_root=runtime, registry_path=registry, + goal_id="goal-a", domain="resume-planning", adapter_kind="generic_project_goal_v0", + state_file=str(state), registered_agents=["agent-a", "agent-b"], quota_allowed_slots=None, + ) + if promoted: + goal = json.loads(registry.read_text())["goals"][0] + fields = parse_active_state_todos(state.read_text(), goal=goal, item_limit=None) + projection = build_todo_runtime_shadow_projection( + goal_id="goal-a", todos=fields["agent_todos"]["items"], handoff_mode="soft_claim", + ) + initialize_canonical_authority(runtime, "goal-a", projection, state_path=state) + state.unlink() + before = state.read_bytes() if state.exists() else None + for capabilities, count in [((), 0), (("--available-capability", "compiler"), 1)]: + code, packet = run_json_cli_result( + "quota", "should-run", "--goal-id", "goal-a", "--agent-id", "agent-a", + "--scan-path", str(tmp_path), *capabilities, registry_path=registry, runtime_root=runtime, + ) + assert code == 0, packet + assert packet["agent_todo_summary"]["current_agent_deferred_resume_count"] == count + if count: + assert packet["effective_action"] == "successor_replan_required" + assert (state.read_bytes() if state.exists() else None) == before diff --git a/tests/control_plane_ts/resume_planning.test.ts b/tests/control_plane_ts/resume_planning.test.ts new file mode 100644 index 0000000000..208cd5d13f --- /dev/null +++ b/tests/control_plane_ts/resume_planning.test.ts @@ -0,0 +1,80 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import type { JsonObject } from "../../loopx/control_plane/effect_program.ts"; +import { projectTodoResumePlanning, RESUME_PLANNING_REQUEST } from "../../loopx/control_plane/todos/resume_planning.ts"; + +function fact(id: string, overrides: JsonObject = {}): JsonObject { + return { + payload: { todo_id: id, index: 1, text: "Wait for work", role: "agent", + task_class: "advancement_task", status: "deferred", resume_when: "capacity_available:compiler", + resume_ready: false }, + id, status: "deferred", claim: null, excluded: [], done: false, + resume: "capacity_available:compiler", ready: false, ready_truthy: false, + priority: 1, index: 1, target_id: null, target_status: null, + target_class: "advancement_task", ...overrides, + }; +} + +function request(overrides: JsonObject = {}): JsonObject { + return { + schema_version: RESUME_PLANNING_REQUEST, agent_id: "agent-a", item_limit: 2, + has_deferred_count: false, has_visible_deferred_count: false, deferred_count: null, + available_capabilities: null, + sources: { items: [], backlog_items: [], first_open_items: [], deferred_items: [], + deferred_resume_candidates: [], resume_blocked_items: [], monitor_open_items: [], + current_agent_claimed_monitor_items: [], claimed_monitor_open_items: [] }, + ...overrides, + }; +} + +test("capacity evaluation and claim lanes use one snapshot without mutating it", () => { + const input = request({ available_capabilities: ["compiler"] }); + (input.sources as JsonObject).items = [fact("todo_capacity", { claim: "agent-a" }), + fact("todo_excluded", { excluded: ["agent-a"] })]; + const before = structuredClone(input); + const result = projectTodoResumePlanning(input); + const lanes = result.deferred_lanes as JsonObject; + assert.equal(lanes.current_agent_deferred_resume_count, 1); + assert.equal(lanes.executor_excluded_self_deferred_resume_count, 1); + assert.equal(lanes.unclaimed_deferred_resume_count, 0); + assert.equal((result.blocked_successor_items as unknown[]).length, 0); + assert.deepEqual(input, before); + input.available_capabilities = []; + const unavailable = projectTodoResumePlanning(input).deferred_lanes as JsonObject; + assert.equal(unavailable.current_agent_deferred_resume_count, 0); +}); + +test("a large wait source retains total counts independently from the display bound", () => { + const input = request(); + const rows = Array.from({ length: 257 }, (_, i) => fact(`todo_wait_${i}`)); + (input.sources as JsonObject).deferred_items = rows; + const result = projectTodoResumePlanning(input); + const lanes = result.deferred_lanes as JsonObject; + assert.equal(lanes.deferred_count, 257); + assert.equal((lanes.deferred_items as unknown[]).length, 2); + assert.equal((result.deferred_items as unknown[]).length, 257); + assert.deepEqual(result.resume_blocked_lanes, {}); +}); + +test("equal priority/index preserves source order, and zero display keeps counts", () => { + const input = request({ item_limit: 0, has_deferred_count: true, deferred_count: 99 }); + (input.sources as JsonObject).items = [fact("todo_zed"), fact("todo_alpha")]; + const result = projectTodoResumePlanning(input); + assert.deepEqual((result.deferred_items as JsonObject[]).map((row) => row.todo_id), + ["todo_zed", "todo_alpha"]); + assert.equal((result.deferred_lanes as JsonObject).deferred_count, 99); + assert.deepEqual((result.deferred_lanes as JsonObject).deferred_items, []); +}); + +test("wire faults fail closed before selection instead of inventing readiness", () => { + for (const override of [{ schema_version: "future" }, { item_limit: "2" }, + { has_deferred_count: null }, { available_capabilities: true }, { sources: {} }]) { + assert.throws(() => projectTodoResumePlanning(request(override))); + } + for (const override of [{ ready: "false" }, { done: 0 }, { excluded: "agent-a" }, + { priority: null }, { payload: {} }]) { + const input = request(); + (input.sources as JsonObject).items = [fact("todo_invalid", override)]; + assert.throws(() => projectTodoResumePlanning(input)); + } +}); diff --git a/tests/control_plane_ts/todo_resume_condition.test.ts b/tests/control_plane_ts/todo_resume_condition.test.ts index a603a17e99..cf08427b31 100644 --- a/tests/control_plane_ts/todo_resume_condition.test.ts +++ b/tests/control_plane_ts/todo_resume_condition.test.ts @@ -7,6 +7,7 @@ import { TODO_RESUME_EVALUATION_REQUEST_SCHEMA_VERSION, TODO_RESUME_NORMALIZE_REQUEST_SCHEMA_VERSION, evaluateTodoResumeConditions, + diagnoseTodoResumeCondition, normalizeTodoResumeWhen, planTodoExternalWaitTransition, } from "../../loopx/control_plane/todos/resume_condition.ts"; @@ -37,6 +38,56 @@ test("resume syntax is normalized by the typed Todo boundary", () => { }), null); }); +test("live monitor completion is invalid, but historical completion remains satisfied", () => { + for (const status of ["open", "blocked", "deferred", "done"]) { + const result = evaluateTodoResumeConditions({ + schema_version: TODO_RESUME_EVALUATION_REQUEST_SCHEMA_VERSION, + items: [todo("todo_waiting", "open", "advancement_task", { resume_when: "todo_done:todo_monitor" })], + source_items: [todo("todo_monitor", status, "continuous_monitor", { archive_state: "archived" })], + }); + const condition = (result.conditions as Array<{ condition: Record }>)[0].condition; + assert.equal(condition.satisfied, status === "done"); + assert.equal(condition.availability_reason, status === "done" ? "resume_condition_satisfied" : "resume_condition_invalid"); + assert.equal(condition.invalid_state, status === "done" ? undefined : "monitor_completion_requires_replan"); + } +}); + +test("self-dependency is not satisfied even when a stale completed row says done", () => { + for (const kind of ["todo_done", "monitor_changed"]) { + const result = evaluateTodoResumeConditions({ + schema_version: TODO_RESUME_EVALUATION_REQUEST_SCHEMA_VERSION, + items: [todo("todo_waiting", "done", "continuous_monitor", { + resume_when: `${kind}:todo_waiting`, resume_monitor_generation: 1, material_change_generation: 2, + })], source_items: [], + }); + const condition = (result.conditions as Array<{ condition: Record }>)[0].condition; + assert.equal(condition.satisfied, false); + assert.equal(condition.invalid_state, "dependency_self_reference"); + } +}); + +test("missing completion target stays pending rather than claiming proof of invalidity", () => { + const result = evaluateTodoResumeConditions({ + schema_version: TODO_RESUME_EVALUATION_REQUEST_SCHEMA_VERSION, + items: [todo("todo_waiting", "open", "advancement_task", { resume_when: "todo_done:todo_missing" })], + source_items: [], + }); + const condition = (result.conditions as Array<{ condition: Record }>)[0].condition; + assert.equal(condition.availability_reason, "resume_condition_pending"); + assert.equal(condition.invalid_state, undefined); +}); + +test("legacy compact diagnosis infers only typed resume syntax, not prose or another kind", () => { + const legacy = { resume_when: "todo_done:todo_monitor", satisfied: false, + target_status: "open", target_task_class: "continuous_monitor" }; + assert.deepEqual(diagnoseTodoResumeCondition(legacy), { + kind: "todo_done", state: "invalid", reason: "monitor_completion_requires_replan", + }); + for (const kind of ["monitor_changed", "capacity_available", "pr_merged"]) { + assert.equal(diagnoseTodoResumeCondition({ ...legacy, kind }).state, "pending"); + } +}); + test("one reducer evaluates Todo, PR, capacity, and monitor resume conditions", () => { const result = evaluateTodoResumeConditions({ schema_version: TODO_RESUME_EVALUATION_REQUEST_SCHEMA_VERSION,