From 8c325299ac412b2464023c006c8e4bffb3048183 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 20 Sep 2026 19:26:20 +0800 Subject: [PATCH 1/2] fix(control-plane): reconcile redirects and watch-only monitors Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- docs/project-agent-todo-contract.md | 19 ++++ loopx/capabilities/issue_fix/README.md | 6 ++ loopx/capabilities/issue_fix/README.zh-CN.md | 5 + loopx/capabilities/issue_fix/pr_lifecycle.py | 24 ++++- .../issue_fix/pr_lifecycle_rollout.py | 14 +++ .../goals/goal_frontier/__init__.py | 52 +++++++--- loopx/control_plane/todos/quota_selection.py | 4 +- loopx/control_plane/todos/quota_selection.ts | 12 ++- loopx/control_plane/todos/quota_summary.py | 27 ++++- loopx/control_plane/todos/summary_lanes.ts | 13 ++- loopx/control_plane/todos/todo_semantics.py | 58 +++++++++++ loopx/control_plane/todos/todo_summary.py | 26 ++--- .../work_items/interaction_contract.py | 47 +++++++++ loopx/control_plane/work_items/work_lane.py | 53 ++++++++-- .../work_items/work_lane_context.py | 18 +++- .../test_issue_fix_pr_gate_reconcile.py | 78 +++++++++++++++ .../test_work_lane_contract_core.py | 98 +++++++++++++++++++ .../control_plane_ts/capability_gate.test.ts | 3 +- .../control_plane_ts/quota_selection.test.ts | 7 +- .../todo_resume_condition.test.ts | 42 ++++++++ .../todo_summary_lanes.test.ts | 17 ++++ 21 files changed, 568 insertions(+), 55 deletions(-) diff --git a/docs/project-agent-todo-contract.md b/docs/project-agent-todo-contract.md index 832586f6c4..d67f47ff2a 100644 --- a/docs/project-agent-todo-contract.md +++ b/docs/project-agent-todo-contract.md @@ -110,6 +110,25 @@ loopx todo add \ --action-kind monitor ``` +`watch_only=true` changes convergence and replan semantics, not schedulability. +A scheduled watch-only monitor remains eligible at `next_due_at`, but it never +creates autonomous replan pressure and never preempts runnable advancement. +When both are present, `interaction_contract` keeps advancement primary and +projects an optional, typed, no-spend `auxiliary_monitor_poll` route. +The canonical watch-only/ordinary-due partition is produced inside the existing +TypeScript Todo summary and quota-planning owners after Agent scope and +capability admission; Python compatibility code only adapts legacy facts and +renders the selected CLI/Lark route. + +`watch_only=true` 改变的是收敛与 replan 语义,而不是可调度性。带 +`next_due_at` 的 watch-only monitor 到期后仍可轮询,但不会制造 autonomous +replan 压力,也不会抢占 runnable advancement;二者同时存在时, +`interaction_contract` 保持 advancement 为主,并投影一条可选、typed、no-spend +的 `auxiliary_monitor_poll` 路由。 +watch-only/普通 due 的权威分区由既有 TypeScript Todo summary 与 quota-planning +owner 在 Agent scope 和 capability admission 之后生成;Python 兼容层只适配旧事实并 +渲染已选中的 CLI/Lark 路由。 + `--action-kind` is a public-safe token. Known generic tokens such as `run_eval`, `validate`, `rebuild`, `writeback`, `monitor`, and `poll` help the CLI project the lane consistently, but explicit `--task-class` is the authority diff --git a/loopx/capabilities/issue_fix/README.md b/loopx/capabilities/issue_fix/README.md index 7f1e709e42..4a5c2b2d90 100644 --- a/loopx/capabilities/issue_fix/README.md +++ b/loopx/capabilities/issue_fix/README.md @@ -424,6 +424,12 @@ should-run` pass can select a matched todo as ordinary runnable work. Replaying the same merged observation reuses the stable event id and creates no second transition. +If GitHub redirects a renamed repository, lifecycle reconciliation treats the +provider-returned PR URL as the canonical repository identity and records the +requested repository as an explicit alias source reference. The alias is valid +only for the same observed PR number; resume evaluation remains repository- +qualified and never matches the same number in an unrelated repository. + This is deliberately event-backed rather than webhook-code coupling: ```text diff --git a/loopx/capabilities/issue_fix/README.zh-CN.md b/loopx/capabilities/issue_fix/README.zh-CN.md index 08d6799a35..9c1f545d2a 100644 --- a/loopx/capabilities/issue_fix/README.zh-CN.md +++ b/loopx/capabilities/issue_fix/README.zh-CN.md @@ -407,6 +407,11 @@ scope、public-safe 且幂等的 `pr_merge` rollout event。Todo resume 投影 后续一次 `status` / `quota should-run` 就能把已匹配的 todo 当作普通 runnable work 选中。相同 merged observation 重放时复用稳定 event id,不会制造第二次 transition。 +如果 GitHub 对改名仓库返回重定向,lifecycle reconciliation 以 provider 返回的 PR URL +作为 canonical repository identity,并把请求时的旧仓库记为显式 alias source ref。 +alias 只对同一次观察到的 PR number 有效;resume evaluation 仍保持 repository-qualified, +不会误匹配其他仓库中的同号 PR。 + 这是一条事件驱动链,而不是 webhook 与业务代码硬耦合: ```text diff --git a/loopx/capabilities/issue_fix/pr_lifecycle.py b/loopx/capabilities/issue_fix/pr_lifecycle.py index c4ffbbdc9d..06358dc06d 100644 --- a/loopx/capabilities/issue_fix/pr_lifecycle.py +++ b/loopx/capabilities/issue_fix/pr_lifecycle.py @@ -945,6 +945,19 @@ def build_issue_fix_pr_lifecycle_monitor_packet( reference, timeout_seconds=fetch_timeout_seconds, ) + canonical_reference = reference + provider_url = payload.get("url") + if isinstance(provider_url, str) and provider_url.strip(): + provider_reference = normalise_github_issue_reference( + repo=repo, + issue_ref=pr_ref, + url=provider_url, + ) + if ( + provider_reference.get("kind") == "pull_request" + and provider_reference.get("number") == reference.get("number") + ): + canonical_reference = provider_reference if not issue_ref: raw_linked_issues = payload.get("closingIssuesReferences") or payload.get( "closing_issues_references" @@ -958,12 +971,17 @@ def build_issue_fix_pr_lifecycle_monitor_packet( issue_ref = f"issues_{number}" break observation = _build_observation( - repo=str(reference["repo"]), - pr_ref=str(reference["issue_ref"]), + repo=str(canonical_reference["repo"]), + pr_ref=str(canonical_reference["issue_ref"]), issue_ref=issue_ref, - reference=reference, + reference=canonical_reference, provider_payload=payload, ) + requested_repo = str(reference["repo"]) + canonical_repo = str(canonical_reference["repo"]) + if canonical_repo != requested_repo: + observation["requested_repo"] = requested_repo + observation["repository_aliases"] = [requested_repo] transition = _decide_transition(observation) maintainer_correction = ( normalise_issue_fix_maintainer_correction_input(maintainer_correction_input) diff --git a/loopx/capabilities/issue_fix/pr_lifecycle_rollout.py b/loopx/capabilities/issue_fix/pr_lifecycle_rollout.py index d778b16245..a7167d7662 100644 --- a/loopx/capabilities/issue_fix/pr_lifecycle_rollout.py +++ b/loopx/capabilities/issue_fix/pr_lifecycle_rollout.py @@ -47,10 +47,23 @@ def append_pr_merge_rollout_event( } pr_ref = f"{repo}#{number}" + repository_aliases = sorted( + { + str(alias).strip().lower() + for alias in observation.get("repository_aliases") or [] + if isinstance(alias, str) + and str(alias).strip() + and str(alias).strip().lower() != repo + } + ) event = build_rollout_event( goal_id=goal_id, event_kind="pr_merge", pr_ref=pr_ref, + source_refs=[ + {"kind": "pull_request", "ref": f"{alias}#{number}"} + for alias in repository_aliases + ], status="merged", summary=f"PR {pr_ref} merged; dependent resume conditions may proceed.", recorded_at=str( @@ -82,4 +95,5 @@ def append_pr_merge_rollout_event( "recorded_at": recorded_event["recorded_at"], "status": recorded_event.get("status"), "pr_ref": pr_ref, + "repository_aliases": repository_aliases, } diff --git a/loopx/control_plane/goals/goal_frontier/__init__.py b/loopx/control_plane/goals/goal_frontier/__init__.py index 610894e389..550cdb0ff1 100644 --- a/loopx/control_plane/goals/goal_frontier/__init__.py +++ b/loopx/control_plane/goals/goal_frontier/__init__.py @@ -578,11 +578,21 @@ def _count_advancement_items(items: Any, *, claimed_by: str | None = None) -> in def _summary_task_counts(summary: dict[str, Any] | None) -> dict[str, int]: open_count = _open_todo_count(summary) if not isinstance(summary, dict): - return {"open": open_count, "advancement": 0, "monitor": 0, "monitor_due": 0} + return { + "open": open_count, + "advancement": 0, + "monitor": 0, + "monitor_due": 0, + "watch_only_monitor_due": 0, + } executable = summary.get("executable_backlog_items") monitor_open = summary.get("monitor_open_items") - watch_only_count = ( - len( + if isinstance(summary.get("watch_only_monitor_count"), int): + watch_only_count = safe_non_negative_int( + summary.get("watch_only_monitor_count") + ) + elif isinstance(monitor_open, list): + watch_only_count = len( [ item for item in monitor_open @@ -591,10 +601,14 @@ def _summary_task_counts(summary: dict[str, Any] | None) -> dict[str, int]: and todo_item_is_watch_only_monitor(item) ] ) - if isinstance(monitor_open, list) - else safe_non_negative_int(summary.get("watch_only_monitor_count")) - ) - open_count = max(0, open_count - watch_only_count) + else: + watch_only_count = 0 + if isinstance(summary.get("convergence_open_count"), int): + open_count = safe_non_negative_int(summary.get("convergence_open_count")) + else: + # Compatibility-only fallback for summaries produced before the typed + # TypeScript lane owner exposed convergence_open_count. + open_count = max(0, open_count - watch_only_count) advancement_count = ( _count_advancement_items(executable) if isinstance(executable, list) @@ -608,8 +622,15 @@ def _summary_task_counts(summary: dict[str, Any] | None) -> dict[str, int]: ] ) ) - monitor_count = ( - len( + work_counts = summary.get("work_counts") + if isinstance(work_counts, dict): + monitor_count = max( + 0, + safe_non_negative_int(work_counts.get("monitor")) + - watch_only_count, + ) + elif isinstance(monitor_open, list): + monitor_count = len( [ item for item in monitor_open @@ -619,9 +640,10 @@ def _summary_task_counts(summary: dict[str, Any] | None) -> dict[str, int]: and not todo_item_is_watch_only_monitor(item) ] ) - if isinstance(monitor_open, list) - else safe_non_negative_int(summary.get("claimed_monitor_open_count")) - ) + else: + monitor_count = safe_non_negative_int( + summary.get("claimed_monitor_open_count") + ) return { "open": open_count, "advancement": advancement_count, @@ -631,6 +653,9 @@ def _summary_task_counts(summary: dict[str, Any] | None) -> dict[str, int]: safe_non_negative_int(summary.get("monitor_due_count")) - safe_non_negative_int(summary.get("watch_only_monitor_due_count")), ), + "watch_only_monitor_due": safe_non_negative_int( + summary.get("watch_only_monitor_due_count") + ), } @@ -1812,6 +1837,9 @@ def build_goal_frontier_projection( "agent_advancement_open_count": agent_counts.get("advancement", 0), "agent_monitor_open_count": agent_counts.get("monitor", 0), "agent_monitor_due_count": agent_counts.get("monitor_due", 0), + "agent_watch_only_monitor_due_count": agent_counts.get( + "watch_only_monitor_due", 0 + ), }, "remaining_advancement_frontier": { "current_agent_claimed_advancement_count": current_agent_claimed_advancement_count, diff --git a/loopx/control_plane/todos/quota_selection.py b/loopx/control_plane/todos/quota_selection.py index 421b00ead8..470a70f045 100644 --- a/loopx/control_plane/todos/quota_selection.py +++ b/loopx/control_plane/todos/quota_selection.py @@ -13,7 +13,8 @@ ) from .todo_semantics import ( todo_item_has_removed_continuation_policy, todo_item_is_actionable_open, - todo_item_is_due_monitor, todo_item_task_class, todo_projection_sort_key, + todo_item_is_due_monitor, todo_item_is_watch_only_monitor, + todo_item_task_class, todo_projection_sort_key, todo_summary_monitor_writeback_supported, ) from .resume_planning import build_todo_resume_planning_request @@ -46,6 +47,7 @@ def encode(item: dict[str, Any]) -> dict[str, Any]: "removed": todo_item_has_removed_continuation_policy(item), "actionable": todo_item_is_actionable_open(item), "due": todo_item_is_due_monitor(item), + "watch_only": todo_item_is_watch_only_monitor(item), "task_class": todo_item_task_class(item), "priority": priority, "index": index, "profile_rank": agent_profile_candidate_rank(item, agent_profile=profile), diff --git a/loopx/control_plane/todos/quota_selection.ts b/loopx/control_plane/todos/quota_selection.ts index 582d728bc5..6bc5295d3f 100644 --- a/loopx/control_plane/todos/quota_selection.ts +++ b/loopx/control_plane/todos/quota_selection.ts @@ -11,7 +11,7 @@ interface Row { payload: JsonObject; display: JsonObject; claim: string | null; bound: string | null; blocks: string | null; excluded: readonly string[]; global: boolean; gate: boolean; removed: boolean; actionable: boolean; - due: boolean; taskClass: string; priority: number; index: number; + due: boolean; watchOnly: boolean; taskClass: string; priority: number; index: number; profileRank: number; missing: readonly string[]; rawClaimed: boolean; } @@ -25,7 +25,8 @@ function decodeRow(value: unknown, available?: readonly string[]): Row { claim: optional("claim"), bound: optional("bound"), blocks: optional("blocks"), excluded: requireStringArray(raw.excluded, "excluded"), global: boolean("global"), gate: boolean("gate"), removed: boolean("removed"), actionable: boolean("actionable"), - due: boolean("due"), taskClass: optional("task_class") ?? "advancement_task", + due: boolean("due"), watchOnly: raw.watch_only === undefined ? false : boolean("watch_only"), + taskClass: optional("task_class") ?? "advancement_task", priority: integer("priority"), index: integer("index"), profileRank: integer("profile_rank"), missing: available === undefined ? requireStringArray(raw.missing, "missing") : missingRequiredCapabilities(requireStringArray(raw.required, "required"), requireStringArray(raw.targets, "targets"), available), @@ -148,6 +149,8 @@ export function projectQuotaSelection(value: unknown): JsonObject { const scope = agent && !userMode ? claimScope(blocking, open, agent, profile, diagnostic) : null; const monitors = open.filter(row => row.actionable && row.taskClass === "continuous_monitor"); const due = supported ? monitors.filter(row => row.due && executableBy(row, agent)) : []; + const admittedDue = due.filter(row => !row.missing.length); + const watchOnlyMonitors = monitors.filter(row => row.watchOnly); const activeVisible = (row: Row) => userMode ? (row.gate ? gateApplies(row, agent) : actionApplies(row, agent)) : executableBy(row, agent); const gateFilter = otherGates.length ? { schema_version: "agent_scoped_user_gate_filter_v0", agent_id: agent, @@ -170,7 +173,10 @@ export function projectQuotaSelection(value: unknown): JsonObject { user_action_agent_scope_filter: actionFilter, other_agent_scoped_items: payloads(otherGates), agent_scope_filter: gateFilter, open_items: payloads(open), claim_scope: scope, executable_items: payloads(open.filter(row => row.actionable && row.taskClass === "advancement_task")), - monitor_items: payloads(monitors), monitor_due_items: payloads(due.filter(row => !row.missing.length)), + monitor_items: payloads(monitors), monitor_due_items: payloads(admittedDue), + watch_only_monitor_items: payloads(watchOnlyMonitors), + watch_only_monitor_due_items: payloads(admittedDue.filter(row => row.watchOnly)), + non_watch_only_monitor_due_items: payloads(admittedDue.filter(row => !row.watchOnly)), monitor_capability_blocked_due_items: due.filter(row => row.missing.length).map(row => ({...row.display, missing_capabilities: [...row.missing]})), claimed_open_items: payloads(blocking.filter(row => row.rawClaimed)), display_open_items: payloads(userMode ? [...open, ...actions] : open), diff --git a/loopx/control_plane/todos/quota_summary.py b/loopx/control_plane/todos/quota_summary.py index 214da11486..bdf94ffd5d 100644 --- a/loopx/control_plane/todos/quota_summary.py +++ b/loopx/control_plane/todos/quota_summary.py @@ -100,6 +100,8 @@ ) QUOTA_PAYLOAD_LANE_LIMITS = { "monitor_due_items": MONITOR_DUE_ITEM_LIMIT, + "watch_only_monitor_due_items": MONITOR_DUE_ITEM_LIMIT, + "non_watch_only_monitor_due_items": MONITOR_DUE_ITEM_LIMIT, "monitor_capability_blocked_due_items": QUOTA_PAYLOAD_DIAGNOSTIC_LANE_LIMIT, "monitor_schedule_gap_items": MONITOR_DUE_ITEM_LIMIT, "first_open_items": 3, @@ -181,6 +183,9 @@ class _QuotaTodoLanes: executable_items: list[dict[str, Any]] monitor_items: list[dict[str, Any]] monitor_due_items: list[dict[str, Any]] + watch_only_monitor_items: list[dict[str, Any]] + watch_only_monitor_due_items: list[dict[str, Any]] + non_watch_only_monitor_due_items: list[dict[str, Any]] monitor_capability_blocked_due_items: list[dict[str, Any]] claimed_open_items: list[dict[str, Any]] display_open_items: list[dict[str, Any]] @@ -401,6 +406,14 @@ def summarize_user_todos_for_quota( "monitor_open_items": lanes.monitor_items, "monitor_due_count": len(lanes.monitor_due_items), "monitor_due_items": lanes.monitor_due_items[:MONITOR_DUE_ITEM_LIMIT], + "watch_only_monitor_count": len(lanes.watch_only_monitor_items), + "watch_only_monitor_due_count": len(lanes.watch_only_monitor_due_items), + "watch_only_monitor_due_items": lanes.watch_only_monitor_due_items[ + :MONITOR_DUE_ITEM_LIMIT + ], + "non_watch_only_monitor_due_items": lanes.non_watch_only_monitor_due_items[ + :MONITOR_DUE_ITEM_LIMIT + ], "monitor_capability_blocked_due_count": len( lanes.monitor_capability_blocked_due_items ), @@ -420,11 +433,7 @@ def summarize_user_todos_for_quota( summary["advancement_frontier_revision_index"] = value[ "advancement_frontier_revision_index" ] - if value.get("watch_only_monitor_count"): - summary["watch_only_monitor_count"] = value["watch_only_monitor_count"] - summary["watch_only_monitor_due_count"] = value.get( - "watch_only_monitor_due_count", 0 - ) + if lanes.watch_only_monitor_items: summary["convergence_open_count"] = value.get("convergence_open_count") if recent_completed_advancement_items: summary["recent_completed_advancement_items"] = recent_completed_advancement_items @@ -821,6 +830,14 @@ def summarize_project_asset_todos_for_quota( "monitor_open_items": lanes.monitor_items, "monitor_due_count": len(lanes.monitor_due_items), "monitor_due_items": lanes.monitor_due_items[:MONITOR_DUE_ITEM_LIMIT], + "watch_only_monitor_count": len(lanes.watch_only_monitor_items), + "watch_only_monitor_due_count": len(lanes.watch_only_monitor_due_items), + "watch_only_monitor_due_items": lanes.watch_only_monitor_due_items[ + :MONITOR_DUE_ITEM_LIMIT + ], + "non_watch_only_monitor_due_items": lanes.non_watch_only_monitor_due_items[ + :MONITOR_DUE_ITEM_LIMIT + ], "monitor_capability_blocked_due_count": len( lanes.monitor_capability_blocked_due_items ), diff --git a/loopx/control_plane/todos/summary_lanes.ts b/loopx/control_plane/todos/summary_lanes.ts index ae2a045798..7051465bb3 100644 --- a/loopx/control_plane/todos/summary_lanes.ts +++ b/loopx/control_plane/todos/summary_lanes.ts @@ -8,7 +8,9 @@ export const TODO_SUMMARY_LANES = [ "open_items", "terminal_items", "deferred_items", "done_items", "projected_open_items", "projected_deferred_items", "budgeted_items", "claimed_open_items", "unclaimed_open_items", "executable_items", "blocker_items", "resume_blocked_items", "monitor_items", - "monitor_due_items", "monitor_schedule_gap_items", "claimed_advancement_items", + "monitor_due_items", "watch_only_monitor_items", "watch_only_monitor_due_items", + "non_watch_only_monitor_due_items", "convergent_open_items", + "monitor_schedule_gap_items", "claimed_advancement_items", "claimed_monitor_items", "active_next_action_items", "active_next_action_executable_items", ] as const; export type TodoSummaryLane = typeof TODO_SUMMARY_LANES[number]; @@ -82,6 +84,9 @@ export function projectTodoSummaryLanes(value: unknown): JsonObject { const monitors = ordered.filter(row => row.actionable && row.taskClass === "continuous_monitor"); const activeMonitor = (row: Row) => row.expiresAt === null || row.expiresAt > now; const due = monitors.filter(row => activeMonitor(row) && row.dueAt !== null && row.dueAt <= now); + const watchOnlyMonitors = monitors.filter(row => row.watchOnly); + const watchOnlyDue = due.filter(row => row.watchOnly); + const nonWatchOnlyDue = due.filter(row => !row.watchOnly); const missing = monitors.filter(row => activeMonitor(row) && !row.watchOnly && row.dueAt === null); const selected = { open_items: open, terminal_items: terminal, deferred_items: deferred, done_items: done, @@ -90,7 +95,11 @@ export function projectTodoSummaryLanes(value: unknown): JsonObject { unclaimed_open_items: ordered.filter(row => !row.claim), executable_items: executable, blocker_items: ordered.filter(row => row.status === "blocked" && row.taskClass === "blocker"), resume_blocked_items: ordered.filter(row => row.resumeBlocked), monitor_items: monitors, - monitor_due_items: due, monitor_schedule_gap_items: missing, + monitor_due_items: due, watch_only_monitor_items: watchOnlyMonitors, + watch_only_monitor_due_items: watchOnlyDue, + non_watch_only_monitor_due_items: nonWatchOnlyDue, + convergent_open_items: open.filter(row => !(row.taskClass === "continuous_monitor" && row.watchOnly)), + monitor_schedule_gap_items: missing, claimed_advancement_items: executable.filter(row => row.claim), claimed_monitor_items: monitors.filter(row => row.claim), active_next_action_items: ordered.filter(row => row.preferred), active_next_action_executable_items: executable.filter(row => row.preferred), diff --git a/loopx/control_plane/todos/todo_semantics.py b/loopx/control_plane/todos/todo_semantics.py index f99519009b..7a4934fd86 100644 --- a/loopx/control_plane/todos/todo_semantics.py +++ b/loopx/control_plane/todos/todo_semantics.py @@ -644,6 +644,64 @@ def todo_summary_monitor_due_items( ) +def todo_summary_watch_only_monitor_due_items( + summary: dict[str, Any] | None, + *, + task_text_keys: tuple[str, ...] = ("title", "text"), + text_mode: str = "label", +) -> list[dict[str, Any]]: + """Consume the typed watch-only partition; legacy summaries fall back safely.""" + + if isinstance(summary, dict) and isinstance( + summary.get("watch_only_monitor_due_items"), list + ): + return _summary_monitor_items( + summary, + projected_key="watch_only_monitor_due_items", + predicate=lambda _item: True, + task_text_keys=task_text_keys, + text_mode=text_mode, + ) + return [ + item + for item in todo_summary_monitor_due_items( + summary, + task_text_keys=task_text_keys, + text_mode=text_mode, + ) + if todo_item_is_watch_only_monitor(item) + ] + + +def todo_summary_non_watch_only_monitor_due_items( + summary: dict[str, Any] | None, + *, + task_text_keys: tuple[str, ...] = ("title", "text"), + text_mode: str = "label", +) -> list[dict[str, Any]]: + """Consume the typed ordinary-due partition; classify only legacy summaries.""" + + if isinstance(summary, dict) and isinstance( + summary.get("non_watch_only_monitor_due_items"), list + ): + return _summary_monitor_items( + summary, + projected_key="non_watch_only_monitor_due_items", + predicate=lambda _item: True, + task_text_keys=task_text_keys, + text_mode=text_mode, + ) + return [ + item + for item in todo_summary_monitor_due_items( + summary, + task_text_keys=task_text_keys, + text_mode=text_mode, + ) + if not todo_item_is_watch_only_monitor(item) + ] + + def todo_summary_monitor_due_count( summary: dict[str, Any] | None, *, diff --git a/loopx/control_plane/todos/todo_summary.py b/loopx/control_plane/todos/todo_summary.py index 168a7c3e20..db4bcda879 100644 --- a/loopx/control_plane/todos/todo_summary.py +++ b/loopx/control_plane/todos/todo_summary.py @@ -116,6 +116,10 @@ class _TodoGroupLanes: resume_blocked_items: list[dict[str, Any]] monitor_items: list[dict[str, Any]] monitor_due_items: list[dict[str, Any]] + watch_only_monitor_items: list[dict[str, Any]] + watch_only_monitor_due_items: list[dict[str, Any]] + non_watch_only_monitor_due_items: list[dict[str, Any]] + convergent_open_items: list[dict[str, Any]] monitor_schedule_gap_items: list[dict[str, Any]] claimed_advancement_items: list[dict[str, Any]] claimed_monitor_items: list[dict[str, Any]] @@ -1096,25 +1100,9 @@ def compact_evaluated_todo_group( item.get("route_continuation_replan_required") is True for item in [*items, *handoff_gates] ) - watch_only_monitor_items = [ - item - for item in lanes.monitor_items - if projection_todo_item_is_watch_only_monitor(item) - ] - watch_only_ids = { - normalize_todo_id(item.get("todo_id")) - for item in watch_only_monitor_items - } - watch_only_monitor_due_items = [ - item - for item in lanes.monitor_due_items - if normalize_todo_id(item.get("todo_id")) in watch_only_ids - ] - convergent_open_items = [ - item - for item in lanes.open_items - if normalize_todo_id(item.get("todo_id")) not in watch_only_ids - ] + watch_only_monitor_items = lanes.watch_only_monitor_items + watch_only_monitor_due_items = lanes.watch_only_monitor_due_items + convergent_open_items = lanes.convergent_open_items summary: dict[str, Any] = { "schema_version": "todo_summary_v0", "source_section": source_section, diff --git a/loopx/control_plane/work_items/interaction_contract.py b/loopx/control_plane/work_items/interaction_contract.py index a5ccbde6de..9b7e036e1b 100644 --- a/loopx/control_plane/work_items/interaction_contract.py +++ b/loopx/control_plane/work_items/interaction_contract.py @@ -1188,6 +1188,15 @@ def _build_interaction_agent_channel( "quiet_noop_allowed": quiet_noop_allowed, } channel.update(build_primary_action_projection(payload, mode=mode)) + work_lane = ( + payload.get("work_lane_contract") + if isinstance(payload.get("work_lane_contract"), Mapping) + else {} + ) + if isinstance(work_lane.get("auxiliary_monitor_poll"), Mapping): + channel["auxiliary_monitor_poll"] = dict( + work_lane["auxiliary_monitor_poll"] + ) if isinstance(payload.get("action_portfolio"), dict): channel["action_portfolio_ref"] = "$.action_portfolio" selection.apply_action_selection_agent_gate(channel, payload) @@ -1308,6 +1317,44 @@ def _build_interaction_cli_channel( spend_after_validation=spend_after_selection, ), } + work_lane = ( + payload.get("work_lane_contract") + if isinstance(payload.get("work_lane_contract"), Mapping) + else {} + ) + auxiliary_monitor = ( + work_lane.get("auxiliary_monitor_poll") + if isinstance(work_lane.get("auxiliary_monitor_poll"), Mapping) + else None + ) + if auxiliary_monitor is not None: + selected_monitor_id = normalize_todo_id( + auxiliary_monitor.get("selected_todo_id") + ) + try: + auxiliary_scheduler_args = render_scheduler_execution_args( + scheduler_execution_context=scheduler_execution_context, + ) + except ValueError: + auxiliary_scheduler_args = "" + if selected_monitor_id: + command_prefix = selection.render_cli_command_prefix( + runtime_root=runtime_root + ) + agent_identity = ( + payload.get("agent_identity") + if isinstance(payload.get("agent_identity"), dict) + else {} + ) + channel["auxiliary_monitor_poll"] = { + **dict(auxiliary_monitor), + "command": ( + f"{command_prefix} quota monitor-poll --goal-id " + f"{str(payload.get('goal_id') or '')}" + f"{_scoped_cli_args(agent_identity, available_capabilities=available_capabilities)}" + f"{auxiliary_scheduler_args} --todo-id {selected_monitor_id} --execute" + ), + } selection.apply_action_selection_cli_gate(channel, payload) if settlement_plan is not None and spend_after_selection: channel["settlement_plan"] = settlement_plan diff --git a/loopx/control_plane/work_items/work_lane.py b/loopx/control_plane/work_items/work_lane.py index ab64eb244d..30498394a2 100644 --- a/loopx/control_plane/work_items/work_lane.py +++ b/loopx/control_plane/work_items/work_lane.py @@ -68,6 +68,7 @@ def observe_work_lane( "target_key", "next_due_at", "expires_at", + "watch_only", "resume_when", "resume_ready", "blocking_monitor_todo_id", @@ -512,6 +513,8 @@ def build_work_lane_contract( todo_counts: dict[str, int], monitor_due_count: int, due_monitor_items: list[dict[str, Any]], + watch_only_due_monitor_items: list[dict[str, Any]] | None = None, + non_watch_only_due_monitor_items: list[dict[str, Any]] | None = None, first_advancement: dict[str, Any] | None, due_monitor_preempts_advancement: bool, outcome_followthrough: dict[str, Any] | None, @@ -544,22 +547,41 @@ def build_work_lane_contract( ) non_runnable_non_monitor_count = max(0, open_count - monitor_count) first_due_monitor = due_monitor_items[0] if due_monitor_items else None + watch_only_due_items = watch_only_due_monitor_items or [] + ordinary_due_items = ( + non_watch_only_due_monitor_items + if non_watch_only_due_monitor_items is not None + else due_monitor_items + ) + first_preemptive_due_monitor = ordinary_due_items[0] if ordinary_due_items else None blocked_by_monitor_items = resume_blocked_by_monitor_items or [] schedule_gap_items = monitor_schedule_gap_items or [] first_schedule_gap = schedule_gap_items[0] if schedule_gap_items else None monitor_debt_backoff_applies = bool( monitor_debt_backoff_active - and first_due_monitor + and first_preemptive_due_monitor and first_advancement - and todo_priority_rank(first_advancement) <= todo_priority_rank(first_due_monitor) + and todo_priority_rank(first_advancement) + <= todo_priority_rank(first_preemptive_due_monitor) ) effective_due_monitor_preemption = due_monitor_preempts_advancement and not ( has_advancement_todos and (monitor_attempt_already_recorded or monitor_debt_backoff_applies) ) - def due_monitor_contract(*, reason_codes: list[str]) -> dict[str, Any]: - selected = first_due_monitor or {} + def due_monitor_contract( + *, reason_codes: list[str], selected_monitor: dict[str, Any] | None = None + ) -> dict[str, Any]: + selected = selected_monitor or first_due_monitor or {} + selected_todo_id = normalize_todo_id(selected.get("todo_id")) + projected_due_items = [ + selected, + *[ + item + for item in due_monitor_items + if normalize_todo_id(item.get("todo_id")) != selected_todo_id + ], + ] return { "schema_version": WORK_LANE_CONTRACT_SCHEMA_VERSION, "lane": "continuous_monitor", @@ -571,7 +593,7 @@ def due_monitor_contract(*, reason_codes: list[str]) -> dict[str, Any]: "monitor_policy": "attempt_due_monitor_once_then_writeback_or_no_spend_if_unchanged", "monitor_due_count": max(0, int(monitor_due_count)), "monitor_due_items": _compact_work_lane_todo_items( - due_monitor_items, + projected_due_items, limit=monitor_due_item_limit, ), "selected_todo_id": selected.get("todo_id"), @@ -585,7 +607,8 @@ def due_monitor_contract(*, reason_codes: list[str]) -> dict[str, Any]: if progress_scope != "dependency_observation": if has_advancement_todos and effective_due_monitor_preemption: return due_monitor_contract( - reason_codes=["monitor_due", "due_monitor_priority_preempts_advancement"] + reason_codes=["monitor_due", "due_monitor_priority_preempts_advancement"], + selected_monitor=first_preemptive_due_monitor, ) if has_advancement_todos and first_advancement is not None: reason_codes = ["open_agent_todo"] @@ -624,6 +647,24 @@ def due_monitor_contract(*, reason_codes: list[str]) -> dict[str, Any]: "monitor_policy": "material_transition_only", "action": action, } + if watch_only_due_items: + selected_watch_only_monitor = watch_only_due_items[0] + contract["auxiliary_monitor_poll"] = { + "schema_version": "auxiliary_monitor_poll_v0", + "required": False, + "preempts_advancement": False, + "spend_policy": "no_spend", + "continuation": "advancement_remains_primary", + "monitor_due_count": len(watch_only_due_items), + "monitor_due_items": _compact_work_lane_todo_items( + watch_only_due_items, + limit=monitor_due_item_limit, + ), + "selected_todo_id": selected_watch_only_monitor.get("todo_id"), + "selected_next_due_at": selected_watch_only_monitor.get( + "next_due_at" + ), + } if outcome_followthrough: contract["outcome_followthrough"] = outcome_followthrough return contract diff --git a/loopx/control_plane/work_items/work_lane_context.py b/loopx/control_plane/work_items/work_lane_context.py index 248eb2c017..7d00e99414 100644 --- a/loopx/control_plane/work_items/work_lane_context.py +++ b/loopx/control_plane/work_items/work_lane_context.py @@ -10,9 +10,11 @@ todo_summary_first_executable_item, todo_summary_monitor_due_count, todo_summary_monitor_due_items, + todo_summary_non_watch_only_monitor_due_items, todo_summary_monitor_schedule_gap_count, todo_summary_monitor_schedule_gap_items, todo_summary_open_task_counts, + todo_summary_watch_only_monitor_due_items, ) from .delivery_history import project_delivery_response from .work_lane import ( @@ -113,6 +115,12 @@ def build_work_lane_context_contract( agent_todo_summary, due_items=due_monitor_items, ) + watch_only_due_monitor_items = todo_summary_watch_only_monitor_due_items( + agent_todo_summary + ) + non_watch_only_due_monitor_items = ( + todo_summary_non_watch_only_monitor_due_items(agent_todo_summary) + ) if not advancement_allowed and due_monitor_count <= 0: # In monitor-only mode, an action phrase such as "observe result" must # not bypass the monitor todo's explicit cadence window. @@ -122,7 +130,11 @@ def build_work_lane_context_contract( agent_todo_summary, gap_items=monitor_schedule_gap_items, ) - first_due_monitor = due_monitor_items[0] if due_monitor_items else None + first_preemptive_due_monitor = ( + non_watch_only_due_monitor_items[0] + if non_watch_only_due_monitor_items + else None + ) first_advancement = ( todo_summary_first_executable_item(agent_todo_summary) if advancement_allowed @@ -143,11 +155,13 @@ def build_work_lane_context_contract( todo_counts=todo_counts, monitor_due_count=due_monitor_count, due_monitor_items=due_monitor_items, + watch_only_due_monitor_items=watch_only_due_monitor_items, + non_watch_only_due_monitor_items=non_watch_only_due_monitor_items, monitor_schedule_gap_count=monitor_schedule_gap_count, monitor_schedule_gap_items=monitor_schedule_gap_items, first_advancement=first_advancement, due_monitor_preempts_advancement=due_monitor_preempts_advancement( - first_due_monitor, + first_preemptive_due_monitor, first_advancement=first_advancement, ), outcome_followthrough=( diff --git a/tests/capabilities/test_issue_fix_pr_gate_reconcile.py b/tests/capabilities/test_issue_fix_pr_gate_reconcile.py index 63b52766a6..62242d37b3 100644 --- a/tests/capabilities/test_issue_fix_pr_gate_reconcile.py +++ b/tests/capabilities/test_issue_fix_pr_gate_reconcile.py @@ -8,6 +8,16 @@ from loopx.capabilities.issue_fix.pr_gate_reconcile import ( reconcile_issue_fix_pr_gate, ) +from loopx.capabilities.issue_fix.pr_lifecycle import ( + build_issue_fix_pr_lifecycle_monitor_packet, +) +from loopx.capabilities.issue_fix.pr_lifecycle_rollout import ( + append_pr_merge_rollout_event, +) +from loopx.control_plane.todos.resume_condition import ( + evaluate_todo_resume_conditions, +) +from loopx.rollout_event_log import load_rollout_events, rollout_event_log_path GOAL_ID = "multi-agent-pr-gate-reconcile" @@ -96,3 +106,71 @@ def test_multi_agent_legacy_merge_gate_without_actor_is_atomic(tmp_path: Path) - _reconcile(registry, state, agent_id=None) assert state.read_text(encoding="utf-8") == before + + +def test_repository_redirect_emits_canonical_merge_with_exact_alias_resume( + tmp_path: Path, +) -> None: + registry, _state = _write_fixture(tmp_path) + runtime_root = tmp_path / "runtime" + packet = build_issue_fix_pr_lifecycle_monitor_packet( + url="https://github.com/huangruiteng/loopx/pull/4344", + provider_payload={ + "state": "MERGED", + "mergedAt": "2026-09-13T13:45:34Z", + "url": "https://github.com/loopx-project/loopx/pull/4344", + }, + generated_at="2026-09-20T10:30:00Z", + ) + + assert packet["observation"]["repo"] == "loopx-project/loopx" + assert packet["observation"]["requested_repo"] == "huangruiteng/loopx" + assert packet["observation"]["repository_aliases"] == [ + "huangruiteng/loopx" + ] + + receipt = append_pr_merge_rollout_event( + payload=packet, + goal_id=GOAL_ID, + registry_path=registry, + runtime_root_arg=str(runtime_root), + ) + assert receipt["pr_ref"] == "loopx-project/loopx#4344" + assert receipt["repository_aliases"] == ["huangruiteng/loopx"] + events = load_rollout_events(rollout_event_log_path(runtime_root, GOAL_ID)) + assert events[0]["code_refs"]["pr_ref"] == "loopx-project/loopx#4344" + assert events[0]["source_refs"] == [ + {"kind": "pull_request", "ref": "huangruiteng/loopx#4344"} + ] + + conditions = evaluate_todo_resume_conditions( + [ + { + "todo_id": "todo_legacy_repository_wait", + "status": "deferred", + "task_repository": "git:github.com/huangruiteng/loopx", + "resume_when": "pr_merged:#4344", + } + ], + source_items=[], + rollout_events=events, + evaluated_at="2026-09-20T10:31:00Z", + ) + assert conditions["todo_legacy_repository_wait"]["satisfied"] is True + assert conditions["todo_legacy_repository_wait"]["matched_pr_ref"] == ( + "huangruiteng/loopx#4344" + ) + + +def test_repository_redirect_alias_does_not_match_another_pr_number() -> None: + packet = build_issue_fix_pr_lifecycle_monitor_packet( + url="https://github.com/huangruiteng/loopx/pull/4344", + provider_payload={ + "state": "MERGED", + "mergedAt": "2026-09-13T13:45:34Z", + "url": "https://github.com/loopx-project/loopx/pull/4345", + }, + ) + + assert packet["observation"]["repo"] == "huangruiteng/loopx" + assert "repository_aliases" not in packet["observation"] diff --git a/tests/control_plane/test_work_lane_contract_core.py b/tests/control_plane/test_work_lane_contract_core.py index 284cc6b41f..ffe80764d4 100644 --- a/tests/control_plane/test_work_lane_contract_core.py +++ b/tests/control_plane/test_work_lane_contract_core.py @@ -92,6 +92,104 @@ def test_due_monitor_preempts_lower_priority_advancement() -> None: assert guard["recommended_action"] == "[P0] Monitor one overdue dependency." +def test_due_watch_only_monitor_is_an_auxiliary_no_spend_route() -> None: + items = _monitor_and_advancement() + items[0].update( + { + "todo_id": "todo_watch_due", + "watch_only": "true", + } + ) + items[1].update( + { + "todo_id": "todo_advancement", + } + ) + payload = _status( + agent_todo_items=items, + next_action="Advance the bounded product slice.", + ) + guard = build_quota_should_run(payload, goal_id=GOAL_ID) + + assert guard["recommended_action"] == "[P1] Advance the bounded product slice." + lane = guard["work_lane_contract"] + assert lane["lane"] == "advancement_task" + assert lane["auxiliary_monitor_poll"] == { + "schema_version": "auxiliary_monitor_poll_v0", + "required": False, + "preempts_advancement": False, + "spend_policy": "no_spend", + "continuation": "advancement_remains_primary", + "monitor_due_count": 1, + "monitor_due_items": [ + { + "index": 1, + "text": "[P0] Monitor one overdue dependency.", + "todo_id": "todo_watch_due", + "status": "open", + "priority": "P0", + "task_class": "continuous_monitor", + "action_kind": "monitor", + "next_due_at": PAST_DUE_AT, + "watch_only": "true", + } + ], + "selected_todo_id": "todo_watch_due", + "selected_next_due_at": PAST_DUE_AT, + } + summary = guard["agent_todo_summary"] + assert summary["watch_only_monitor_due_count"] == 1 + assert summary["watch_only_monitor_due_items"][0]["todo_id"] == ( + "todo_watch_due" + ) + interaction = guard["interaction_contract"] + assert interaction["agent_channel"]["auxiliary_monitor_poll"][ + "required" + ] is False + auxiliary_cli = interaction["cli_channel"]["auxiliary_monitor_poll"] + assert auxiliary_cli["spend_policy"] == "no_spend" + assert "--todo-id todo_watch_due --execute" in auxiliary_cli["command"] + + +def test_watch_only_priority_cannot_hide_an_ordinary_due_monitor() -> None: + items = _monitor_and_advancement() + items[0].update( + { + "todo_id": "todo_watch_due", + "watch_only": "true", + } + ) + items.insert( + 1, + { + "index": 3, + "todo_id": "todo_ordinary_due", + "text": "[P1] Poll the ordinary due dependency.", + "role": "agent", + "status": "open", + "priority": "P1", + "task_class": "continuous_monitor", + "action_kind": "monitor", + "next_due_at": PAST_DUE_AT, + }, + ) + items[-1].update( + { + "priority": "P2", + "text": "[P2] Advance the bounded product slice.", + } + ) + payload = _status(agent_todo_items=items) + + guard = build_quota_should_run(payload, goal_id=GOAL_ID) + lane = guard["work_lane_contract"] + + assert lane["lane"] == "continuous_monitor" + assert lane["selected_todo_id"] == "todo_ordinary_due" + assert lane["monitor_due_items"][0]["todo_id"] == "todo_ordinary_due" + assert lane.get("auxiliary_monitor_poll") is None + + def test_receipt_bound_advancement_retains_auxiliary_due_monitor_context() -> None: due_monitor = { "todo_id": "todo_due_monitor", diff --git a/tests/control_plane_ts/capability_gate.test.ts b/tests/control_plane_ts/capability_gate.test.ts index 91e128d8b8..42dcc01a92 100644 --- a/tests/control_plane_ts/capability_gate.test.ts +++ b/tests/control_plane_ts/capability_gate.test.ts @@ -73,7 +73,8 @@ test("production-scale requirements use the same rule for every retained record test("quota v1 recomputes missing capabilities, never trusts stale Python results", () => { const candidate = {payload: {todo_id: "monitor"}, claim: null, bound: null, blocks: null, excluded: [], global: false, gate: false, removed: false, actionable: true, due: true, task_class: "continuous_monitor", - priority: 1, index: 1, profile_rank: 1, missing: [], raw_claimed: false, required: ["network"], targets: []}; + priority: 1, index: 1, profile_rank: 1, missing: [], raw_claimed: false, + watch_only: false, required: ["network"], targets: []}; const request = {items: [candidate], active_items: [], active_executable_items: [], agent_id: "agent-a", user_gate_scope: false, monitor_supported: true, diagnostic_limit: 3, backlog_limit: 8, visibility_limit: 16, profile: null, source_open_count: 1, available: []}; diff --git a/tests/control_plane_ts/quota_selection.test.ts b/tests/control_plane_ts/quota_selection.test.ts index eafec322e9..f9aa7ff10c 100644 --- a/tests/control_plane_ts/quota_selection.test.ts +++ b/tests/control_plane_ts/quota_selection.test.ts @@ -7,6 +7,7 @@ import { productionScaleCoordinationFixture } from "./production_scale_coordinat function row(id: string, fields: JsonObject = {}): JsonObject { return {payload: {todo_id: id}, claim: null, bound: null, blocks: null, excluded: [], global: false, gate: false, removed: false, actionable: true, due: false, + watch_only: false, task_class: "advancement_task", priority: 1, index: 1, profile_rank: 1, missing: [], raw_claimed: false, ...fields}; } @@ -50,10 +51,14 @@ test("execution scope is shared with active-next-action, including removed-polic test("monitor eligibility preserves provider writeback and capability fences", () => { const items = [row("due", {task_class: "continuous_monitor", due: true}), + row("watch", {task_class: "continuous_monitor", due: true, watch_only: true}), row("missing", {task_class: "continuous_monitor", due: true, missing: ["network"]}), row("future", {task_class: "continuous_monitor"})]; const lanes = projectQuotaSelection(request(items)).lanes as JsonObject; - assert.deepEqual(ids(lanes.monitor_due_items), ["due"]); + assert.deepEqual(ids(lanes.monitor_due_items), ["due", "watch"]); + assert.deepEqual(ids(lanes.watch_only_monitor_items), ["watch"]); + assert.deepEqual(ids(lanes.watch_only_monitor_due_items), ["watch"]); + assert.deepEqual(ids(lanes.non_watch_only_monitor_due_items), ["due"]); assert.deepEqual(ids(lanes.monitor_capability_blocked_due_items), ["missing"]); assert.deepEqual(ids(lanes.executable_items), []); const unsupported = projectQuotaSelection(request(items, {monitor_supported: false})).lanes as JsonObject; diff --git a/tests/control_plane_ts/todo_resume_condition.test.ts b/tests/control_plane_ts/todo_resume_condition.test.ts index efbfa32c1b..537356f894 100644 --- a/tests/control_plane_ts/todo_resume_condition.test.ts +++ b/tests/control_plane_ts/todo_resume_condition.test.ts @@ -203,6 +203,48 @@ test("one reducer evaluates Todo, PR, capacity, monitor, and date resume conditi assert.equal(conditions.todo_wait_date.material_change_generation, 1); }); +test("PR merge source aliases bridge a repository redirect without weakening repository binding", () => { + const request = { + schema_version: TODO_RESUME_EVALUATION_REQUEST_SCHEMA_VERSION, + items: [ + todo("todo_old_repo", "deferred", "advancement_task", { + resume_when: "pr_merged:#4344", + task_repository: "git:github.com/huangruiteng/loopx", + }), + todo("todo_unrelated_repo", "deferred", "advancement_task", { + resume_when: "pr_merged:#4344", + task_repository: "git:github.com/example/loopx", + }), + ], + source_items: [], + rollout_events: [{ + event_id: "event-redirected-merge-4344", + event_kind: "pr_merge", + pr_ref: "loopx-project/loopx#4344", + source_refs: [ + {kind: "pull_request", ref: "huangruiteng/loopx#4344"}, + {kind: "pull_request", ref: "example/loopx#9999"}, + ], + recorded_at: "2026-09-13T13:45:34Z", + }], + available_capabilities: [], + evaluated_at: "2026-09-20T00:00:00Z", + }; + const result = evaluateTodoResumeConditions(request); + const conditions = Object.fromEntries( + (result.conditions as Array>).map((row) => [ + row.todo_id, + row.condition, + ]), + ) as Record>; + assert.equal(conditions.todo_old_repo.satisfied, true); + assert.equal( + conditions.todo_old_repo.matched_pr_ref, + "huangruiteng/loopx#4344", + ); + assert.equal(conditions.todo_unrelated_repo.satisfied, false); +}); + test("monitor resume is generation-fenced and fail-closed without a baseline", () => { const result = evaluateTodoResumeConditions({ schema_version: TODO_RESUME_EVALUATION_REQUEST_SCHEMA_VERSION, diff --git a/tests/control_plane_ts/todo_summary_lanes.test.ts b/tests/control_plane_ts/todo_summary_lanes.test.ts index 29feaf084f..5e0d196a72 100644 --- a/tests/control_plane_ts/todo_summary_lanes.test.ts +++ b/tests/control_plane_ts/todo_summary_lanes.test.ts @@ -34,10 +34,27 @@ test("one observation time fences due, expiry, missing schedule and watch-only m monitor({due_at: 10, expires_at: 100}), monitor({}), monitor({watch_only: true}), monitor({due_at: 90, acceptance_blocked: true})]).lanes as JsonObject; assert.deepEqual(lanes.monitor_due_items, [0]); + assert.deepEqual(lanes.watch_only_monitor_items, [4]); + assert.deepEqual(lanes.watch_only_monitor_due_items, []); + assert.deepEqual(lanes.non_watch_only_monitor_due_items, [0]); + assert.deepEqual(lanes.convergent_open_items, [0, 1, 2, 3, 5]); assert.deepEqual(lanes.monitor_schedule_gap_items, [3]); assert.deepEqual(lanes.monitor_items, [0, 1, 2, 3, 4]); }); +test("watch-only due monitors remain schedulable but are partitioned from ordinary due work", () => { + const monitor = (fields: JsonObject) => row({task_class: "continuous_monitor", ...fields}); + const lanes = project([ + monitor({due_at: 90, watch_only: true}), + monitor({due_at: 80}), + row(), + ]).lanes as JsonObject; + assert.deepEqual(lanes.monitor_due_items, [0, 1]); + assert.deepEqual(lanes.watch_only_monitor_due_items, [0]); + assert.deepEqual(lanes.non_watch_only_monitor_due_items, [1]); + assert.deepEqual(lanes.convergent_open_items, [1, 2]); +}); + test("display ordering preserves stable legacy ties and Python Unicode ordering", () => { const lanes = project([row({sort: [1, 2, "", ""]}), row({sort: [0, 9, "", ""]}), row({sort: [1, 2, "", ""]}), row({sort: [1, 999999, "", "\u{10000}"]}), From 6942ad6671928303c912b28801e08bf68edb50ed Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Mon, 21 Sep 2026 01:54:49 +0800 Subject: [PATCH 2/2] fix(control-plane): bind auxiliary monitor observations Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- docs/project-agent-todo-contract.md | 8 +++ .../work_items/interaction_contract.py | 56 +++++++++++++++++-- .../test_monitor_observation_admission.py | 36 +++++------- .../test_quota_settlement_cli.py | 3 + .../test_work_lane_contract_core.py | 54 +++++++++++++++++- 5 files changed, 128 insertions(+), 29 deletions(-) diff --git a/docs/project-agent-todo-contract.md b/docs/project-agent-todo-contract.md index f41af665e2..d19d05ff37 100644 --- a/docs/project-agent-todo-contract.md +++ b/docs/project-agent-todo-contract.md @@ -154,6 +154,11 @@ A scheduled watch-only monitor remains eligible at `next_due_at`, but it never creates autonomous replan pressure and never preempts runnable advancement. When both are present, `interaction_contract` keeps advancement primary and projects an optional, typed, no-spend `auxiliary_monitor_poll` route. +That CLI route is available only when it is bound to the current Turn. It +requires the caller to place the fresh observation digest in +`LOOPX_MONITOR_RESULT_HASH` and exposes separate unchanged and material-change +commands; omitting either the Turn binding or result digest fails closed before +monitor writeback. The canonical watch-only/ordinary-due partition is produced inside the existing TypeScript Todo summary and quota-planning owners after Agent scope and capability admission; Python compatibility code only adapts legacy facts and @@ -164,6 +169,9 @@ renders the selected CLI/Lark route. replan 压力,也不会抢占 runnable advancement;二者同时存在时, `interaction_contract` 保持 advancement 为主,并投影一条可选、typed、no-spend 的 `auxiliary_monitor_poll` 路由。 +该 CLI 路由仅在绑定当前 Turn 时可用;调用方必须把本次新鲜 observation digest +写入 `LOOPX_MONITOR_RESULT_HASH`,并在 unchanged 与 material-change 两条命令中 +明确选择。缺少 Turn 绑定或 result digest 时,monitor writeback 会在写入前失败关闭。 watch-only/普通 due 的权威分区由既有 TypeScript Todo summary 与 quota-planning owner 在 Agent scope 和 capability admission 之后生成;Python 兼容层只适配旧事实并 渲染已选中的 CLI/Lark 路由。 diff --git a/loopx/control_plane/work_items/interaction_contract.py b/loopx/control_plane/work_items/interaction_contract.py index 9b7e036e1b..06a13e1028 100644 --- a/loopx/control_plane/work_items/interaction_contract.py +++ b/loopx/control_plane/work_items/interaction_contract.py @@ -57,6 +57,11 @@ INTERACTION_RESPONSE_PLAN_SCHEMA_VERSION = "interaction_response_plan_v0" PROTOCOL_ACTION_PACKET_SCHEMA_VERSION = "protocol_action_packet_v0" PROTOCOL_ACTION_PACKET_LLM_POLICY = "no_api" +AUXILIARY_MONITOR_POLL_CLI_SCHEMA_VERSION = "auxiliary_monitor_poll_cli_v0" +AUXILIARY_MONITOR_OBSERVATION_INPUT_SCHEMA_VERSION = ( + "auxiliary_monitor_observation_input_v0" +) +AUXILIARY_MONITOR_RESULT_HASH_ENV = "LOOPX_MONITOR_RESULT_HASH" class _InteractionContractRequired(typing.TypedDict): @@ -1346,15 +1351,54 @@ def _build_interaction_cli_channel( if isinstance(payload.get("agent_identity"), dict) else {} ) - channel["auxiliary_monitor_poll"] = { + safe_turn_instance_id = str(turn_instance_id or "").strip() + auxiliary_projection: dict[str, Any] = { **dict(auxiliary_monitor), - "command": ( + "schema_version": AUXILIARY_MONITOR_POLL_CLI_SCHEMA_VERSION, + "input_contract": { + "schema_version": ( + AUXILIARY_MONITOR_OBSERVATION_INPUT_SCHEMA_VERSION + ), + "result_hash": { + "required": True, + "environment_variable": AUXILIARY_MONITOR_RESULT_HASH_ENV, + "source": "fresh_external_observation_digest", + }, + "material_change": { + "required": True, + "unchanged_command_key": "command", + "changed_command_key": "material_change_command", + }, + }, + } + if not safe_turn_instance_id: + auxiliary_projection.update( + { + "availability": "turn_binding_required", + "reason_code": "auxiliary_monitor_turn_instance_id_missing", + } + ) + else: + command = ( f"{command_prefix} quota monitor-poll --goal-id " - f"{str(payload.get('goal_id') or '')}" + f"{shlex.quote(str(payload.get('goal_id') or ''))}" f"{_scoped_cli_args(agent_identity, available_capabilities=available_capabilities)}" - f"{auxiliary_scheduler_args} --todo-id {selected_monitor_id} --execute" - ), - } + f"{auxiliary_scheduler_args} --turn-instance-id " + f"{shlex.quote(safe_turn_instance_id)} --todo-id " + f"{shlex.quote(selected_monitor_id)} --result-hash " + f'"${{{AUXILIARY_MONITOR_RESULT_HASH_ENV}:?}}"' + ) + auxiliary_projection.update( + { + "availability": "ready", + "turn_instance_id": safe_turn_instance_id, + "command": f"{command} --execute", + "material_change_command": ( + f"{command} --material-change --execute" + ), + } + ) + channel["auxiliary_monitor_poll"] = auxiliary_projection selection.apply_action_selection_cli_gate(channel, payload) if settlement_plan is not None and spend_after_selection: channel["settlement_plan"] = settlement_plan diff --git a/tests/control_plane/test_monitor_observation_admission.py b/tests/control_plane/test_monitor_observation_admission.py index 8549e29c06..e5aa1225ba 100644 --- a/tests/control_plane/test_monitor_observation_admission.py +++ b/tests/control_plane/test_monitor_observation_admission.py @@ -9,6 +9,7 @@ TODO_ID, _append_newly_due_monitor, _classification_count, + _projected_cli_args, _run_cli, _spend_run_count, _write_fixture, @@ -37,7 +38,7 @@ def test_receipt_bound_advancement_allows_one_auxiliary_due_monitor_receipt( assert first_rc == 0, first assert first["heartbeat_receipt"]["settlement_identity"]["todo_id"] == TODO_ID - _append_newly_due_monitor(project) + _append_newly_due_monitor(project, watch_only=True) capability_args = ( "--available-capability", "network", @@ -54,27 +55,20 @@ def test_receipt_bound_advancement_allows_one_auxiliary_due_monitor_receipt( assert replay["selected_todo"]["todo_id"] == TODO_ID assert "due_monitor_context" in replay["work_lane_contract"]["reason_codes"] - poll_args = ( - "quota", - "monitor-poll", - "--codex-app", - "--goal-id", - GOAL_ID, - "--agent-id", - AGENT_ID, - "--turn-instance-id", - turn_instance_id, - "--todo-id", - DUE_MONITOR_TODO_ID, - "--target-key", - "due-monitor-fixture", - "--result-hash", - "unchanged-auxiliary-monitor", - *capability_args, - "--execute", - "--scan-path", - str(project), + auxiliary_cli = replay["interaction_contract"]["cli_channel"][ + "auxiliary_monitor_poll" + ] + assert auxiliary_cli["turn_instance_id"] == turn_instance_id + poll_args = _projected_cli_args( + auxiliary_cli["command"], + turn_instance_id=turn_instance_id, ) + poll_args = tuple( + "unchanged-auxiliary-monitor" + if token == "${LOOPX_MONITOR_RESULT_HASH:?}" + else token + for token in poll_args + ) + ("--scan-path", str(project)) poll_rc, poll = _run_cli(registry_path, runtime, *poll_args) poll_replay_rc, poll_replay = _run_cli( registry_path, diff --git a/tests/control_plane/test_quota_settlement_cli.py b/tests/control_plane/test_quota_settlement_cli.py index c8f9ba7107..b3170f79ed 100644 --- a/tests/control_plane/test_quota_settlement_cli.py +++ b/tests/control_plane/test_quota_settlement_cli.py @@ -330,9 +330,11 @@ def _append_newly_due_monitor( project: Path, *, priority: str = "P0-monitor", + watch_only: bool = False, ) -> None: state_path = project / f".codex/goals/{GOAL_ID}/ACTIVE_GOAL_STATE.md" state_text = state_path.read_text(encoding="utf-8") + watch_only_field = "watch_only=true " if watch_only else "" state_path.write_text( state_text.replace( "## Agent Todo\n\n", @@ -342,6 +344,7 @@ def _append_newly_due_monitor( "task_class=continuous_monitor action_kind=observe " f"claimed_by={AGENT_ID} target_key=due-monitor-fixture " "required_capabilities=network%2Cexternal_evidence_poll " + f"{watch_only_field}" "cadence=1m next_due_at=2000-01-01T00%3A00%3A00Z -->\n", ), encoding="utf-8", diff --git a/tests/control_plane/test_work_lane_contract_core.py b/tests/control_plane/test_work_lane_contract_core.py index ffe80764d4..556c738663 100644 --- a/tests/control_plane/test_work_lane_contract_core.py +++ b/tests/control_plane/test_work_lane_contract_core.py @@ -109,7 +109,12 @@ def test_due_watch_only_monitor_is_an_auxiliary_no_spend_route() -> None: agent_todo_items=items, next_action="Advance the bounded product slice.", ) - guard = build_quota_should_run(payload, goal_id=GOAL_ID) + turn_instance_id = "turn-watch-only-auxiliary" + guard = build_quota_should_run( + payload, + goal_id=GOAL_ID, + turn_instance_id=turn_instance_id, + ) assert guard["recommended_action"] == "[P1] Advance the bounded product slice." lane = guard["work_lane_contract"] @@ -147,8 +152,53 @@ def test_due_watch_only_monitor_is_an_auxiliary_no_spend_route() -> None: "required" ] is False auxiliary_cli = interaction["cli_channel"]["auxiliary_monitor_poll"] + assert auxiliary_cli["schema_version"] == "auxiliary_monitor_poll_cli_v0" + assert auxiliary_cli["availability"] == "ready" + assert auxiliary_cli["turn_instance_id"] == turn_instance_id assert auxiliary_cli["spend_policy"] == "no_spend" - assert "--todo-id todo_watch_due --execute" in auxiliary_cli["command"] + assert f"--turn-instance-id {turn_instance_id}" in auxiliary_cli["command"] + assert "--todo-id todo_watch_due" in auxiliary_cli["command"] + assert '--result-hash "${LOOPX_MONITOR_RESULT_HASH:?}"' in auxiliary_cli[ + "command" + ] + assert auxiliary_cli["command"].endswith("--execute") + assert "--material-change --execute" in auxiliary_cli[ + "material_change_command" + ] + assert auxiliary_cli["input_contract"] == { + "schema_version": "auxiliary_monitor_observation_input_v0", + "result_hash": { + "required": True, + "environment_variable": "LOOPX_MONITOR_RESULT_HASH", + "source": "fresh_external_observation_digest", + }, + "material_change": { + "required": True, + "unchanged_command_key": "command", + "changed_command_key": "material_change_command", + }, + } + + +def test_due_watch_only_monitor_without_turn_binding_is_not_executable() -> None: + items = _monitor_and_advancement() + items[0].update({"todo_id": "todo_watch_due", "watch_only": "true"}) + items[1]["todo_id"] = "todo_advancement" + + guard = build_quota_should_run( + _status(agent_todo_items=items), + goal_id=GOAL_ID, + ) + + auxiliary_cli = guard["interaction_contract"]["cli_channel"][ + "auxiliary_monitor_poll" + ] + assert auxiliary_cli["availability"] == "turn_binding_required" + assert auxiliary_cli["reason_code"] == ( + "auxiliary_monitor_turn_instance_id_missing" + ) + assert "command" not in auxiliary_cli + assert "material_change_command" not in auxiliary_cli def test_watch_only_priority_cannot_hide_an_ordinary_due_monitor() -> None: