Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
43 changes: 43 additions & 0 deletions loopx/chat_action_store.py
Original file line number Diff line number Diff line change
Expand Up @@ -1036,3 +1036,46 @@ def apply(
proposal["updated_at"] = now
self._write(payload)
return proposal

def record_team_plan_recovery(
self, proposal_id: str, *, receipt: Mapping[str, Any]
) -> dict[str, Any]:
"""Replace the receipt of an already applied plan with its recovery.

A recovery finishes lanes a confirmed plan left unstaffed, so it can
never be a first apply: the proposal has to be `applied` and its stored
receipt has to be the one the recovery is continuing from. Replacing
that receipt wholesale is what keeps the plan's readback single-sourced
-- the lane identities, the outcome, the remaining gap count and the
attempt history all move together under one lock.
"""

token = _opaque_id(proposal_id, field="proposal_id")
with exclusive_file_lock(
self.path,
agent_id="loopx-chat",
operation="record_team_plan_recovery",
):
payload = self._read()
proposal = payload["proposals"].get(token)
if not isinstance(proposal, dict):
raise KeyError("typed Chat action proposal was not found")
if str(proposal.get("status") or "") != "applied":
raise ActionConflictError(
"only an applied team plan can record a recovery"
)
if not isinstance(proposal.get("receipt"), Mapping):
raise ActionConflictError(
"an applied team plan needs its receipt before a recovery"
)
safe_receipt = _safe_json_value(dict(receipt), path="receipt")
if not isinstance(safe_receipt, dict):
raise ValueError("receipt must be an object")
_opaque_id(safe_receipt.get("receipt_id"), field="receipt.receipt_id")
_opaque_id(safe_receipt.get("outcome"), field="receipt.outcome")
if safe_receipt.get("projection_verified") is not True:
raise ValueError("receipt.projection_verified must be true")
proposal["receipt"] = safe_receipt
proposal["updated_at"] = _utc_now()
self._write(payload)
return proposal
196 changes: 14 additions & 182 deletions loopx/chat_actions.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,9 +16,10 @@
from .chat_goal_lifecycle_actions import ChatGoalLifecycleActionMixin
from .chat_monitor_actions import ChatMonitorActionMixin
from .chat_store import ChatSessionStore
from .chat_team_plan_actions import ChatTeamPlanActionMixin
from .chat_todo_actions import ChatTodoActionMixin
from .configure_goal import configure_goal
from .control_plane.runtime.time import now_utc, parse_timestamp
from .control_plane.runtime.time import now_utc, now_utc_iso, parse_timestamp
from .control_plane.scheduler.monitor_todo import monitor_next_due_at
from .history import load_registry
from .host_loop_activation import build_host_loop_activation_packet
Expand Down Expand Up @@ -177,6 +178,7 @@ def _monitor_text(parameters: Mapping[str, Any]) -> str | None:

class ChatActionService(
ChatActionNormalizationMixin,
ChatTeamPlanActionMixin,
ChatGoalLifecycleActionMixin,
ChatMonitorActionMixin,
ChatTodoActionMixin,
Expand Down Expand Up @@ -260,53 +262,6 @@ def _registry_fingerprint(self) -> str:
raise ValueError("the active LoopX registry is unavailable") from exc
return hashlib.sha256(content).hexdigest()

def _team_plan_state_fingerprint(
self, goal_id: str, plan: Mapping[str, Any]
) -> str:
"""Bind every fact a confirmed team plan was reviewed against.

Registry bytes are not enough. A plan is reviewed against the Goal's own
intent -- the objective its work advances -- and that intent lives in the
active-state document and in the canonical source basis the lanes would
be created against, neither of which the registry bytes cover. Changing
the objective therefore used to leave the confirmed plan applicable,
because nothing the preview bound had moved.

An unreadable fact is bound as its own explicit absence rather than
dropped from the digest, so the precondition fails closed in both
directions: a Goal whose intent becomes readable after the preview asks
the owner to confirm again instead of silently dropping the check.
"""

from .control_plane.work_items.governed_transition_proposal import (
steward_team_plan_intent_basis,
)

goal = self._goal(goal_id)
project = Path(str(goal.get("repo") or "")).expanduser()
state_file = Path(str(goal.get("state_file") or ""))
if not state_file.is_absolute():
state_file = project / state_file
try:
state_digest: str | None = hashlib.sha256(
state_file.read_bytes()
).hexdigest()
except OSError:
state_digest = None
return _digest(
{
"registry": self._registry_fingerprint(),
"goal_id": goal_id,
"active_state": state_digest,
"intent_basis": steward_team_plan_intent_basis(
goal_id=goal_id,
goal=goal,
registry_path=self.registry_path,
plan=plan,
),
}
)

def _agent_eligibility(
self,
agent_id: str,
Expand Down Expand Up @@ -1014,140 +969,6 @@ def _apply_monitor_create(
)
return {"proposal": stored, "turn": None}

def _apply_team_plan(
self, proposal_id: str, proposal: dict[str, Any], parameters: dict[str, Any]
) -> dict[str, Any]:
"""Create each ready lane's first bounded Todo through the Todo owner."""

from .control_plane.work_items.governed_transition_proposal import (
GovernedTransitionSettlementPhase,
settle_governed_transition_proposals,
)

goal_id = str(parameters["goal_id"])
plan = parameters.get("plan")
if not isinstance(plan, Mapping):
raise ValueError("team plan proposal is malformed")
# The preview bound this Goal's registration facts, its active-state
# intent and the canonical basis its lanes would advance; re-read them
# here so a plan confirmed against one objective cannot become work
# under another, and so a registry change still asks for confirmation.
current_fingerprint = self._team_plan_state_fingerprint(goal_id, plan)
if current_fingerprint != proposal.get("expected_state_fingerprint"):
stale = self.store.apply(
proposal_id,
current_state_fingerprint=current_fingerprint,
receipt={},
)
return {"proposal": stale, "turn": None}
# The governed transition owner re-validates the plan with the host's own
# facts and owns the settlement phase, so this action never becomes a
# second writer of lanes.
settlements = settle_governed_transition_proposals(
registry_path=self.registry_path,
goal_id=goal_id,
agent_id=str(parameters.get("requested_by") or "owner"),
effect_id=proposal_id,
proposals=[{**dict(plan), "proposal_id": proposal_id}],
existing_receipts=[],
checkpoint=lambda _receipts: None,
phase=GovernedTransitionSettlementPhase.PRE_SETTLEMENT,
)
settlement = settlements[0]
lane_todo_ids = [str(item) for item in (settlement.get("lane_todo_ids") or [])]
intent_basis = str(settlement.get("intent_basis") or "")
gap_count = int(settlement.get("gap_count") or 0)
if not lane_todo_ids:
# Every lane stayed a gap, so this confirmation created nothing and
# reused nothing. The old path wrote a receipt that reported
# "lanes already present" with a verified projection and an empty
# Todo id, which reads as success where the readback finds no work.
# A confirmation that can only create nothing is recorded as the
# typed failure it is, and the plan's lanes and reasons stay in the
# card the owner confirmed.
return {
"proposal": self.store.mark_failed(
proposal_id,
error_code="team_plan_no_staffable_lane",
message=(
f"none of the plan's {gap_count} lane(s) can be staffed by "
"this host, so confirming it created no work"
),
),
"turn": None,
}
# The outcome is read from what the settlement actually produced, not
# from "the action was not a creation": a plan that created lanes beside
# a gap is a partial application, and reporting it as a full success
# told the owner the commitment was kept when part of it was not.
lane_failure = settlement.get("lane_failure")
if lane_failure:
# A lane failed after earlier lanes were written. The plan did not
# apply, so it is not reported as applied; the identities that do
# exist are recorded with the failure so the retry reconciles
# against them instead of creating a second copy of the same lane.
return {
"proposal": self.store.mark_failed(
proposal_id,
error_code="team_plan_lane_write_failed",
message=(
f"lane {lane_failure['lane_id']} could not be created; "
f"{len(lane_todo_ids)} lane Todo(s) from this plan already exist"
),
details={
"goal_id": goal_id,
"lane_todo_ids": lane_todo_ids,
"lane_settlements": [
dict(item)
for item in (settlement.get("lane_settlements") or [])
],
"failed_lane_id": str(lane_failure["lane_id"]),
"failed_lane_reason_code": str(
lane_failure["reason_code"]
),
},
),
"turn": None,
}
if str(settlement.get("action") or "") == "reused":
outcome = "team_plan_lanes_already_present"
elif gap_count:
outcome = "team_plan_partially_applied"
else:
outcome = "team_plan_applied"
receipt = {
"receipt_id": _digest(
{
"proposal_id": proposal_id,
"goal_id": goal_id,
"lane_todo_ids": lane_todo_ids,
}
)[:32],
"outcome": outcome,
"projection_verified": True,
"resource_ids": {
"goal_id": goal_id,
"todo_id": str(settlement.get("todo_id") or ""),
"lane_todo_ids": lane_todo_ids,
},
}
lane_settlements = settlement.get("lane_settlements")
if lane_settlements:
# Which lane each created Todo is, who runs it, the priority it
# carries and the acceptance it was confirmed to end on, so the
# owner's readback still names the commitment and not just the work.
receipt["lanes"] = [dict(item) for item in lane_settlements]
if gap_count:
receipt["gap_count"] = gap_count
if intent_basis:
# The canonical revision these lanes were created against, so the
# owner's readback can name what the work advances.
receipt["intent_basis"] = intent_basis
stored = self.store.apply(
proposal_id, current_state_fingerprint=current_fingerprint, receipt=receipt
)
return {"proposal": stored, "turn": None}

def preview(self, request: Mapping[str, Any]) -> dict[str, Any]:
unknown = set(request) - {
"action_kind",
Expand Down Expand Up @@ -1351,6 +1172,17 @@ def apply(self, proposal_id: str) -> dict[str, Any]:
if proposal is None:
raise KeyError("typed Chat action proposal was not found")
if proposal.get("status") == "applied":
receipt = proposal.get("receipt")
cursor = (
receipt.get("recovery_cursor")
if isinstance(receipt, Mapping)
else None
)
if isinstance(cursor, Mapping) and (cursor.get("gap_lane_ids") or []):
# An applied plan that still owns unstaffed lanes is the one
# case where applying again is not a replay: the confirmation
# already happened, and what is left is finishing it.
return self._recover_team_plan(proposal_id, proposal)
return {
"proposal": proposal,
"turn": self._turn_from_receipt(proposal.get("receipt")),
Expand Down
Loading
Loading