From 32082e97d4f01db0cffe17b1d7943d543caa2c8c Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Thu, 17 Sep 2026 18:18:09 +0800 Subject: [PATCH] fix(goal_frontier): carry agent-owned frontier identity into writeback ACKs The long-chain ACK fence already matches on the identity over the rows this agent owns (`frontier_owned_identity`, added in #4610), but the typed-delta writeback path never carried that identity: - `LongTodoChainObservation.to_trigger()` omitted it, so the recorded obligation trigger only exposed `frontier_revision`; - `replan_obligation_trigger_checkpoints()` rebuilds the ACK checkpoints from those triggers, so writeback ACKs recorded a revision-only checkpoint. The revision covers the whole selectable set, which also contains rows nobody has claimed yet. Any peer lane claiming or editing such a row moves the revision, so a writeback ACK re-armed the pending `long_todo_chain` replan obligation on the next wake even though this agent's own work basis had not changed. Observed live: the 15:24 ACK for `replan-132c272361cab7f1` recorded `trigger_checkpoints: [{kind, frontier_revision}]` only, and the next guard reported `rearmed_after_obligation_id: replan-132c272361cab7f1`. The prompt-path ACK (`replan_successor_transition_ack`) was unaffected because it composes `long_todo_chain_source_checkpoint`, which already carries the identity; this change makes both ACK paths record the same fence. Validation: `tests/control_plane/test_goal_frontier_replan_rules.py::test_writeback_ack_owned_identity_absorbs_a_peer_claim` reproduces the peer-claim re-arm and fails on either dropped hop (verified by reverting each half independently); 151 focused control-plane tests pass. Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../goals/goal_frontier/long_todo_chain.py | 7 ++ .../work_items/progress_observation.py | 17 ++- .../test_goal_frontier_replan_rules.py | 104 +++++++++++++++++- 3 files changed, 121 insertions(+), 7 deletions(-) diff --git a/loopx/control_plane/goals/goal_frontier/long_todo_chain.py b/loopx/control_plane/goals/goal_frontier/long_todo_chain.py index dd506e0fff..99e2fa11ad 100644 --- a/loopx/control_plane/goals/goal_frontier/long_todo_chain.py +++ b/loopx/control_plane/goals/goal_frontier/long_todo_chain.py @@ -52,6 +52,13 @@ def to_trigger(self) -> dict[str, Any]: } if self.frontier_revision_complete and self.frontier_revision: trigger["frontier_revision"] = self.frontier_revision + if self.frontier_owned_identity: + # Refs #4610: the ACK fence also matches on the identity of the rows + # this agent owns, so another lane claiming or editing an unclaimed + # row must not re-arm the obligation. The typed-delta writeback + # records this trigger verbatim, so dropping the identity here would + # leave that fence inert for every writeback-path ACK. + trigger["frontier_owned_identity"] = self.frontier_owned_identity return trigger diff --git a/loopx/control_plane/work_items/progress_observation.py b/loopx/control_plane/work_items/progress_observation.py index 512521688e..ef1e16736e 100644 --- a/loopx/control_plane/work_items/progress_observation.py +++ b/loopx/control_plane/work_items/progress_observation.py @@ -404,12 +404,17 @@ def replan_obligation_trigger_checkpoints( frontier_revision = str(trigger.get("frontier_revision") or "").strip() if not kind or not frontier_revision: continue - checkpoints.append( - { - "kind": kind, - "frontier_revision": frontier_revision, - } - ) + checkpoint = { + "kind": kind, + "frontier_revision": frontier_revision, + } + # Refs #4610: carry the identity over the rows this agent owns, so the + # ACK fence can keep a writeback-path ACK valid when another lane claims + # or edits an unclaimed row without changing this agent's own chain. + owned_identity = str(trigger.get("frontier_owned_identity") or "").strip() + if owned_identity: + checkpoint["frontier_owned_identity"] = owned_identity + checkpoints.append(checkpoint) return checkpoints diff --git a/tests/control_plane/test_goal_frontier_replan_rules.py b/tests/control_plane/test_goal_frontier_replan_rules.py index fb5572a078..5d234b70f9 100644 --- a/tests/control_plane/test_goal_frontier_replan_rules.py +++ b/tests/control_plane/test_goal_frontier_replan_rules.py @@ -29,6 +29,9 @@ build_interaction_contract, interaction_next_cli_actions, ) +from loopx.control_plane.work_items.progress_observation import ( + replan_obligation_trigger_checkpoints, +) @pytest.mark.parametrize( @@ -228,12 +231,13 @@ def _accepted_long_chain_ack(obligation: dict[str, object]) -> dict[str, object] def _derive_long_chain( source_items: list[dict[str, object]], *, + agent_todo_summary: dict[str, object] | None = None, latest_replan_ack: dict[str, object] | None = None, current_transition_replan_ack: dict[str, object] | None = None, ) -> dict[str, object] | None: return derive_goal_frontier_replan_obligation_from_summaries( user_todo_summary={"open_count": 0}, - agent_todo_summary=_long_chain_summary(source_items), + agent_todo_summary=agent_todo_summary or _long_chain_summary(source_items), agent_todo_source_items=source_items, work_lane_contract={"lane": "advancement_task", "must_attempt_work": True}, agent_id="current-agent", @@ -271,6 +275,104 @@ def test_long_todo_chain_checkpoint_is_edge_triggered_and_rearms_on_change() -> ) +def test_writeback_ack_owned_identity_absorbs_a_peer_claim() -> None: + """A writeback ACK must survive a peer claim on an unclaimed row. + + The writeback path records ``replan_obligation_trigger_checkpoints``, so the + identity over this agent's own rows has to survive the trigger -> checkpoint + hop. Without it the recorded ACK only carries the frontier revision, which a + peer lane moves by claiming an unclaimed row, and every such claim re-arms + the obligation even though this agent's own selectable chain is unchanged. + """ + + claimed_items = [ + { + **_advancement(f"todo_{index:012x}", "current-agent"), + "updated_at": "2026-08-22T09:00:00+08:00", + } + for index in range(15) + ] + unclaimed_item = { + **_advancement("todo_0000000000ff", ""), + "updated_at": "2026-08-22T09:00:00+08:00", + } + + def summary( + unclaimed: list[dict[str, object]], + other_agents: list[dict[str, object]], + ) -> dict[str, object]: + return { + "open_count": 16, + "current_agent_claimed_open_count": 15, + "current_agent_claimed_advancement_count": 15, + "unclaimed_open_count": len(unclaimed), + "unclaimed_priority_open_items": unclaimed, + "executable_backlog_items": [*claimed_items, *unclaimed, *other_agents], + "claim_scope": {"other_agent_claimed_items": other_agents}, + } + + obligation = _derive_long_chain( + [*claimed_items, unclaimed_item], + agent_todo_summary=summary([unclaimed_item], []), + ) + assert obligation is not None + trigger = obligation["triggers"][0] + assert trigger["frontier_owned_identity"].startswith( + "todo_frontier_revision_v0:owned:" + ) + + checkpoints = replan_obligation_trigger_checkpoints(obligation) + assert [row["kind"] for row in checkpoints] == ["long_todo_chain"] + assert checkpoints[0]["frontier_owned_identity"] == ( + trigger["frontier_owned_identity"] + ) + + ack = { + "schema_version": "autonomous_replan_ack_v0", + "recorded": True, + "generated_at": "2026-08-22T09:05:00+08:00", + "semantic_delta": { + "schema_version": "replan_semantic_delta_v0", + "accepted": True, + "outcomes": ["new_surface"], + "satisfying_outcomes": ["new_surface"], + "trigger_kinds": ["long_todo_chain"], + "trigger_checkpoints": checkpoints, + "obligation_id": obligation["obligation_id"], + }, + } + + peer_claimed_item = { + **unclaimed_item, + "claimed_by": "peer-agent", + "updated_at": "2026-08-22T09:20:00+08:00", + } + peer_items = [*claimed_items, peer_claimed_item] + peer_summary = summary([], [peer_claimed_item]) + + rearmed = _derive_long_chain(peer_items, agent_todo_summary=peer_summary) + assert rearmed is not None + assert rearmed["obligation_id"] != obligation["obligation_id"] + assert rearmed["triggers"][0]["frontier_revision"] != ( + trigger["frontier_revision"] + ) + assert rearmed["triggers"][0]["frontier_owned_identity"] == ( + trigger["frontier_owned_identity"] + ) + assert _derive_long_chain( + peer_items, agent_todo_summary=peer_summary, latest_replan_ack=ack + ) is None + + revision_only_ack = deepcopy(ack) + for row in revision_only_ack["semantic_delta"]["trigger_checkpoints"]: + row.pop("frontier_owned_identity") + assert _derive_long_chain( + peer_items, + agent_todo_summary=peer_summary, + latest_replan_ack=revision_only_ack, + ) is not None + + def test_frontier_revision_index_preserves_complete_agent_lane_semantics() -> None: source_items = [ {