diff --git a/loopx/control_plane/goals/goal_frontier/__init__.py b/loopx/control_plane/goals/goal_frontier/__init__.py index f71765ccff..74c7146a7c 100644 --- a/loopx/control_plane/goals/goal_frontier/__init__.py +++ b/loopx/control_plane/goals/goal_frontier/__init__.py @@ -47,6 +47,14 @@ autonomous_replan_ack_satisfies_obligation, replan_successor_transition_ack, ) +from .fallback_disposition import ( + VISION_FRONTIER_TODO_DELTA_ACTIONS, # noqa: F401 + FallbackDeclaration, # noqa: F401 + agent_scoped_selectable_advancement_todo_ids, # noqa: F401 + declared_fallback_gap_from_agent_vision, + parse_fallback_declarations, # noqa: F401 + parse_vision_todo_delta_entries, +) from .long_todo_chain import ( LONG_TODO_CHAIN_TRIGGER, classify_long_todo_chain_ack, @@ -91,13 +99,6 @@ TODO_SUCCESSION_GAP_TRIGGER = TODO_SUCCESSION_WARNING_REASON_CODE TODO_TASK_CLASS_ADVANCEMENT = "advancement_task" TODO_TASK_CLASS_MONITOR = "continuous_monitor" -VISION_FRONTIER_TODO_DELTA_ACTIONS = { - "activate", - "create", - "reopen", - "resume", - "retain", -} def safe_non_negative_int(value: Any) -> int: @@ -393,11 +394,9 @@ def acceptance_gaps_from_agent_vision( gap["acceptance_summary"] = acceptance vision_todo_ids = [ todo_id - for value in (agent_vision.get("todo_delta") or []) - if isinstance(value, str) - and (parts := value.strip().partition(":"))[1] - and parts[0].strip().lower() in VISION_FRONTIER_TODO_DELTA_ACTIONS - and (todo_id := _compact_projection_text(parts[2], limit=120)) + for _, todo_id in parse_vision_todo_delta_entries( + agent_vision.get("todo_delta") + ) ] if vision_todo_ids: gap["vision_todo_ids"] = list(dict.fromkeys(vision_todo_ids)) @@ -1416,6 +1415,7 @@ def build_goal_frontier_projection_from_summaries( replan_obligation: dict[str, Any] | None, acceptance_gaps: list[dict[str, Any]] | None = None, vision_wait_state: dict[str, Any] | None = None, + fallback_gaps: list[dict[str, Any]] | None = None, ) -> dict[str, Any]: user_counts = _summary_task_counts(user_todo_summary) agent_counts = _summary_task_counts(agent_todo_summary) @@ -1453,6 +1453,7 @@ def build_goal_frontier_projection_from_summaries( replan_obligation=replan_obligation, acceptance_gaps=acceptance_gaps, vision_wait_state=vision_wait_state, + fallback_gaps=fallback_gaps, deferred_successors=_deferred_successors( agent_todo_summary, agent_id=agent_id, @@ -1602,6 +1603,17 @@ def build_goal_frontier_projection_context_from_status( ), ) acceptance_gaps = [] if vision_wait_state else source_acceptance_gaps + declared_fallback_gaps = [ + gap + for gap in ( + declared_fallback_gap_from_agent_vision( + latest_agent_vision, + agent_todo_summary=agent_todo_summary, + agent_id=agent_id, + ), + ) + if isinstance(gap, dict) + ] projected_replan_ack = projected_autonomous_replan_ack_for_agent( item, project_asset, @@ -1716,6 +1728,7 @@ def build_goal_frontier_projection_context_from_status( replan_obligation=replan_obligation, acceptance_gaps=acceptance_gaps, vision_wait_state=vision_wait_state, + fallback_gaps=declared_fallback_gaps, ) if latest_replan_ack_feedback: goal_frontier_projection["replan_ack_feedback"] = ( @@ -1821,6 +1834,7 @@ def build_goal_frontier_projection( acceptance_gaps: list[dict[str, Any]] | None = None, deferred_successors: dict[str, Any] | None = None, vision_wait_state: dict[str, Any] | None = None, + fallback_gaps: list[dict[str, Any]] | None = None, ) -> dict[str, Any]: replan_required = autonomous_replan_is_required(replan_obligation) blockers: list[str] = [] @@ -1883,6 +1897,11 @@ def build_goal_frontier_projection( } if vision_continuation_audit: projection["vision_continuation_audit"] = vision_continuation_audit + # Advisory-only field: unlike acceptance_gaps it is never cleared by the + # blocked-successor wait state, which is exactly when a declared fallback + # would otherwise disappear silently. + if fallback_gaps: + projection["fallback_gaps"] = fallback_gaps[:1] if isinstance(vision_wait_state, dict): projection["vision_wait_state"] = vision_wait_state if replan_required and isinstance(replan_obligation, dict): diff --git a/loopx/control_plane/goals/goal_frontier/fallback_disposition.py b/loopx/control_plane/goals/goal_frontier/fallback_disposition.py new file mode 100644 index 0000000000..c27cf0f28c --- /dev/null +++ b/loopx/control_plane/goals/goal_frontier/fallback_disposition.py @@ -0,0 +1,289 @@ +from __future__ import annotations + +from dataclasses import dataclass +from typing import Any + +from ...todos.contract import ( + normalize_todo_id, +) +from ...todos.deferred_resume import todo_summary_blocked_successor_items +from ...todos.projection import ( + agent_scoped_selectable_advancement_todo_ids, +) +from ..goal_vision_state import goal_vision_state_is_closed + +# Single owner of the vision todo_delta action contract shared by the +# acceptance-gap projection and this module. +VISION_FRONTIER_TODO_DELTA_ACTIONS = frozenset( + {"activate", "create", "reopen", "resume", "retain"} +) +# create/reopen entries are bounded successor declarations and resolve the +# 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_TODO_DELTA_ID_LIMIT = 120 +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_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 " + "successor, or record an explicit terminal no-follow-up disposition; " + "do not invent a user gate" +) + + +@dataclass(frozen=True) +class FallbackDeclaration: + """Structured declaration of a fallback direction and its associated work.""" + + declaration_id: str + target_todo_id: str | None = None + successor_todo_id: str | None = None + + @property + def candidate_todo_ids(self) -> set[str]: + return { + todo_id + 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 + + +def _compact_text(value: Any, *, limit: int) -> str: + return " ".join(str(value or "").strip().split())[:limit] + + +def parse_vision_todo_delta_entries(entries: Any) -> list[tuple[str, str]]: + """Parse ``action:todo_id`` vision todo_delta entries once for consumers.""" + + parsed: list[tuple[str, str]] = [] + for value in entries or []: + if not isinstance(value, str): + continue + action, separator, raw_todo_id = value.strip().partition(":") + todo_id = _compact_text(raw_todo_id, limit=VISION_TODO_DELTA_ID_LIMIT) + normalized_action = action.strip().lower() + if ( + separator + and todo_id + and normalized_action in (VISION_FRONTIER_TODO_DELTA_ACTIONS) + ): + parsed.append((normalized_action, todo_id)) + return parsed + + +def parse_fallback_declarations( + agent_vision: dict[str, Any] | None, +) -> list[FallbackDeclaration]: + """Parse typed fallback declarations written by the TS Vision contract. + + The only supported authoring path is ``agent_vision.fallback_declarations`` + as validated and persisted by the TS-owned ``goal.vision_checkpoint`` + prepare (and mirrored through the status/shared-runtime compact read + model). Prose mentions, generic ``todo_delta`` actions, and legacy alias + shapes are not declarations. + """ + + declarations: list[FallbackDeclaration] = [] + if not isinstance(agent_vision, dict): + return declarations + source = agent_vision.get("fallback_declarations") + if not isinstance(source, list): + return declarations + + seen: set[tuple[str, str | None, str | None]] = set() + for raw in source[:VISION_FALLBACK_DECLARATION_ENTRY_LIMIT]: + if not isinstance(raw, dict): + continue + declaration_id = _compact_text( + raw.get("declaration_id"), + limit=VISION_TODO_DELTA_ID_LIMIT, + ) + if not declaration_id: + continue + target_todo_id = normalize_todo_id(raw.get("target_todo_id")) + successor_todo_id = normalize_todo_id(raw.get("successor_todo_id")) + key = (declaration_id, target_todo_id, successor_todo_id) + if key in seen: + continue + seen.add(key) + declarations.append( + FallbackDeclaration( + declaration_id=declaration_id, + target_todo_id=target_todo_id, + successor_todo_id=successor_todo_id, + ) + ) + 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 todo_summary_blocked_successor_items( + agent_todo_summary, + agent_id=agent_id, + ) + 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( + todo_summary_blocked_successor_items( + agent_todo_summary, + agent_id=agent_id, + ) + ) + + +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 declared_fallback_gap_from_agent_vision( + agent_vision: dict[str, Any] | None, + *, + agent_todo_summary: dict[str, Any] | None, + agent_id: str | 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 + + declarations = parse_fallback_declarations(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, + ) + 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() + for declaration in declarations: + candidate_ids = declaration.candidate_todo_ids - waiting_todo_ids + if not candidate_ids: + 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 ( + declaration.successor_todo_id + and declaration.successor_todo_id in created_or_reopened_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: + return None + + gap: dict[str, Any] = { + "kind": 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, + } + 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 + generated_at = _compact_text(agent_vision.get("generated_at"), limit=80) + if generated_at: + gap["generated_at"] = generated_at + return {key: value for key, value in gap.items() if value is not None} diff --git a/loopx/control_plane/goals/goal_frontier/semantic_history.py b/loopx/control_plane/goals/goal_frontier/semantic_history.py index bf92e520ca..796b7c49f0 100644 --- a/loopx/control_plane/goals/goal_frontier/semantic_history.py +++ b/loopx/control_plane/goals/goal_frontier/semantic_history.py @@ -167,6 +167,8 @@ def latest_agent_vision_from_runs( } if isinstance(vision.get("path_delta"), dict): result["path_delta"] = vision["path_delta"] + if isinstance(vision.get("fallback_declarations"), list): + result["fallback_declarations"] = vision["fallback_declarations"] return result return None diff --git a/loopx/control_plane/goals/goal_vision.py b/loopx/control_plane/goals/goal_vision.py index b44235a46a..1e8e33af52 100644 --- a/loopx/control_plane/goals/goal_vision.py +++ b/loopx/control_plane/goals/goal_vision.py @@ -44,6 +44,11 @@ "total_limit", "total_usage", ) +# Mirrors the TS-owned prepare contract for bounded typed fallback +# declarations so compaction cannot drop a declared fallback direction. +GOAL_VISION_FALLBACK_DECLARATION_ENTRY_LIMIT = 4 +GOAL_VISION_FALLBACK_DECLARATION_ID_LIMIT = 120 +GOAL_VISION_FALLBACK_DECLARATION_FIELDS = ("target_todo_id", "successor_todo_id") def _compact_public_text(value: Any, *, limit: int) -> str | None: @@ -82,6 +87,33 @@ def _compact_goal_path_delta(value: Any) -> dict[str, Any] | None: return compact if len(compact) > 1 else None +def _compact_fallback_declarations(value: Any) -> list[dict[str, str]]: + if not isinstance(value, list): + return [] + declarations: list[dict[str, str]] = [] + seen: set[str] = set() + for raw in value[:GOAL_VISION_FALLBACK_DECLARATION_ENTRY_LIMIT]: + if not isinstance(raw, dict): + continue + declaration_id = _compact_public_text( + raw.get("declaration_id"), + limit=GOAL_VISION_FALLBACK_DECLARATION_ID_LIMIT, + ) + if not declaration_id or declaration_id in seen: + continue + seen.add(declaration_id) + entry = {"declaration_id": declaration_id} + for field in GOAL_VISION_FALLBACK_DECLARATION_FIELDS: + text = _compact_public_text( + raw.get(field), + limit=GOAL_VISION_FALLBACK_DECLARATION_ID_LIMIT, + ) + if text: + entry[field] = text + declarations.append(entry) + return declarations + + def compact_goal_vision_packet(value: Any) -> dict[str, Any] | None: """Return the public read-path shape of an agent goal-vision packet.""" @@ -93,7 +125,9 @@ def compact_goal_vision_packet(value: Any) -> dict[str, Any] | None: if text: compact[field] = text - patch = value.get("vision_patch") if isinstance(value.get("vision_patch"), dict) else {} + patch = ( + value.get("vision_patch") if isinstance(value.get("vision_patch"), dict) else {} + ) compact_patch: dict[str, str] = {} for field, limit in GOAL_VISION_FIELD_LIMITS.items(): text = _compact_public_text(patch.get(field), limit=limit) @@ -116,7 +150,15 @@ def compact_goal_vision_packet(value: Any) -> dict[str, Any] | None: if todo_delta: compact["todo_delta"] = todo_delta - budget = value.get("vision_budget") if isinstance(value.get("vision_budget"), dict) else {} + declarations = _compact_fallback_declarations(value.get("fallback_declarations")) + if declarations: + compact["fallback_declarations"] = declarations + + budget = ( + value.get("vision_budget") + if isinstance(value.get("vision_budget"), dict) + else {} + ) compact_budget = { field: budget[field] for field in GOAL_VISION_BUDGET_COMPACT_FIELDS @@ -125,7 +167,9 @@ def compact_goal_vision_packet(value: Any) -> dict[str, Any] | None: if compact_budget: compact["vision_budget"] = compact_budget - validation = value.get("validation") if isinstance(value.get("validation"), dict) else {} + validation = ( + value.get("validation") if isinstance(value.get("validation"), dict) else {} + ) compact_validation = { field: validation[field] for field in ("budget_checked", "budget_status", "write_correctness_checked") diff --git a/loopx/control_plane/goals/vision_checkpoint.ts b/loopx/control_plane/goals/vision_checkpoint.ts index 043a1c9aba..e175a22271 100644 --- a/loopx/control_plane/goals/vision_checkpoint.ts +++ b/loopx/control_plane/goals/vision_checkpoint.ts @@ -58,6 +58,14 @@ const GOAL_PATH_DELTA_LIST_LIMITS = { unresolved_questions: [2, 140], evidence_refs: [4, 140], } as const; +// Bounded typed fallback declarations survive prepare unchanged so the +// declared direction cannot disappear behind later read-model compaction. +const VISION_FALLBACK_DECLARATION_ENTRY_LIMIT = 4; +const VISION_FALLBACK_DECLARATION_ID_LIMIT = 120; +const VISION_FALLBACK_DECLARATION_FIELDS = [ + "target_todo_id", + "successor_todo_id", +] as const; const GOAL_VISION_STATE_ALIASES: Readonly> = { closed: "vision_closed", satisfied: "vision_closed", @@ -359,6 +367,63 @@ function normalizeGoalPathDelta( return [normalized, fieldUsage]; } +function normalizeFallbackDeclarations( + value: unknown, +): [JsonObject[], Record] | null { + if (value === null || value === undefined) return null; + if (!Array.isArray(value)) { + throw new EffectRuntimeRequestError( + "agent_vision.fallback_declarations must be a JSON array", + ); + } + if (value.length === 0) return null; + if (value.length > VISION_FALLBACK_DECLARATION_ENTRY_LIMIT) { + throw new EffectRuntimeRequestError( + `agent_vision.fallback_declarations has ${value.length} items; limit is ${VISION_FALLBACK_DECLARATION_ENTRY_LIMIT}`, + ); + } + const declarations: JsonObject[] = []; + const seenIds = new Set(); + const fieldUsage: Record = {}; + value.forEach((raw, index) => { + const entry = requiredObject( + raw, + `agent_vision.fallback_declarations[${index}]`, + ); + const declarationId = boundedPublicText( + `fallback_declarations[${index}].declaration_id`, + entry.declaration_id ?? null, + VISION_FALLBACK_DECLARATION_ID_LIMIT, + ); + if (!declarationId) { + throw new EffectRuntimeRequestError( + `agent_vision.fallback_declarations[${index}] requires a non-empty declaration_id`, + ); + } + if (seenIds.has(declarationId)) { + throw new EffectRuntimeRequestError( + `agent_vision.fallback_declarations repeats declaration_id ${JSON.stringify(declarationId)}`, + ); + } + seenIds.add(declarationId); + const declaration: JsonObject = { declaration_id: declarationId }; + fieldUsage[`fallback_declarations[${index}].declaration_id`] = + declarationId.length; + for (const field of VISION_FALLBACK_DECLARATION_FIELDS) { + const text = boundedPublicText( + `fallback_declarations[${index}].${field}`, + entry[field], + VISION_FALLBACK_DECLARATION_ID_LIMIT, + ); + if (!text) continue; + declaration[field] = text; + fieldUsage[`fallback_declarations[${index}].${field}`] = text.length; + } + declarations.push(declaration); + }); + return [declarations, fieldUsage]; +} + function decodePrepareRequest(request: JsonObject): VisionRefreshPrepareRequest { return { phase: "prepare", @@ -438,6 +503,10 @@ function prepareVisionRefresh(request: VisionRefreshPrepareRequest): JsonObject const [pathDelta, pathDeltaUsage] = normalizeGoalPathDelta(updatePacket.path_delta); Object.assign(fieldUsage, pathDeltaUsage); + const [fallbackDeclarations, fallbackUsage] = normalizeFallbackDeclarations( + updatePacket.fallback_declarations, + ) ?? [null, {}]; + Object.assign(fieldUsage, fallbackUsage); const totalUsage = Object.values(fieldUsage).reduce((total, used) => total + used, 0); if (totalUsage > GOAL_VISION_TOTAL_LIMIT) { throw visionBudgetError( @@ -475,6 +544,8 @@ function prepareVisionRefresh(request: VisionRefreshPrepareRequest): JsonObject for (const [field, [, itemLimit]] of Object.entries(GOAL_PATH_DELTA_LIST_LIMITS)) { fieldLimits[`path_delta.${field}[]`] = itemLimit; } + fieldLimits["fallback_declarations"] = VISION_FALLBACK_DECLARATION_ENTRY_LIMIT; + fieldLimits["fallback_declarations[]"] = VISION_FALLBACK_DECLARATION_ID_LIMIT; const agentVision: JsonObject = { schema_version: GOAL_VISION_REPLAN_SCHEMA_VERSION, @@ -494,6 +565,9 @@ function prepareVisionRefresh(request: VisionRefreshPrepareRequest): JsonObject validation, }; if (pathDelta !== null) agentVision.path_delta = pathDelta; + if (fallbackDeclarations !== null) { + agentVision.fallback_declarations = fallbackDeclarations; + } if ( request.require_path_delta_for_durable_change && diff --git a/loopx/control_plane/todos/projection.py b/loopx/control_plane/todos/projection.py index c6b570ed26..048c0fa04e 100644 --- a/loopx/control_plane/todos/projection.py +++ b/loopx/control_plane/todos/projection.py @@ -149,9 +149,9 @@ def todo_claimed_visibility_items( if len(selected) >= limit: break - return sorted(selected, key=lambda item: original_index.get(id(item), TODO_MISSING_INDEX))[ - :limit - ] + return sorted( + selected, key=lambda item: original_index.get(id(item), TODO_MISSING_INDEX) + )[:limit] def todo_item_task_text( @@ -160,9 +160,7 @@ def todo_item_task_text( keys: tuple[str, ...] = ("title", "text"), ) -> str: return " ".join( - str(item.get(key) or "") - for key in keys - if str(item.get(key) or "").strip() + str(item.get(key) or "") for key in keys if str(item.get(key) or "").strip() ) @@ -197,7 +195,9 @@ def todo_item_expires_at(item: dict[str, Any]) -> datetime | None: return monitor_todo_expires_at(item) -def todo_item_is_expired_monitor(item: dict[str, Any], *, now: datetime | None = None) -> bool: +def todo_item_is_expired_monitor( + item: dict[str, Any], *, now: datetime | None = None +) -> bool: return monitor_todo_is_expired(item, now=now) @@ -245,38 +245,33 @@ def todo_item_claimed_by_agent_or_unclaimed( return not claimed_by or claimed_by == normalized_agent_id -def todo_advancement_frontier_counts( +def todo_advancement_frontier_items( summary: dict[str, Any] | None, *, agent_id: str | None, -) -> dict[str, int]: - """Classify the durable advancement frontier by exact claim ownership.""" - +) -> dict[str, list[dict[str, Any]]]: + """Return the authoritative advancement frontier items grouped by claim ownership. + + Preserves the slot precedence of executable backlog first, falling back to + unclaimed priority and claimed advancement open items when the executable backlog + is omitted. Peer-claimed items are tracked separately and excluded from the current + agent's selectable advancement frontier. + """ + + empty: dict[str, list[dict[str, Any]]] = { + "current_agent_claimed_items": [], + "unclaimed_items": [], + "other_agent_claimed_items": [], + } if not isinstance(summary, dict): - return { - "current_agent_claimed_advancement_count": 0, - "unclaimed_advancement_count": 0, - "other_agent_claimed_advancement_count": 0, - } + return empty + normalized_agent_id = normalize_todo_claimed_by(agent_id) - claim_scope = summary.get("claim_scope") - other_items = ( - claim_scope.get("other_agent_claimed_items") - if isinstance(claim_scope, dict) - else [] - ) - diagnostic_other_count = sum( - 1 - for value in other_items or [] - if isinstance(value, dict) - and todo_item_is_actionable_open(value) - and todo_item_task_class(value) == TODO_TASK_CLASS_ADVANCEMENT - ) executable_items = summary.get("executable_backlog_items") if isinstance(executable_items, list): - current_count = 0 - unclaimed_count = 0 - other_count = 0 + current_items: list[dict[str, Any]] = [] + unclaimed_items: list[dict[str, Any]] = [] + other_items: list[dict[str, Any]] = [] for value in executable_items: if not isinstance(value, dict): continue @@ -287,44 +282,123 @@ def todo_advancement_frontier_counts( claimed_by = normalize_todo_claimed_by(value.get("claimed_by")) if claimed_by: if normalized_agent_id and claimed_by == normalized_agent_id: - if not todo_item_excludes_agent(value, agent_id=normalized_agent_id): - current_count += 1 + if not todo_item_excludes_agent( + value, agent_id=normalized_agent_id + ): + current_items.append(value) elif normalized_agent_id: - other_count += 1 + other_items.append(value) else: - current_count += 1 + current_items.append(value) continue if not todo_item_excludes_agent(value, agent_id=normalized_agent_id): - unclaimed_count += 1 + unclaimed_items.append(value) return { - "current_agent_claimed_advancement_count": max( - current_count, - _positive_int(summary.get("current_agent_claimed_advancement_count")), - ), - "unclaimed_advancement_count": unclaimed_count, - # Agent-scoped executable backlogs intentionally omit peer-owned - # work. Preserve that diagnostic lane from claim_scope without - # letting it contribute to the current Agent's selectable count. - "other_agent_claimed_advancement_count": max( - other_count, - diagnostic_other_count, - ), + "current_agent_claimed_items": current_items, + "unclaimed_items": unclaimed_items, + "other_agent_claimed_items": other_items, } - unclaimed_count = sum( - 1 + unclaimed_items = [ + value for value in summary.get("unclaimed_priority_open_items") or [] if isinstance(value, dict) and todo_item_is_actionable_open(value) and todo_item_task_class(value) == TODO_TASK_CLASS_ADVANCEMENT and not todo_item_excludes_agent(value, agent_id=normalized_agent_id) + ] + current_items = [ + value + for value in summary.get("claimed_advancement_open_items") or [] + if isinstance(value, dict) + and todo_item_is_actionable_open(value) + and todo_item_task_class(value) == TODO_TASK_CLASS_ADVANCEMENT + and ( + not normalized_agent_id + or normalize_todo_claimed_by(value.get("claimed_by")) == normalized_agent_id + ) + and not todo_item_excludes_agent(value, agent_id=normalized_agent_id) + ] + other_items = [ + value + for value in summary.get("claimed_advancement_open_items") or [] + if isinstance(value, dict) + and todo_item_is_actionable_open(value) + and todo_item_task_class(value) == TODO_TASK_CLASS_ADVANCEMENT + and normalized_agent_id + and normalize_todo_claimed_by(value.get("claimed_by")) + and normalize_todo_claimed_by(value.get("claimed_by")) != normalized_agent_id + ] + return { + "current_agent_claimed_items": current_items, + "unclaimed_items": unclaimed_items, + "other_agent_claimed_items": other_items, + } + + +def agent_scoped_selectable_advancement_todo_ids( + agent_todo_summary: dict[str, Any] | None, + *, + agent_id: str | None, +) -> set[str]: + """Return the ids the agent-scoped selectable advancement frontier holds. + + Derived directly from the authoritative ``todo_advancement_frontier_items`` + helper so that slot precedence and claim ownership predicates never diverge + from the frontier counter. + """ + + frontier_items = todo_advancement_frontier_items( + agent_todo_summary, + agent_id=agent_id, + ) + selectable: set[str] = set() + for item in ( + frontier_items["current_agent_claimed_items"] + + frontier_items["unclaimed_items"] + ): + if todo_id := normalize_todo_id(item.get("todo_id")): + selectable.add(todo_id) + return selectable + + +def todo_advancement_frontier_counts( + summary: dict[str, Any] | None, + *, + agent_id: str | None, +) -> dict[str, int]: + """Classify the durable advancement frontier by exact claim ownership.""" + + if not isinstance(summary, dict): + return { + "current_agent_claimed_advancement_count": 0, + "unclaimed_advancement_count": 0, + "other_agent_claimed_advancement_count": 0, + } + frontier_items = todo_advancement_frontier_items(summary, agent_id=agent_id) + claim_scope = summary.get("claim_scope") + other_items = ( + claim_scope.get("other_agent_claimed_items") + if isinstance(claim_scope, dict) + else [] + ) + diagnostic_other_count = sum( + 1 + for value in other_items or [] + if isinstance(value, dict) + and todo_item_is_actionable_open(value) + and todo_item_task_class(value) == TODO_TASK_CLASS_ADVANCEMENT ) return { - "current_agent_claimed_advancement_count": _positive_int( - summary.get("current_agent_claimed_advancement_count") + "current_agent_claimed_advancement_count": max( + len(frontier_items["current_agent_claimed_items"]), + _positive_int(summary.get("current_agent_claimed_advancement_count")), + ), + "unclaimed_advancement_count": len(frontier_items["unclaimed_items"]), + "other_agent_claimed_advancement_count": max( + len(frontier_items["other_agent_claimed_items"]), + diagnostic_other_count, ), - "unclaimed_advancement_count": unclaimed_count, - "other_agent_claimed_advancement_count": diagnostic_other_count, } @@ -344,7 +418,8 @@ def todo_item_excludes_agent( normalized_agent_id = normalize_todo_claimed_by(agent_id) return bool( normalized_agent_id - and normalized_agent_id in normalize_todo_excluded_agents(item.get("excluded_agents")) + and normalized_agent_id + in normalize_todo_excluded_agents(item.get("excluded_agents")) ) @@ -674,7 +749,9 @@ def todo_summary_open_task_counts(summary: dict[str, Any] | None) -> dict[str, i } -def todo_summary_has_only_future_scoped_monitor_work(summary: dict[str, Any] | None) -> bool: +def todo_summary_has_only_future_scoped_monitor_work( + summary: dict[str, Any] | None, +) -> bool: """Return true when the scoped agent has only non-due monitor work left.""" agent_id = todo_summary_claim_scope_agent_id(summary) diff --git a/loopx/state_refresh.py b/loopx/state_refresh.py index 4630d393cf..11e61e95c3 100644 --- a/loopx/state_refresh.py +++ b/loopx/state_refresh.py @@ -536,7 +536,7 @@ def _build_state_refresh_output_projections( field: agent_vision.get(field) for field in ( "schema_version", "agent_id", "state", "vision_patch", - "todo_delta", "vision_budget", + "todo_delta", "fallback_declarations", "vision_budget", ) } if isinstance(agent_vision.get("path_delta"), dict): diff --git a/tests/control_plane/test_goal_frontier_fallback_disposition.py b/tests/control_plane/test_goal_frontier_fallback_disposition.py new file mode 100644 index 0000000000..c02190e43c --- /dev/null +++ b/tests/control_plane/test_goal_frontier_fallback_disposition.py @@ -0,0 +1,557 @@ +from __future__ import annotations + +from typing import Any + +import pytest + +from loopx.control_plane.goals.goal_frontier import ( + VISION_FRONTIER_TODO_DELTA_ACTIONS, + agent_scoped_selectable_advancement_todo_ids, + build_goal_frontier_projection_context_from_status, +) +from loopx.control_plane.goals.goal_vision import ( + compact_goal_vision_packet, + normalize_goal_vision_packet, +) +from loopx.control_plane.scheduler.execution_context import ( + GENERIC_CLI_OUTER_CONTROLLER_SCHEDULER_CONTEXT, +) +from loopx.control_plane.testing.quota_fixtures import ( + quota_status_payload, + quota_todo_item, + quota_todo_summary, +) +from loopx.control_plane.todos.projection import todo_advancement_frontier_counts +from loopx.quota import build_quota_should_run + +GOAL_ID = "vision-fallback-disposition-fixture" +AGENT_ID = "codex-fallback-agent" +PRIMARY_AGENT = "codex-primary-agent" +PREREQ_ID = "todo_primary_prereq" +PRIMARY_WAIT_ID = "todo_primary_successor" +FALLBACK_ID = "todo_declared_fallback" +DECLARED_FALLBACK_ACCEPTANCE = ( + "Deliver the primary successor; if the primary stays blocked, " + "deliver the declared fallback direction instead." +) + + +def _fallback_vision_run( + *, + state: str = "vision_drift_detected", + todo_delta: list[str] | None = None, + acceptance_summary: str = DECLARED_FALLBACK_ACCEPTANCE, + path_outcome: str | None = None, + fallback_declarations: list[Any] | None = None, +) -> dict: + """Persist a caller packet through the production write/readback chain. + + The caller packet goes through the real TS ``goal.vision_checkpoint`` + prepare (the executor entry) and the compact read-model projection (the + status/shared-runtime entry) before it becomes a run-history record, so + the tests can only pass when the typed declaration survives the same + chain a real caller uses. + """ + + packet: dict = { + "goal_id": GOAL_ID, + "agent_id": AGENT_ID, + "state": state, + "todo_delta": todo_delta + if todo_delta is not None + else [f"retain:{PRIMARY_WAIT_ID}"], + "vision_patch": { + "acceptance_summary": acceptance_summary, + "replan_trigger_summary": "The primary acceptance remains open.", + "advancement_policy": "repeat_until_closed", + }, + } + if fallback_declarations is not None: + packet["fallback_declarations"] = fallback_declarations + elif acceptance_summary == DECLARED_FALLBACK_ACCEPTANCE: + packet["fallback_declarations"] = [ + { + "declaration_id": "declared_fallback_direction", + "target_todo_id": FALLBACK_ID, + "successor_todo_id": FALLBACK_ID, + } + ] + if path_outcome is not None: + packet["path_delta"] = { + "outcome": path_outcome, + "prior_assumption": "The primary route stays viable after the " + "prerequisite clears.", + "observed_reality": "The declared fallback direction remains the " + "bounded alternative path.", + "stopped": ["Continue only the primary route."], + } + prepared = normalize_goal_vision_packet(packet, goal_id=GOAL_ID, agent_id=AGENT_ID) + compact = compact_goal_vision_packet(prepared) + assert compact is not None + return { + "classification": "vision_fallback_disposition_fixture", + "generated_at": "2026-09-05T00:00:00+00:00", + "agent_id": AGENT_ID, + "progress_scope": "agent_lane", + "agent_vision": compact, + } + + +def _agent_todos(*, fallback_runnable: bool) -> dict: + prereq = quota_todo_item( + todo_id=PREREQ_ID, + index=1, + text="[P0] Complete the primary prerequisite.", + claimed_by=PRIMARY_AGENT, + ) + waiting = quota_todo_item( + todo_id=PRIMARY_WAIT_ID, + index=2, + text="[P0] Resume the primary successor.", + status="deferred", + claimed_by=AGENT_ID, + resume_when=f"todo_done:{PREREQ_ID}", + ) + items = [prereq, waiting] + if fallback_runnable: + items.append( + quota_todo_item( + todo_id=FALLBACK_ID, + index=3, + text="[P1] Deliver the declared fallback direction.", + claimed_by=AGENT_ID, + ) + ) + return quota_todo_summary(items, role="agent") + + +def _status_payload( + *, + fallback_runnable: bool, + latest_runs: list[dict], +) -> dict: + return quota_status_payload( + goal_id=GOAL_ID, + status="active", + recommended_action="Resolve the declared fallback direction.", + agent_todos=_agent_todos(fallback_runnable=fallback_runnable), + coordination={ + "agent_model": "peer_v1", + "registered_agents": [PRIMARY_AGENT, AGENT_ID], + }, + latest_runs=latest_runs, + ) + + +def _frontier_projection(payload: dict) -> dict: + item = payload["attention_queue"]["items"][0] + context = build_goal_frontier_projection_context_from_status( + goal_id=GOAL_ID, + agent_id=AGENT_ID, + status_payload=payload, + item=item, + project_asset=item["project_asset"], + user_todo_summary=item["user_todos"], + agent_todo_summary=item["agent_todos"], + work_lane_contract=None, + neutral_replan_ack_classifications=set(), + registered_agent_ids=[PRIMARY_AGENT, AGENT_ID], + goal_status="active", + ) + return context["goal_frontier_projection"] + + +def test_blocked_primary_with_runnable_fallback_projects_todo_selectable() -> None: + payload = _status_payload( + fallback_runnable=True, + latest_runs=[_fallback_vision_run(todo_delta=[f"retain:{FALLBACK_ID}"])], + ) + + frontier = _frontier_projection(payload) + assert "fallback_gaps" not in frontier + assert "vision_wait_state" not in frontier + remaining = frontier["remaining_advancement_frontier"] + assert remaining["current_agent_claimed_advancement_count"] == 1 + assert remaining["unclaimed_advancement_count"] == 0 + + decision = build_quota_should_run( + payload, + goal_id=GOAL_ID, + agent_id=AGENT_ID, + scheduler_execution_context=(GENERIC_CLI_OUTER_CONTROLLER_SCHEDULER_CONTEXT), + ) + assert decision["decision"] == "run" + assert decision["should_run"] is True + assert decision["selected_todo"]["todo_id"] == FALLBACK_ID + assert "fallback_gaps" not in decision["goal_frontier_projection"] + + +def test_declared_fallback_survives_prepare_compact_and_readback() -> None: + # The typed declaration must survive the real TS prepare, the compact + # read model, and the history readback before any gap can be projected. + run = _fallback_vision_run( + todo_delta=[f"retain:{PRIMARY_WAIT_ID}", f"retain:{FALLBACK_ID}"] + ) + vision = run["agent_vision"] + + assert vision["fallback_declarations"] == [ + { + "declaration_id": "declared_fallback_direction", + "target_todo_id": FALLBACK_ID, + "successor_todo_id": FALLBACK_ID, + } + ] + from loopx.control_plane.goals.goal_frontier.semantic_history import ( + latest_agent_vision_from_runs, + ) + + readback = latest_agent_vision_from_runs([run], goal_id=GOAL_ID, agent_id=AGENT_ID) + assert readback is not None + assert readback["fallback_declarations"] == vision["fallback_declarations"] + + +def test_declared_fallback_without_resolution_projects_single_gap() -> None: + # The structured declaration links the fallback direction to a Todo id, + # but no runnable Todo with that id exists on this agent's frontier. + payload = _status_payload( + fallback_runnable=False, + latest_runs=[ + _fallback_vision_run( + todo_delta=[f"retain:{PRIMARY_WAIT_ID}", f"retain:{FALLBACK_ID}"] + ) + ], + ) + + frontier = _frontier_projection(payload) + + # The blocked-successor wait state clears ordinary acceptance gaps; the + # declared fallback would disappear silently without the dedicated field. + assert frontier["acceptance_gaps"] == [] + wait = frontier["vision_wait_state"] + assert wait["reason_code"] == "exact_blocked_successor" + assert wait["selected_todo_id"] == PRIMARY_WAIT_ID + gaps = frontier["fallback_gaps"] + assert len(gaps) == 1 + gap = gaps[0] + assert gap["kind"] == "vision_fallback_unresolved" + assert gap["reason_code"] == "declared_fallback_without_runnable_or_terminal" + assert gap["agent_id"] == AGENT_ID + assert gap["unresolved_todo_ids"] == [FALLBACK_ID] + assert "fallback" in gap["recommended_action"] + assert "do not invent a user gate" in gap["recommended_action"] + assert frontier["replan_required"] is False + + decision = build_quota_should_run( + payload, + goal_id=GOAL_ID, + agent_id=AGENT_ID, + 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] + + +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. + payload = _status_payload( + fallback_runnable=False, + latest_runs=[ + _fallback_vision_run( + todo_delta=[f"retain:{PRIMARY_WAIT_ID}"], + fallback_declarations=[], + ) + ], + ) + + frontier = _frontier_projection(payload) + assert "fallback_gaps" not in frontier + + +@pytest.mark.parametrize( + "acceptance_summary", + [ + # Owner probe 1: explicit negation must not project a fallback gap. + "No fallback is authorized; wait for the primary prerequisite.", + # Owner probe 2: a non-English prose declaration has no structured + # declaration channel, so the conservative projection yields no gap. + "主路径阻塞时,执行已声明的备用方案。", + # English prose mentioning a fallback is equally non-declarative. + "Deliver the primary successor after its prerequisite clears; " + "the fallback wording lives in prose only.", + ], + ids=["negated-english", "chinese-prose", "english-prose"], +) +def test_prose_text_alone_never_declares_a_fallback( + acceptance_summary: str, +) -> None: + payload = _status_payload( + fallback_runnable=False, + latest_runs=[ + _fallback_vision_run( + acceptance_summary=acceptance_summary, + todo_delta=[f"retain:{PRIMARY_WAIT_ID}"], + fallback_declarations=[], + ) + ], + ) + + frontier = _frontier_projection(payload) + + assert "fallback_gaps" not in frontier + + +def test_other_agent_primary_todo_without_fallback_declaration_generates_no_gap() -> ( + None +): + # Maintainer blocker 2: generic primary-path retain (e.g. peer-held prerequisite) + # is not a fallback declaration; without a structured declaration, no gap is invented. + payload = _status_payload( + fallback_runnable=False, + latest_runs=[ + _fallback_vision_run( + todo_delta=[f"retain:{PRIMARY_WAIT_ID}", f"retain:{PREREQ_ID}"], + fallback_declarations=[], + ) + ], + ) + + frontier = _frontier_projection(payload) + + assert "fallback_gaps" not in frontier + + +def test_other_agent_primary_todo_does_not_resolve_the_gap() -> None: + # A peer-claimed prerequisite retained on the primary path is not on this + # agent's selectable frontier and is not linked to the declared fallback, + # so the declared fallback gap survives for the declared fallback Todo. + payload = _status_payload( + fallback_runnable=False, + latest_runs=[ + _fallback_vision_run( + todo_delta=[f"retain:{PRIMARY_WAIT_ID}", f"retain:{PREREQ_ID}"], + fallback_declarations=[ + { + "declaration_id": "declared_fallback_direction", + "target_todo_id": FALLBACK_ID, + "successor_todo_id": FALLBACK_ID, + } + ], + ) + ], + ) + + frontier = _frontier_projection(payload) + + assert frontier["acceptance_gaps"] == [] + gaps = frontier["fallback_gaps"] + assert len(gaps) == 1 + assert gaps[0]["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. + payload = _status_payload( + fallback_runnable=False, + latest_runs=[_fallback_vision_run(todo_delta=[f"retain:{FALLBACK_ID}"])], + ) + item = payload["attention_queue"]["items"][0] + monitor = quota_todo_item( + todo_id=FALLBACK_ID, + index=3, + text="[P2] Watch the declared fallback direction.", + claimed_by=AGENT_ID, + ) + monitor["task_class"] = "continuous_monitor" + summary = _agent_todos(fallback_runnable=False) + for slot in ("executable_backlog_items", "backlog_items"): + summary[slot] = list(summary.get(slot) or []) + [monitor] + item["agent_todos"] = summary + + frontier = _frontier_projection(payload) + + gaps = frontier["fallback_gaps"] + assert len(gaps) == 1 + assert gaps[0]["unresolved_todo_ids"] == [FALLBACK_ID] + + +@pytest.mark.parametrize( + ("state", "path_outcome"), + [ + ("no_followup", None), + ("vision_drift_detected", "stop"), + ], + ids=["closed-state", "terminal-path-outcome"], +) +def test_terminal_disposition_closes_fallback_gap_without_regenerating( + state: str, + path_outcome: str | None, +) -> None: + payload = _status_payload( + fallback_runnable=False, + latest_runs=[_fallback_vision_run(state=state, path_outcome=path_outcome)], + ) + + first = _frontier_projection(payload) + assert "fallback_gaps" not in first + assert first["acceptance_gaps"] == [] + + second = _frontier_projection(payload) + assert "fallback_gaps" not in second + assert second["acceptance_gaps"] == [] + + +def test_declared_bounded_successor_delta_resolves_the_gap() -> None: + payload = _status_payload( + fallback_runnable=False, + latest_runs=[_fallback_vision_run(todo_delta=[f"create:{FALLBACK_ID}"])], + ) + + frontier = _frontier_projection(payload) + + assert "fallback_gaps" not in frontier + + +def test_unrelated_create_does_not_resolve_fallback_gap() -> None: + # Maintainer blocker 1: An unrelated create/reopen action must not + # resolve the declared fallback direction. + payload = _status_payload( + fallback_runnable=False, + latest_runs=[ + _fallback_vision_run( + todo_delta=[ + f"retain:{PRIMARY_WAIT_ID}", + "create:todo_unrelated_maintenance", + ], + fallback_declarations=[ + { + "declaration_id": "fallback_direction", + "target_todo_id": FALLBACK_ID, + "successor_todo_id": FALLBACK_ID, + } + ], + ) + ], + ) + + frontier = _frontier_projection(payload) + + gaps = frontier["fallback_gaps"] + assert len(gaps) == 1 + assert gaps[0]["unresolved_todo_ids"] == [FALLBACK_ID] + + +def test_typed_declaration_successor_relation_resolves_gap() -> None: + # A typed declaration-to-successor relation resolves when its declared + # bounded successor is created in todo_delta. + payload = _status_payload( + fallback_runnable=False, + latest_runs=[ + _fallback_vision_run( + todo_delta=[f"retain:{PRIMARY_WAIT_ID}", "create:todo_typed_successor"], + fallback_declarations=[ + { + "declaration_id": "fallback_direction", + "successor_todo_id": "todo_typed_successor", + } + ], + ) + ], + ) + + frontier = _frontier_projection(payload) + + assert "fallback_gaps" not in frontier + + +def test_legacy_alias_shapes_do_not_declare_a_fallback() -> None: + # Unsupported alias shapes (string arrows, patch-level fields, renamed + # keys) have no production author; the reader must not resurrect them. + run = _fallback_vision_run( + todo_delta=[f"retain:{PRIMARY_WAIT_ID}"], + fallback_declarations=[], + ) + vision = dict(run["agent_vision"]) + vision["fallback_relationships"] = [ + {"fallback_id": FALLBACK_ID, "successor_id": FALLBACK_ID} + ] + vision["vision_patch"] = { + **vision["vision_patch"], + "fallback_declarations": [f"{FALLBACK_ID}->{FALLBACK_ID}"], + } + payload = _status_payload( + fallback_runnable=False, + latest_runs=[{**run, "agent_vision": vision}], + ) + + frontier = _frontier_projection(payload) + + assert "fallback_gaps" not in frontier + + +def test_compact_read_model_mirrors_the_bounded_declaration_contract() -> None: + # Compaction is the defensive mirror of the TS prepare contract: at most + # four entries, one per unique declaration_id, typed fields only, and + # anything past the bound is truncated exactly like the write contract + # rejects it. + compact = compact_goal_vision_packet( + { + "schema_version": "goal_vision_replan_contract_v0", + "goal_id": GOAL_ID, + "agent_id": AGENT_ID, + "state": "vision_drift_detected", + "vision_patch": {"vision_summary": "Bounded route."}, + "fallback_declarations": [ + { + "declaration_id": "first_direction", + "target_todo_id": FALLBACK_ID, + "legacy_field": "dropped", + }, + {"declaration_id": "first_direction", "target_todo_id": "todo_dup"}, + {"target_todo_id": "todo_missing_id"}, + "not-an-object", + {"declaration_id": "truncated_direction"}, + ], + } + ) + + assert compact is not None + assert compact["fallback_declarations"] == [ + {"declaration_id": "first_direction", "target_todo_id": FALLBACK_ID}, + ] + + +def test_selectable_frontier_ids_mirror_the_authoritative_counts() -> None: + # The completion-evidence id set must be the same agent-scoped frontier + # the authoritative advancement counter projects. + summary = _agent_todos(fallback_runnable=True) + peer_only = quota_todo_item( + todo_id="todo_peer_owned_direction", + index=4, + text="[P1] Advance the peer-owned direction.", + claimed_by=PRIMARY_AGENT, + ) + for slot in ("executable_backlog_items", "backlog_items"): + summary[slot] = list(summary.get(slot) or []) + [peer_only] + + selectable_ids = agent_scoped_selectable_advancement_todo_ids( + summary, + agent_id=AGENT_ID, + ) + counts = todo_advancement_frontier_counts(summary, agent_id=AGENT_ID) + + assert selectable_ids == {FALLBACK_ID} + assert counts["current_agent_claimed_advancement_count"] == 1 + assert counts["unclaimed_advancement_count"] == 0 + assert counts["other_agent_claimed_advancement_count"] == 2 + assert PREREQ_ID not in selectable_ids + assert PRIMARY_WAIT_ID not in selectable_ids + + +def test_vision_todo_delta_actions_contract_stays_the_shared_owner() -> None: + # Both the acceptance-gap projection and the fallback disposition must + # consume one action contract; create/reopen stay the successor subset. + assert VISION_FRONTIER_TODO_DELTA_ACTIONS == frozenset( + {"activate", "create", "reopen", "resume", "retain"} + ) diff --git a/tests/control_plane_ts/vision_checkpoint.test.ts b/tests/control_plane_ts/vision_checkpoint.test.ts index 51f2f82240..b9cd27499c 100644 --- a/tests/control_plane_ts/vision_checkpoint.test.ts +++ b/tests/control_plane_ts/vision_checkpoint.test.ts @@ -224,6 +224,121 @@ test("prepare rejects incomplete and over-wide path deltas", () => { ); }); +test("prepare validates and carries bounded unique fallback declarations", () => { + const result = buildVisionCheckpoint(prepareRequest({ + agent_vision_packet: { + vision_summary: "Ship one bounded route.", + fallback_declarations: [ + { + declaration_id: " declared_fallback_direction ", + target_todo_id: "todo_declared_fallback", + successor_todo_id: null, + }, + { declaration_id: "typed_successor", successor_todo_id: "todo_successor" }, + ], + }, + })); + + const vision = result.agent_vision as Record; + assert.deepEqual(vision.fallback_declarations, [ + { + declaration_id: "declared_fallback_direction", + target_todo_id: "todo_declared_fallback", + }, + { declaration_id: "typed_successor", successor_todo_id: "todo_successor" }, + ]); + const budget = vision.vision_budget as Record; + const fieldLimits = budget.field_limits as Record; + const fieldUsage = budget.field_usage as Record; + assert.equal(fieldLimits["fallback_declarations"], 4); + assert.equal(fieldLimits["fallback_declarations[]"], 120); + assert.equal( + fieldUsage["fallback_declarations[0].declaration_id"], + "declared_fallback_direction".length, + ); + + const empty = buildVisionCheckpoint(prepareRequest({ + agent_vision_packet: { + vision_summary: "Ship one bounded route.", + fallback_declarations: [], + }, + })) as Record; + assert.equal( + (empty.agent_vision as Record).fallback_declarations, + undefined, + ); +}); + +test("prepare rejects unbounded, duplicated, and unsafe fallback declarations", () => { + const declaration = (id: string) => ({ declaration_id: id }); + assert.throws( + () => buildVisionCheckpoint(prepareRequest({ + agent_vision_packet: { + vision_summary: "Ship one bounded route.", + fallback_declarations: "declared_fallback_direction", + }, + })), + /fallback_declarations must be a JSON array/, + ); + assert.throws( + () => buildVisionCheckpoint(prepareRequest({ + agent_vision_packet: { + vision_summary: "Ship one bounded route.", + fallback_declarations: [1, 2, 3, 4, 5].map(() => declaration("one")), + }, + })), + /fallback_declarations has 5 items; limit is 4/, + ); + assert.throws( + () => buildVisionCheckpoint(prepareRequest({ + agent_vision_packet: { + vision_summary: "Ship one bounded route.", + fallback_declarations: [ + declaration("declared_fallback_direction"), + declaration("declared_fallback_direction"), + ], + }, + })), + /repeats declaration_id "declared_fallback_direction"/, + ); + assert.throws( + () => buildVisionCheckpoint(prepareRequest({ + agent_vision_packet: { + vision_summary: "Ship one bounded route.", + fallback_declarations: [{ target_todo_id: "todo_declared_fallback" }], + }, + })), + /requires a non-empty declaration_id/, + ); + assert.throws( + () => buildVisionCheckpoint(prepareRequest({ + agent_vision_packet: { + vision_summary: "Ship one bounded route.", + fallback_declarations: [ + { declaration_id: "route", target_todo_id: "todo_with_fallback" }, + { declaration_id: "token = leaked", successor_todo_id: "todo_x" }, + ], + }, + })), + /contains a private-looking value/, + ); + assert.throws( + () => buildVisionCheckpoint(prepareRequest({ + agent_vision_packet: { + vision_summary: "Ship one bounded route.", + fallback_declarations: [ + { declaration_id: "x".repeat(121) }, + ], + }, + })), + (error: unknown) => { + const rejection = error as { code?: string; message?: string }; + assert.equal(rejection.code, "vision_budget_exceeded"); + return true; + }, + ); +}); + test("pure prepare and finalize reductions replay deterministically", () => { const prepare = prepareRequest(); assert.deepEqual( diff --git a/tests/test_state_refresh_projections.py b/tests/test_state_refresh_projections.py index 6e77217389..b411ac344d 100644 --- a/tests/test_state_refresh_projections.py +++ b/tests/test_state_refresh_projections.py @@ -4,7 +4,10 @@ def test_refresh_output_projections_preserve_optional_run_contracts() -> None: - delta_contract = {"schema_version": "repair_delta_contract_v0", "delta_present": True} + delta_contract = { + "schema_version": "repair_delta_contract_v0", + "delta_present": True, + } agent_vision = { "schema_version": "goal_vision_v0", "agent_id": "quality-agent", @@ -12,6 +15,7 @@ def test_refresh_output_projections_preserve_optional_run_contracts() -> None: "vision_patch": {"acceptance_summary": "Keep simplification measurable."}, "todo_delta": ["Simplify the projection owner."], "vision_budget": {"status": "within_budget"}, + "fallback_declarations": None, "path_delta": {"changed": ["projection assembly"]}, "validation": {"budget_checked": True}, } @@ -25,7 +29,10 @@ def test_refresh_output_projections_preserve_optional_run_contracts() -> None: "state": { "path": "/ignored/project/state.md", "sha256_16": "0123456789abcdef", - "frontmatter": {"updated_at": "2026-07-30T23:59:00+00:00", "status": "active"}, + "frontmatter": { + "updated_at": "2026-07-30T23:59:00+00:00", + "status": "active", + }, "next_action": ["continue"], }, "runtime_projection_route": {"status": "resolved", "route_id": "route-1"}, @@ -71,6 +78,7 @@ def test_refresh_output_projections_preserve_optional_run_contracts() -> None: "state", "vision_patch", "todo_delta", + "fallback_declarations", "vision_budget", "path_delta", ) @@ -95,3 +103,59 @@ def test_refresh_output_projections_preserve_optional_run_contracts() -> None: assert payload["active_state_next_action_update"] == {"updated": True} assert payload["agent_vision"] is agent_vision assert payload["state"] is record["state"] + + +def test_run_index_agent_vision_keeps_fallback_declarations(tmp_path): + """The run index must persist typed fallback declarations for later reads.""" + + from loopx.state_refresh import _build_state_refresh_output_projections + + record = { + "schema_version": "loopx_goal_run_v0", + "generated_at": "2026-09-07T00:00:00+00:00", + "goal_id": "goal-1", + "agent_id": "agent-a", + "turn_instance_id": "turn-1", + "observed_at": "2026-09-07T00:00:00Z", + "classification": "validated_progress", + "recommended_action": "continue", + "recommended_action_source": "explicit_arg", + "health_check": "state_file 1/1", + "state": {"frontmatter": {"updated_at": "2026-09-06T23:59:00+00:00"}}, + "runtime_projection_route": {"status": "resolved", "route_id": "route-1"}, + "delivery_batch_scale": "single_surface", + "delivery_outcome": "outcome_progress", + "delivery_workspace": {"workspace_kind": "independent_git_worktree"}, + "agent_vision": { + "schema_version": "loopx_goal_vision_packet_v0", + "agent_id": "agent-a", + "state": "active", + "vision_patch": "Primary path; fallback direction declared.", + "todo_delta": [], + "fallback_declarations": [ + {"declaration_id": "decl_fallback_1", "target_todo_id": "todo_fb1"} + ], + "vision_budget": {}, + }, + } + registry = tmp_path / "registry.global.json" + registry.write_text("{}", encoding="utf-8") + runtime_root = tmp_path / "runtime" + json_path = runtime_root / "goals" / "goal-1" / "runs" / "run.json" + markdown_path = runtime_root / "goals" / "goal-1" / "runs" / "run.md" + index_path = runtime_root / "goals" / "goal-1" / "runs" / "index.jsonl" + _record, index_record = _build_state_refresh_output_projections( + record=record, + registry_path=registry, + runtime_root=runtime_root, + project=tmp_path, + json_path=json_path, + markdown_path=markdown_path, + index_path=index_path, + dry_run=True, + autonomous_replan_recorded_requested=False, + ) + indexed_vision = index_record.get("agent_vision") or {} + assert indexed_vision.get("fallback_declarations") == [ + {"declaration_id": "decl_fallback_1", "target_todo_id": "todo_fb1"} + ]