Skip to content
43 changes: 31 additions & 12 deletions loopx/control_plane/goals/goal_frontier/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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))
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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"] = (
Expand Down Expand Up @@ -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] = []
Expand Down Expand Up @@ -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):
Expand Down
289 changes: 289 additions & 0 deletions loopx/control_plane/goals/goal_frontier/fallback_disposition.py
Original file line number Diff line number Diff line change
@@ -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}
Loading