From 8e0476d5b3146eaec179b2c7aea934e11e8e8d66 Mon Sep 17 00:00:00 2001 From: "lusendong.6789" Date: Mon, 7 Sep 2026 22:31:55 +0800 Subject: [PATCH 1/7] fix(quota): resolve fallback from canonical todos Signed-off-by: lusendong.6789 Co-authored-by: TRAE CLI --- .../goals/goal_frontier/__init__.py | 15 ++ .../goal_frontier/fallback_disposition.py | 177 +++++++++++- loopx/control_plane/quota/live_decision.py | 100 +++++++ loopx/control_plane/quota/should_run.py | 2 + .../control_plane/quota/should_run_prepare.py | 6 + loopx/quota.py | 2 + .../test_fallback_wait_admission.py | 251 ++++++++++++++++++ ...test_goal_frontier_fallback_disposition.py | 128 ++++++++- 8 files changed, 665 insertions(+), 16 deletions(-) create mode 100644 tests/control_plane/test_fallback_wait_admission.py diff --git a/loopx/control_plane/goals/goal_frontier/__init__.py b/loopx/control_plane/goals/goal_frontier/__init__.py index 44080125be..80d06d26c0 100644 --- a/loopx/control_plane/goals/goal_frontier/__init__.py +++ b/loopx/control_plane/goals/goal_frontier/__init__.py @@ -1482,6 +1482,9 @@ def build_goal_frontier_projection_context_from_status( work_lane_contract: dict[str, Any] | None, neutral_replan_ack_classifications: set[str], agent_todo_source_items: list[dict[str, Any]] | None = None, + fallback_todo_source_items: list[dict[str, Any]] | None = None, + fallback_todo_source_authoritative: bool | None = None, + available_capabilities: Any = None, registered_agent_ids: list[str] | None = None, goal_status: str | None = None, agent_profile: dict[str, Any] | None = None, @@ -1612,6 +1615,18 @@ def build_goal_frontier_projection_context_from_status( latest_agent_vision, agent_todo_summary=agent_todo_summary, agent_id=agent_id, + agent_todo_source_items=( + fallback_todo_source_items + if fallback_todo_source_authoritative is True + else None + if fallback_todo_source_authoritative is False + else agent_todo_source_items + ), + rollout_events=latest_runs_for_goal( + status_payload, + goal_id=goal_id, + ), + available_capabilities=available_capabilities, ), ) if isinstance(gap, dict) diff --git a/loopx/control_plane/goals/goal_frontier/fallback_disposition.py b/loopx/control_plane/goals/goal_frontier/fallback_disposition.py index c27cf0f28c..3c43a11653 100644 --- a/loopx/control_plane/goals/goal_frontier/fallback_disposition.py +++ b/loopx/control_plane/goals/goal_frontier/fallback_disposition.py @@ -4,12 +4,21 @@ from typing import Any from ...todos.contract import ( + TODO_STATUS_DEFERRED, + TODO_STATUS_OPEN, + TODO_TASK_CLASS_ADVANCEMENT, normalize_todo_id, + normalize_todo_resume_when, + normalize_todo_status, ) from ...todos.deferred_resume import todo_summary_blocked_successor_items from ...todos.projection import ( agent_scoped_selectable_advancement_todo_ids, + todo_item_claimed_by_agent_or_unclaimed, + todo_item_is_actionable_open, + todo_item_task_class, ) +from ...todos.resume_condition import evaluate_todo_resume_conditions from ..goal_vision_state import goal_vision_state_is_closed # Single owner of the vision todo_delta action contract shared by the @@ -29,6 +38,10 @@ VISION_FALLBACK_DECLARATION_FIELDS = ("target_todo_id", "successor_todo_id") VISION_FALLBACK_GAP_TRIGGER = "vision_fallback_unresolved" VISION_FALLBACK_GAP_REASON_CODE = "declared_fallback_without_runnable_or_terminal" +VISION_FALLBACK_LOOKUP_UNCERTAIN_TRIGGER = "vision_fallback_lookup_uncertain" +VISION_FALLBACK_LOOKUP_UNCERTAIN_REASON_CODE = ( + "declared_fallback_authoritative_lookup_unavailable" +) VISION_FALLBACK_TERMINAL_PATH_OUTCOME = "stop" VISION_FALLBACK_RUNNABLE_ITEM_LIMIT = 3 VISION_FALLBACK_RECOMMENDED_ACTION = ( @@ -37,6 +50,10 @@ "successor, or record an explicit terminal no-follow-up disposition; " "do not invent a user gate" ) +VISION_FALLBACK_LOOKUP_UNCERTAIN_ACTION = ( + "retry the fallback disposition from the complete canonical Todo source; " + "do not infer absence from a bounded presentation lane or invent a user gate" +) @dataclass(frozen=True) @@ -54,14 +71,13 @@ def candidate_todo_ids(self) -> set[str]: for todo_id in ( self.target_todo_id, self.successor_todo_id, - self.declaration_id, ) if todo_id } @property def unresolved_todo_id(self) -> str: - return self.target_todo_id or self.declaration_id + return self.target_todo_id or self.successor_todo_id or self.declaration_id def _compact_text(value: Any, *, limit: int) -> str: @@ -188,11 +204,99 @@ def _vision_has_terminal_disposition(agent_vision: dict[str, Any]) -> bool: ) +def _authoritative_fallback_disposition_ids( + declarations: list[FallbackDeclaration], + *, + agent_todo_source_items: list[dict[str, Any]], + agent_id: str | None, + rollout_events: list[dict[str, Any]] | None, + available_capabilities: Any, +) -> tuple[set[str], set[str], set[str]]: + """Return exact declared Todo ids that are runnable or validly waiting. + + The planning source is complete canonical state, not a presentation lane. + Existing Todo predicates retain ownership, exclusion, task-class, and + lifecycle semantics; the TS resume evaluator remains the authority for a + linked external wait. + """ + + declared_ids = { + todo_id + for declaration in declarations + for todo_id in declaration.candidate_todo_ids + } + matched_items = [ + item + for item in agent_todo_source_items + if isinstance(item, dict) + and normalize_todo_id(item.get("todo_id")) in declared_ids + ] + resume_items = [ + item + for item in matched_items + if normalize_todo_resume_when(item.get("resume_when")) + ] + resume_conditions = ( + evaluate_todo_resume_conditions( + resume_items, + source_items=agent_todo_source_items, + rollout_events=rollout_events, + available_capabilities=available_capabilities, + ) + if resume_items + else {} + ) + + runnable_ids: set[str] = set() + waiting_ids: set[str] = set() + uncertain_ids: set[str] = set() + for item in matched_items: + todo_id = normalize_todo_id(item.get("todo_id")) + if not todo_id: + continue + if todo_item_task_class(item) != TODO_TASK_CLASS_ADVANCEMENT: + continue + if not todo_item_claimed_by_agent_or_unclaimed(item, agent_id=agent_id): + continue + + resume_when = normalize_todo_resume_when(item.get("resume_when")) + if resume_when: + status = normalize_todo_status(item.get("status")) or TODO_STATUS_OPEN + if status not in {TODO_STATUS_OPEN, TODO_STATUS_DEFERRED}: + continue + condition = resume_conditions.get(todo_id) + if not isinstance(condition, dict): + uncertain_ids.add(todo_id) + continue + if condition.get("invalid_target") is True or condition.get( + "invalid_state" + ): + continue + if condition.get("provider_required") is True: + # Capability evidence was unavailable, so this lookup is + # uncertain rather than proof that the fallback is absent. + uncertain_ids.add(todo_id) + continue + if condition.get("satisfied") is True: + runnable_ids.add(todo_id) + elif condition.get("satisfied") is False: + waiting_ids.add(todo_id) + continue + + if todo_item_is_actionable_open(item): + runnable_ids.add(todo_id) + + return runnable_ids, waiting_ids, uncertain_ids + + def declared_fallback_gap_from_agent_vision( agent_vision: dict[str, Any] | None, *, agent_todo_summary: dict[str, Any] | None, agent_id: str | None, + agent_todo_source_items: list[dict[str, Any]] | None = None, + rollout_events: list[dict[str, Any]] | None = None, + available_capabilities: Any = None, ) -> dict[str, Any] | None: """Project one advisory gap for an unresolved declared fallback. @@ -231,14 +335,32 @@ def declared_fallback_gap_from_agent_vision( if not declarations: return None - selectable_ids = agent_scoped_selectable_advancement_todo_ids( - agent_todo_summary, - agent_id=agent_id, - ) - waiting_todo_ids = _blocked_successor_todo_ids( - agent_todo_summary, - agent_id=agent_id, - ) + source_is_authoritative = agent_todo_source_items is not None + uncertain_todo_ids: set[str] = set() + if source_is_authoritative: + ( + selectable_ids, + waiting_todo_ids, + uncertain_todo_ids, + ) = _authoritative_fallback_disposition_ids( + declarations, + agent_todo_source_items=agent_todo_source_items or [], + agent_id=agent_id, + rollout_events=rollout_events, + available_capabilities=available_capabilities, + ) + else: + # Legacy direct callers may only have a compact display summary. It + # remains valid positive evidence, but omission from a bounded lane is + # uncertainty and must not be projected as authoritative absence. + selectable_ids = agent_scoped_selectable_advancement_todo_ids( + agent_todo_summary, + agent_id=agent_id, + ) + waiting_todo_ids = _blocked_successor_todo_ids( + agent_todo_summary, + agent_id=agent_id, + ) todo_delta = parse_vision_todo_delta_entries(agent_vision.get("todo_delta")) created_or_reopened_ids = { todo_id @@ -247,9 +369,12 @@ def declared_fallback_gap_from_agent_vision( } unresolved_ids: set[str] = set() + lookup_uncertain_ids: set[str] = set(uncertain_todo_ids) for declaration in declarations: candidate_ids = declaration.candidate_todo_ids - waiting_todo_ids if not candidate_ids: + if not declaration.candidate_todo_ids: + unresolved_ids.add(declaration.unresolved_todo_id) continue # Disposition 1: Runnable on authoritative selectable advancement frontier if candidate_ids & selectable_ids: @@ -257,32 +382,56 @@ def declared_fallback_gap_from_agent_vision( # Disposition 2: Bounded successor created/reopened specifically for this fallback if candidate_ids & created_or_reopened_ids: continue + if candidate_ids & lookup_uncertain_ids: + continue if ( declaration.successor_todo_id and declaration.successor_todo_id in created_or_reopened_ids ): continue + if not source_is_authoritative: + lookup_uncertain_ids.update(candidate_ids) + continue unresolved_id = declaration.unresolved_todo_id if unresolved_id not in waiting_todo_ids: unresolved_ids.add(unresolved_id) - if not unresolved_ids: + if not unresolved_ids and not lookup_uncertain_ids: return None + lookup_uncertain_only = not unresolved_ids + gap: dict[str, Any] = { - "kind": VISION_FALLBACK_GAP_TRIGGER, + "kind": ( + VISION_FALLBACK_LOOKUP_UNCERTAIN_TRIGGER + if lookup_uncertain_only + else VISION_FALLBACK_GAP_TRIGGER + ), "source": "latest_agent_vision", "agent_id": agent_vision.get("agent_id"), "state": agent_vision.get("state"), - "reason_code": VISION_FALLBACK_GAP_REASON_CODE, - "recommended_action": VISION_FALLBACK_RECOMMENDED_ACTION, + "reason_code": ( + VISION_FALLBACK_LOOKUP_UNCERTAIN_REASON_CODE + if lookup_uncertain_only + else VISION_FALLBACK_GAP_REASON_CODE + ), + "recommended_action": ( + VISION_FALLBACK_LOOKUP_UNCERTAIN_ACTION + if lookup_uncertain_only + else VISION_FALLBACK_RECOMMENDED_ACTION + ), } unresolved_todo_ids = [todo_id for todo_id in sorted(unresolved_ids) if todo_id][ :VISION_FALLBACK_RUNNABLE_ITEM_LIMIT ] if unresolved_todo_ids: gap["unresolved_todo_ids"] = unresolved_todo_ids + uncertain_todo_ids = [ + todo_id for todo_id in sorted(lookup_uncertain_ids) if todo_id + ][:VISION_FALLBACK_RUNNABLE_ITEM_LIMIT] + if uncertain_todo_ids: + gap["lookup_uncertain_todo_ids"] = uncertain_todo_ids generated_at = _compact_text(agent_vision.get("generated_at"), limit=80) if generated_at: gap["generated_at"] = generated_at diff --git a/loopx/control_plane/quota/live_decision.py b/loopx/control_plane/quota/live_decision.py index 01be74d04d..e7b1ee2cdf 100644 --- a/loopx/control_plane/quota/live_decision.py +++ b/loopx/control_plane/quota/live_decision.py @@ -6,6 +6,16 @@ from typing import Any from ...quota import build_quota_should_run +from ...todos import list_goal_todos +from ..coordination.local_authority import LocalCoordinationAuthorityUnavailable +from ..effect_runtime import EffectRuntimeRemoteError +from ..goals.goal_frontier.fallback_disposition import ( + parse_fallback_declarations, +) +from ..goals.goal_frontier.semantic_history import ( + latest_agent_vision_from_status_payload, +) +from ..todos.contract import normalize_todo_id, normalize_todo_resume_when from ..capability_hooks import ( InteractionProjectionHookRegistration, dispatch_interaction_projection_hooks, @@ -27,6 +37,88 @@ BoundedResearchFrontierProjector = Callable[..., Mapping[str, Any] | None] +def _fallback_authority_todo_ids( + status_payload: dict[str, Any], + *, + goal_id: str, + agent_id: str | None, +) -> set[str]: + vision = latest_agent_vision_from_status_payload( + status_payload, + goal_id=goal_id, + agent_id=agent_id, + ) + return { + todo_id + for declaration in parse_fallback_declarations(vision) + for todo_id in declaration.candidate_todo_ids + } + + +def _live_fallback_authority_items( + status_payload: dict[str, Any], + *, + registry_path: Path, + runtime_root: Path, + goal_id: str, + agent_id: str | None, +) -> list[dict[str, Any]] | None: + """Read only the exact canonical Todos needed by fallback disposition. + + ``status`` is deliberately presentation-bounded, so omission from it can + never prove that a declared fallback is absent. The live CLI path owns the + registry/runtime authority needed for exact reads. Resume dependencies are + included only when referenced by one of the declared Todos so the existing + TypeScript evaluator retains state-transition authority. + """ + + requested_ids = _fallback_authority_todo_ids( + status_payload, + goal_id=goal_id, + agent_id=agent_id, + ) + if not requested_ids: + return None + + items: dict[str, dict[str, Any]] = {} + pending_ids = set(requested_ids) + while pending_ids: + todo_id = pending_ids.pop() + try: + projection = list_goal_todos( + registry_path=registry_path, + goal_id=goal_id, + todo_id=todo_id, + runtime_root_arg=str(runtime_root), + limit=None, + ) + except ( + EffectRuntimeRemoteError, + LocalCoordinationAuthorityUnavailable, + OSError, + ValueError, + ): + return None + item = projection.get("todo") + if not isinstance(item, dict): + continue + normalized_id = normalize_todo_id(item.get("todo_id")) + if not normalized_id: + return None + items[normalized_id] = dict(item) + resume_when = normalize_todo_resume_when(item.get("resume_when")) + if resume_when: + resume_kind, _, target = resume_when.partition(":") + dependency_id = ( + normalize_todo_id(target) + if resume_kind in {"todo_done", "monitor_changed"} + else None + ) + if dependency_id and dependency_id not in items: + pending_ids.add(dependency_id) + return list(items.values()) + + def _fresh_read_covers_all_pending_material( dispatch: Mapping[str, Any] | None, projector: Callable[..., dict[str, Any]] | None, @@ -440,6 +532,13 @@ def build_live_quota_should_run_decision( fresh_operator_inbox_read = _fresh_operator_inbox_read_required( turn_start_hook_dispatch ) + authoritative_fallback_todo_items = _live_fallback_authority_items( + decision_status_payload, + registry_path=registry_path, + runtime_root=runtime_root, + goal_id=goal_id, + agent_id=agent_id, + ) payload = build_quota_should_run( decision_status_payload, goal_id=goal_id, @@ -465,6 +564,7 @@ def build_live_quota_should_run_decision( receipt_bound_replan_obligation_id=receipt_bound_replan_obligation_id, turn_instance_id=turn_instance_id, runtime_root=runtime_root, + authoritative_fallback_todo_items=authoritative_fallback_todo_items, ) if route_source.startswith("loopx_turn_"): payload["runtime_root"] = str(runtime_root) diff --git a/loopx/control_plane/quota/should_run.py b/loopx/control_plane/quota/should_run.py index a8f4a265ab..4fa1e728fc 100644 --- a/loopx/control_plane/quota/should_run.py +++ b/loopx/control_plane/quota/should_run.py @@ -258,6 +258,7 @@ def build_quota_should_run( receipt_bound_replan_obligation_id: str | None = None, turn_instance_id: str | None = None, runtime_root: str | Path | None = None, + authoritative_fallback_todo_items: list[dict[str, Any]] | None = None, ) -> dict[str, Any]: safe_goal_id = str(goal_id or "").strip() resolved_scheduler_context = resolve_scheduler_execution_context( @@ -326,6 +327,7 @@ def build_quota_should_run( receipt_bound_monitor_phase=receipt_bound_monitor_phase, receipt_bound_replay_phase=receipt_bound_replay_phase, receipt_bound_replan_obligation_id=receipt_bound_replan_obligation_id, + authoritative_fallback_todo_items=authoritative_fallback_todo_items, ) route = _resolve_quota_route_with_settled_replay_precedence(prepared) route = _apply_selected_todo_guards(prepared, route) diff --git a/loopx/control_plane/quota/should_run_prepare.py b/loopx/control_plane/quota/should_run_prepare.py index f075066714..3d42c07867 100644 --- a/loopx/control_plane/quota/should_run_prepare.py +++ b/loopx/control_plane/quota/should_run_prepare.py @@ -454,6 +454,7 @@ def _prepare_quota_should_run_item( receipt_bound_monitor_phase: ReceiptBoundMonitorPhase | None, receipt_bound_replay_phase: ReceiptBoundReplayPhase | None, receipt_bound_replan_obligation_id: str | None, + authoritative_fallback_todo_items: list[dict[str, Any]] | None = None, ) -> _QuotaDecisionPreparation: quota = item.get("quota") if isinstance(item.get("quota"), dict) else {} state = str(quota.get("state") or "unknown") @@ -749,6 +750,11 @@ def _prepare_quota_should_run_item( user_todo_summary=user_todo_summary, agent_todo_summary=agent_todo_summary, agent_todo_source_items=agent_todo_source_items, + fallback_todo_source_items=authoritative_fallback_todo_items, + fallback_todo_source_authoritative=( + authoritative_fallback_todo_items is not None + ), + available_capabilities=effective_available_capabilities, work_lane_contract=work_lane_contract, neutral_replan_ack_classifications=AUTONOMOUS_REPLAN_ACK_NEUTRAL_CLASSIFICATIONS, registered_agent_ids=registered_agent_ids, diff --git a/loopx/quota.py b/loopx/quota.py index 91d39a9789..6d7c7c1720 100644 --- a/loopx/quota.py +++ b/loopx/quota.py @@ -877,6 +877,7 @@ def build_quota_should_run( receipt_bound_replan_obligation_id: str | None = None, turn_instance_id: str | None = None, runtime_root: str | Path | None = None, + authoritative_fallback_todo_items: list[dict[str, Any]] | None = None, ) -> dict[str, Any]: from .control_plane.quota.should_run import ( build_quota_should_run as _build_quota_should_run, @@ -901,6 +902,7 @@ def build_quota_should_run( receipt_bound_replan_obligation_id=receipt_bound_replan_obligation_id, turn_instance_id=turn_instance_id, runtime_root=runtime_root, + authoritative_fallback_todo_items=authoritative_fallback_todo_items, ) diff --git a/tests/control_plane/test_fallback_wait_admission.py b/tests/control_plane/test_fallback_wait_admission.py new file mode 100644 index 0000000000..d0a87ff590 --- /dev/null +++ b/tests/control_plane/test_fallback_wait_admission.py @@ -0,0 +1,251 @@ +"""Fallback wait admission must not depend on compact Todo lane capacity.""" + +from __future__ import annotations + +import json +from pathlib import Path + +import pytest + +from loopx.cli import main as cli_main +from loopx.control_plane.goals.goal_vision import normalize_goal_vision_packet + +GOAL_ID = "fallback-wait-capacity-fixture" +AGENT_ID = "worker" +PRIMARY_TODO_ID = "todo_source_a" +FALLBACK_TODO_ID = "todo_source_b" +PREREQUISITE_TODO_ID = "todo_recovery" + + +def _state_text( + *, + unrelated_deferred_count: int, + fallback_present: bool = True, + prerequisite_status: str = "open", +) -> str: + unrelated = "".join( + f""" +- [ ] [P1] Wait for unrelated prerequisite {index}. + +""" + for index in range(unrelated_deferred_count) + ) + fallback = ( + f""" +- [ ] [P1] Read authorized source B after recovery. + +""" + if fallback_present + else "" + ) + prerequisite_marker = "x" if prerequisite_status == "done" else " " + return f"""# Active Goal State + +## Agent Todo + +- [{prerequisite_marker}] [P0] Observe recovery of source A. + +- [ ] [P0] Read source A after recovery. + +{unrelated} +{fallback} +""" + + +def _vision_run() -> dict: + vision = normalize_goal_vision_packet( + { + "state": "vision_drift_detected", + "vision_patch": { + "acceptance_summary": "Complete the bounded source check.", + "replan_trigger_summary": "Source A is unavailable.", + }, + "todo_delta": [f"retain:{PRIMARY_TODO_ID}"], + "fallback_declarations": [ + { + "declaration_id": "source-b", + "target_todo_id": FALLBACK_TODO_ID, + } + ], + }, + goal_id=GOAL_ID, + agent_id=AGENT_ID, + ) + return { + "classification": "source_path_checkpoint", + "agent_id": AGENT_ID, + "generated_at": "2026-09-01T00:00:00+00:00", + "agent_vision": vision, + } + + +def _write_fixture( + tmp_path: Path, + *, + unrelated_deferred_count: int, + fallback_present: bool = True, + prerequisite_status: str = "open", +) -> list[str]: + state_file = tmp_path / "ACTIVE_GOAL_STATE.md" + state_file.write_text( + _state_text( + unrelated_deferred_count=unrelated_deferred_count, + fallback_present=fallback_present, + prerequisite_status=prerequisite_status, + ) + ) + registry = tmp_path / "registry.json" + registry.write_text( + json.dumps( + { + "goals": [ + { + "id": GOAL_ID, + "status": "active", + "domain": "engineering", + "waiting_on": "codex", + "state_file": str(state_file), + "repo": str(tmp_path), + "adapter": { + "kind": "fixture_connected_delivery_v0", + "status": "connected-delivery", + }, + "quota": { + "compute": 1.0, + "window_hours": 24, + "slot_minutes": 1, + }, + "coordination": { + "agent_model": "peer_v1", + "registered_agents": [AGENT_ID, "observer"], + }, + } + ] + } + ) + ) + runtime = tmp_path / "runtime" + runs = runtime / "goals" / GOAL_ID / "runs" + runs.mkdir(parents=True) + run = _vision_run() + json_path = runs / "source-checkpoint.json" + markdown_path = runs / "source-checkpoint.md" + json_path.write_text(json.dumps(run) + "\n") + markdown_path.write_text("# Source checkpoint\n") + (runs / "index.jsonl").write_text( + json.dumps( + { + **run, + "json_path": str(json_path), + "markdown_path": str(markdown_path), + } + ) + + "\n" + ) + return [ + "--format", + "json", + "--registry", + str(registry), + "--runtime-root", + str(runtime), + "quota", + "should-run", + "--verbose", + "--goal-id", + GOAL_ID, + "--agent-id", + AGENT_ID, + "--runtime-profile", + "ark_managed_agent_goal", + ] + + +@pytest.mark.parametrize("unrelated_deferred_count", [0, 8, 20]) +def test_real_cli_wait_admission_is_stable_across_compact_lane_capacity( + tmp_path: Path, + capsys: pytest.CaptureFixture[str], + unrelated_deferred_count: int, +) -> None: + args = _write_fixture( + tmp_path, + unrelated_deferred_count=unrelated_deferred_count, + ) + + assert cli_main(args) == 0 + result = json.loads(capsys.readouterr().out) + assert result["should_run"] is False + assert result["execution_obligation"]["must_attempt_work"] is False + assert "fallback_gaps" not in result["goal_frontier_projection"] + if unrelated_deferred_count: + assert ( + result["agent_todo_summary"]["payload_compaction"]["omitted_lanes"][ + "deferred_items" + ] + > 0 + ) + assert ( + result["scheduler_hint"]["goal_runtime_continuation"]["disposition"] + == "defer" + ) + + +def test_real_cli_missing_declared_fallback_projects_true_unresolved_gap( + tmp_path: Path, + capsys: pytest.CaptureFixture[str], +) -> None: + args = _write_fixture( + tmp_path, + unrelated_deferred_count=20, + fallback_present=False, + ) + + assert cli_main(args) == 0 + result = json.loads(capsys.readouterr().out) + gap = result["goal_frontier_projection"]["fallback_gaps"][0] + assert gap["kind"] == "vision_fallback_unresolved" + assert gap["unresolved_todo_ids"] == [FALLBACK_TODO_ID] + assert "lookup_uncertain_todo_ids" not in gap + + +def test_real_cli_authority_read_failure_projects_uncertainty( + tmp_path: Path, + capsys: pytest.CaptureFixture[str], + monkeypatch: pytest.MonkeyPatch, +) -> None: + args = _write_fixture(tmp_path, unrelated_deferred_count=20) + + def fail_authority_read(**_kwargs: object) -> dict: + raise OSError("fixture canonical authority unavailable") + + monkeypatch.setattr( + "loopx.control_plane.quota.live_decision.list_goal_todos", + fail_authority_read, + ) + + assert cli_main(args) == 0 + result = json.loads(capsys.readouterr().out) + gap = result["goal_frontier_projection"]["fallback_gaps"][0] + assert gap["kind"] == "vision_fallback_lookup_uncertain" + assert gap["lookup_uncertain_todo_ids"] == [FALLBACK_TODO_ID] + assert "unresolved_todo_ids" not in gap + + +def test_real_cli_uses_dependency_state_from_exact_canonical_read( + tmp_path: Path, + capsys: pytest.CaptureFixture[str], +) -> None: + args = _write_fixture( + tmp_path, + unrelated_deferred_count=20, + prerequisite_status="done", + ) + + assert cli_main(args) == 0 + result = json.loads(capsys.readouterr().out) + assert result["should_run"] is True + assert result["selected_todo"]["todo_id"] in { + PRIMARY_TODO_ID, + FALLBACK_TODO_ID, + } + assert "fallback_gaps" not in result["goal_frontier_projection"] diff --git a/tests/control_plane/test_goal_frontier_fallback_disposition.py b/tests/control_plane/test_goal_frontier_fallback_disposition.py index c02190e43c..015b039eb1 100644 --- a/tests/control_plane/test_goal_frontier_fallback_disposition.py +++ b/tests/control_plane/test_goal_frontier_fallback_disposition.py @@ -22,6 +22,7 @@ quota_todo_summary, ) from loopx.control_plane.todos.projection import todo_advancement_frontier_counts +from loopx.control_plane.todos.summary_item import todo_planning_source_items from loopx.quota import build_quota_should_run GOAL_ID = "vision-fallback-disposition-fixture" @@ -143,8 +144,11 @@ def _status_payload( ) -def _frontier_projection(payload: dict) -> dict: +def _frontier_projection(payload: dict, *, include_source: bool = True) -> dict: item = payload["attention_queue"]["items"][0] + source_items = ( + todo_planning_source_items(item["agent_todos"]) if include_source else None + ) context = build_goal_frontier_projection_context_from_status( goal_id=GOAL_ID, agent_id=AGENT_ID, @@ -153,6 +157,7 @@ def _frontier_projection(payload: dict) -> dict: project_asset=item["project_asset"], user_todo_summary=item["user_todos"], agent_todo_summary=item["agent_todos"], + agent_todo_source_items=source_items, work_lane_contract=None, neutral_replan_ack_classifications=set(), registered_agent_ids=[PRIMARY_AGENT, AGENT_ID], @@ -248,7 +253,39 @@ def test_declared_fallback_without_resolution_projects_single_gap() -> None: scheduler_execution_context=(GENERIC_CLI_OUTER_CONTROLLER_SCHEDULER_CONTEXT), ) quota_projection = decision["goal_frontier_projection"] - assert quota_projection["fallback_gaps"][0]["unresolved_todo_ids"] == [FALLBACK_ID] + gap = quota_projection["fallback_gaps"][0] + assert gap["kind"] == "vision_fallback_lookup_uncertain" + assert gap["lookup_uncertain_todo_ids"] == [FALLBACK_ID] + + authoritative_decision = build_quota_should_run( + payload, + goal_id=GOAL_ID, + agent_id=AGENT_ID, + authoritative_fallback_todo_items=[], + scheduler_execution_context=(GENERIC_CLI_OUTER_CONTROLLER_SCHEDULER_CONTEXT), + ) + authoritative_gap = authoritative_decision["goal_frontier_projection"][ + "fallback_gaps" + ][0] + assert authoritative_gap["kind"] == "vision_fallback_unresolved" + assert authoritative_gap["unresolved_todo_ids"] == [FALLBACK_ID] + + +def test_missing_authoritative_source_projects_uncertainty_not_absence() -> None: + payload = _status_payload( + fallback_runnable=False, + latest_runs=[_fallback_vision_run(todo_delta=[f"retain:{FALLBACK_ID}"])], + ) + + frontier = _frontier_projection(payload, include_source=False) + + gap = frontier["fallback_gaps"][0] + assert gap["kind"] == "vision_fallback_lookup_uncertain" + assert gap["reason_code"] == ( + "declared_fallback_authoritative_lookup_unavailable" + ) + assert gap["lookup_uncertain_todo_ids"] == [FALLBACK_ID] + assert "unresolved_todo_ids" not in gap def test_retaining_only_the_blocked_primary_successor_is_no_declaration() -> None: @@ -349,6 +386,41 @@ def test_other_agent_primary_todo_does_not_resolve_the_gap() -> None: assert gaps[0]["unresolved_todo_ids"] == [FALLBACK_ID] +@pytest.mark.parametrize( + "ownership_metadata", + [ + {"claimed_by": PRIMARY_AGENT}, + {"excluded_agents": [AGENT_ID]}, + ], + ids=["peer-claimed", "current-agent-excluded"], +) +def test_authoritative_fallback_keeps_existing_ownership_and_exclusion_semantics( + ownership_metadata: dict[str, Any], +) -> None: + payload = _status_payload( + fallback_runnable=False, + latest_runs=[_fallback_vision_run(todo_delta=[f"retain:{FALLBACK_ID}"])], + ) + fallback = quota_todo_item( + todo_id=FALLBACK_ID, + index=3, + text="[P1] Deliver the declared fallback direction.", + **ownership_metadata, + ) + + decision = build_quota_should_run( + payload, + goal_id=GOAL_ID, + agent_id=AGENT_ID, + authoritative_fallback_todo_items=[fallback], + scheduler_execution_context=(GENERIC_CLI_OUTER_CONTROLLER_SCHEDULER_CONTEXT), + ) + + gap = decision["goal_frontier_projection"]["fallback_gaps"][0] + assert gap["kind"] == "vision_fallback_unresolved" + assert gap["unresolved_todo_ids"] == [FALLBACK_ID] + + def test_declared_fallback_linked_to_monitor_todo_keeps_the_gap() -> None: # A linked Todo that is not advancement work is not a runnable fallback # successor, so the declaration stays unresolved. @@ -367,6 +439,9 @@ def test_declared_fallback_linked_to_monitor_todo_keeps_the_gap() -> None: summary = _agent_todos(fallback_runnable=False) for slot in ("executable_backlog_items", "backlog_items"): summary[slot] = list(summary.get(slot) or []) + [monitor] + summary["monitor_open_items"] = list(summary.get("monitor_open_items") or []) + [ + monitor + ] item["agent_todos"] = summary frontier = _frontier_projection(payload) @@ -376,6 +451,55 @@ def test_declared_fallback_linked_to_monitor_todo_keeps_the_gap() -> None: assert gaps[0]["unresolved_todo_ids"] == [FALLBACK_ID] +@pytest.mark.parametrize( + ("resume_monitor_generation", "expected_gap_kind"), + [ + (3, None), + (None, "vision_fallback_unresolved"), + ], + ids=["valid-pending-monitor-wait", "invalid-missing-generation-baseline"], +) +def test_fallback_monitor_wait_uses_typed_resume_evaluator( + resume_monitor_generation: int | None, + expected_gap_kind: str | None, +) -> None: + payload = _status_payload( + fallback_runnable=False, + latest_runs=[_fallback_vision_run(todo_delta=[f"retain:{FALLBACK_ID}"])], + ) + fallback = quota_todo_item( + todo_id=FALLBACK_ID, + index=3, + text="[P1] Resume the declared fallback after observation changes.", + claimed_by=AGENT_ID, + resume_when="monitor_changed:todo_fallback_monitor", + resume_monitor_generation=resume_monitor_generation, + ) + monitor = quota_todo_item( + todo_id="todo_fallback_monitor", + index=4, + text="[P2] Observe the fallback dependency.", + task_class="continuous_monitor", + claimed_by=AGENT_ID, + material_change_generation=3, + ) + + decision = build_quota_should_run( + payload, + goal_id=GOAL_ID, + agent_id=AGENT_ID, + authoritative_fallback_todo_items=[fallback, monitor], + scheduler_execution_context=(GENERIC_CLI_OUTER_CONTROLLER_SCHEDULER_CONTEXT), + ) + + gaps = decision["goal_frontier_projection"].get("fallback_gaps", []) + if expected_gap_kind is None: + assert gaps == [] + else: + assert gaps[0]["kind"] == expected_gap_kind + assert gaps[0]["unresolved_todo_ids"] == [FALLBACK_ID] + + @pytest.mark.parametrize( ("state", "path_outcome"), [ From acd302cedafe2b9c12d624ae0b92371f78cb92ff Mon Sep 17 00:00:00 2001 From: "lusendong.6789" Date: Tue, 8 Sep 2026 12:49:26 +0800 Subject: [PATCH 2/7] fix(quota): bound fallback reads to direct dependencies Signed-off-by: lusendong.6789 --- .../goal-vision-replan-contract-v0.md | 8 ++ .../goals/goal_frontier/__init__.py | 9 +- loopx/control_plane/quota/live_decision.py | 75 +++++++++-------- .../control_plane/quota/should_run_prepare.py | 3 - .../work_items/semantic_replan_writeback.py | 1 + .../test_fallback_wait_admission.py | 84 +++++++++++++++++++ ...test_goal_frontier_fallback_disposition.py | 1 + 7 files changed, 134 insertions(+), 47 deletions(-) diff --git a/docs/reference/protocols/goal-vision-replan-contract-v0.md b/docs/reference/protocols/goal-vision-replan-contract-v0.md index 7700f9b742..4b0d6dcc9f 100644 --- a/docs/reference/protocols/goal-vision-replan-contract-v0.md +++ b/docs/reference/protocols/goal-vision-replan-contract-v0.md @@ -449,6 +449,14 @@ rule uses existing acceptance/lineage facts regardless of the optional advisory relationships, and it grants no additional authority. Ownership, exclusions, capabilities, user gates, and quota remain independent execution constraints. +For optional fallback advice, live quota reads the declared target/successor +Todos and one layer of direct resume dependencies, at most 16 exact reads. +Deeper dependency chains do not expand that lookup. Existing typed resume, +ownership, and lifecycle rules decide whether each declared path is available. +An unavailable or mismatched canonical read produces +`vision_fallback_lookup_uncertain`, rather than treating a missing display row +as a missing Todo. `fallback_gaps` remains advisory and adds no replan obligation. + 等待资格现在逐项检查已有 acceptance 的 Todo 关联,并在展示裁剪前从完整来源计算。 A 的等待不能遮住尚未落实的 B;有可执行工作则继续,相关工作都具有合法等待证据才暂缓。 无需另外维护 fallback 声明;无关 Todo 的数量和顺序不应改变决策。 diff --git a/loopx/control_plane/goals/goal_frontier/__init__.py b/loopx/control_plane/goals/goal_frontier/__init__.py index e6f98f4d90..f771028de1 100644 --- a/loopx/control_plane/goals/goal_frontier/__init__.py +++ b/loopx/control_plane/goals/goal_frontier/__init__.py @@ -1406,7 +1406,6 @@ def build_goal_frontier_projection_context_from_status( neutral_replan_ack_classifications: set[str], agent_todo_source_items: list[dict[str, Any]] | None = None, fallback_todo_source_items: list[dict[str, Any]] | None = None, - fallback_todo_source_authoritative: bool | None = None, available_capabilities: Any = None, registered_agent_ids: list[str] | None = None, goal_status: str | None = None, @@ -1539,13 +1538,7 @@ def build_goal_frontier_projection_context_from_status( latest_agent_vision, agent_todo_summary=agent_todo_summary, agent_id=agent_id, - agent_todo_source_items=( - fallback_todo_source_items - if fallback_todo_source_authoritative is True - else None - if fallback_todo_source_authoritative is False - else agent_todo_source_items - ), + agent_todo_source_items=fallback_todo_source_items, rollout_events=latest_runs_for_goal( status_payload, goal_id=goal_id, diff --git a/loopx/control_plane/quota/live_decision.py b/loopx/control_plane/quota/live_decision.py index e7b1ee2cdf..f93795e58c 100644 --- a/loopx/control_plane/quota/live_decision.py +++ b/loopx/control_plane/quota/live_decision.py @@ -68,8 +68,10 @@ def _live_fallback_authority_items( ``status`` is deliberately presentation-bounded, so omission from it can never prove that a declared fallback is absent. The live CLI path owns the registry/runtime authority needed for exact reads. Resume dependencies are - included only when referenced by one of the declared Todos so the existing - TypeScript evaluator retains state-transition authority. + included only when directly referenced by a declared Todo. Four declarations + can name eight targets/successors and eight direct dependencies: at most 16 + reads. The TypeScript evaluator inspects a dependency's current state, not + its own resume chain, and retains state-transition authority. """ requested_ids = _fallback_authority_todo_ids( @@ -82,40 +84,41 @@ def _live_fallback_authority_items( items: dict[str, dict[str, Any]] = {} pending_ids = set(requested_ids) - while pending_ids: - todo_id = pending_ids.pop() - try: - projection = list_goal_todos( - registry_path=registry_path, - goal_id=goal_id, - todo_id=todo_id, - runtime_root_arg=str(runtime_root), - limit=None, - ) - except ( - EffectRuntimeRemoteError, - LocalCoordinationAuthorityUnavailable, - OSError, - ValueError, - ): - return None - item = projection.get("todo") - if not isinstance(item, dict): - continue - normalized_id = normalize_todo_id(item.get("todo_id")) - if not normalized_id: - return None - items[normalized_id] = dict(item) - resume_when = normalize_todo_resume_when(item.get("resume_when")) - if resume_when: - resume_kind, _, target = resume_when.partition(":") - dependency_id = ( - normalize_todo_id(target) - if resume_kind in {"todo_done", "monitor_changed"} - else None - ) - if dependency_id and dependency_id not in items: - pending_ids.add(dependency_id) + for include_dependencies in (True, False): + dependency_ids: set[str] = set() + for todo_id in sorted(pending_ids): + try: + projection = list_goal_todos( + registry_path=registry_path, + goal_id=goal_id, + todo_id=todo_id, + runtime_root_arg=str(runtime_root), + limit=None, + ) + except ( + EffectRuntimeRemoteError, + LocalCoordinationAuthorityUnavailable, + OSError, + ValueError, + ): + return None + item = projection.get("todo") + if item is None: + continue + if not isinstance(item, dict) or normalize_todo_id(item.get("todo_id")) != todo_id: + return None + items[todo_id] = dict(item) + resume_when = normalize_todo_resume_when(item.get("resume_when")) + if include_dependencies and resume_when: + resume_kind, _, target = resume_when.partition(":") + dependency_id = ( + normalize_todo_id(target) + if resume_kind in {"todo_done", "monitor_changed"} + else None + ) + if dependency_id: + dependency_ids.add(dependency_id) + pending_ids = dependency_ids - requested_ids return list(items.values()) diff --git a/loopx/control_plane/quota/should_run_prepare.py b/loopx/control_plane/quota/should_run_prepare.py index 8d7d20f3ac..2b679827e6 100644 --- a/loopx/control_plane/quota/should_run_prepare.py +++ b/loopx/control_plane/quota/should_run_prepare.py @@ -755,9 +755,6 @@ def _prepare_quota_should_run_item( include_terminal=True, ), fallback_todo_source_items=authoritative_fallback_todo_items, - fallback_todo_source_authoritative=( - authoritative_fallback_todo_items is not None - ), available_capabilities=effective_available_capabilities, work_lane_contract=work_lane_contract, neutral_replan_ack_classifications=AUTONOMOUS_REPLAN_ACK_NEUTRAL_CLASSIFICATIONS, diff --git a/loopx/control_plane/work_items/semantic_replan_writeback.py b/loopx/control_plane/work_items/semantic_replan_writeback.py index 48db421e62..d70e62fc88 100644 --- a/loopx/control_plane/work_items/semantic_replan_writeback.py +++ b/loopx/control_plane/work_items/semantic_replan_writeback.py @@ -274,6 +274,7 @@ def qualify_replan_writeback( user_todo_summary=user_todos, agent_todo_summary=agent_todos, agent_todo_source_items=agent_todo_source_items, + fallback_todo_source_items=agent_todo_source_items, work_lane_contract=build_work_lane_context_contract( {"progress_scope": "primary_goal"}, agent_todo_summary=agent_todos, diff --git a/tests/control_plane/test_fallback_wait_admission.py b/tests/control_plane/test_fallback_wait_admission.py index d0a87ff590..20c68a2540 100644 --- a/tests/control_plane/test_fallback_wait_admission.py +++ b/tests/control_plane/test_fallback_wait_admission.py @@ -9,6 +9,7 @@ from loopx.cli import main as cli_main from loopx.control_plane.goals.goal_vision import normalize_goal_vision_packet +from loopx.control_plane.quota import live_decision GOAL_ID = "fallback-wait-capacity-fixture" AGENT_ID = "worker" @@ -249,3 +250,86 @@ def test_real_cli_uses_dependency_state_from_exact_canonical_read( FALLBACK_TODO_ID, } assert "fallback_gaps" not in result["goal_frontier_projection"] + + +@pytest.mark.parametrize("depth", [0, 4, 64, 512]) +def test_real_cli_reads_only_direct_fallback_dependencies( + tmp_path: Path, + capsys: pytest.CaptureFixture[str], + monkeypatch: pytest.MonkeyPatch, + depth: int, +) -> None: + args = _write_fixture(tmp_path, unrelated_deferred_count=20) + state_file = tmp_path / "ACTIVE_GOAL_STATE.md" + state = state_file.read_text() + if depth: + state = state.replace( + f"todo_id={PREREQUISITE_TODO_ID} status=open", + f"todo_id={PREREQUISITE_TODO_ID} resume_when=todo_done:todo_chain_0 status=deferred", + ) + for index in range(depth): + state += ( + f"\n- [ ] [P2] Observe prerequisite {index}.\n" + f" \n" + ) + state_file.write_text(state) + reads: list[str] = [] + original_read = live_decision.list_goal_todos + + def record_read(**kwargs: object) -> dict: + reads.append(str(kwargs["todo_id"])) + return original_read(**kwargs) + + monkeypatch.setattr(live_decision, "list_goal_todos", record_read) + assert cli_main(args) == 0 + result = json.loads(capsys.readouterr().out) + assert result["should_run"] is False + assert "fallback_gaps" not in result["goal_frontier_projection"] + # The direct prerequisite is deferred regardless of its own dependency. + # Its chain cannot alter this wait or add canonical lookup calls. + assert reads == [FALLBACK_TODO_ID, PREREQUISITE_TODO_ID] + + +def test_real_cli_mismatched_authority_identity_is_uncertain( + tmp_path: Path, + capsys: pytest.CaptureFixture[str], + monkeypatch: pytest.MonkeyPatch, +) -> None: + args = _write_fixture(tmp_path, unrelated_deferred_count=20) + original_read = live_decision.list_goal_todos + + def mismatched_read(**kwargs: object) -> dict: + projection = original_read(**kwargs) + return {**projection, "todo": {**projection["todo"], "todo_id": "todo_other"}} + + monkeypatch.setattr(live_decision, "list_goal_todos", mismatched_read) + assert cli_main(args) == 0 + result = json.loads(capsys.readouterr().out) + gap = result["goal_frontier_projection"]["fallback_gaps"][0] + assert gap["kind"] == "vision_fallback_lookup_uncertain" + assert gap["lookup_uncertain_todo_ids"] == [FALLBACK_TODO_ID] + assert "unresolved_todo_ids" not in gap + + +def test_real_cli_without_declarations_does_not_read_fallback_authority( + tmp_path: Path, + capsys: pytest.CaptureFixture[str], + monkeypatch: pytest.MonkeyPatch, +) -> None: + args = _write_fixture(tmp_path, unrelated_deferred_count=20) + index = tmp_path / "runtime" / "goals" / GOAL_ID / "runs" / "index.jsonl" + run = json.loads(index.read_text()) + run["agent_vision"].pop("fallback_declarations") + index.write_text(json.dumps(run) + "\n") + Path(run["json_path"]).write_text(json.dumps(run) + "\n") + + def unexpected_read(**_kwargs: object) -> dict: + pytest.fail("No fallback declaration authorizes an exact fallback lookup") + + monkeypatch.setattr(live_decision, "list_goal_todos", unexpected_read) + assert cli_main(args) == 0 + result = json.loads(capsys.readouterr().out) + assert result["should_run"] is False + assert "fallback_gaps" not in result["goal_frontier_projection"] diff --git a/tests/control_plane/test_goal_frontier_fallback_disposition.py b/tests/control_plane/test_goal_frontier_fallback_disposition.py index efcc6cabc8..3b23e03dfd 100644 --- a/tests/control_plane/test_goal_frontier_fallback_disposition.py +++ b/tests/control_plane/test_goal_frontier_fallback_disposition.py @@ -158,6 +158,7 @@ def _frontier_projection(payload: dict, *, include_source: bool = True) -> dict: user_todo_summary=item["user_todos"], agent_todo_summary=item["agent_todos"], agent_todo_source_items=source_items, + fallback_todo_source_items=source_items, work_lane_contract=None, neutral_replan_ack_classifications=set(), registered_agent_ids=[PRIMARY_AGENT, AGENT_ID], From b27a7a000b8b98d164a81c708172a4a3634a6d10 Mon Sep 17 00:00:00 2001 From: "lusendong.6789" Date: Tue, 8 Sep 2026 22:21:47 +0800 Subject: [PATCH 3/7] fix(quota): preserve fallback wait eligibility and uncertainty Signed-off-by: lusendong.6789 --- .../goal-vision-replan-contract-v0.md | 9 ++- .../goals/goal_frontier/__init__.py | 3 +- .../goal_frontier/fallback_disposition.py | 29 +++++-- loopx/control_plane/quota/live_decision.py | 14 +++- loopx/control_plane/quota/should_run.py | 7 +- .../control_plane/quota/should_run_prepare.py | 7 +- loopx/control_plane/todos/deferred_resume.py | 19 +++-- loopx/quota.py | 7 +- .../test_fallback_wait_admission.py | 81 ++++++++++++++++++- ...test_goal_frontier_fallback_disposition.py | 9 +++ 10 files changed, 157 insertions(+), 28 deletions(-) diff --git a/docs/reference/protocols/goal-vision-replan-contract-v0.md b/docs/reference/protocols/goal-vision-replan-contract-v0.md index 4b0d6dcc9f..4355117500 100644 --- a/docs/reference/protocols/goal-vision-replan-contract-v0.md +++ b/docs/reference/protocols/goal-vision-replan-contract-v0.md @@ -453,9 +453,12 @@ For optional fallback advice, live quota reads the declared target/successor Todos and one layer of direct resume dependencies, at most 16 exact reads. Deeper dependency chains do not expand that lookup. Existing typed resume, ownership, and lifecycle rules decide whether each declared path is available. -An unavailable or mismatched canonical read produces -`vision_fallback_lookup_uncertain`, rather than treating a missing display row -as a missing Todo. `fallback_gaps` remains advisory and adds no replan obligation. +An unavailable, ambiguous, or mismatched canonical read produces +`vision_fallback_lookup_uncertain`; compact display evidence cannot override a +failed authority read. Only explicit canonical not-found proves absence. +Pending `todo_done:` keeps an unresolved fallback; a valid +`monitor_changed` generation condition may establish a wait. `fallback_gaps` +remains advisory and adds no replan obligation. 等待资格现在逐项检查已有 acceptance 的 Todo 关联,并在展示裁剪前从完整来源计算。 A 的等待不能遮住尚未落实的 B;有可执行工作则继续,相关工作都具有合法等待证据才暂缓。 diff --git a/loopx/control_plane/goals/goal_frontier/__init__.py b/loopx/control_plane/goals/goal_frontier/__init__.py index f771028de1..d9f6af8342 100644 --- a/loopx/control_plane/goals/goal_frontier/__init__.py +++ b/loopx/control_plane/goals/goal_frontier/__init__.py @@ -56,6 +56,7 @@ from .fallback_disposition import ( VISION_FRONTIER_TODO_DELTA_ACTIONS, # noqa: F401 FallbackDeclaration, # noqa: F401 + FallbackTodoSource, agent_scoped_selectable_advancement_todo_ids, # noqa: F401 declared_fallback_gap_from_agent_vision, parse_fallback_declarations, # noqa: F401 @@ -1405,7 +1406,7 @@ def build_goal_frontier_projection_context_from_status( work_lane_contract: dict[str, Any] | None, neutral_replan_ack_classifications: set[str], agent_todo_source_items: list[dict[str, Any]] | None = None, - fallback_todo_source_items: list[dict[str, Any]] | None = None, + fallback_todo_source_items: FallbackTodoSource = None, available_capabilities: Any = None, registered_agent_ids: list[str] | None = None, goal_status: str | None = None, diff --git a/loopx/control_plane/goals/goal_frontier/fallback_disposition.py b/loopx/control_plane/goals/goal_frontier/fallback_disposition.py index ce1d2bb16e..38337ec922 100644 --- a/loopx/control_plane/goals/goal_frontier/fallback_disposition.py +++ b/loopx/control_plane/goals/goal_frontier/fallback_disposition.py @@ -1,6 +1,7 @@ from __future__ import annotations from dataclasses import dataclass +from enum import Enum from typing import Any from ...todos.contract import ( @@ -11,7 +12,10 @@ normalize_todo_resume_when, normalize_todo_status, ) -from ...todos.deferred_resume import todo_summary_blocked_successor_items +from ...todos.deferred_resume import ( + todo_resume_condition_is_non_monitor_wait, + todo_summary_blocked_successor_items, +) from ...todos.projection import ( agent_scoped_selectable_advancement_todo_ids, todo_item_claimed_by_agent_or_unclaimed, @@ -55,6 +59,15 @@ ) +class FallbackTodoReadState(Enum): + UNAVAILABLE = "unavailable" + + +# None is an omitted source from a legacy caller; UNAVAILABLE is a failed +# authority read, which compact display evidence must not override. +FallbackTodoSource = list[dict[str, Any]] | FallbackTodoReadState | None + + @dataclass(frozen=True) class FallbackDeclaration: """Structured declaration of a fallback direction and its associated work.""" @@ -259,7 +272,10 @@ def _authoritative_fallback_disposition_ids( continue if condition.get("satisfied") is True: runnable_ids.add(todo_id) - elif condition.get("satisfied") is False: + elif condition.get("satisfied") is False and ( + condition.get("kind") == "monitor_changed" + or todo_resume_condition_is_non_monitor_wait(condition) + ): waiting_ids.add(todo_id) continue @@ -274,7 +290,7 @@ def declared_fallback_gap_from_agent_vision( *, agent_todo_summary: dict[str, Any] | None, agent_id: str | None, - agent_todo_source_items: list[dict[str, Any]] | None = None, + agent_todo_source_items: FallbackTodoSource = None, rollout_events: list[dict[str, Any]] | None = None, available_capabilities: Any = None, ) -> dict[str, Any] | None: @@ -315,9 +331,9 @@ def declared_fallback_gap_from_agent_vision( if not declarations: return None - source_is_authoritative = agent_todo_source_items is not None + source_is_authoritative = isinstance(agent_todo_source_items, list) uncertain_todo_ids: set[str] = set() - if source_is_authoritative: + if isinstance(agent_todo_source_items, list): ( selectable_ids, waiting_todo_ids, @@ -329,6 +345,9 @@ def declared_fallback_gap_from_agent_vision( rollout_events=rollout_events, available_capabilities=available_capabilities, ) + elif agent_todo_source_items is FallbackTodoReadState.UNAVAILABLE: + selectable_ids = set() + waiting_todo_ids = set() else: # Legacy direct callers may only have a compact display summary. It # remains valid positive evidence, but omission from a bounded lane is diff --git a/loopx/control_plane/quota/live_decision.py b/loopx/control_plane/quota/live_decision.py index f93795e58c..49156bbc3f 100644 --- a/loopx/control_plane/quota/live_decision.py +++ b/loopx/control_plane/quota/live_decision.py @@ -10,6 +10,8 @@ from ..coordination.local_authority import LocalCoordinationAuthorityUnavailable from ..effect_runtime import EffectRuntimeRemoteError from ..goals.goal_frontier.fallback_disposition import ( + FallbackTodoReadState, + FallbackTodoSource, parse_fallback_declarations, ) from ..goals.goal_frontier.semantic_history import ( @@ -62,7 +64,7 @@ def _live_fallback_authority_items( runtime_root: Path, goal_id: str, agent_id: str | None, -) -> list[dict[str, Any]] | None: +) -> FallbackTodoSource: """Read only the exact canonical Todos needed by fallback disposition. ``status`` is deliberately presentation-bounded, so omission from it can @@ -101,12 +103,16 @@ def _live_fallback_authority_items( OSError, ValueError, ): - return None + return FallbackTodoReadState.UNAVAILABLE + if projection.get("ambiguous") is True: + return FallbackTodoReadState.UNAVAILABLE item = projection.get("todo") if item is None: - continue + if projection.get("not_found") is True: + continue + return FallbackTodoReadState.UNAVAILABLE if not isinstance(item, dict) or normalize_todo_id(item.get("todo_id")) != todo_id: - return None + return FallbackTodoReadState.UNAVAILABLE items[todo_id] = dict(item) resume_when = normalize_todo_resume_when(item.get("resume_when")) if include_dependencies and resume_when: diff --git a/loopx/control_plane/quota/should_run.py b/loopx/control_plane/quota/should_run.py index 4fa1e728fc..b248b88bb5 100644 --- a/loopx/control_plane/quota/should_run.py +++ b/loopx/control_plane/quota/should_run.py @@ -2,7 +2,10 @@ from collections.abc import Callable, Mapping from pathlib import Path -from typing import Any +from typing import TYPE_CHECKING, Any + +if TYPE_CHECKING: + from ..goals.goal_frontier.fallback_disposition import FallbackTodoSource from ...quota import ( _build_quota_plan_for_goal, @@ -258,7 +261,7 @@ def build_quota_should_run( receipt_bound_replan_obligation_id: str | None = None, turn_instance_id: str | None = None, runtime_root: str | Path | None = None, - authoritative_fallback_todo_items: list[dict[str, Any]] | None = None, + authoritative_fallback_todo_items: FallbackTodoSource = None, ) -> dict[str, Any]: safe_goal_id = str(goal_id or "").strip() resolved_scheduler_context = resolve_scheduler_execution_context( diff --git a/loopx/control_plane/quota/should_run_prepare.py b/loopx/control_plane/quota/should_run_prepare.py index 2b679827e6..deb451ce81 100644 --- a/loopx/control_plane/quota/should_run_prepare.py +++ b/loopx/control_plane/quota/should_run_prepare.py @@ -3,7 +3,10 @@ from collections.abc import Callable, Mapping from dataclasses import dataclass from pathlib import Path -from typing import Any +from typing import TYPE_CHECKING, Any + +if TYPE_CHECKING: + from ..goals.goal_frontier.fallback_disposition import FallbackTodoSource from ...quota import ( AUTONOMOUS_REPLAN_ACK_NEUTRAL_CLASSIFICATIONS, @@ -454,7 +457,7 @@ def _prepare_quota_should_run_item( receipt_bound_monitor_phase: ReceiptBoundMonitorPhase | None, receipt_bound_replay_phase: ReceiptBoundReplayPhase | None, receipt_bound_replan_obligation_id: str | None, - authoritative_fallback_todo_items: list[dict[str, Any]] | None = None, + authoritative_fallback_todo_items: FallbackTodoSource = None, ) -> _QuotaDecisionPreparation: quota = item.get("quota") if isinstance(item.get("quota"), dict) else {} state = str(quota.get("state") or "unknown") diff --git a/loopx/control_plane/todos/deferred_resume.py b/loopx/control_plane/todos/deferred_resume.py index 10572a4e90..87a7838cc7 100644 --- a/loopx/control_plane/todos/deferred_resume.py +++ b/loopx/control_plane/todos/deferred_resume.py @@ -299,6 +299,17 @@ def todo_summary_monitor_blocked_resume_items( return sorted(_dedupe_todo_items(candidates), key=todo_projection_sort_key) +def todo_resume_condition_is_non_monitor_wait(condition: dict[str, Any]) -> bool: + """A pending condition cannot use a standing monitor as a done prerequisite. + + Generation-fenced monitor waits are validated separately by the TS resume + evaluator; this preserves the existing ordinary successor wait boundary. + """ + return condition.get("satisfied") is False and normalize_todo_task_class( + condition.get("target_task_class"), text="" + ) != TODO_TASK_CLASS_MONITOR + + def todo_summary_blocked_successor_items( value: dict[str, Any], *, @@ -336,13 +347,7 @@ def todo_summary_blocked_successor_items( 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: + if not resume_when or not todo_resume_condition_is_non_monitor_wait(condition): continue compact = dict(item) compact["resume_when"] = resume_when diff --git a/loopx/quota.py b/loopx/quota.py index 6d7c7c1720..afc31852a2 100644 --- a/loopx/quota.py +++ b/loopx/quota.py @@ -3,7 +3,10 @@ from collections.abc import Callable, Mapping from datetime import datetime, timedelta, timezone from pathlib import Path -from typing import Any +from typing import TYPE_CHECKING, Any + +if TYPE_CHECKING: + from .control_plane.goals.goal_frontier.fallback_disposition import FallbackTodoSource from .control_plane import compact_control_plane_policy from .control_plane.goals.activation import goal_is_stopped @@ -877,7 +880,7 @@ def build_quota_should_run( receipt_bound_replan_obligation_id: str | None = None, turn_instance_id: str | None = None, runtime_root: str | Path | None = None, - authoritative_fallback_todo_items: list[dict[str, Any]] | None = None, + authoritative_fallback_todo_items: FallbackTodoSource = None, ) -> dict[str, Any]: from .control_plane.quota.should_run import ( build_quota_should_run as _build_quota_should_run, diff --git a/tests/control_plane/test_fallback_wait_admission.py b/tests/control_plane/test_fallback_wait_admission.py index 20c68a2540..d77f7d37c6 100644 --- a/tests/control_plane/test_fallback_wait_admission.py +++ b/tests/control_plane/test_fallback_wait_admission.py @@ -209,12 +209,14 @@ def test_real_cli_missing_declared_fallback_projects_true_unresolved_gap( assert "lookup_uncertain_todo_ids" not in gap +@pytest.mark.parametrize("unrelated_deferred_count", [0, 20]) def test_real_cli_authority_read_failure_projects_uncertainty( tmp_path: Path, capsys: pytest.CaptureFixture[str], monkeypatch: pytest.MonkeyPatch, + unrelated_deferred_count: int, ) -> None: - args = _write_fixture(tmp_path, unrelated_deferred_count=20) + args = _write_fixture(tmp_path, unrelated_deferred_count=unrelated_deferred_count) def fail_authority_read(**_kwargs: object) -> dict: raise OSError("fixture canonical authority unavailable") @@ -292,12 +294,14 @@ def record_read(**kwargs: object) -> dict: assert reads == [FALLBACK_TODO_ID, PREREQUISITE_TODO_ID] +@pytest.mark.parametrize("unrelated_deferred_count", [0, 20]) def test_real_cli_mismatched_authority_identity_is_uncertain( tmp_path: Path, capsys: pytest.CaptureFixture[str], monkeypatch: pytest.MonkeyPatch, + unrelated_deferred_count: int, ) -> None: - args = _write_fixture(tmp_path, unrelated_deferred_count=20) + args = _write_fixture(tmp_path, unrelated_deferred_count=unrelated_deferred_count) original_read = live_decision.list_goal_todos def mismatched_read(**kwargs: object) -> dict: @@ -333,3 +337,76 @@ def unexpected_read(**_kwargs: object) -> dict: result = json.loads(capsys.readouterr().out) assert result["should_run"] is False assert "fallback_gaps" not in result["goal_frontier_projection"] + + +@pytest.mark.parametrize("fallback_status", ["open", "deferred"]) +@pytest.mark.parametrize( + ("resume_kind", "generation", "expected_gap"), + [ + ("todo_done", None, "vision_fallback_unresolved"), + ("monitor_changed", 3, None), + ("monitor_changed", None, "vision_fallback_unresolved"), + ], + ids=["monitor-completion-is-invalid", "generation-wait", "missing-baseline"], +) +def test_real_cli_fallback_monitor_wait_requires_generation_condition( + tmp_path: Path, + capsys: pytest.CaptureFixture[str], + fallback_status: str, + resume_kind: str, + generation: int | None, + expected_gap: str | None, +) -> None: + args = _write_fixture(tmp_path, unrelated_deferred_count=0) + state_file = tmp_path / "ACTIVE_GOAL_STATE.md" + generation_metadata = ( + f" resume_monitor_generation={generation}" if generation is not None else "" + ) + state = state_file.read_text().replace( + f"todo_id={FALLBACK_TODO_ID} status=deferred " + f"task_class=advancement_task claimed_by={AGENT_ID} " + f"resume_when=todo_done:{PREREQUISITE_TODO_ID}", + f"todo_id={FALLBACK_TODO_ID} status={fallback_status} " + f"task_class=advancement_task claimed_by={AGENT_ID} " + f"resume_when={resume_kind}:todo_monitor{generation_metadata}", + ) + state += ( + "\n- [ ] [P2] Observe the fallback dependency.\n" + " \n" + ) + state_file.write_text(state) + + assert cli_main(args) == 0 + result = json.loads(capsys.readouterr().out) + gaps = result["goal_frontier_projection"].get("fallback_gaps", []) + if expected_gap is None: + assert gaps == [] + else: + assert gaps[0]["kind"] == expected_gap + assert gaps[0]["unresolved_todo_ids"] == [FALLBACK_TODO_ID] + + +@pytest.mark.parametrize("duplicate_id", [FALLBACK_TODO_ID, PREREQUISITE_TODO_ID]) +def test_real_cli_ambiguous_canonical_fallback_read_is_uncertain( + tmp_path: Path, + capsys: pytest.CaptureFixture[str], + duplicate_id: str, +) -> None: + args = _write_fixture(tmp_path, unrelated_deferred_count=0) + state_file = tmp_path / "ACTIVE_GOAL_STATE.md" + state_file.write_text( + state_file.read_text() + + "\n## User Todo\n\n- [ ] [P2] Resolve a separate user action.\n" + + f" \n" + ) + + # The existing duplicate-id input gate remains an error. Its advisory + # must also retain the ambiguity, rather than claiming canonical absence. + assert cli_main(args) == 1 + result = json.loads(capsys.readouterr().out) + gap = result["goal_frontier_projection"]["fallback_gaps"][0] + assert gap["kind"] == "vision_fallback_lookup_uncertain" + assert gap["lookup_uncertain_todo_ids"] == [FALLBACK_TODO_ID] + assert "unresolved_todo_ids" not in gap diff --git a/tests/control_plane/test_goal_frontier_fallback_disposition.py b/tests/control_plane/test_goal_frontier_fallback_disposition.py index 3b23e03dfd..fcb5fca353 100644 --- a/tests/control_plane/test_goal_frontier_fallback_disposition.py +++ b/tests/control_plane/test_goal_frontier_fallback_disposition.py @@ -287,6 +287,15 @@ def test_missing_authoritative_source_projects_uncertainty_not_absence() -> None assert "unresolved_todo_ids" not in gap +def test_legacy_omitted_source_retains_positive_display_evidence() -> None: + payload = _status_payload( + fallback_runnable=True, + latest_runs=[_fallback_vision_run()], + ) + + assert "fallback_gaps" not in _frontier_projection(payload, include_source=False) + + def test_retaining_only_the_blocked_primary_successor_is_no_declaration() -> None: # The primary successor is the wait state's own object; retaining it does # not declare a fallback direction, so no gap is invented. From 661c4e67a8ef2c39b906c984f4a0215e2bdfe6f6 Mon Sep 17 00:00:00 2001 From: "lusendong.6789" Date: Fri, 11 Sep 2026 18:07:07 +0800 Subject: [PATCH 4/7] fix(quota): reject waits on completed fallback monitors Signed-off-by: lusendong.6789 Co-authored-by: TRAE CLI --- .../goal_frontier/fallback_disposition.py | 25 +++++++++++++++---- ...test_goal_frontier_fallback_disposition.py | 25 +++++++++++++++---- 2 files changed, 40 insertions(+), 10 deletions(-) diff --git a/loopx/control_plane/goals/goal_frontier/fallback_disposition.py b/loopx/control_plane/goals/goal_frontier/fallback_disposition.py index 38337ec922..5dc73900a7 100644 --- a/loopx/control_plane/goals/goal_frontier/fallback_disposition.py +++ b/loopx/control_plane/goals/goal_frontier/fallback_disposition.py @@ -8,6 +8,7 @@ TODO_STATUS_DEFERRED, TODO_STATUS_OPEN, TODO_TASK_CLASS_ADVANCEMENT, + TODO_TASK_CLASS_MONITOR, normalize_todo_id, normalize_todo_resume_when, normalize_todo_status, @@ -272,11 +273,25 @@ def _authoritative_fallback_disposition_ids( continue if condition.get("satisfied") is True: runnable_ids.add(todo_id) - elif condition.get("satisfied") is False and ( - condition.get("kind") == "monitor_changed" - or todo_resume_condition_is_non_monitor_wait(condition) - ): - waiting_ids.add(todo_id) + elif condition.get("satisfied") is False: + if condition.get("kind") == "monitor_changed": + # An unchanged generation is a durable wait only while + # its continuous-monitor target remains open. A terminal + # monitor cannot produce another generation, so accepting + # that condition as pending would hide the fallback gap + # forever. A changed generation is handled above and + # remains runnable even if the monitor closed concurrently. + target_status = normalize_todo_status( + condition.get("target_status") + ) + if ( + target_status == TODO_STATUS_OPEN + and condition.get("target_task_class") + == TODO_TASK_CLASS_MONITOR + ): + waiting_ids.add(todo_id) + elif todo_resume_condition_is_non_monitor_wait(condition): + waiting_ids.add(todo_id) continue if todo_item_is_actionable_open(item): diff --git a/tests/control_plane/test_goal_frontier_fallback_disposition.py b/tests/control_plane/test_goal_frontier_fallback_disposition.py index fcb5fca353..5b74048630 100644 --- a/tests/control_plane/test_goal_frontier_fallback_disposition.py +++ b/tests/control_plane/test_goal_frontier_fallback_disposition.py @@ -460,15 +460,29 @@ def test_declared_fallback_linked_to_monitor_todo_keeps_the_gap() -> None: @pytest.mark.parametrize( - ("resume_monitor_generation", "expected_gap_kind"), + ( + "resume_monitor_generation", + "monitor_generation", + "monitor_status", + "expected_gap_kind", + ), [ - (3, None), - (None, "vision_fallback_unresolved"), + (3, 3, "open", None), + (None, 3, "open", "vision_fallback_unresolved"), + (3, 3, "done", "vision_fallback_unresolved"), + (3, 4, "done", None), + ], + ids=[ + "valid-pending-monitor-wait", + "invalid-missing-generation-baseline", + "completed-monitor-cannot-remain-pending", + "changed-generation-remains-runnable-after-monitor-completes", ], - ids=["valid-pending-monitor-wait", "invalid-missing-generation-baseline"], ) def test_fallback_monitor_wait_uses_typed_resume_evaluator( resume_monitor_generation: int | None, + monitor_generation: int, + monitor_status: str, expected_gap_kind: str | None, ) -> None: payload = _status_payload( @@ -489,7 +503,8 @@ def test_fallback_monitor_wait_uses_typed_resume_evaluator( text="[P2] Observe the fallback dependency.", task_class="continuous_monitor", claimed_by=AGENT_ID, - material_change_generation=3, + status=monitor_status, + material_change_generation=monitor_generation, ) decision = build_quota_should_run( From 7ae1b177ffaaaa26b558edf774d4788cba50af09 Mon Sep 17 00:00:00 2001 From: huangruiteng Date: Sat, 12 Sep 2026 11:30:46 +0800 Subject: [PATCH 5/7] refactor(quota): compose typed fallback disposition from one source snapshot Signed-off-by: huangruiteng --- .../control_plane/effect_runtime_handlers.ts | 2 + .../goal_frontier/fallback_disposition.py | 404 ++++-------------- .../goal_frontier/fallback_disposition.ts | 125 ++++++ .../goals/goal_frontier/fallback_source.py | 79 ++++ loopx/control_plane/quota/live_decision.py | 104 +---- loopx/control_plane/todos/resume_condition.py | 22 +- .../test_fallback_wait_admission.py | 94 +++- ...test_goal_frontier_fallback_disposition.py | 32 ++ .../fallback_disposition.test.ts | 98 +++++ 9 files changed, 518 insertions(+), 442 deletions(-) create mode 100644 loopx/control_plane/goals/goal_frontier/fallback_disposition.ts create mode 100644 loopx/control_plane/goals/goal_frontier/fallback_source.py create mode 100644 tests/control_plane_ts/fallback_disposition.test.ts diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index 0122154cee..36d2172122 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -77,6 +77,7 @@ import { } from "./todos/resume_condition.ts"; import { evaluateSchedulerStateTransition } from "./scheduler/state_transition_rules.ts"; import { projectTodoResumePlanning } from "./todos/resume_planning.ts"; +import { projectFallbackDisposition } from "./goals/goal_frontier/fallback_disposition.ts"; import { projectTodoQuotaPlanning } from "./todos/quota_selection.ts"; import { evaluateSchedulerStateOperation, @@ -416,6 +417,7 @@ export function createEffectRuntimeHandlers( ["todo.resume_condition.normalize", normalizeTodoResumeWhen], ["todo.resume_condition.evaluate", evaluateTodoResumeConditions], ["todo.resume_planning.project", projectTodoResumePlanning], + ["goal.fallback_disposition.project", projectFallbackDisposition], ["todo.quota_planning.project", projectTodoQuotaPlanning], ["todo.external_wait.plan", planTodoExternalWaitTransition], ["scheduler.state_transition.evaluate", evaluateSchedulerStateTransition], diff --git a/loopx/control_plane/goals/goal_frontier/fallback_disposition.py b/loopx/control_plane/goals/goal_frontier/fallback_disposition.py index 1fe3db0af7..cb214f84e5 100644 --- a/loopx/control_plane/goals/goal_frontier/fallback_disposition.py +++ b/loopx/control_plane/goals/goal_frontier/fallback_disposition.py @@ -5,22 +5,23 @@ from typing import Any from ...todos.contract import ( - TODO_STATUS_DEFERRED, - TODO_STATUS_OPEN, - TODO_TASK_CLASS_ADVANCEMENT, - TODO_TASK_CLASS_MONITOR, + normalize_todo_claimed_by, + normalize_todo_excluded_agents, normalize_todo_id, normalize_todo_resume_when, normalize_todo_status, ) -from ...todos.resume_planning import project_todo_resume_planning +from ...effect_runtime import effect_runtime_result +from ...todos.resume_planning import build_todo_resume_planning_request from ...todos.projection import ( agent_scoped_selectable_advancement_todo_ids, - todo_item_claimed_by_agent_or_unclaimed, - todo_item_is_actionable_open, + todo_item_has_removed_continuation_policy, todo_item_task_class, ) -from ...todos.resume_condition import evaluate_todo_resume_conditions +from ...todos.resume_condition import ( + build_todo_resume_evaluation_request, + normalize_todo_generation, +) from ..goal_vision_read_model import ( VISION_FRONTIER_TODO_DELTA_ACTIONS as VISION_FRONTIER_TODO_DELTA_ACTIONS, VISION_TODO_DELTA_ID_LIMIT as VISION_TODO_DELTA_ID_LIMIT, @@ -32,11 +33,7 @@ # fallback disposition on their own; activate/resume/retain entries only link # the vision to existing Todos and still need a selectable frontier match. VISION_TODO_DELTA_SUCCESSOR_ACTIONS = frozenset({"create", "reopen"}) -VISION_TODO_DELTA_LINKAGE_ACTIONS = frozenset( - VISION_FRONTIER_TODO_DELTA_ACTIONS - VISION_TODO_DELTA_SUCCESSOR_ACTIONS -) VISION_FALLBACK_DECLARATION_ENTRY_LIMIT = 4 -VISION_FALLBACK_DECLARATION_FIELDS = ("target_todo_id", "successor_todo_id") VISION_FALLBACK_GAP_TRIGGER = "vision_fallback_unresolved" VISION_FALLBACK_GAP_REASON_CODE = "declared_fallback_without_runnable_or_terminal" VISION_FALLBACK_LOOKUP_UNCERTAIN_TRIGGER = "vision_fallback_lookup_uncertain" @@ -44,7 +41,6 @@ "declared_fallback_authoritative_lookup_unavailable" ) VISION_FALLBACK_TERMINAL_PATH_OUTCOME = "stop" -VISION_FALLBACK_RUNNABLE_ITEM_LIMIT = 3 VISION_FALLBACK_RECOMMENDED_ACTION = ( "resolve the declared fallback direction: link or retain a runnable " "successor Todo referencing it, declare a bounded create/reopen " @@ -139,185 +135,6 @@ def parse_fallback_declarations( return declarations -def _blocked_successor_todo_ids( - agent_todo_summary: dict[str, Any] | None, - *, - agent_id: str | None, -) -> set[str]: - """Ids the blocked-successor wait state itself is waiting on.""" - - if not isinstance(agent_todo_summary, dict): - return set() - return { - todo_id - for todo_id in ( - normalize_todo_id(item.get("todo_id")) - for item in project_todo_resume_planning( - agent_todo_summary, - agent_id=agent_id, - )["blocked_successor_items"] - if isinstance(item, dict) - ) - if todo_id - } - - -def _blocked_primary_waiting( - agent_todo_summary: dict[str, Any] | None, - *, - agent_id: str | None, -) -> bool: - """Reuse the blocked-successor wait scope as the primary-blocked signal.""" - - if not isinstance(agent_todo_summary, dict): - return False - blocker_items = agent_todo_summary.get("current_agent_blocker_items") - if isinstance(blocker_items, list) and blocker_items: - return True - return bool( - project_todo_resume_planning( - agent_todo_summary, - agent_id=agent_id, - )["blocked_successor_items"] - ) - - -def _vision_has_terminal_disposition(agent_vision: dict[str, Any]) -> bool: - """Terminal evidence: closed-family state or path_delta.outcome=stop.""" - - if goal_vision_state_is_closed(agent_vision.get("state")): - return True - path_delta = agent_vision.get("path_delta") - path_delta = path_delta if isinstance(path_delta, dict) else {} - return ( - str(path_delta.get("outcome") or "").strip().lower() - == VISION_FALLBACK_TERMINAL_PATH_OUTCOME - ) - - -def _authoritative_fallback_disposition_ids( - declarations: list[FallbackDeclaration], - *, - agent_todo_source_items: list[dict[str, Any]], - agent_id: str | None, - rollout_events: list[dict[str, Any]] | None, - available_capabilities: Any, -) -> tuple[set[str], set[str], set[str]]: - """Return exact declared Todo ids that are runnable or validly waiting. - - The planning source is complete canonical state, not a presentation lane. - Existing Todo predicates retain ownership, exclusion, task-class, and - lifecycle semantics; the TS resume evaluator remains the authority for a - linked external wait. - """ - - declared_ids = { - todo_id - for declaration in declarations - for todo_id in declaration.candidate_todo_ids - } - matched_items = [ - item - for item in agent_todo_source_items - if isinstance(item, dict) - and normalize_todo_id(item.get("todo_id")) in declared_ids - ] - resume_items = [ - item - for item in matched_items - if normalize_todo_resume_when(item.get("resume_when")) - ] - resume_conditions = ( - evaluate_todo_resume_conditions( - resume_items, - source_items=agent_todo_source_items, - rollout_events=rollout_events, - available_capabilities=available_capabilities, - ) - if resume_items - else {} - ) - evaluated_source_items: list[dict[str, Any]] = [] - for source_item in agent_todo_source_items: - evaluated_item = dict(source_item) - source_todo_id = normalize_todo_id(source_item.get("todo_id")) - condition = resume_conditions.get(source_todo_id or "") - if isinstance(condition, dict): - evaluated_item["resume_condition"] = condition - evaluated_item["resume_ready"] = condition.get("satisfied") is True - evaluated_source_items.append(evaluated_item) - ordinary_waiting_ids = { - todo_id - for todo_id in ( - normalize_todo_id(item.get("todo_id")) - for item in project_todo_resume_planning( - {"items": evaluated_source_items}, - agent_id=agent_id, - )["blocked_successor_items"] - if isinstance(item, dict) - ) - if todo_id - } - - runnable_ids: set[str] = set() - waiting_ids: set[str] = set() - uncertain_ids: set[str] = set() - for item in matched_items: - todo_id = normalize_todo_id(item.get("todo_id")) - if not todo_id: - continue - if todo_item_task_class(item) != TODO_TASK_CLASS_ADVANCEMENT: - continue - if not todo_item_claimed_by_agent_or_unclaimed(item, agent_id=agent_id): - continue - - resume_when = normalize_todo_resume_when(item.get("resume_when")) - if resume_when: - status = normalize_todo_status(item.get("status")) or TODO_STATUS_OPEN - if status not in {TODO_STATUS_OPEN, TODO_STATUS_DEFERRED}: - continue - condition = resume_conditions.get(todo_id) - if not isinstance(condition, dict): - uncertain_ids.add(todo_id) - continue - if condition.get("invalid_target") is True or condition.get( - "invalid_state" - ): - continue - if condition.get("provider_required") is True: - # Capability evidence was unavailable, so this lookup is - # uncertain rather than proof that the fallback is absent. - uncertain_ids.add(todo_id) - continue - if condition.get("satisfied") is True: - runnable_ids.add(todo_id) - elif condition.get("satisfied") is False: - if condition.get("kind") == "monitor_changed": - # An unchanged generation is a durable wait only while - # its continuous-monitor target remains open. A terminal - # monitor cannot produce another generation, so accepting - # that condition as pending would hide the fallback gap - # forever. A changed generation is handled above and - # remains runnable even if the monitor closed concurrently. - target_status = normalize_todo_status( - condition.get("target_status") - ) - if ( - target_status == TODO_STATUS_OPEN - and condition.get("target_task_class") - == TODO_TASK_CLASS_MONITOR - ): - waiting_ids.add(todo_id) - elif todo_id in ordinary_waiting_ids: - waiting_ids.add(todo_id) - continue - - if todo_item_is_actionable_open(item): - runnable_ids.add(todo_id) - - return runnable_ids, waiting_ids, uncertain_ids - - def declared_fallback_gap_from_agent_vision( agent_vision: dict[str, Any] | None, *, @@ -327,143 +144,82 @@ def declared_fallback_gap_from_agent_vision( rollout_events: list[dict[str, Any]] | None = None, available_capabilities: Any = None, ) -> dict[str, Any] | None: - """Project one advisory gap for an unresolved declared fallback. - - A fallback direction is declared structurally via the agent vision's - typed ``fallback_declarations`` contract, which the TS-owned Vision - prepare validates and persists. Prose mentions never declare a - fallback, and generic ``todo_delta`` actions are not fallback - declarations on their own. - - The declared direction is resolved when one of: - 1. A linked Todo sits on the authoritative agent-scoped selectable - advancement frontier (peer-claimed primary-path Todos do not); - 2. A bounded successor Todo is created or reopened specifically for - this fallback direction; or - 3. The vision records an explicit terminal disposition (closed-family state - or path_delta.outcome=stop). - - When the primary path is blocked and none of the resolutions holds, the - declared fallback would otherwise disappear silently behind the - blocked-successor wait state, which clears the ordinary acceptance gaps. - This advisory gap stays in the independent ``fallback_gaps`` projection - field and never enters the acceptance-gap replan stream. - """ - - if not isinstance(agent_vision, dict): - return None - if _vision_has_terminal_disposition(agent_vision): - return None - if not _blocked_primary_waiting( - agent_todo_summary, - agent_id=agent_id, - ): - return None - + """Encode persisted facts; TypeScript owns the advisory disposition.""" declarations = parse_fallback_declarations(agent_vision) if not declarations: return None - - source_is_authoritative = isinstance(agent_todo_source_items, list) - uncertain_todo_ids: set[str] = set() - if isinstance(agent_todo_source_items, list): - ( - selectable_ids, - waiting_todo_ids, - uncertain_todo_ids, - ) = _authoritative_fallback_disposition_ids( - declarations, - agent_todo_source_items=agent_todo_source_items or [], - agent_id=agent_id, - rollout_events=rollout_events, - available_capabilities=available_capabilities, - ) - elif agent_todo_source_items is FallbackTodoReadState.UNAVAILABLE: - selectable_ids = set() - waiting_todo_ids = set() - else: - # Legacy direct callers may only have a compact display summary. It - # remains valid positive evidence, but omission from a bounded lane is - # uncertainty and must not be projected as authoritative absence. - selectable_ids = agent_scoped_selectable_advancement_todo_ids( - agent_todo_summary, - agent_id=agent_id, - ) - waiting_todo_ids = _blocked_successor_todo_ids( - agent_todo_summary, - agent_id=agent_id, - ) - todo_delta = parse_vision_todo_delta_entries(agent_vision.get("todo_delta")) - created_or_reopened_ids = { - todo_id - for action, todo_id in todo_delta - if action in VISION_TODO_DELTA_SUCCESSOR_ACTIONS - } - - unresolved_ids: set[str] = set() - lookup_uncertain_ids: set[str] = set(uncertain_todo_ids) - for declaration in declarations: - candidate_ids = declaration.candidate_todo_ids - waiting_todo_ids - if not candidate_ids: - if not declaration.candidate_todo_ids: - unresolved_ids.add(declaration.unresolved_todo_id) - continue - # Disposition 1: Runnable on authoritative selectable advancement frontier - if candidate_ids & selectable_ids: - continue - # Disposition 2: Bounded successor created/reopened specifically for this fallback - if candidate_ids & created_or_reopened_ids: - continue - if candidate_ids & lookup_uncertain_ids: - continue - if ( - declaration.successor_todo_id - and declaration.successor_todo_id in created_or_reopened_ids - ): - continue - if not source_is_authoritative: - lookup_uncertain_ids.update(candidate_ids) - continue - - unresolved_id = declaration.unresolved_todo_id - if unresolved_id not in waiting_todo_ids: - unresolved_ids.add(unresolved_id) - - if not unresolved_ids and not lookup_uncertain_ids: + assert isinstance(agent_vision, dict) + summary = agent_todo_summary if isinstance(agent_todo_summary, dict) else {} + source = agent_todo_source_items + items = [ + { + **item, + "todo_id": normalize_todo_id(item.get("todo_id")), + "status": normalize_todo_status(item.get("status")) or "open", + "task_class": todo_item_task_class(item), + "claimed_by": normalize_todo_claimed_by(item.get("claimed_by")), + "excluded_agents": normalize_todo_excluded_agents(item.get("excluded_agents")), + "done": item.get("done") is True, + "removed_continuation": todo_item_has_removed_continuation_policy(item), + "resume_when": normalize_todo_resume_when(item.get("resume_when")), + "resume_monitor_generation": normalize_todo_generation(item.get("resume_monitor_generation")), + } + for item in source if isinstance(item, dict) and normalize_todo_id(item.get("todo_id")) + ] if isinstance(source, list) else [] + resume_request = build_todo_resume_evaluation_request( + items, source_items=items, rollout_events=rollout_events, + available_capabilities=available_capabilities, + ) + # Reuse the resume codec's bounded evidence, never send arbitrary Todo prose + # or raw rollout events into this read-only projection. + facts = [ + {**row, **{key: item[key] for key in + ("excluded_agents", "done", "removed_continuation")}} + for row, item in zip(resume_request["source_items"], items, strict=True) + ] + path = agent_vision.get("path_delta") + path = path if isinstance(path, dict) else {} + result = effect_runtime_result("goal.fallback_disposition.project", { + "schema_version": "goal_fallback_disposition_request_v0", + "terminal": goal_vision_state_is_closed(agent_vision.get("state")) + or str(path.get("outcome") or "").strip().lower() == VISION_FALLBACK_TERMINAL_PATH_OUTCOME, + "agent_id": normalize_todo_claimed_by(agent_id), + "blocker_present": bool(summary.get("current_agent_blocker_items")), + "resume_planning": build_todo_resume_planning_request(summary, agent_id=agent_id), + "source_state": "complete" if isinstance(source, list) else + "unavailable" if source is FallbackTodoReadState.UNAVAILABLE else "omitted", + "items": facts, + "resume_evaluation": { + key: value for key, value in resume_request.items() + if key not in {"items", "source_items"} + }, + "legacy_selectable_ids": sorted(agent_scoped_selectable_advancement_todo_ids( + summary, agent_id=agent_id, + )) if source is None else [], + "declarations": [ + {"candidate_ids": sorted(entry.candidate_todo_ids), + "unresolved_id": entry.unresolved_todo_id} + for entry in declarations + ], + "created_or_reopened_ids": sorted({ + todo_id for action, todo_id in parse_vision_todo_delta_entries(agent_vision.get("todo_delta")) + if action in VISION_TODO_DELTA_SUCCESSOR_ACTIONS + }), + }) + if not isinstance(result, dict) or result.get("schema_version") != "goal_fallback_disposition_v0": + raise RuntimeError("TypeScript fallback disposition shape mismatch") + if result["kind"] == "resolved": return None - - lookup_uncertain_only = not unresolved_ids - - gap: dict[str, Any] = { - "kind": ( - VISION_FALLBACK_LOOKUP_UNCERTAIN_TRIGGER - if lookup_uncertain_only - else VISION_FALLBACK_GAP_TRIGGER - ), - "source": "latest_agent_vision", - "agent_id": agent_vision.get("agent_id"), - "state": agent_vision.get("state"), - "reason_code": ( - VISION_FALLBACK_LOOKUP_UNCERTAIN_REASON_CODE - if lookup_uncertain_only - else VISION_FALLBACK_GAP_REASON_CODE - ), - "recommended_action": ( - VISION_FALLBACK_LOOKUP_UNCERTAIN_ACTION - if lookup_uncertain_only - else VISION_FALLBACK_RECOMMENDED_ACTION - ), + uncertain = result["kind"] == VISION_FALLBACK_LOOKUP_UNCERTAIN_TRIGGER + gap = { + "kind": result["kind"], "source": "latest_agent_vision", + "agent_id": agent_vision.get("agent_id"), "state": agent_vision.get("state"), + "reason_code": VISION_FALLBACK_LOOKUP_UNCERTAIN_REASON_CODE if uncertain else VISION_FALLBACK_GAP_REASON_CODE, + "recommended_action": VISION_FALLBACK_LOOKUP_UNCERTAIN_ACTION if uncertain else VISION_FALLBACK_RECOMMENDED_ACTION, } - unresolved_todo_ids = [todo_id for todo_id in sorted(unresolved_ids) if todo_id][ - :VISION_FALLBACK_RUNNABLE_ITEM_LIMIT - ] - if unresolved_todo_ids: - gap["unresolved_todo_ids"] = unresolved_todo_ids - uncertain_todo_ids = [ - todo_id for todo_id in sorted(lookup_uncertain_ids) if todo_id - ][:VISION_FALLBACK_RUNNABLE_ITEM_LIMIT] - if uncertain_todo_ids: - gap["lookup_uncertain_todo_ids"] = uncertain_todo_ids + for key in ("unresolved_todo_ids", "lookup_uncertain_todo_ids"): + if result[key]: + gap[key] = result[key] generated_at = _compact_text(agent_vision.get("generated_at"), limit=80) if generated_at: gap["generated_at"] = generated_at diff --git a/loopx/control_plane/goals/goal_frontier/fallback_disposition.ts b/loopx/control_plane/goals/goal_frontier/fallback_disposition.ts new file mode 100644 index 0000000000..894ed9d1ec --- /dev/null +++ b/loopx/control_plane/goals/goal_frontier/fallback_disposition.ts @@ -0,0 +1,125 @@ +import type { JsonObject } from "../../effect_program.ts"; +import { EffectRuntimeRequestError } from "../../effect_runtime_errors.ts"; +import { + requireJsonObject, requireStringArray, requireBoolean, + requireNonEmptyString, optionalNonEmptyString, requireStringLiteral, +} from "../../runtime_decode.ts"; +import { + evaluateTodoResumeConditions, diagnoseTodoResumeCondition, + resumeConditionHasKnownPendingTarget, +} from "../../todos/resume_condition.ts"; +import { projectTodoResumePlanning } from "../../todos/resume_planning.ts"; + +type Disposition = "runnable" | "waiting" | "unresolved" | "uncertain"; +interface Declaration { candidates: string[]; unresolved: string } +const RESULT = "goal_fallback_disposition_v0"; + +function objects(value: unknown, label: string): JsonObject[] { + if (!Array.isArray(value)) throw new EffectRuntimeRequestError(`${label} must be an array`); + return value.map((item) => requireJsonObject(item, label)); +} + +function disposition( + item: JsonObject, agent: string | null, condition: JsonObject | undefined, +): Disposition { + const claim = optionalNonEmptyString(item.claimed_by, "claimed_by"); + const excluded = requireStringArray(item.excluded_agents, "excluded_agents"); + if (item.task_class !== "advancement_task" || + (item.role !== undefined && item.role !== "agent") || + (item.archive_state !== undefined && item.archive_state !== "active") || + requireBoolean(item.done, "done") || + requireBoolean(item.removed_continuation, "removed_continuation") || + (agent && (excluded.includes(agent) || (claim !== null && claim !== agent)))) return "unresolved"; + if (!["open", "deferred"].includes(String(item.status))) return "unresolved"; + if (!item.resume_when) return item.status === "open" ? "runnable" : "unresolved"; + if (!condition) return "uncertain"; + const diagnosis = diagnoseTodoResumeCondition(condition, String(item.todo_id)); + if (diagnosis.state === "invalid") return "unresolved"; + if (condition.provider_required === true) return "uncertain"; + if (diagnosis.state === "satisfied") return "runnable"; + // A false ready bit alone is not a durable wait. Reuse the same positive + // target/generation/repository proof as supervision and settlement. + return resumeConditionHasKnownPendingTarget(condition, item) ? "waiting" : "unresolved"; +} + +function authoritativeDispositions( + request: JsonObject, agent: string | null, declarations: Declaration[], +): Map { + const items = objects(request.items, "items"); + const declared = new Set(declarations.flatMap((entry) => entry.candidates)); + const evaluation = requireJsonObject(request.resume_evaluation, "resume_evaluation"); + // The same source rows feed target classification and condition evaluation; + // callers cannot supply a contradictory second set of dependency facts. + const evaluated = evaluateTodoResumeConditions({ + ...evaluation, source_items: items, + items: items.filter((item) => declared.has(String(item.todo_id))), + }); + const conditions = new Map(objects(evaluated.conditions, "conditions").map((row) => + [String(row.todo_id), requireJsonObject(row.condition, "condition")] as const)); + const byId = new Map(); + const counts = new Map(); + for (const item of items) { + const id = requireNonEmptyString(item.todo_id, "todo_id"); + counts.set(id, (counts.get(id) ?? 0) + 1); + } + for (const item of items) { + const id = String(item.todo_id); + if (!declared.has(id)) continue; + const condition = conditions.get(id); + const target = condition?.target_todo_id; + // Duplicate target/dependency identity is uncertainty, never last-row-wins. + byId.set(id, counts.get(id) !== 1 || + (typeof target === "string" && (counts.get(target) ?? 0) > 1) + ? "uncertain" : disposition(item, agent, condition)); + } + return byId; +} + +export function projectFallbackDisposition(value: unknown): JsonObject { + const request = requireJsonObject(value, "fallback disposition"); + if (request.schema_version !== "goal_fallback_disposition_request_v0") + throw new EffectRuntimeRequestError("unsupported fallback disposition schema"); + const source = requireStringLiteral(request.source_state, + ["complete", "unavailable", "omitted"], "source_state"); + const agent = optionalNonEmptyString(request.agent_id, "agent_id"); + const declarations: Declaration[] = objects(request.declarations, "declarations").map((entry) => ({ + candidates: requireStringArray(entry.candidate_ids, "candidate_ids"), + unresolved: requireNonEmptyString(entry.unresolved_id, "unresolved_id"), + })); + if (declarations.length > 4 || declarations.some((entry) => entry.candidates.length > 2)) + throw new EffectRuntimeRequestError("fallback declaration bound exceeded"); + const resolved: JsonObject = { schema_version: RESULT, kind: "resolved", + unresolved_todo_ids: [], lookup_uncertain_todo_ids: [] }; + if (requireBoolean(request.terminal, "terminal") || declarations.length === 0) return resolved; + const planning = projectTodoResumePlanning(request.resume_planning); + const blocked = objects(planning.blocked_successor_items, "blocked_successor_items"); + if (!requireBoolean(request.blocker_present, "blocker_present") && blocked.length === 0) return resolved; + const states = source === "complete" ? authoritativeDispositions(request, agent, declarations) + : new Map(); + if (source === "omitted") { + for (const id of requireStringArray(request.legacy_selectable_ids, "legacy_selectable_ids")) + states.set(id, "runnable"); + for (const row of blocked) states.set(String(row.todo_id), "waiting"); + } + const created = new Set(requireStringArray(request.created_or_reopened_ids, "created_or_reopened_ids")); + const unresolved = new Set(); + const uncertain = new Set(); + for (const entry of declarations) { + // Alternatives resolve a declaration symmetrically. Uncertainty belongs + // only to declarations without a positively resolved alternative. + if (entry.candidates.some((id) => created.has(id) || + states.get(id) === "runnable" || states.get(id) === "waiting")) continue; + const unknown = entry.candidates.filter((id) => states.get(id) === "uncertain" || + (source !== "complete" && !states.has(id))); + if (unknown.length) { + for (const id of unknown) uncertain.add(id); + } else unresolved.add(entry.unresolved); + } + return { + schema_version: RESULT, + kind: unresolved.size ? "vision_fallback_unresolved" + : uncertain.size ? "vision_fallback_lookup_uncertain" : "resolved", + unresolved_todo_ids: [...unresolved].sort().slice(0, 3), + lookup_uncertain_todo_ids: [...uncertain].sort().slice(0, 3), + }; +} diff --git a/loopx/control_plane/goals/goal_frontier/fallback_source.py b/loopx/control_plane/goals/goal_frontier/fallback_source.py new file mode 100644 index 0000000000..18945b9a7c --- /dev/null +++ b/loopx/control_plane/goals/goal_frontier/fallback_source.py @@ -0,0 +1,79 @@ +"""One source snapshot, bounded transport, no fallback-specific state authority.""" + +from __future__ import annotations + +from pathlib import Path +from typing import Any + +from ....history import load_registry +from ....state_refresh import resolve_goal_state +from ...coordination.local_authority import ( + LocalCoordinationAuthorityUnavailable, + read_canonical_todos_if_promoted, +) +from ...effect_runtime import EffectRuntimeRemoteError +from ...todos.active_state_todo_parser import parse_todo_source +from ...todos.contract import normalize_todo_id, normalize_todo_resume_when +from .fallback_disposition import ( + FallbackTodoReadState, FallbackTodoSource, parse_fallback_declarations, +) +from .semantic_history import latest_agent_vision_from_status_payload + + +def read_fallback_source_snapshot( + *, registry_path: Path, runtime_root: Path, goal_id: str, +) -> list[dict[str, Any]]: + canonical = read_canonical_todos_if_promoted(runtime_root=runtime_root, goal_id=goal_id) + if canonical is not None: + # Failure after promotion must not consult the Markdown projection. + return canonical["todos"] + goal, _, state_file = resolve_goal_state( + registry=load_registry(registry_path), goal_id=goal_id, + project_override=None, state_file_override=None, + ) + if goal is None: + raise ValueError("fallback source goal is absent") + groups, archived, _ = parse_todo_source(state_file.read_text(encoding="utf-8")) + return [*groups["user"], *groups["agent"], *archived] + + +def live_fallback_authority_items( + status_payload: dict[str, Any], *, registry_path: Path, runtime_root: Path, + goal_id: str, agent_id: str | None, +) -> FallbackTodoSource: + vision = latest_agent_vision_from_status_payload( + status_payload, goal_id=goal_id, agent_id=agent_id, + ) + requested = { + todo_id for declaration in parse_fallback_declarations(vision) + for todo_id in declaration.candidate_todo_ids + } + if not requested: + return None + try: + source = read_fallback_source_snapshot( + registry_path=registry_path, runtime_root=runtime_root, goal_id=goal_id, + ) + except (EffectRuntimeRemoteError, LocalCoordinationAuthorityUnavailable, OSError, ValueError): + return FallbackTodoReadState.UNAVAILABLE + if not isinstance(source, list) or any(not isinstance(item, dict) for item in source): + return FallbackTodoReadState.UNAVAILABLE + by_id: dict[str, list[dict[str, Any]]] = {} + for item in source: + todo_id = normalize_todo_id(item.get("todo_id")) + if todo_id: + by_id.setdefault(todo_id, []).append(item) + selected = set(requested) + for todo_id in requested: + for item in by_id.get(todo_id, []): + resume = normalize_todo_resume_when(item.get("resume_when")) + kind, _, target = (resume or "").partition(":") + if kind in {"todo_done", "monitor_changed"}: + dependency = normalize_todo_id(target) + if dependency: + selected.add(dependency) + # Four declarations name at most eight alternatives and eight direct + # dependencies. Do not follow dependency chains or re-read another revision. + if any(len(by_id.get(todo_id, [])) > 1 for todo_id in selected): + return FallbackTodoReadState.UNAVAILABLE + return [by_id[todo_id][0] for todo_id in sorted(selected) if todo_id in by_id] diff --git a/loopx/control_plane/quota/live_decision.py b/loopx/control_plane/quota/live_decision.py index aa1e795d18..05a9a536eb 100644 --- a/loopx/control_plane/quota/live_decision.py +++ b/loopx/control_plane/quota/live_decision.py @@ -6,19 +6,8 @@ from typing import Any from ...quota import build_quota_should_run -from ...todos import list_goal_todos from ..agent_context import project_agent_context -from ..coordination.local_authority import LocalCoordinationAuthorityUnavailable -from ..effect_runtime import EffectRuntimeRemoteError -from ..goals.goal_frontier.fallback_disposition import ( - FallbackTodoReadState, - FallbackTodoSource, - parse_fallback_declarations, -) -from ..goals.goal_frontier.semantic_history import ( - latest_agent_vision_from_status_payload, -) -from ..todos.contract import normalize_todo_id, normalize_todo_resume_when +from ..goals.goal_frontier.fallback_source import live_fallback_authority_items from ..capability_hooks import ( InteractionProjectionHookRegistration, dispatch_interaction_projection_hooks, @@ -40,95 +29,6 @@ BoundedResearchFrontierProjector = Callable[..., Mapping[str, Any] | None] -def _fallback_authority_todo_ids( - status_payload: dict[str, Any], - *, - goal_id: str, - agent_id: str | None, -) -> set[str]: - vision = latest_agent_vision_from_status_payload( - status_payload, - goal_id=goal_id, - agent_id=agent_id, - ) - return { - todo_id - for declaration in parse_fallback_declarations(vision) - for todo_id in declaration.candidate_todo_ids - } - - -def _live_fallback_authority_items( - status_payload: dict[str, Any], - *, - registry_path: Path, - runtime_root: Path, - goal_id: str, - agent_id: str | None, -) -> FallbackTodoSource: - """Read only the exact canonical Todos needed by fallback disposition. - - ``status`` is deliberately presentation-bounded, so omission from it can - never prove that a declared fallback is absent. The live CLI path owns the - registry/runtime authority needed for exact reads. Resume dependencies are - included only when directly referenced by a declared Todo. Four declarations - can name eight targets/successors and eight direct dependencies: at most 16 - reads. The TypeScript evaluator inspects a dependency's current state, not - its own resume chain, and retains state-transition authority. - """ - - requested_ids = _fallback_authority_todo_ids( - status_payload, - goal_id=goal_id, - agent_id=agent_id, - ) - if not requested_ids: - return None - - items: dict[str, dict[str, Any]] = {} - pending_ids = set(requested_ids) - for include_dependencies in (True, False): - dependency_ids: set[str] = set() - for todo_id in sorted(pending_ids): - try: - projection = list_goal_todos( - registry_path=registry_path, - goal_id=goal_id, - todo_id=todo_id, - runtime_root_arg=str(runtime_root), - limit=None, - ) - except ( - EffectRuntimeRemoteError, - LocalCoordinationAuthorityUnavailable, - OSError, - ValueError, - ): - return FallbackTodoReadState.UNAVAILABLE - if projection.get("ambiguous") is True: - return FallbackTodoReadState.UNAVAILABLE - item = projection.get("todo") - if item is None: - if projection.get("not_found") is True: - continue - return FallbackTodoReadState.UNAVAILABLE - if not isinstance(item, dict) or normalize_todo_id(item.get("todo_id")) != todo_id: - return FallbackTodoReadState.UNAVAILABLE - items[todo_id] = dict(item) - resume_when = normalize_todo_resume_when(item.get("resume_when")) - if include_dependencies and resume_when: - resume_kind, _, target = resume_when.partition(":") - dependency_id = ( - normalize_todo_id(target) - if resume_kind in {"todo_done", "monitor_changed"} - else None - ) - if dependency_id: - dependency_ids.add(dependency_id) - pending_ids = dependency_ids - requested_ids - return list(items.values()) - - def _fresh_read_covers_all_pending_material( dispatch: Mapping[str, Any] | None, projector: Callable[..., dict[str, Any]] | None, @@ -542,7 +442,7 @@ def build_live_quota_should_run_decision( fresh_operator_inbox_read = _fresh_operator_inbox_read_required( turn_start_hook_dispatch ) - authoritative_fallback_todo_items = _live_fallback_authority_items( + authoritative_fallback_todo_items = live_fallback_authority_items( decision_status_payload, registry_path=registry_path, runtime_root=runtime_root, diff --git a/loopx/control_plane/todos/resume_condition.py b/loopx/control_plane/todos/resume_condition.py index f76015f0c0..0bcffb3580 100644 --- a/loopx/control_plane/todos/resume_condition.py +++ b/loopx/control_plane/todos/resume_condition.py @@ -297,15 +297,15 @@ def normalize_todo_resume_when_via_runtime(value: Any) -> str | None: return result -def evaluate_todo_resume_conditions( +def build_todo_resume_evaluation_request( items: list[dict[str, Any]], *, source_items: list[dict[str, Any]], rollout_events: list[dict[str, Any]] | None = None, available_capabilities: Any = None, kinds: list[str] | None = None, -) -> dict[str, dict[str, Any]]: - """Return TS-owned resume conditions keyed by the waiting Todo id.""" +) -> dict[str, Any]: + """Encode bounded resume facts for direct or composed typed projections.""" request: dict[str, Any] = { "schema_version": TODO_RESUME_EVALUATION_REQUEST_SCHEMA_VERSION, @@ -328,6 +328,22 @@ def evaluate_todo_resume_conditions( ) if kinds is not None: request["kinds"] = kinds + return request + + +def evaluate_todo_resume_conditions( + items: list[dict[str, Any]], + *, + source_items: list[dict[str, Any]], + rollout_events: list[dict[str, Any]] | None = None, + available_capabilities: Any = None, + kinds: list[str] | None = None, +) -> dict[str, dict[str, Any]]: + """Return TS-owned resume conditions keyed by the waiting Todo id.""" + request = build_todo_resume_evaluation_request( + items, source_items=source_items, rollout_events=rollout_events, + available_capabilities=available_capabilities, kinds=kinds, + ) try: result = effect_runtime_result("todo.resume_condition.evaluate", request) except EffectRuntimeRejected as exc: diff --git a/tests/control_plane/test_fallback_wait_admission.py b/tests/control_plane/test_fallback_wait_admission.py index d77f7d37c6..2eef513ab3 100644 --- a/tests/control_plane/test_fallback_wait_admission.py +++ b/tests/control_plane/test_fallback_wait_admission.py @@ -9,7 +9,10 @@ from loopx.cli import main as cli_main from loopx.control_plane.goals.goal_vision import normalize_goal_vision_packet -from loopx.control_plane.quota import live_decision +from loopx.control_plane.goals.goal_frontier import fallback_disposition, fallback_source +from canonical_authority_fixture import initialize_canonical_authority, isolate_sqlite_runtime +from loopx.control_plane.coordination.runtime_shadow import build_todo_runtime_shadow_projection +from loopx.control_plane.todos.todo_summary import TODO_ITEM_SCHEMA_VERSION GOAL_ID = "fallback-wait-capacity-fixture" AGENT_ID = "worker" @@ -222,7 +225,7 @@ def fail_authority_read(**_kwargs: object) -> dict: raise OSError("fixture canonical authority unavailable") monkeypatch.setattr( - "loopx.control_plane.quota.live_decision.list_goal_todos", + "loopx.control_plane.goals.goal_frontier.fallback_source.read_fallback_source_snapshot", fail_authority_read, ) @@ -278,37 +281,47 @@ def test_real_cli_reads_only_direct_fallback_dependencies( ) state_file.write_text(state) reads: list[str] = [] - original_read = live_decision.list_goal_todos + original_read = fallback_source.read_fallback_source_snapshot - def record_read(**kwargs: object) -> dict: - reads.append(str(kwargs["todo_id"])) + def record_read(**kwargs: object) -> list[dict]: + reads.append(str(kwargs["goal_id"])) return original_read(**kwargs) - monkeypatch.setattr(live_decision, "list_goal_todos", record_read) + monkeypatch.setattr(fallback_source, "read_fallback_source_snapshot", record_read) + transported: list[int] = [] + original_projection = fallback_disposition.effect_runtime_result + + def record_projection(method: str, params: dict) -> object: + if method == "goal.fallback_disposition.project": + transported.append(len(params["items"])) + return original_projection(method, params) + + monkeypatch.setattr(fallback_disposition, "effect_runtime_result", record_projection) assert cli_main(args) == 0 result = json.loads(capsys.readouterr().out) assert result["should_run"] is False assert "fallback_gaps" not in result["goal_frontier_projection"] # The direct prerequisite is deferred regardless of its own dependency. # Its chain cannot alter this wait or add canonical lookup calls. - assert reads == [FALLBACK_TODO_ID, PREREQUISITE_TODO_ID] + assert reads == [GOAL_ID] + assert transported and set(transported) == {2} @pytest.mark.parametrize("unrelated_deferred_count", [0, 20]) -def test_real_cli_mismatched_authority_identity_is_uncertain( +def test_real_cli_malformed_authority_snapshot_is_uncertain( tmp_path: Path, capsys: pytest.CaptureFixture[str], monkeypatch: pytest.MonkeyPatch, unrelated_deferred_count: int, ) -> None: args = _write_fixture(tmp_path, unrelated_deferred_count=unrelated_deferred_count) - original_read = live_decision.list_goal_todos + original_read = fallback_source.read_fallback_source_snapshot - def mismatched_read(**kwargs: object) -> dict: + def mismatched_read(**kwargs: object) -> list: projection = original_read(**kwargs) - return {**projection, "todo": {**projection["todo"], "todo_id": "todo_other"}} + return [*projection, None] - monkeypatch.setattr(live_decision, "list_goal_todos", mismatched_read) + monkeypatch.setattr(fallback_source, "read_fallback_source_snapshot", mismatched_read) assert cli_main(args) == 0 result = json.loads(capsys.readouterr().out) gap = result["goal_frontier_projection"]["fallback_gaps"][0] @@ -332,7 +345,7 @@ def test_real_cli_without_declarations_does_not_read_fallback_authority( def unexpected_read(**_kwargs: object) -> dict: pytest.fail("No fallback declaration authorizes an exact fallback lookup") - monkeypatch.setattr(live_decision, "list_goal_todos", unexpected_read) + monkeypatch.setattr(fallback_source, "read_fallback_source_snapshot", unexpected_read) assert cli_main(args) == 0 result = json.loads(capsys.readouterr().out) assert result["should_run"] is False @@ -410,3 +423,58 @@ def test_real_cli_ambiguous_canonical_fallback_read_is_uncertain( assert gap["kind"] == "vision_fallback_lookup_uncertain" assert gap["lookup_uncertain_todo_ids"] == [FALLBACK_TODO_ID] assert "unresolved_todo_ids" not in gap + + +@pytest.mark.parametrize("provider", ["file", "sqlite"]) +@pytest.mark.parametrize("display", ["stale", "missing"]) +@pytest.mark.parametrize("dependency_status", ["open", "done"]) +def test_real_provider_cli_uses_one_snapshot_including_archived_dependency( + tmp_path: Path, capsys: pytest.CaptureFixture[str], monkeypatch: pytest.MonkeyPatch, + provider: str, display: str, dependency_status: str, +) -> None: + if provider == "sqlite": + isolate_sqlite_runtime(tmp_path, monkeypatch) + args = _write_fixture(tmp_path, unrelated_deferred_count=20, + prerequisite_status=dependency_status) + state_file = tmp_path / "ACTIVE_GOAL_STATE.md" + runtime_root = tmp_path / "runtime" + rows = fallback_source.read_fallback_source_snapshot( + registry_path=tmp_path / "registry.json", runtime_root=runtime_root, goal_id=GOAL_ID, + ) + for row in rows: + row["schema_version"] = TODO_ITEM_SCHEMA_VERSION + if dependency_status == "done": + for row in rows: + if row["todo_id"] == PREREQUISITE_TODO_ID: + row["archive_state"] = "archive" + projection = build_todo_runtime_shadow_projection( + goal_id=GOAL_ID, todos=rows, handoff_mode="soft_claim", + ) + initialize_canonical_authority(runtime_root, GOAL_ID, projection, + state_path=state_file, provider=provider) + before = fallback_source.read_canonical_todos_if_promoted( + runtime_root=runtime_root, goal_id=GOAL_ID, + ) + if display == "missing": + state_file.unlink() + else: + state_file.write_text("# Stale display\n\n## Agent Todo\n") + saved_display = state_file.read_bytes() if state_file.exists() else None + snapshots: list[int] = [] + original = fallback_source.read_fallback_source_snapshot + + def read_once(**kwargs: object) -> list[dict]: + snapshot = original(**kwargs) + snapshots.append(len(snapshot)) + return snapshot + + monkeypatch.setattr(fallback_source, "read_fallback_source_snapshot", read_once) + assert cli_main(args) == 0 + result = json.loads(capsys.readouterr().out) + assert "fallback_gaps" not in result["goal_frontier_projection"] + assert snapshots == [len(rows)] + after = fallback_source.read_canonical_todos_if_promoted( + runtime_root=runtime_root, goal_id=GOAL_ID, + ) + assert before == after + assert (state_file.read_bytes() if state_file.exists() else None) == saved_display diff --git a/tests/control_plane/test_goal_frontier_fallback_disposition.py b/tests/control_plane/test_goal_frontier_fallback_disposition.py index 5b74048630..a8eea06324 100644 --- a/tests/control_plane/test_goal_frontier_fallback_disposition.py +++ b/tests/control_plane/test_goal_frontier_fallback_disposition.py @@ -37,6 +37,38 @@ ) +@pytest.mark.parametrize("resolution", ["waiting", "runnable"]) +@pytest.mark.parametrize("reverse", [False, True]) +def test_one_resolved_alternative_discharges_only_its_declaration( + resolution: str, reverse: bool, +) -> None: + from loopx.control_plane.goals.goal_frontier.fallback_disposition import ( + declared_fallback_gap_from_agent_vision, + ) + + ready = quota_todo_item(todo_id="todo_ready", index=1, text="Deliver alternative") + source = [ready] + if resolution == "waiting": + ready.update(status="deferred", resume_when="todo_done:todo_dependency") + source.append(quota_todo_item(todo_id="todo_dependency", index=2, text="Dependency")) + else: + source.append(quota_todo_item( + todo_id="todo_uncertain", index=2, text="Other alternative", + status="deferred", resume_when="capacity_available:delivery", + )) + alternatives = ["todo_missing" if resolution == "waiting" else "todo_uncertain", "todo_ready"] + if reverse: + alternatives.reverse() + vision = {"state": "vision_drift_detected", "fallback_declarations": [{ + "declaration_id": "direction", "target_todo_id": alternatives[0], + "successor_todo_id": alternatives[1], + }]} + assert declared_fallback_gap_from_agent_vision( + vision, agent_todo_summary={"current_agent_blocker_items": [{}]}, + agent_id=AGENT_ID, agent_todo_source_items=source, + ) is None + + def _fallback_vision_run( *, state: str = "vision_drift_detected", diff --git a/tests/control_plane_ts/fallback_disposition.test.ts b/tests/control_plane_ts/fallback_disposition.test.ts new file mode 100644 index 0000000000..251f475d94 --- /dev/null +++ b/tests/control_plane_ts/fallback_disposition.test.ts @@ -0,0 +1,98 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import type { JsonObject } from "../../loopx/control_plane/effect_program.ts"; +import { projectFallbackDisposition } from "../../loopx/control_plane/goals/goal_frontier/fallback_disposition.ts"; + +function todo(id: string, overrides: JsonObject = {}): JsonObject { + return { todo_id: id, role: "agent", task_class: "advancement_task", status: "open", + archive_state: "active", done: false, removed_continuation: false, + excluded_agents: [], ...overrides }; +} +function request(items: JsonObject[], overrides: JsonObject = {}): JsonObject { + return { schema_version: "goal_fallback_disposition_request_v0", terminal: false, + source_state: "complete", agent_id: "worker", blocker_present: true, + declarations: [{ candidate_ids: ["todo_fallback"], unresolved_id: "todo_fallback" }], + items, legacy_selectable_ids: [], created_or_reopened_ids: [], + resume_evaluation: { schema_version: "todo_resume_evaluation_request_v0", rollout_events: [] }, + resume_planning: { schema_version: "todo_resume_planning_request_v0", + agent_id: "worker", item_limit: 5, has_deferred_count: false, + has_visible_deferred_count: false, deferred_count: null, available_capabilities: null, + sources: Object.fromEntries(["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"].map((key) => [key, []])), + }, ...overrides }; +} + +for (const [name, patch] of Object.entries({ + "peer claim": { claimed_by: "peer" }, "exclusion": { excluded_agents: ["worker"] }, + "archived": { archive_state: "archive" }, "terminal": { status: "done" }, + "contradictory done": { done: true }, "removed continuation": { removed_continuation: true }, + "monitor": { task_class: "continuous_monitor" }, "user role": { role: "user" }, +})) { + test(`fallback cannot resolve from ${name}`, () => { + assert.equal(projectFallbackDisposition(request([todo("todo_fallback", patch)])).kind, + "vision_fallback_unresolved"); + }); +} + +test("dependency facts, not cached readiness or supplied secondary rows, decide waits", () => { + const waiter = todo("todo_fallback", { status: "deferred", + resume_when: "todo_done:todo_dependency", resume_ready: true }); + const input = request([waiter]); + input.resume_evaluation = { ...(input.resume_evaluation as JsonObject), + source_items: [todo("todo_dependency", { status: "done" })] }; + assert.equal(projectFallbackDisposition(input).kind, "vision_fallback_unresolved"); + input.items = [waiter, todo("todo_dependency")]; + assert.equal(projectFallbackDisposition(input).kind, "resolved"); + input.items = [waiter, todo("todo_dependency", { status: "done", archive_state: "archive" })]; + assert.equal(projectFallbackDisposition(input).kind, "resolved"); +}); + +test("only repository-bound PR waits count as known pending evidence", () => { + const waiting = todo("todo_fallback", { resume_when: "pr_merged:#23" }); + assert.equal(projectFallbackDisposition(request([waiting])).kind, "vision_fallback_unresolved"); + waiting.task_repository = "git:github.com/example/project"; + assert.equal(projectFallbackDisposition(request([waiting])).kind, "resolved"); +}); + +test("duplicate alternatives or prerequisites are uncertain, independent of row order", () => { + const waiting = todo("todo_fallback", { resume_when: "todo_done:todo_dependency" }); + for (const items of [ + [waiting, waiting], [waiting, todo("todo_dependency"), todo("todo_dependency", { status: "done" })], + ]) { + for (const rows of [items, [...items].reverse()]) { + assert.equal(projectFallbackDisposition(request(rows)).kind, "vision_fallback_lookup_uncertain"); + } + } +}); + +test("resolving one declaration never clears another, and projection is read-only", () => { + const input = request([todo("todo_fallback")], { + declarations: [ + { candidate_ids: ["todo_fallback", "todo_missing"], unresolved_id: "todo_missing" }, + { candidate_ids: ["todo_other"], unresolved_id: "todo_other" }, + ], + }); + const before = structuredClone(input); + const result = projectFallbackDisposition(input); + assert.deepEqual(result.unresolved_todo_ids, ["todo_other"]); + assert.deepEqual(input, before); + assert.equal("execution_obligation" in result, false); +}); + +test("unavailable source cannot be overridden by display positives; omitted source can", () => { + for (const [state, expected] of [ + ["unavailable", "vision_fallback_lookup_uncertain"], ["omitted", "resolved"], + ]) assert.equal(projectFallbackDisposition(request([], { + source_state: state, legacy_selectable_ids: ["todo_fallback"], + })).kind, expected); +}); + +test("no-declaration and terminal projections are elided, malformed wire data rejected", () => { + assert.equal(projectFallbackDisposition(request([], { terminal: true })).kind, "resolved"); + assert.equal(projectFallbackDisposition(request([], { declarations: [] })).kind, "resolved"); + for (const patch of [{ schema_version: "future" }, { source_state: "partial" }, + { terminal: "false" }, { declarations: Array(5).fill({ candidate_ids: [], unresolved_id: "direction" }) }, + { items: [todo("todo_fallback", { excluded_agents: "worker" })] }]) + assert.throws(() => projectFallbackDisposition(request([], patch))); +}); From a432b578b978df2a513744ddb3856e04ecf49b6e Mon Sep 17 00:00:00 2001 From: huangruiteng Date: Sat, 12 Sep 2026 11:31:24 +0800 Subject: [PATCH 6/7] docs(rfc): record typed fallback consumer closure and semantic corrections Signed-off-by: huangruiteng --- ...shared-goal-authority-state-provider-v0.md | 8 +++++ ...-goal-authority-state-provider-v0.zh-CN.md | 6 ++++ .../typescript-control-plane-migration-v0.md | 11 +++++++ ...script-control-plane-migration-v0.zh-CN.md | 8 +++++ .../goal-vision-replan-contract-v0.md | 32 +++++++++++++++---- 5 files changed, 59 insertions(+), 6 deletions(-) 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 9ce7e05802..faa524445d 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md @@ -2680,6 +2680,14 @@ It does not qualify D1/D2, alter the default provider, or relax D3 promotion hol Markdown remains the permanent one-way display; the remaining execution cards below are unchanged. +Optional fallback advice now consumes a single source snapshot, including direct +archived prerequisites, and delegates dispositions to the existing typed resume +semantics. This retires per-Todo source reads and Python wait/aggregation rules; +it does not add a provider, relax a fence, qualify durability, or make Markdown +authoritative after promotion. Read failure remains uncertainty, never permission +to consult stale display. See the TS RFC's T3 consumer notes for the intentional +alternative-resolution corrections; D1–D3 and the permanent projection plan stay. + Use the [TS execution cards](typescript-control-plane-migration-v0.md#execution-cards-after-the-current-stack) for command inventory, update/monitor transactions and consumer deletion. Do not repeat that plan in a second implementation or treat a merged read-policy PR 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 197fcb6a5b..821fb0a2a3 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 @@ -2137,6 +2137,12 @@ T3 decision-dependency 读取策略现由同一 TS owner 解释 scope coverage 不改变默认 provider 或放宽 D3 promotion hold。Markdown 继续作为永久单向展示, 后续执行卡与退役条件保留。 +可选 fallback 提示改为消费一次来源快照(含直接归档依赖),由 TS 复用既有 resume +语义判断去向,退役逐 Todo 来源读取与 Python 等待/聚合规则。没有新增 provider、 +放宽 fence 或完成持久化资格化;promotion 后 Markdown 不恢复 authority。读取失败 +仍为 uncertainty,不能回退陈旧展示。候选去向的有意修正见 TS RFC 的 T3 说明; +D1–D3 与永久投影规划不变。 + **D1 — 资格化永久投影交付,可与 T1/T2 重叠推进。** 能力缺口 consumer 在 legacy/canonical 输入上共用 TS requirement/resolution owner, diff --git a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md index 48d12a57bc..114facce13 100644 --- a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md +++ b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md @@ -577,6 +577,17 @@ Unidentifiable referenced history still needs explicit repair, not a post-promot Markdown fallback. This closes the demonstrated dependency omission, not all history import, provider qualification, soak, or D3 cutover requirements. +Optional Vision fallback advice now composes resume evaluation and positive +pending-target proof in `goal_frontier/fallback_disposition.ts`. Python retains +persisted-input codecs, source I/O and presentation, not a second Monitor-wait or +declaration aggregation rule. Live quota reads one complete canonical/legacy +snapshot and selects at most 16 target/direct-dependency records, including +archive; it no longer performs repeated per-Todo queries. Intentional corrections: +alternatives resolve their own declaration symmetrically, uncertainty is scoped +to unresolved declarations, and invalid/ambiguous wait evidence cannot suppress +advice. This closes one T3 consumer, not all quota snapshot consistency or T1/T2; +fallback declarations remain optional advisory input, never execution authority. + - Audit Turn/quota, Dashboard, standing decisions, shared-goal alignment and amendment revision inputs. Reuse #4117's canonical source adapter and pass one snapshot through a decision; do not build another Todo inventory. 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 a33d6740ef..8611cad34c 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 @@ -441,6 +441,14 @@ user role。历史节点不会进入活动工作或 lease lane。无法识别的 不得恢复 promotion 后的 Markdown fallback。本批闭合已复现的依赖遗漏,不代表所有 历史导入、provider 资格化、soak 或 D3 cutover 条件均完成。 +可选 Vision fallback 提示在 `goal_frontier/fallback_disposition.ts` 内组合已有 +resume evaluator 与合法 pending target 正向证明。Python 仅保留持久输入 codec、 +来源 I/O 和展示,删除独立 Monitor 等待分支与声明聚合规则。Live quota 一次读取 +完整 canonical/legacy 快照,再选择最多 16 个目标/直接依赖记录(含归档),不再 +逐 Todo 重复查询。有意修正:候选对称地消除所属声明,uncertainty 仅属于尚未解决 +的声明,非法/含混等待证据不能隐藏提示。这只闭合一个 T3 consumer,不代表整个 +quota 的快照一致性或 T1/T2 完成;声明仍是可选 advisory,不是执行 authority。 + - 分别审计 Turn/quota、Dashboard、standing decision、shared-goal alignment、 amendment revision 输入。复用 #4117 canonical source adapter,一次决策传递一份 snapshot,不新增 Todo inventory。 diff --git a/docs/reference/protocols/goal-vision-replan-contract-v0.md b/docs/reference/protocols/goal-vision-replan-contract-v0.md index a5298df203..57b3357da0 100644 --- a/docs/reference/protocols/goal-vision-replan-contract-v0.md +++ b/docs/reference/protocols/goal-vision-replan-contract-v0.md @@ -462,17 +462,37 @@ rule uses existing acceptance/lineage facts regardless of the optional advisory relationships, and it grants no additional authority. Ownership, exclusions, capabilities, user gates, and quota remain independent execution constraints. -For optional fallback advice, live quota reads the declared target/successor -Todos and one layer of direct resume dependencies, at most 16 exact reads. -Deeper dependency chains do not expand that lookup. Existing typed resume, -ownership, and lifecycle rules decide whether each declared path is available. -An unavailable, ambiguous, or mismatched canonical read produces +For optional fallback advice, live quota takes one complete Todo source snapshot: +canonical provider records after promotion, otherwise one legacy source read. +Only the declared targets/successors and one layer of direct resume dependencies +are transported to the typed projection (at most 16 records, including archived +prerequisites). Deeper chains do not expand the selection or add reads. No +declarations means no additional source read. This does not promise an atomic +snapshot across the entire quota/status packet. +Existing typed resume, ownership, and lifecycle rules decide availability. +An unavailable, ambiguous, or malformed source read produces `vision_fallback_lookup_uncertain`; compact display evidence cannot override a -failed authority read. Only explicit canonical not-found proves absence. +failed authority read. Only absence from a complete source proves not-found. Pending `todo_done:` keeps an unresolved fallback; a valid `monitor_changed` generation condition may establish a wait. `fallback_gaps` remains advisory and adds no replan obligation. +Disposition is per declaration, not per candidate: a runnable, validly waiting, +or bounded create/reopen alternative resolves its own declaration regardless of +the other candidate's order, absence or uncertainty. It cannot clear a different +declaration. Unavailable source/provider evidence stays uncertain; it is not proof of a +valid wait. Waiting uses the shared positive target/generation/repository proof; +an unbound PR reference, archived fallback itself, or contradictory done flag +cannot suppress the advisory. Archived completed prerequisites remain valid +evidence. No fallback declaration is required to discover canonical work, and +none grants execution, user-decision or settlement authority. + +可选 fallback 提示从一次完整来源快照读取,最多输送 16 个直接相关记录,保留归档 +依赖;无声明不额外读取。按整条声明判断:一个可执行或有合法等待证据的候选即可 +消除该声明的提示,但不能替其他声明清账。候选顺序、另一个候选缺失或不确定不再 +误伤有效去向。等待复用 TS 的目标、generation 和仓库绑定正向证明;这仍是 advisory, +不要求 Agent 为发现 canonical 工作额外维护声明,也不授予执行、审批或结算权限。 + 等待资格现在逐项检查已有 acceptance 的 Todo 关联,并在展示裁剪前从完整来源计算。 A 的等待不能遮住尚未落实的 B;有可执行工作则继续,相关工作都具有合法等待证据才暂缓。 无需另外维护 fallback 声明;无关 Todo 的数量和顺序不应改变决策。 From 8c520fec18d35f8c6220a8b7c417c8eecae430f9 Mon Sep 17 00:00:00 2001 From: huangruiteng Date: Sat, 12 Sep 2026 11:36:26 +0800 Subject: [PATCH 7/7] fix(frontier): bound preloaded fallback sources at the shared codec boundary Signed-off-by: huangruiteng --- .../goal_frontier/fallback_disposition.py | 31 +++++++++++++++++++ .../goals/goal_frontier/fallback_source.py | 24 ++------------ ...test_goal_frontier_fallback_disposition.py | 22 +++++++++++++ 3 files changed, 55 insertions(+), 22 deletions(-) diff --git a/loopx/control_plane/goals/goal_frontier/fallback_disposition.py b/loopx/control_plane/goals/goal_frontier/fallback_disposition.py index cb214f84e5..db1b106aac 100644 --- a/loopx/control_plane/goals/goal_frontier/fallback_disposition.py +++ b/loopx/control_plane/goals/goal_frontier/fallback_disposition.py @@ -135,6 +135,33 @@ def parse_fallback_declarations( return declarations +def select_fallback_source_items( + source: list[dict[str, Any]], requested: set[str], +) -> FallbackTodoSource: + """Bound transport for every caller, including already-loaded writeback sources.""" + if not isinstance(source, list) or any(not isinstance(item, dict) for item in source): + return FallbackTodoReadState.UNAVAILABLE + by_id: dict[str, list[dict[str, Any]]] = {} + for item in source: + todo_id = normalize_todo_id(item.get("todo_id")) + if todo_id: + by_id.setdefault(todo_id, []).append(item) + selected = set(requested) + for todo_id in requested: + for item in by_id.get(todo_id, []): + resume = normalize_todo_resume_when(item.get("resume_when")) + kind, _, target = (resume or "").partition(":") + if kind in {"todo_done", "monitor_changed"}: + dependency = normalize_todo_id(target) + if dependency: + selected.add(dependency) + # Four declarations name at most eight alternatives and eight direct + # dependencies. Do not follow dependency chains or re-read another revision. + if any(len(by_id.get(todo_id, [])) > 1 for todo_id in selected): + return FallbackTodoReadState.UNAVAILABLE + return [by_id[todo_id][0] for todo_id in sorted(selected) if todo_id in by_id] + + def declared_fallback_gap_from_agent_vision( agent_vision: dict[str, Any] | None, *, @@ -151,6 +178,10 @@ def declared_fallback_gap_from_agent_vision( assert isinstance(agent_vision, dict) summary = agent_todo_summary if isinstance(agent_todo_summary, dict) else {} source = agent_todo_source_items + if isinstance(source, list): + source = select_fallback_source_items(source, { + todo_id for entry in declarations for todo_id in entry.candidate_todo_ids + }) items = [ { **item, diff --git a/loopx/control_plane/goals/goal_frontier/fallback_source.py b/loopx/control_plane/goals/goal_frontier/fallback_source.py index 18945b9a7c..0a6fcc8317 100644 --- a/loopx/control_plane/goals/goal_frontier/fallback_source.py +++ b/loopx/control_plane/goals/goal_frontier/fallback_source.py @@ -13,9 +13,9 @@ ) from ...effect_runtime import EffectRuntimeRemoteError from ...todos.active_state_todo_parser import parse_todo_source -from ...todos.contract import normalize_todo_id, normalize_todo_resume_when from .fallback_disposition import ( FallbackTodoReadState, FallbackTodoSource, parse_fallback_declarations, + select_fallback_source_items, ) from .semantic_history import latest_agent_vision_from_status_payload @@ -56,24 +56,4 @@ def live_fallback_authority_items( ) except (EffectRuntimeRemoteError, LocalCoordinationAuthorityUnavailable, OSError, ValueError): return FallbackTodoReadState.UNAVAILABLE - if not isinstance(source, list) or any(not isinstance(item, dict) for item in source): - return FallbackTodoReadState.UNAVAILABLE - by_id: dict[str, list[dict[str, Any]]] = {} - for item in source: - todo_id = normalize_todo_id(item.get("todo_id")) - if todo_id: - by_id.setdefault(todo_id, []).append(item) - selected = set(requested) - for todo_id in requested: - for item in by_id.get(todo_id, []): - resume = normalize_todo_resume_when(item.get("resume_when")) - kind, _, target = (resume or "").partition(":") - if kind in {"todo_done", "monitor_changed"}: - dependency = normalize_todo_id(target) - if dependency: - selected.add(dependency) - # Four declarations name at most eight alternatives and eight direct - # dependencies. Do not follow dependency chains or re-read another revision. - if any(len(by_id.get(todo_id, [])) > 1 for todo_id in selected): - return FallbackTodoReadState.UNAVAILABLE - return [by_id[todo_id][0] for todo_id in sorted(selected) if todo_id in by_id] + return select_fallback_source_items(source, requested) diff --git a/tests/control_plane/test_goal_frontier_fallback_disposition.py b/tests/control_plane/test_goal_frontier_fallback_disposition.py index a8eea06324..8d4a7bc7ed 100644 --- a/tests/control_plane/test_goal_frontier_fallback_disposition.py +++ b/tests/control_plane/test_goal_frontier_fallback_disposition.py @@ -69,6 +69,28 @@ def test_one_resolved_alternative_discharges_only_its_declaration( ) is None +def test_loaded_writeback_source_is_bounded_before_effect_transport(monkeypatch) -> None: + from loopx.control_plane.goals.goal_frontier import fallback_disposition as module + + source = [quota_todo_item(todo_id=f"todo_unrelated_{index}", index=index + 1, + text="Unrelated work") for index in range(4096)] + source.append(quota_todo_item(todo_id=FALLBACK_ID, index=4097, text="Fallback")) + original = module.effect_runtime_result + counts = [] + + def capture(method, request): + counts.append(len(request["items"])) + return original(method, request) + + monkeypatch.setattr(module, "effect_runtime_result", capture) + assert module.declared_fallback_gap_from_agent_vision( + {"fallback_declarations": [{"declaration_id": "direction", "target_todo_id": FALLBACK_ID}]}, + agent_todo_summary={"current_agent_blocker_items": [{}]}, + agent_id=AGENT_ID, agent_todo_source_items=source, + ) is None + assert counts == [1] + + def _fallback_vision_run( *, state: str = "vision_drift_detected",