From 3ef552989f1eda258e05f4793d9ed86c8bbb8638 Mon Sep 17 00:00:00 2001 From: Duang777 Date: Thu, 8 Oct 2026 23:51:26 +0800 Subject: [PATCH 1/3] fix: isolate acceptance across goal recreation Signed-off-by: Duang777 --- loopx/control_plane/goals/acceptance.py | 57 ++- .../goals/acceptance_authority.ts | 18 + .../goals/acceptance_lifecycle.ts | 6 +- .../goals/source_session_recreation.py | 320 ++++++++------ .../project_registry_io_manifest_v1.json | 14 +- .../test_goal_acceptance_source_recreation.py | 404 ++++++++++++++++++ .../goal_acceptance_authority.test.ts | 45 ++ 7 files changed, 729 insertions(+), 135 deletions(-) create mode 100644 tests/control_plane/test_goal_acceptance_source_recreation.py diff --git a/loopx/control_plane/goals/acceptance.py b/loopx/control_plane/goals/acceptance.py index f0fcb1eba5..13814a8045 100644 --- a/loopx/control_plane/goals/acceptance.py +++ b/loopx/control_plane/goals/acceptance.py @@ -14,13 +14,19 @@ from uuid import uuid4 from ...agent_registry import load_goal_from_registry, registered_agent_ids_for_goal +from ...paths import resolve_runtime_root from ..coordination.local_authority import local_authority_is_promoted from ..coordination.local_authority_shadow_adapter import effective_runtime_root -from ..effect_runtime import effect_runtime_result +from ..effect_runtime import ( + CANONICAL_AUTHORITY_WRITE_TIMEOUT_SECONDS, + effect_runtime_result, +) +from ..projects.registry_codec import load_project_registry from ..todos.completion_validation import ( _resolve_completion_validation_workspace, run_declared_completion_validation_effect, ) +from .goal_ref_validation import exact_goal_ref _INSPECT_METHOD = "goal.acceptance.inspect" @@ -41,11 +47,27 @@ def _routing( raise ValueError( "Goal acceptance requires an existing canonical authority; activation never promotes a provider" ) - return {"runtime_root": str(root.resolve()), "goal_id": goal_id} + route: dict[str, Any] = { + "runtime_root": str(root.resolve()), + "goal_id": goal_id, + } + goal_instance_id = goal.get("goal_instance_id") + if goal_instance_id is not None: + route["goal_ref"] = exact_goal_ref(goal_id, goal_instance_id) + return route -def _result(method: str, request: Mapping[str, Any]) -> dict[str, Any]: - value = effect_runtime_result(method, dict(request)) +def _result( + method: str, + request: Mapping[str, Any], + *, + timeout: float | None = None, +) -> dict[str, Any]: + value = effect_runtime_result( + method, + dict(request), + **({"timeout": timeout} if timeout is not None else {}), + ) if not isinstance(value, dict): raise TypeError("Goal acceptance authority returned an invalid result") if value.get("status") not in { @@ -66,6 +88,33 @@ def _result(method: str, request: Mapping[str, Any]) -> dict[str, Any]: return value +def transition_goal_acceptance_lifecycle( + *, + registry_path: Path, + goal_id: str, + transition: Mapping[str, Any], + operation_id: str, +) -> dict[str, Any] | None: + """Commit one source-owned acceptance lifecycle transition if promoted.""" + root = resolve_runtime_root( + load_project_registry(registry_path), + registry_path=registry_path, + ) + if not local_authority_is_promoted(runtime_root=root, goal_id=goal_id): + return None + return _result( + "goal.acceptance.lifecycle.transition", + { + "runtime_root": str(root.resolve()), + "goal_id": goal_id, + "actor_agent_id": None, + "operation_id": operation_id, + "transition": dict(transition), + }, + timeout=CANONICAL_AUTHORITY_WRITE_TIMEOUT_SECONDS, + ) + + def inspect_goal_acceptance( *, registry_path: Path, diff --git a/loopx/control_plane/goals/acceptance_authority.ts b/loopx/control_plane/goals/acceptance_authority.ts index e0d217a50e..e62be27ee7 100644 --- a/loopx/control_plane/goals/acceptance_authority.ts +++ b/loopx/control_plane/goals/acceptance_authority.ts @@ -310,6 +310,20 @@ function planLifecycleTransition( reason: "A retiring Goal instance cannot be rebound as active."}; } if (current === null) { + if (transition.kind === "reconcile_recreated") { + const acceptance = readGoalAcceptance(head, goalId); + return { + kind: "commit", + lifecycle: { + schema_version: GOAL_ACCEPTANCE_LIFECYCLE_SCHEMA, + state: "active", + goal_ref: transition.goal_ref, + }, + acceptance: acceptance === null + ? null + : bindGoalAcceptanceStateOwner(acceptance, transition.retired_goal_ref), + }; + } return {kind: "failure", reason_code: "goal_acceptance_lifecycle_unbound", reason: "Bind the existing Goal instance before changing its acceptance lifecycle."}; } @@ -323,6 +337,10 @@ function planLifecycleTransition( : {kind: "commit", lifecycle: {...current, state: "retiring"}, acceptance: readGoalAcceptance(head, goalId)}; } + if (transition.kind === "reconcile_recreated" + && sameExactGoalRef(current.goal_ref, transition.goal_ref)) { + return {kind: "no_change", lifecycle: current}; + } if (current.state !== "retiring" || !sameExactGoalRef(current.goal_ref, transition.retired_goal_ref)) { return lifecycleMismatch("Goal acceptance successor activation requires its exact retiring predecessor."); diff --git a/loopx/control_plane/goals/acceptance_lifecycle.ts b/loopx/control_plane/goals/acceptance_lifecycle.ts index 886dce3d29..ab4a35294d 100644 --- a/loopx/control_plane/goals/acceptance_lifecycle.ts +++ b/loopx/control_plane/goals/acceptance_lifecycle.ts @@ -19,7 +19,7 @@ export type GoalAcceptanceLifecycleTransition = | Readonly<{kind: "bind_existing"; goal_ref: WireExactGoalRef}> | Readonly<{kind: "retire"; goal_ref: WireExactGoalRef}> | Readonly<{ - kind: "activate_successor"; + kind: "activate_successor" | "reconcile_recreated"; retired_goal_ref: WireExactGoalRef; goal_ref: WireExactGoalRef; }>; @@ -72,7 +72,7 @@ export function parseGoalAcceptanceLifecycleTransition(value: unknown): GoalAcce "goal acceptance lifecycle transition fields are missing or unsupported"); return {kind: raw.kind, goal_ref: parseWireExactGoalRef(raw.goal_ref, "transition goal_ref")}; } - requireLifecycle(raw.kind === "activate_successor", + requireLifecycle(raw.kind === "activate_successor" || raw.kind === "reconcile_recreated", "unsupported goal acceptance lifecycle transition"); requireLifecycle(Object.keys(raw).length === 3 && Object.hasOwn(raw, "retired_goal_ref") && Object.hasOwn(raw, "goal_ref"), @@ -83,5 +83,5 @@ export function parseGoalAcceptanceLifecycleTransition(value: unknown): GoalAcce "goal acceptance successor must preserve the Goal alias"); requireLifecycle(!sameExactGoalRef(retiredGoalRef, goalRef), "goal acceptance successor must use a new Goal instance"); - return {kind: "activate_successor", retired_goal_ref: retiredGoalRef, goal_ref: goalRef}; + return {kind: raw.kind, retired_goal_ref: retiredGoalRef, goal_ref: goalRef}; } diff --git a/loopx/control_plane/goals/source_session_recreation.py b/loopx/control_plane/goals/source_session_recreation.py index 8c10b8e51c..0ddfeeb6a6 100644 --- a/loopx/control_plane/goals/source_session_recreation.py +++ b/loopx/control_plane/goals/source_session_recreation.py @@ -1,7 +1,6 @@ from __future__ import annotations import copy -from contextlib import ExitStack from dataclasses import dataclass import hashlib import json @@ -21,6 +20,7 @@ source_session_registry_transaction, ) from ..runtime.time import now_local_iso +from .acceptance import transition_goal_acceptance_lifecycle from .source_session_registry_state import ( GOAL_INSTANCE_ID, alias_digest, @@ -69,6 +69,53 @@ def _canonical_writer_guard_path(registry_path: Path, goal_id: str) -> Path: return shadow_maintenance_lock_target(runtime_root.resolve(), goal_id) +def _acceptance_lifecycle_operation_id( + request: RecreateGoalRequest, + phase: str, +) -> str: + digest = hashlib.sha256( + f"{request.operation_id}\0{phase}".encode("utf-8") + ).hexdigest() + return f"goal-acceptance-lifecycle:{phase}:{digest}" + + +def _retire_acceptance_lifecycle( + request: RecreateGoalRequest, + *, + requested_goal_ref: dict[str, str], +) -> None: + transition_goal_acceptance_lifecycle( + registry_path=request.registry_path, + goal_id=request.goal_id, + operation_id=_acceptance_lifecycle_operation_id(request, "bind"), + transition={"kind": "bind_existing", "goal_ref": requested_goal_ref}, + ) + transition_goal_acceptance_lifecycle( + registry_path=request.registry_path, + goal_id=request.goal_id, + operation_id=_acceptance_lifecycle_operation_id(request, "retire"), + transition={"kind": "retire", "goal_ref": requested_goal_ref}, + ) + + +def _activate_recreated_acceptance_lifecycle( + request: RecreateGoalRequest, + *, + retired_goal_ref: dict[str, str], + goal_ref: dict[str, str], +) -> None: + transition_goal_acceptance_lifecycle( + registry_path=request.registry_path, + goal_id=request.goal_id, + operation_id=_acceptance_lifecycle_operation_id(request, "activate"), + transition={ + "kind": "reconcile_recreated", + "retired_goal_ref": retired_goal_ref, + "goal_ref": goal_ref, + }, + ) + + def _recreation_journal_path( registry_path: Path, *, @@ -345,6 +392,12 @@ def recreate_goal_instance(request: RecreateGoalRequest) -> dict[str, Any]: "reserved_at": replay_receipt["committed_at"], } write_journal(journal_path, {**journal, "phase": "published"}) + if state.active_goal_ref == state.reserved_goal_ref: + _activate_recreated_acceptance_lifecycle( + request, + retired_goal_ref=requested_goal_ref, + goal_ref=state.reserved_goal_ref, + ) return _recreation_result( request, receipt=replay_receipt, @@ -372,6 +425,10 @@ def recreate_goal_instance(request: RecreateGoalRequest) -> dict[str, Any]: goal_id=request.goal_id, gate=closing_gate, ) + _retire_acceptance_lifecycle( + request, + requested_goal_ref=requested_goal_ref, + ) drain_result = drain_releasable_source_turn_effects( registry_path=request.registry_path, @@ -386,140 +443,153 @@ def recreate_goal_instance(request: RecreateGoalRequest) -> dict[str, Any]: changed=gate_changed or drain_result.changed, ) - with ExitStack() as locks: - locks.enter_context(exclusive_cross_runtime_file_lock( - guard, - operation="source_session_goal_lifetime_publish", - )) - locks.enter_context(exclusive_cross_runtime_file_lock( + with exclusive_cross_runtime_file_lock( + guard, + operation="source_session_goal_lifetime_publish", + ): + with exclusive_cross_runtime_file_lock( _canonical_writer_guard_path(request.registry_path, request.goal_id), operation="source_session_goal_canonical_publish", timeout_seconds=CANONICAL_AUTHORITY_WRITE_TIMEOUT_SECONDS, - )) - journal = _read_recreation_journal(journal_path) - if journal is None: - raise RuntimeError("Goal recreation reservation journal is missing") - with source_session_registry_transaction( - request.registry_path, - operation="source_session_goal_recreate", - ) as transaction: - state = _evaluate_recreation( - transaction.payload_copy(), - request, - requested_goal_ref=requested_goal_ref, - request_digest=request_digest, - journal=journal, - ) - if state.decision["kind"] == "replay": - replay_receipt = state.decision.get("receipt") - if not isinstance(replay_receipt, dict): - raise RuntimeError("Goal recreation replay omitted its receipt") - next_gate = decide_source_turn_effect_repair_locked( - registry_path=request.registry_path, - goal_id=request.goal_id, + ): + journal = _read_recreation_journal(journal_path) + if journal is None: + raise RuntimeError("Goal recreation reservation journal is missing") + with source_session_registry_transaction( + request.registry_path, + operation="source_session_goal_recreate", + ) as transaction: + state = _evaluate_recreation( + transaction.payload_copy(), + request, requested_goal_ref=requested_goal_ref, - current_goal_ref=state.active_goal_ref, - reserved_goal_ref=state.reserved_goal_ref, - operation_id=request.operation_id, request_digest=request_digest, + journal=journal, ) - if next_gate is not None: - write_source_turn_effect_gate_locked( + if state.decision["kind"] == "replay": + replay_receipt = state.decision.get("receipt") + if not isinstance(replay_receipt, dict): + raise RuntimeError("Goal recreation replay omitted its receipt") + next_gate = decide_source_turn_effect_repair_locked( registry_path=request.registry_path, goal_id=request.goal_id, - gate=next_gate, + requested_goal_ref=requested_goal_ref, + current_goal_ref=state.active_goal_ref, + reserved_goal_ref=state.reserved_goal_ref, + operation_id=request.operation_id, + request_digest=request_digest, ) - write_journal(journal_path, {**journal, "phase": "published"}) - return _recreation_result( - request, - receipt=replay_receipt, - replayed=True, - ) - - next_gate = decide_source_turn_effect_publish_locked( - registry_path=request.registry_path, - goal_id=request.goal_id, - requested_goal_ref=requested_goal_ref, - current_goal_ref=state.active_goal_ref, - reserved_goal_ref=state.reserved_goal_ref, - operation_id=request.operation_id, - request_digest=request_digest, - ) - committed_at = now_local_iso() - retired_session_ids = sorted( - str(binding["session_id"]) for binding in state.retiring_bindings - ) - receipt = { - "schema_version": "loopx_goal_recreation_receipt_v1", - "operation_id": request.operation_id, - "request_digest": request_digest, - "retired_goal_ref": copy.deepcopy(requested_goal_ref), - "new_goal_ref": copy.deepcopy(state.reserved_goal_ref), - "retired_session_ids": retired_session_ids, - "committed_at": committed_at, - } - state.registry["goals"] = [ - { - **candidate, - "goal_instance_id": state.reserved_goal_ref["goal_instance_id"], - "execution_authority": False, - } - if candidate.get("id") == request.goal_id - else candidate - for candidate in required_list(state.registry, "goals") - ] - state.registry["session_bindings"] = [ - binding - for binding in state.bindings - if binding.get("foreground_goal_ref") != requested_goal_ref - ] - session_receipts = state.session_receipts - if retired_session_ids: - session_receipts = [ - *session_receipts, - { - "schema_version": ( - "loopx_source_session_retirement_receipt_v1" - ), - "operation": "retire_bindings", + if next_gate is not None: + write_source_turn_effect_gate_locked( + registry_path=request.registry_path, + goal_id=request.goal_id, + gate=next_gate, + ) + write_journal(journal_path, {**journal, "phase": "published"}) + result = _recreation_result( + request, + receipt=replay_receipt, + replayed=True, + ) + else: + next_gate = decide_source_turn_effect_publish_locked( + registry_path=request.registry_path, + goal_id=request.goal_id, + requested_goal_ref=requested_goal_ref, + current_goal_ref=state.active_goal_ref, + reserved_goal_ref=state.reserved_goal_ref, + operation_id=request.operation_id, + request_digest=request_digest, + ) + committed_at = now_local_iso() + retired_session_ids = sorted( + str(binding["session_id"]) + for binding in state.retiring_bindings + ) + receipt = { + "schema_version": "loopx_goal_recreation_receipt_v1", "operation_id": request.operation_id, "request_digest": request_digest, "retired_goal_ref": copy.deepcopy(requested_goal_ref), - "session_ids": retired_session_ids, + "new_goal_ref": copy.deepcopy(state.reserved_goal_ref), + "retired_session_ids": retired_session_ids, "committed_at": committed_at, - }, - ] - state.registry["session_receipts"] = session_receipts - retired = state.registry.get("retired_goal_instances", []) - if not isinstance(retired, list) or any( - not isinstance(item, dict) for item in retired - ): - raise ValueError( - "source-session retired_goal_instances must be a list of objects" - ) - state.registry["retired_goal_instances"] = [ - *retired, - { - "goal_ref": copy.deepcopy(requested_goal_ref), - "successor_goal_ref": copy.deepcopy(state.reserved_goal_ref), - "operation_id": request.operation_id, - "retired_at": committed_at, - }, - ] - state.registry["lifetime_receipts"] = [ - *state.lifetime_receipts, - receipt, - ] - state.registry["updated_at"] = committed_at - transaction.commit(state.registry) - write_source_turn_effect_gate_locked( - registry_path=request.registry_path, - goal_id=request.goal_id, - gate=next_gate, - ) - write_journal(journal_path, {**journal, "phase": "published"}) - return _recreation_result( - request, - receipt=receipt, - replayed=False, - ) + } + state.registry["goals"] = [ + { + **candidate, + "goal_instance_id": ( + state.reserved_goal_ref["goal_instance_id"] + ), + "execution_authority": False, + } + if candidate.get("id") == request.goal_id + else candidate + for candidate in required_list(state.registry, "goals") + ] + state.registry["session_bindings"] = [ + binding + for binding in state.bindings + if binding.get("foreground_goal_ref") != requested_goal_ref + ] + session_receipts = state.session_receipts + if retired_session_ids: + session_receipts = [ + *session_receipts, + { + "schema_version": ( + "loopx_source_session_retirement_receipt_v1" + ), + "operation": "retire_bindings", + "operation_id": request.operation_id, + "request_digest": request_digest, + "retired_goal_ref": copy.deepcopy( + requested_goal_ref + ), + "session_ids": retired_session_ids, + "committed_at": committed_at, + }, + ] + state.registry["session_receipts"] = session_receipts + retired = state.registry.get("retired_goal_instances", []) + if not isinstance(retired, list) or any( + not isinstance(item, dict) for item in retired + ): + raise ValueError( + "source-session retired_goal_instances must be a list of objects" + ) + state.registry["retired_goal_instances"] = [ + *retired, + { + "goal_ref": copy.deepcopy(requested_goal_ref), + "successor_goal_ref": copy.deepcopy( + state.reserved_goal_ref + ), + "operation_id": request.operation_id, + "retired_at": committed_at, + }, + ] + state.registry["lifetime_receipts"] = [ + *state.lifetime_receipts, + receipt, + ] + state.registry["updated_at"] = committed_at + transaction.commit(state.registry) + write_source_turn_effect_gate_locked( + registry_path=request.registry_path, + goal_id=request.goal_id, + gate=next_gate, + ) + write_journal(journal_path, {**journal, "phase": "published"}) + result = _recreation_result( + request, + receipt=receipt, + replayed=False, + ) + + _activate_recreated_acceptance_lifecycle( + request, + retired_goal_ref=requested_goal_ref, + goal_ref=state.reserved_goal_ref, + ) + return result diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index 420fc348cd..55ce7f338f 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -1141,6 +1141,14 @@ "api": "load_project_registry", "classification": "codec_api" }, + { + "site": "loopx/control_plane/goals/acceptance.py::.transition_goal_acceptance_lifecycle::codec_read:load_project_registry#1", + "line": 100, + "column": 9, + "kind": "codec_read", + "api": "load_project_registry", + "classification": "codec_api" + }, { "site": "loopx/control_plane/goals/activation_service.py::._projected_source_identity::codec_read:load_registry#1", "line": 90, @@ -1455,7 +1463,7 @@ }, { "site": "loopx/control_plane/goals/source_session_recreation.py::.recreate_goal_instance::codec_transaction:source_session_registry_transaction#1", - "line": 307, + "line": 354, "column": 14, "kind": "codec_transaction", "api": "source_session_registry_transaction", @@ -1463,8 +1471,8 @@ }, { "site": "loopx/control_plane/goals/source_session_recreation.py::.recreate_goal_instance::codec_transaction:source_session_registry_transaction#2", - "line": 402, - "column": 14, + "line": 458, + "column": 18, "kind": "codec_transaction", "api": "source_session_registry_transaction", "classification": "codec_api" diff --git a/tests/control_plane/test_goal_acceptance_source_recreation.py b/tests/control_plane/test_goal_acceptance_source_recreation.py new file mode 100644 index 0000000000..468850ac61 --- /dev/null +++ b/tests/control_plane/test_goal_acceptance_source_recreation.py @@ -0,0 +1,404 @@ +from __future__ import annotations + +from contextlib import contextmanager +import hashlib +from pathlib import Path +from typing import Any + +import pytest +from canonical_authority_fixture import ( + initialize_canonical_authority, + isolate_sqlite_runtime, +) + +from loopx.control_plane.coordination.local_authority_shadow_projection import ( + canonical_bytes, +) +from loopx.control_plane.coordination.runtime_shadow import ( + build_todo_runtime_shadow_projection, +) +from loopx.control_plane.effect_runtime import effect_runtime_result +from loopx.control_plane.goals import source_session_recreation +from loopx.control_plane.goals.acceptance import ( + configure_goal_acceptance, + inspect_goal_acceptance, +) +from loopx.control_plane.goals.source_session_recreation import ( + RecreateGoalRequest, + recreate_goal_instance, +) +from loopx.control_plane.projects.registry_codec import ( + load_project_registry, + source_session_registry_transaction, +) +from loopx.control_plane.testing.canary_harness import write_fixture_registry + + +GOAL_ID = "goal-acceptance" +INSTANCE_A = "ginst_aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" + + +def _acceptance_document() -> dict[str, Any]: + return { + "scope": {"kind": "all_advancement"}, + "objective": "Keep acceptance bound to its originating Goal instance", + "non_goals": [], + "criteria": [ + { + "id": "proof", + "description": "The independent check passes", + "validation_argv": ["true"], + "validation_timeout_seconds": 5, + "validation_files": [], + } + ], + "bindings": [], + } + + +def _canonical_projection(*, legacy_acceptance: bool) -> dict[str, Any]: + projection = build_todo_runtime_shadow_projection( + goal_id=GOAL_ID, + todos=[], + ) + if legacy_acceptance: + document = _acceptance_document() + projection["goal_acceptance"] = { + "schema_version": "loopx_goal_acceptance_v0", + "enabled": True, + "revision": 1, + "digest": hashlib.sha256(canonical_bytes(document)).hexdigest(), + "document": document, + "bindings": [], + "verification": None, + } + return projection + + +def _write_source_registry(path: Path) -> None: + payload = { + "schema_version": "0.2", + "registry_role": "project-local", + "profile_id": "source_session_v1", + "common_runtime_root": str(path.parent), + "projects": [], + "goals": [ + { + "id": GOAL_ID, + "goal_instance_id": INSTANCE_A, + "status": "active", + "execution_authority": False, + } + ], + "session_bindings": [], + "session_receipts": [], + "lifetime_receipts": [], + "retired_goal_instances": [], + } + with source_session_registry_transaction( + path, + operation="acceptance_source_recreation_fixture", + create=lambda: payload, + ) as transaction: + transaction.commit(transaction.payload_copy()) + + +def _promote( + *, + tmp_path: Path, + registry_path: Path, + provider: str, + legacy_acceptance: bool, +) -> Path: + runtime_root = Path( + str(load_project_registry(registry_path)["common_runtime_root"]) + ) + state_path = tmp_path / "acceptance-source-state.md" + state_path.write_text("# Acceptance source fixture\n", encoding="utf-8") + initialize_canonical_authority( + runtime_root, + GOAL_ID, + _canonical_projection(legacy_acceptance=legacy_acceptance), + state_path=state_path, + provider=provider, + ) + return runtime_root + + +def _inspect(runtime_root: Path, goal_ref: dict[str, str]) -> dict[str, Any]: + result = effect_runtime_result( + "goal.acceptance.inspect", + { + "runtime_root": str(runtime_root.resolve()), + "goal_id": GOAL_ID, + "goal_ref": goal_ref, + }, + ) + assert isinstance(result, dict) + return result + + +@pytest.mark.parametrize("provider", ["file", "sqlite"]) +def test_python_acceptance_routing_uses_registered_goal_instance( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, + provider: str, +) -> None: + if provider == "sqlite": + isolate_sqlite_runtime(tmp_path, monkeypatch) + project = tmp_path / "project" + project.mkdir() + state_path = project / "state.md" + state_path.write_text("# Goal\n", encoding="utf-8") + runtime_root = tmp_path / "runtime" + registry_path = tmp_path / "registry.json" + write_fixture_registry( + project=project, + runtime_root=runtime_root, + registry_path=registry_path, + goal_id=GOAL_ID, + domain="acceptance", + adapter_kind="generic_project_goal_v0", + state_file=str(state_path), + extra_goal_fields={"goal_instance_id": INSTANCE_A}, + ) + initialize_canonical_authority( + runtime_root, + GOAL_ID, + _canonical_projection(legacy_acceptance=False), + state_path=state_path, + provider=provider, + ) + + before = inspect_goal_acceptance( + registry_path=registry_path, + goal_id=GOAL_ID, + ) + applied = configure_goal_acceptance( + registry_path=registry_path, + goal_id=GOAL_ID, + expected_provider_revision=str(before["provider_revision"]), + document=_acceptance_document(), + operation_id="configure-exact-acceptance", + execute=True, + ) + + assert applied["status"] == "applied" + assert applied["goal_acceptance_contract"]["goal_ref"] == { + "goal_id": GOAL_ID, + "goal_instance_id": INSTANCE_A, + } + assert applied["goal_acceptance_contract"]["lifecycle_state"] == "active" + + +@pytest.mark.parametrize("provider", ["file", "sqlite"]) +def test_promoted_source_recreation_fences_old_acceptance( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, + provider: str, +) -> None: + if provider == "sqlite": + isolate_sqlite_runtime(tmp_path, monkeypatch) + registry_path = tmp_path / "project" / ".loopx" / "registry.json" + _write_source_registry(registry_path) + runtime_root = _promote( + tmp_path=tmp_path, + registry_path=registry_path, + provider=provider, + legacy_acceptance=True, + ) + request = RecreateGoalRequest( + registry_path=registry_path, + goal_id=GOAL_ID, + goal_instance_id=INSTANCE_A, + operation_id="recreate-acceptance-a-to-b", + ) + + recreated = recreate_goal_instance(request) + goal_ref = recreated["goal_ref"] + current = _inspect(runtime_root, goal_ref) + + assert recreated["ok"] is True + assert current["status"] == "loaded" + assert current["contract"] is None + assert current["goal_acceptance_contract"] == { + "enabled": False, + "lifecycle_state": "active", + "goal_ref": goal_ref, + } + + +@pytest.mark.parametrize("provider", ["file", "sqlite"]) +def test_recreation_retry_repairs_acceptance_after_source_publication( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, + provider: str, +) -> None: + if provider == "sqlite": + isolate_sqlite_runtime(tmp_path, monkeypatch) + registry_path = tmp_path / "project" / ".loopx" / "registry.json" + _write_source_registry(registry_path) + runtime_root = _promote( + tmp_path=tmp_path, + registry_path=registry_path, + provider=provider, + legacy_acceptance=True, + ) + request = RecreateGoalRequest( + registry_path=registry_path, + goal_id=GOAL_ID, + goal_instance_id=INSTANCE_A, + operation_id="recreate-with-activation-retry", + ) + transition = source_session_recreation.transition_goal_acceptance_lifecycle + lock = source_session_recreation.exclusive_cross_runtime_file_lock + held_operations: list[str] = [] + activation_lock_snapshots: list[tuple[str, ...]] = [] + activation_failed = False + + @contextmanager + def track_lock(path: Path, **kwargs: Any): + operation = str(kwargs.get("operation")) + with lock(path, **kwargs) as lock_path: + held_operations.append(operation) + try: + yield lock_path + finally: + held_operations.remove(operation) + + def fail_first_activation(**kwargs: Any) -> dict[str, Any] | None: + nonlocal activation_failed + if kwargs["transition"]["kind"] == "reconcile_recreated": + activation_lock_snapshots.append(tuple(held_operations)) + if not activation_failed: + activation_failed = True + raise RuntimeError("injected acceptance activation failure") + return transition(**kwargs) + + monkeypatch.setattr( + source_session_recreation, + "exclusive_cross_runtime_file_lock", + track_lock, + ) + monkeypatch.setattr( + source_session_recreation, + "transition_goal_acceptance_lifecycle", + fail_first_activation, + ) + with pytest.raises(RuntimeError, match="injected acceptance activation failure"): + recreate_goal_instance(request) + + registry = load_project_registry(registry_path) + goal_ref = { + "goal_id": GOAL_ID, + "goal_instance_id": registry["goals"][0]["goal_instance_id"], + } + assert goal_ref["goal_instance_id"] != INSTANCE_A + retiring = _inspect( + runtime_root, + {"goal_id": GOAL_ID, "goal_instance_id": INSTANCE_A}, + ) + assert retiring["goal_acceptance_contract"]["lifecycle_state"] == "retiring" + + replayed = recreate_goal_instance(request) + current = _inspect(runtime_root, goal_ref) + + assert replayed["replayed"] is True + assert activation_lock_snapshots == [ + ("source_session_goal_lifetime_publish",), + ("source_session_goal_lifetime_close",), + ] + assert current["status"] == "loaded" + assert current["goal_acceptance_contract"] == { + "enabled": False, + "lifecycle_state": "active", + "goal_ref": goal_ref, + } + + +@pytest.mark.parametrize("provider", ["file", "sqlite"]) +def test_historical_recreation_migrates_legacy_acceptance_to_retired_instance( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, + provider: str, +) -> None: + if provider == "sqlite": + isolate_sqlite_runtime(tmp_path, monkeypatch) + registry_path = tmp_path / "project" / ".loopx" / "registry.json" + _write_source_registry(registry_path) + request = RecreateGoalRequest( + registry_path=registry_path, + goal_id=GOAL_ID, + goal_instance_id=INSTANCE_A, + operation_id="historical-recreation-a-to-b", + ) + recreated = recreate_goal_instance(request) + goal_ref = recreated["goal_ref"] + runtime_root = _promote( + tmp_path=tmp_path, + registry_path=registry_path, + provider=provider, + legacy_acceptance=True, + ) + + replayed = recreate_goal_instance(request) + current = _inspect(runtime_root, goal_ref) + + assert replayed["replayed"] is True + assert current["status"] == "loaded" + assert current["contract"] is None + assert current["goal_acceptance_contract"] == { + "enabled": False, + "lifecycle_state": "active", + "goal_ref": goal_ref, + } + + +@pytest.mark.parametrize("provider", ["file", "sqlite"]) +def test_historical_replay_does_not_replace_a_later_acceptance_successor( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, + provider: str, +) -> None: + if provider == "sqlite": + isolate_sqlite_runtime(tmp_path, monkeypatch) + registry_path = tmp_path / "project" / ".loopx" / "registry.json" + _write_source_registry(registry_path) + runtime_root = _promote( + tmp_path=tmp_path, + registry_path=registry_path, + provider=provider, + legacy_acceptance=True, + ) + first_request = RecreateGoalRequest( + registry_path=registry_path, + goal_id=GOAL_ID, + goal_instance_id=INSTANCE_A, + operation_id="recreate-a-to-b", + ) + first = recreate_goal_instance(first_request) + goal_b = first["goal_ref"] + second = recreate_goal_instance( + RecreateGoalRequest( + registry_path=registry_path, + goal_id=GOAL_ID, + goal_instance_id=goal_b["goal_instance_id"], + operation_id="recreate-b-to-c", + ) + ) + goal_c = second["goal_ref"] + + replayed = recreate_goal_instance(first_request) + current = _inspect(runtime_root, goal_c) + + assert replayed["replayed"] is True + assert replayed["goal_ref"] == goal_b + assert load_project_registry(registry_path)["goals"][0]["goal_instance_id"] == ( + goal_c["goal_instance_id"] + ) + assert current["status"] == "loaded" + assert current["goal_acceptance_contract"] == { + "enabled": False, + "lifecycle_state": "active", + "goal_ref": goal_c, + } diff --git a/tests/control_plane_ts/goal_acceptance_authority.test.ts b/tests/control_plane_ts/goal_acceptance_authority.test.ts index 04828750e0..53842453ed 100644 --- a/tests/control_plane_ts/goal_acceptance_authority.test.ts +++ b/tests/control_plane_ts/goal_acceptance_authority.test.ts @@ -247,6 +247,51 @@ for (const provider of providers) { assert.deepEqual((current.head.goal_acceptance as JsonObject).owner_goal_ref, goalA); assert.equal((await inspectGoalAcceptance(store, goal, undefined, goalA)).status, "loaded"); }); + test(`${provider}: historical recreation migration keeps legacy acceptance with the retired instance`, options, async t => { + const store = await fixture(t, provider); await seed(store); + assert.equal((await configureGoalAcceptance(store, await configureRequest(store))).status, "applied"); + assert.equal((await lifecycleTransition(store, { + kind: "reconcile_recreated", + retired_goal_ref: goalA, + goal_ref: goalB, + })).status, "applied"); + const migrated = await head(store); + assert.deepEqual(migrated.head.goal_acceptance_lifecycle, { + schema_version: "loopx_goal_acceptance_lifecycle_v0", + state: "active", + goal_ref: goalB, + }); + assert.deepEqual((migrated.head.goal_acceptance as JsonObject).owner_goal_ref, goalA); + assert.deepEqual(projectGoalAcceptance(migrated.head, goal), { + enabled: false, + lifecycle_state: "active", + goal_ref: goalB, + }); + assert.equal(acceptanceWorkGuard(migrated.head, goal, "todo_first"), null); + }); + test(`${provider}: historical recreation replay does not reactivate a retiring successor`, options, async t => { + const store = await fixture(t, provider); await seed(store); + assert.equal((await lifecycleTransition(store, { + kind: "reconcile_recreated", + retired_goal_ref: goalA, + goal_ref: goalB, + })).status, "applied"); + assert.equal((await lifecycleTransition(store, { + kind: "retire", + goal_ref: goalB, + })).status, "applied"); + const replay = await lifecycleTransition(store, { + kind: "reconcile_recreated", + retired_goal_ref: goalA, + goal_ref: goalB, + }); + assert.equal(replay.status, "no_change"); + assert.deepEqual(replay.goal_acceptance_lifecycle, { + schema_version: "loopx_goal_acceptance_lifecycle_v0", + state: "retiring", + goal_ref: goalB, + }); + }); test(`${provider}: exact lifecycle migration fences retirement and same-alias replacement`, options, async t => { const store = await fixture(t, provider); await seed(store); const legacyConfigure = await configureRequest(store); From 98f3fc3bb3068d09b0f5895bc670133b8e7de852 Mon Sep 17 00:00:00 2001 From: Duang777 Date: Fri, 9 Oct 2026 09:15:30 +0800 Subject: [PATCH 2/3] test: observe acceptance fence during recreation Signed-off-by: Duang777 --- .../test_shadow_native_todo_update_e2e.py | 36 ++++++++++++------- 1 file changed, 24 insertions(+), 12 deletions(-) diff --git a/tests/control_plane/test_shadow_native_todo_update_e2e.py b/tests/control_plane/test_shadow_native_todo_update_e2e.py index 38db0a66c4..8fed031452 100644 --- a/tests/control_plane/test_shadow_native_todo_update_e2e.py +++ b/tests/control_plane/test_shadow_native_todo_update_e2e.py @@ -7,7 +7,6 @@ """ from __future__ import annotations -from contextlib import contextmanager from dataclasses import dataclass import hashlib import json @@ -409,23 +408,34 @@ def test_goal_recreation_waits_for_an_admitted_canonical_update( authority_reason=None, ) writer = start(node_command("update_paused", request)) - publish_started = threading.Event() + fence_started = threading.Event() recreation_finished = threading.Event() recreated: list[dict] = [] recreation_errors: list[BaseException] = [] - actual_lock = source_session_recreation.exclusive_cross_runtime_file_lock + actual_transition = ( + source_session_recreation.transition_goal_acceptance_lifecycle + ) - @contextmanager - def observed_lock(path: Path, **options: object): - if options.get("operation") == "source_session_goal_lifetime_publish": - publish_started.set() - with actual_lock(path, **options): - yield + def observed_transition( + *, + registry_path: Path, + goal_id: str, + transition: dict[str, object], + operation_id: str, + ) -> dict[str, object] | None: + if transition.get("kind") == "bind_existing": + fence_started.set() + return actual_transition( + registry_path=registry_path, + goal_id=goal_id, + transition=transition, + operation_id=operation_id, + ) monkeypatch.setattr( source_session_recreation, - "exclusive_cross_runtime_file_lock", - observed_lock, + "transition_goal_acceptance_lifecycle", + observed_transition, ) def recreate() -> None: @@ -447,7 +457,9 @@ def recreate() -> None: try: expect_barrier(writer, "commit") thread.start() - assert publish_started.wait(timeout=5), "Goal recreation did not reach publication" + assert fence_started.wait(timeout=5), ( + "Goal recreation did not reach acceptance fencing" + ) assert not recreation_finished.wait(timeout=0.2), ( "Goal B published while Goal A's admitted canonical update was paused " "before provider commit" From b32ba92c9939f4e912b320265974ba7aed7a6c50 Mon Sep 17 00:00:00 2001 From: Duang777 Date: Fri, 9 Oct 2026 13:18:50 +0800 Subject: [PATCH 3/3] fix(acceptance): route lifecycle transitions through source owner Signed-off-by: Duang777 --- loopx/control_plane/goals/acceptance.py | 9 ++------- .../goals/source_session_recreation.py | 19 +++++++++++++------ .../project_registry_io_manifest_v1.json | 14 +++----------- .../test_shadow_native_todo_update_e2e.py | 4 ++-- 4 files changed, 20 insertions(+), 26 deletions(-) diff --git a/loopx/control_plane/goals/acceptance.py b/loopx/control_plane/goals/acceptance.py index 13814a8045..64f9aad348 100644 --- a/loopx/control_plane/goals/acceptance.py +++ b/loopx/control_plane/goals/acceptance.py @@ -14,14 +14,12 @@ from uuid import uuid4 from ...agent_registry import load_goal_from_registry, registered_agent_ids_for_goal -from ...paths import resolve_runtime_root from ..coordination.local_authority import local_authority_is_promoted from ..coordination.local_authority_shadow_adapter import effective_runtime_root from ..effect_runtime import ( CANONICAL_AUTHORITY_WRITE_TIMEOUT_SECONDS, effect_runtime_result, ) -from ..projects.registry_codec import load_project_registry from ..todos.completion_validation import ( _resolve_completion_validation_workspace, run_declared_completion_validation_effect, @@ -90,16 +88,13 @@ def _result( def transition_goal_acceptance_lifecycle( *, - registry_path: Path, + runtime_root: Path, goal_id: str, transition: Mapping[str, Any], operation_id: str, ) -> dict[str, Any] | None: """Commit one source-owned acceptance lifecycle transition if promoted.""" - root = resolve_runtime_root( - load_project_registry(registry_path), - registry_path=registry_path, - ) + root = runtime_root.resolve() if not local_authority_is_promoted(runtime_root=root, goal_id=goal_id): return None return _result( diff --git a/loopx/control_plane/goals/source_session_recreation.py b/loopx/control_plane/goals/source_session_recreation.py index 0ddfeeb6a6..6a91f84ac9 100644 --- a/loopx/control_plane/goals/source_session_recreation.py +++ b/loopx/control_plane/goals/source_session_recreation.py @@ -63,10 +63,16 @@ class _RecreationState: decision: dict[str, Any] -def _canonical_writer_guard_path(registry_path: Path, goal_id: str) -> Path: +def _canonical_runtime_root(registry_path: Path) -> Path: registry = load_project_registry(registry_path) - runtime_root = resolve_runtime_root(registry, registry_path=registry_path) - return shadow_maintenance_lock_target(runtime_root.resolve(), goal_id) + return resolve_runtime_root(registry, registry_path=registry_path).resolve() + + +def _canonical_writer_guard_path(registry_path: Path, goal_id: str) -> Path: + return shadow_maintenance_lock_target( + _canonical_runtime_root(registry_path), + goal_id, + ) def _acceptance_lifecycle_operation_id( @@ -84,14 +90,15 @@ def _retire_acceptance_lifecycle( *, requested_goal_ref: dict[str, str], ) -> None: + runtime_root = _canonical_runtime_root(request.registry_path) transition_goal_acceptance_lifecycle( - registry_path=request.registry_path, + runtime_root=runtime_root, goal_id=request.goal_id, operation_id=_acceptance_lifecycle_operation_id(request, "bind"), transition={"kind": "bind_existing", "goal_ref": requested_goal_ref}, ) transition_goal_acceptance_lifecycle( - registry_path=request.registry_path, + runtime_root=runtime_root, goal_id=request.goal_id, operation_id=_acceptance_lifecycle_operation_id(request, "retire"), transition={"kind": "retire", "goal_ref": requested_goal_ref}, @@ -105,7 +112,7 @@ def _activate_recreated_acceptance_lifecycle( goal_ref: dict[str, str], ) -> None: transition_goal_acceptance_lifecycle( - registry_path=request.registry_path, + runtime_root=_canonical_runtime_root(request.registry_path), goal_id=request.goal_id, operation_id=_acceptance_lifecycle_operation_id(request, "activate"), transition={ diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index 8767965f24..cfb50d0f68 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -1141,14 +1141,6 @@ "api": "load_project_registry", "classification": "codec_api" }, - { - "site": "loopx/control_plane/goals/acceptance.py::.transition_goal_acceptance_lifecycle::codec_read:load_project_registry#1", - "line": 100, - "column": 9, - "kind": "codec_read", - "api": "load_project_registry", - "classification": "codec_api" - }, { "site": "loopx/control_plane/goals/activation_service.py::._projected_source_identity::codec_read:load_registry#1", "line": 90, @@ -1454,7 +1446,7 @@ "classification": "codec_api" }, { - "site": "loopx/control_plane/goals/source_session_recreation.py::._canonical_writer_guard_path::codec_read:load_project_registry#1", + "site": "loopx/control_plane/goals/source_session_recreation.py::._canonical_runtime_root::codec_read:load_project_registry#1", "line": 67, "column": 16, "kind": "codec_read", @@ -1463,7 +1455,7 @@ }, { "site": "loopx/control_plane/goals/source_session_recreation.py::.recreate_goal_instance::codec_transaction:source_session_registry_transaction#1", - "line": 354, + "line": 361, "column": 14, "kind": "codec_transaction", "api": "source_session_registry_transaction", @@ -1471,7 +1463,7 @@ }, { "site": "loopx/control_plane/goals/source_session_recreation.py::.recreate_goal_instance::codec_transaction:source_session_registry_transaction#2", - "line": 458, + "line": 465, "column": 18, "kind": "codec_transaction", "api": "source_session_registry_transaction", diff --git a/tests/control_plane/test_shadow_native_todo_update_e2e.py b/tests/control_plane/test_shadow_native_todo_update_e2e.py index 8fed031452..9a75ec2578 100644 --- a/tests/control_plane/test_shadow_native_todo_update_e2e.py +++ b/tests/control_plane/test_shadow_native_todo_update_e2e.py @@ -418,7 +418,7 @@ def test_goal_recreation_waits_for_an_admitted_canonical_update( def observed_transition( *, - registry_path: Path, + runtime_root: Path, goal_id: str, transition: dict[str, object], operation_id: str, @@ -426,7 +426,7 @@ def observed_transition( if transition.get("kind") == "bind_existing": fence_started.set() return actual_transition( - registry_path=registry_path, + runtime_root=runtime_root, goal_id=goal_id, transition=transition, operation_id=operation_id,