diff --git a/loopx/chat_action_store.py b/loopx/chat_action_store.py index 9cfcb39632..c5294776a5 100644 --- a/loopx/chat_action_store.py +++ b/loopx/chat_action_store.py @@ -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 diff --git a/loopx/chat_actions.py b/loopx/chat_actions.py index f128dda6fd..79e6277b9c 100644 --- a/loopx/chat_actions.py +++ b/loopx/chat_actions.py @@ -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 @@ -177,6 +178,7 @@ def _monitor_text(parameters: Mapping[str, Any]) -> str | None: class ChatActionService( ChatActionNormalizationMixin, + ChatTeamPlanActionMixin, ChatGoalLifecycleActionMixin, ChatMonitorActionMixin, ChatTodoActionMixin, @@ -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, @@ -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", @@ -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")), diff --git a/loopx/chat_team_plan_actions.py b/loopx/chat_team_plan_actions.py new file mode 100644 index 0000000000..e03300d951 --- /dev/null +++ b/loopx/chat_team_plan_actions.py @@ -0,0 +1,522 @@ +"""Typed Chat actions for confirming a steward team plan. + +A team plan is the one typed action that commits a *set* of lanes, so it owns +two things the general action service does not: the facts a confirmation was +reviewed against, and what happens when that confirmation could only be applied +partly. Both live here so `chat_actions` stays the router rather than the +settlement. +""" + +from __future__ import annotations + +import hashlib +from pathlib import Path +from typing import Any, Mapping + +from .agent_registry import registered_agent_ids_for_goal +from .control_plane.runtime.time import now_utc_iso + + +def _team_plan_gap_lane_ids( + plan: Mapping[str, Any], settled: Mapping[str, Any] +) -> list[str]: + """Name the confirmed lanes a settlement did not staff. + + The lanes that stayed unstaffed are the plan's own lanes minus the ones the + settlement reported, read here rather than carried as a second list, so a + cursor cannot disagree with the plan about which lane is missing. + """ + + settled_lane_ids = { + str(item.get("lane_id") or "") + for item in settled.get("lane_settlements") or [] + if isinstance(item, Mapping) + } + return [ + str(lane.get("lane_id") or "") + for lane in (plan.get("lanes") or []) + if isinstance(lane, Mapping) + and str(lane.get("lane_id") or "") + and str(lane.get("lane_id") or "") not in settled_lane_ids + ] + + + + +class ChatTeamPlanActionMixin: + """Keep team-plan settlement and recovery out of the action router.""" + + def _team_plan_state_fingerprint( + self, goal_id: str, plan: Mapping[str, Any] + ) -> str: + from .chat_actions import _digest + + """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, + ) + + 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 _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.""" + + 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} + settled = self._settle_team_plan_lanes( + proposal_id=proposal_id, + goal_id=goal_id, + plan=plan, + requested_by=str(parameters.get("requested_by") or "owner"), + ) + if settled.get("error") == "team_plan_no_staffable_lane": + # 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 {settled['gap_count']} lane(s) can be " + "staffed by this host, so confirming it created no work" + ), + ), + "turn": None, + } + lane_failure = settled.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(settled['lane_todo_ids'])} lane Todo(s) from this " + "plan already exist" + ), + details={ + "goal_id": goal_id, + "lane_todo_ids": [str(item) for item in settled["lane_todo_ids"]], + "lane_settlements": [ + dict(item) for item in settled["lane_settlements"] + ], + "failed_lane_id": str(lane_failure["lane_id"]), + "failed_lane_reason_code": str(lane_failure["reason_code"]), + }, + ), + "turn": None, + } + receipt = self._team_plan_receipt( + proposal_id=proposal_id, + goal_id=goal_id, + settled=settled, + ) + if settled["gap_count"]: + # A partially committed plan keeps the facts a recovery has to + # satisfy, so the lanes it could not create stay recoverable + # instead of becoming a receipt detail nobody can act on. + receipt["recovery_cursor"] = self._team_plan_recovery_cursor( + plan=plan, + settled=settled, + ) + stored = self.store.apply( + proposal_id, current_state_fingerprint=current_fingerprint, receipt=receipt + ) + return {"proposal": stored, "turn": None} + + def _team_plan_staffing_gap_lane_ids( + self, goal_id: str, plan: Mapping[str, Any] + ) -> list[str]: + """Read the host's staffing verdict for a plan without creating work. + + A recovery has to decide whether it may act before the settlement + creates anything, so the validation the settlement performs is read here + as a verdict only: which of the plan's lanes this host cannot staff now. + """ + + from .control_plane.todos.contract import ( + TODO_ACTION_KIND_ADVANCEMENT_VALUES, + ) + from .control_plane.work_items.governed_transition_proposal import ( + validate_steward_team_plan_preview, + ) + + goal = self._goal(goal_id) + verdict = validate_steward_team_plan_preview( + plan, + registered_agent_ids=registered_agent_ids_for_goal(goal), + supported_action_kinds=sorted(TODO_ACTION_KIND_ADVANCEMENT_VALUES), + ) + return [ + str(lane.get("lane_id") or "") + for lane in verdict["lanes"] + if str(lane.get("staffing") or "") != "ready" + ] + + def _settle_team_plan_lanes( + self, + *, + proposal_id: str, + goal_id: str, + plan: Mapping[str, Any], + requested_by: str, + ) -> dict[str, Any]: + """Ask the governed transition owner to ensure the plan's lane Todos. + + The settlement re-validates the plan against the host's own facts every + time it runs, so both the first apply and a recovery go through this one + call: a lane whose Todo already exists comes back as a reuse, and a lane + this host still cannot staff comes back as a gap. + """ + + from .control_plane.work_items.governed_transition_proposal import ( + GovernedTransitionSettlementPhase, + settle_governed_transition_proposals, + ) + + settlements = settle_governed_transition_proposals( + registry_path=self.registry_path, + goal_id=goal_id, + agent_id=requested_by, + 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 [])] + lane_settlements = [ + dict(item) for item in (settlement.get("lane_settlements") or []) + ] + return { + "error": "" if lane_todo_ids else "team_plan_no_staffable_lane", + "action": str(settlement.get("action") or ""), + "todo_id": str(settlement.get("todo_id") or ""), + "lane_todo_ids": lane_todo_ids, + # A lane write that fails after earlier lanes exist is reported by + # the settlement owner rather than raised, so the apply can record a + # retry-safe failure that still names the identities it created. + "lane_failure": settlement.get("lane_failure"), + # A lane the settlement created is a lane this plan had not yet + # committed; the settlement records that per lane, so the recovery + # reads the same fact rather than re-deriving it from Todo ids. + "created_lane_ids": [ + str(item.get("lane_id") or "") + for item in lane_settlements + if str(item.get("disposition") or "") == "created" + ], + "lane_settlements": lane_settlements, + "intent_basis": str(settlement.get("intent_basis") or ""), + "gap_count": int(settlement.get("gap_count") or 0), + } + + def _team_plan_receipt( + self, + *, + proposal_id: str, + goal_id: str, + settled: Mapping[str, Any], + ) -> dict[str, Any]: + from .chat_actions import _digest + + """Read what the settlement produced as the receipt the card shows.""" + + lane_todo_ids = [str(item) for item in settled["lane_todo_ids"]] + gap_count = int(settled["gap_count"]) + # 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. + if str(settled["action"]) == "reused": + outcome = "team_plan_lanes_already_present" + elif gap_count: + outcome = "team_plan_partially_applied" + else: + outcome = "team_plan_applied" + receipt: dict[str, Any] = { + "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(settled["todo_id"]), + "lane_todo_ids": lane_todo_ids, + }, + } + if settled["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 settled["lane_settlements"]] + if gap_count: + receipt["gap_count"] = gap_count + if settled["intent_basis"]: + # The canonical revision these lanes were created against, so the + # owner's readback can name what the work advances. + receipt["intent_basis"] = str(settled["intent_basis"]) + return receipt + + def _team_plan_recovery_cursor( + self, + *, + plan: Mapping[str, Any], + settled: Mapping[str, Any], + ) -> dict[str, Any]: + from .chat_actions import _digest + + """Record what a later recovery of this plan still owes. + + A recovery finishes the lanes a confirmed plan could not staff, so the + cursor carries the two facts that decision needs: the plan the owner + confirmed (its digest, so a recovery can prove it is completing the same + commitment) and the lanes that confirmation left unstaffed, named. The + remaining lanes are derived from the plan itself rather than from a + second list, so the cursor cannot disagree with the plan about which + lane is missing. + """ + + return { + "schema_version": "team_plan_recovery_cursor_v0", + "plan_digest": _digest(dict(plan)), + "gap_lane_ids": _team_plan_gap_lane_ids(plan, settled), + "attempts": [], + } + + def _record_team_plan_recovery_attempt( + self, + cursor: Mapping[str, Any], + *, + outcome: str, + recovered_lane_ids: Mapping[str, Any] = (), + unstaffable_lane_ids: Mapping[str, Any] = (), + remaining_gap_lane_ids: Mapping[str, Any] | None = None, + ) -> dict[str, Any]: + """Append one bounded attempt to the cursor's own history.""" + + attempts = [dict(item) for item in (cursor.get("attempts") or [])] + attempts.append( + { + "outcome": str(outcome), + "recovered_lane_ids": [str(item) for item in recovered_lane_ids], + "unstaffable_lane_ids": [str(item) for item in unstaffable_lane_ids], + "remaining_gap_lane_ids": ( + [str(item) for item in (cursor.get("gap_lane_ids") or [])] + if remaining_gap_lane_ids is None + else [str(item) for item in remaining_gap_lane_ids] + ), + "recorded_at": now_utc_iso(), + } + ) + # A plan can only be recovered as often as it has lanes, so the history + # stays a readback of what happened rather than an unbounded log. + return { + **dict(cursor), + "gap_lane_ids": [ + str(item) + for item in ( + cursor.get("gap_lane_ids") or [] + if remaining_gap_lane_ids is None + else remaining_gap_lane_ids + ) + ], + "attempts": attempts[-3:], + } + + def _recover_team_plan( + self, proposal_id: str, proposal: Mapping[str, Any] + ) -> dict[str, Any]: + from .chat_actions import _digest + + """Finish the lanes a confirmed plan left unstaffed. + + This is the re-entrant apply the roadmap's R1 exit names: the owner does + not confirm the plan again, because the commitment is already theirs and + the lanes are already recorded. Two things must hold, and both are read + from the plan the owner actually confirmed: + + - the stored plan is still the plan that was confirmed (its digest), so + a recovery can never complete a different commitment; and + - the settlement still staffs every lane the plan already committed, so + a recovery may only *add* the lanes that are staffable now, never + quietly replace or drop one that already has a Todo. + + A refusal is recorded on the plan instead of being turned into a failure + of an apply that already happened: the lanes that exist stay exactly as + they are, and the cursor says what the attempt found. + """ + + parameters = proposal.get("normalized_parameters") + parameters = parameters if isinstance(parameters, Mapping) else {} + goal_id = str(parameters.get("goal_id") or "") + plan = parameters.get("plan") + if not goal_id or not isinstance(plan, Mapping): + raise ValueError("team plan proposal is malformed") + receipt = proposal.get("receipt") + receipt = dict(receipt) if isinstance(receipt, Mapping) else {} + cursor = receipt.get("recovery_cursor") + if not isinstance(cursor, Mapping) or not (cursor.get("gap_lane_ids") or []): + raise ValueError("this team plan has no lane left to recover") + if str(cursor.get("plan_digest") or "") != _digest(dict(plan)): + # The stored plan is not the one this cursor was written for, so no + # recovery can claim to be completing the confirmed commitment. + receipt["recovery_cursor"] = self._record_team_plan_recovery_attempt( + cursor, outcome="confirmed_plan_changed" + ) + stored = self.store.record_team_plan_recovery( + proposal_id, receipt=receipt + ) + return {"proposal": stored, "turn": None} + # The staffability verdict is read before the settlement runs, because + # the settlement creates the lanes it finds staffable. A recovery that + # would strand a lane the plan already committed has to refuse *before* + # it writes anything, not after. + committed_lane_ids = [ + str(item.get("lane_id") or "") + for item in (receipt.get("lanes") or []) + if isinstance(item, Mapping) + ] + verdict_gaps = self._team_plan_staffing_gap_lane_ids(goal_id, plan) + regressed = sorted( + lane_id for lane_id in committed_lane_ids if lane_id in verdict_gaps + ) + if regressed: + # A lane the plan already committed is unstaffable on the host now. + # Finishing the plan would commit work the host cannot run, so this + # is a refusal to act rather than a partial success: the committed + # lanes stand where they are, and the plan says which one the host + # can no longer staff instead of leaving it to the next reader. + receipt["recovery_cursor"] = self._record_team_plan_recovery_attempt( + cursor, + outcome="committed_lane_unstaffable", + unstaffable_lane_ids=regressed, + remaining_gap_lane_ids=verdict_gaps, + ) + stored = self.store.record_team_plan_recovery( + proposal_id, receipt=receipt + ) + return {"proposal": stored, "turn": None} + settled = self._settle_team_plan_lanes( + proposal_id=proposal_id, + goal_id=goal_id, + plan=plan, + requested_by=str(parameters.get("requested_by") or "owner"), + ) + recovered = [ + lane_id + for lane_id in settled["created_lane_ids"] + if lane_id and lane_id not in committed_lane_ids + ] + unsettled_lane_ids = _team_plan_gap_lane_ids(plan, settled) + merged = self._team_plan_receipt( + proposal_id=proposal_id, + goal_id=goal_id, + settled=settled, + ) + # The lanes this plan ensured before are still this plan's lanes, so the + # recovery reports the committed set rather than only what it touched. + previous_lanes = { + str(item.get("lane_id") or ""): dict(item) + for item in (receipt.get("lanes") or []) + if isinstance(item, Mapping) + } + for item in merged.get("lanes") or []: + previous_lanes[str(item.get("lane_id") or "")] = dict(item) + if previous_lanes: + merged["lanes"] = list(previous_lanes.values()) + merged["outcome"] = ( + "team_plan_partially_applied" + if merged.get("gap_count") + else "team_plan_applied" + ) + merged["recovered_lane_ids"] = sorted( + { + *[ + str(lane_id) + for attempt in (cursor.get("attempts") or []) + if isinstance(attempt, Mapping) + for lane_id in (attempt.get("recovered_lane_ids") or []) + ], + *recovered, + } + ) + merged["recovery_cursor"] = self._record_team_plan_recovery_attempt( + cursor, + outcome="recovered" if recovered else "no_progress", + recovered_lane_ids=recovered, + remaining_gap_lane_ids=unsettled_lane_ids, + ) + stored = self.store.record_team_plan_recovery(proposal_id, receipt=merged) + return {"proposal": stored, "turn": None} diff --git a/tests/test_chat_team_plan_action.py b/tests/test_chat_team_plan_action.py index c8e414709a..812a144078 100644 --- a/tests/test_chat_team_plan_action.py +++ b/tests/test_chat_team_plan_action.py @@ -323,6 +323,169 @@ def test_a_lane_that_fails_before_any_work_is_an_error_not_a_partial( service.apply(_preview(service, _two_lane_plan())["proposal_id"]) assert "loopx:todo " not in _todos(project) +def _unstaffed_second_lane_plan(*, second_agent: str = "agent-beta") -> dict: + """One side of the recovery cases: a plan whose second lane is unstaffed. + + Main's `_two_lane_plan` staffs both lanes on the same Agent, so it cannot + express the F4 gap. This variant assigns the second lane to another Agent, + which makes it a staffing gap on this host until that Agent is registered. + """ + + plan = _plan() + plan["lanes"].append( + { + "lane_id": "lane-beta", + "agent_id": second_agent, + "acceptance": "The review lane's receipt is recorded", + "first_todo": { + "text": "Advance the review lane", + "priority": "P1", + "task_class": "advancement_task", + "action_kind": "validate", + }, + } + ) + return plan + + +def _register_agent(registry_path: Path, agent_id: str) -> None: + """Register one more Agent for the Goal, exactly as a host change would.""" + + registry = json.loads(registry_path.read_text(encoding="utf-8")) + registered = registry["goals"][0]["coordination"]["registered_agents"] + registry["goals"][0]["coordination"]["registered_agents"] = [*registered, agent_id] + registry_path.write_text(json.dumps(registry), encoding="utf-8") + + +def test_a_partially_applied_plan_can_be_finished_without_confirming_again( + tmp_path: Path, +) -> None: + """The lanes a confirmed plan left unstaffed stay recoverable. + + R1's exit names this: a re-entrant apply has to finish a partial plan + without a fresh owner confirmation. The confirmation already happened, the + lanes are already recorded, and what moved is the host fact that stopped one + of them -- so the plan records the commitment it still has to meet and the + next apply completes it instead of replaying an empty success. + """ + + project, registry_path, service = _fixture(tmp_path, agents=(AGENT_ID,)) + preview = _preview(service, _unstaffed_second_lane_plan()) + + applied = service.apply(preview["proposal_id"])["proposal"] + assert applied["status"] == "applied" + receipt = applied["receipt"] + assert receipt["outcome"] == "team_plan_partially_applied" + assert receipt["gap_count"] == 1 + first_lane_todo_ids = receipt["resource_ids"]["lane_todo_ids"] + assert len(first_lane_todo_ids) == 1 + cursor = receipt["recovery_cursor"] + assert cursor["schema_version"] == "team_plan_recovery_cursor_v0" + # The cursor names the lane the confirmation left unstaffed and the plan it + # belongs to, so a later recovery can prove it is completing this plan. + assert cursor["gap_lane_ids"] == ["lane-beta"] + assert len(cursor["plan_digest"]) == 64 + assert cursor["attempts"] == [] + assert _todos(project).count("loopx:todo ") == 1 + + # The Agent the second lane needs is registered, which is the change that + # makes the lane staffable. It does not invalidate the commitment the owner + # already confirmed, so recovering it needs no second confirmation. + _register_agent(registry_path, "agent-beta") + + recovered = service.apply(preview["proposal_id"])["proposal"] + assert recovered["status"] == "applied" + receipt = recovered["receipt"] + assert receipt["outcome"] == "team_plan_applied" + assert "gap_count" not in receipt + lane_todo_ids = receipt["resource_ids"]["lane_todo_ids"] + assert len(lane_todo_ids) == 2 + # The lane that already had its Todo keeps it; the recovery created only the + # lane that was missing. + assert first_lane_todo_ids[0] in lane_todo_ids + assert len(receipt["lanes"]) == 2 + assert [lane["disposition"] for lane in receipt["lanes"]] == ["reused", "created"] + attempt = receipt["recovery_cursor"]["attempts"][-1] + assert attempt["outcome"] == "recovered" + assert attempt["remaining_gap_lane_ids"] == [] + assert attempt["recovered_lane_ids"] == ["lane-beta"] + assert receipt["recovery_cursor"]["gap_lane_ids"] == [] + assert receipt["recovered_lane_ids"] == ["lane-beta"] + assert _todos(project).count("loopx:todo ") == 2 + + +def test_a_recovery_that_staffs_nothing_creates_nothing_and_says_so( + tmp_path: Path, +) -> None: + """Re-applying a plan whose host still cannot staff the lane is a no-op.""" + + project, _registry_path, service = _fixture(tmp_path, agents=(AGENT_ID,)) + preview = _preview(service, _unstaffed_second_lane_plan()) + applied = service.apply(preview["proposal_id"])["proposal"] + before = _todos(project) + + again = service.apply(preview["proposal_id"])["proposal"] + + receipt = again["receipt"] + assert receipt["outcome"] == "team_plan_partially_applied" + assert receipt["gap_count"] == 1 + assert len(receipt["resource_ids"]["lane_todo_ids"]) == 1 + assert _todos(project) == before + attempt = receipt["recovery_cursor"]["attempts"][-1] + assert attempt["outcome"] == "no_progress" + assert attempt["recovered_lane_ids"] == [] + assert attempt["remaining_gap_lane_ids"] == ["lane-beta"] + + +def test_a_recovery_cannot_silently_drop_a_committed_lane( + tmp_path: Path, +) -> None: + """A recovery may only add lanes; it may never replace one that exists. + + If the host stops registering the Agent of a lane the plan already + committed, finishing the plan would commit work the host cannot run. The + refusal is recorded on the plan the owner already confirmed instead of + failing an apply that already happened, so the committed lanes stay exactly + as they are. + """ + + project, registry_path, service = _fixture(tmp_path, agents=(AGENT_ID,)) + preview = _preview(service, _unstaffed_second_lane_plan()) + applied = service.apply(preview["proposal_id"])["proposal"] + assert applied["receipt"]["gap_count"] == 1 + before = _todos(project) + + # The Agent of the lane that already has its Todo is no longer registered, + # and the Agent the missing lane needs is. Recovering now would finish a plan + # whose committed lane this host can no longer run. + registry = json.loads(registry_path.read_text(encoding="utf-8")) + registry["goals"][0]["coordination"]["registered_agents"] = ["agent-beta"] + registry_path.write_text(json.dumps(registry), encoding="utf-8") + + recovered = service.apply(preview["proposal_id"])["proposal"] + assert recovered["status"] == "applied" + receipt = recovered["receipt"] + assert receipt["gap_count"] == 1 + assert len(receipt["resource_ids"]["lane_todo_ids"]) == 1 + assert _todos(project) == before + attempt = receipt["recovery_cursor"]["attempts"][-1] + assert attempt["outcome"] == "committed_lane_unstaffable" + assert attempt["unstaffable_lane_ids"] == ["lane-alpha"] + + +def test_a_fully_applied_plan_is_not_recoverable(tmp_path: Path) -> None: + """A plan with no gap stays a replay, not a recovery.""" + + project, _registry_path, service = _fixture(tmp_path) + preview = _preview(service) + first = service.apply(preview["proposal_id"])["proposal"] + assert first["receipt"]["outcome"] == "team_plan_applied" + assert "recovery_cursor" not in first["receipt"] + + again = service.apply(preview["proposal_id"])["proposal"] + + assert again["receipt"] == first["receipt"] + assert _todos(project).count("loopx:todo ") == 1 def _validated(preview_plan: dict) -> dict: