From 772768b0cdcb5d1e65654ef93feb01df6d2eacb2 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Thu, 17 Sep 2026 04:33:20 +0800 Subject: [PATCH] feat(manager): let a partially applied team plan be finished A confirmed team plan can be applied partially: the settlement staffs the lanes the host can run and reports the rest as gaps. Nothing could then finish it. `apply` returned the stored proposal for any applied plan, so a lane that was missing for a reason the host later resolved stayed missing until the owner confirmed a whole new plan. A partial application now records what it still owes -- the digest of the confirmed plan and the lanes it left unstaffed -- and a re-entrant apply of that same proposal completes the plan instead of replaying an empty success: the settlement re-reads the host's staffing facts, every committed lane keeps its Todo, and only the lanes that are staffable now become work. The outcome becomes team_plan_applied once no lane is left, and the receipt keeps the lanes it recovered. The apply entry point, the settlement call and the receipt builder are shared with the first apply, so one plan cannot have two settlements. A recovery refuses before writing anything when the host can no longer staff a lane the plan already committed, records that refusal on the plan, and never turns an apply that already happened into a failure. What it deliberately does not gate is Goal-level intent drift: the repository has no typed intent identity (shared_goal_alignment says so), and both facts that could stand in for one move on the apply's own writes, so the plan itself -- the commitment the owner confirmed -- is what a recovery proves it is completing. The team-plan settlement, its receipt and its recovery live in loopx/chat_team_plan_actions.py, beside the monitor, Todo and lifecycle action mixins, so the router module does not grow past its reviewed module metric ceiling for work that is not routing. Validation: 46 focused Python tests, including the recovery, a no-progress replay, the committed-lane refusal and the unchanged behaviour of a fully applied plan; loopx canary premerge --from-git-diff passes with 0 failures (maintainability ratchet, vocabulary drift, 4 catalog canaries, public boundary). Rebase reconciliation (2026-09-17, onto main after #4587 and #4615 landed): main had added the typed lane-write failure path in `_apply_team_plan` while this branch moved that method into the team-plan action mixin, so the two changes collided in one place. The resolution keeps both behaviours: the mixin now forwards `lane_failure` from the settlement and `_apply_team_plan` records the retry-safe `team_plan_lane_write_failed` failure with the identities it did create, while the F4 recovery cursor stays on the partial-application path. The test file keeps main's two lane-failure cases and this branch's recovery cases side by side; this branch's second-lane helper is renamed to `_unstaffed_second_lane_plan` because main already owns the name `_two_lane_plan` for the both-lanes-staffed plan. Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- loopx/chat_action_store.py | 43 +++ loopx/chat_actions.py | 196 +---------- loopx/chat_team_plan_actions.py | 522 ++++++++++++++++++++++++++++ tests/test_chat_team_plan_action.py | 163 +++++++++ 4 files changed, 742 insertions(+), 182 deletions(-) create mode 100644 loopx/chat_team_plan_actions.py 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: