diff --git a/loopx/global_registry.py b/loopx/global_registry.py index 9afec2ec7e..72a5260f8c 100644 --- a/loopx/global_registry.py +++ b/loopx/global_registry.py @@ -10,6 +10,10 @@ from typing import Any, Callable from .authority import compact_authority_registry +from .control_plane.coordination.shadow_management import ( + shadow_maintenance_lock_target, +) +from .control_plane.effect_runtime import CANONICAL_AUTHORITY_WRITE_TIMEOUT_SECONDS from .control_plane.projects.contract import validate_project_record_bindings from .control_plane.projects.registry_codec import load_registry from .control_plane.runtime.time import now_local_iso @@ -578,16 +582,29 @@ def retire_global_registry_goals( "recommended_action": writability.get("recommended_action"), } - mutation = mutate_global_registry( - global_path, - "retire_global_registry_goals", - lambda current: _retire_global_registry_reduction( - current, - requested_ids=requested_ids, - global_path=global_path, - updated_at=updated_at, - ), - ) + canonical_runtime_root = global_path.expanduser().resolve().parent + with ExitStack() as stack: + for goal_id in sorted(requested_ids): + stack.enter_context( + exclusive_cross_runtime_file_lock( + shadow_maintenance_lock_target( + canonical_runtime_root, + goal_id, + ), + operation="retire_global_registry_goal_canonical", + timeout_seconds=CANONICAL_AUTHORITY_WRITE_TIMEOUT_SECONDS, + ) + ) + mutation = mutate_global_registry( + global_path, + "retire_global_registry_goals", + lambda current: _retire_global_registry_reduction( + current, + requested_ids=requested_ids, + global_path=global_path, + updated_at=updated_at, + ), + ) receipt = mutation["receipt"] backup_path = mutation["backup_path"] diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index 7c6f406ce8..0f42dcc912 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -1975,7 +1975,7 @@ }, { "site": "loopx/global_registry.py::._sync_project_registry_to_global_once::codec_read:load_registry#1", - "line": 649, + "line": 666, "column": 24, "kind": "codec_read", "api": "load_registry", @@ -1983,7 +1983,7 @@ }, { "site": "loopx/global_registry.py::.sync_project_registry_to_global::codec_read:load_registry#1", - "line": 857, + "line": 874, "column": 22, "kind": "codec_read", "api": "load_registry", diff --git a/tests/control_plane/test_global_goal_retirement_writer_fence.py b/tests/control_plane/test_global_goal_retirement_writer_fence.py new file mode 100644 index 0000000000..6930653d98 --- /dev/null +++ b/tests/control_plane/test_global_goal_retirement_writer_fence.py @@ -0,0 +1,249 @@ +from __future__ import annotations + +import hashlib +import json +from pathlib import Path +import select +import subprocess +import threading +from typing import Any + +import pytest + +from canonical_authority_fixture import initialize_canonical_authority +from loopx import global_registry as global_registry_module +from loopx.control_plane.coordination.coordination_state_contract import ( + TODO_DOMAIN_ITEM_SCHEMA_VERSION, + TODO_DOMAIN_READ_RECORD_SCHEMA_VERSION, + TODO_DOMAIN_RECORD_FIELDS, +) +from loopx.control_plane.coordination.local_authority_shadow_projection import ( + canonical_bytes, +) +from loopx.global_registry import ( + retire_global_registry_goals, + sync_project_registry_to_global, +) +from loopx.paths import global_registry_path +from loopx.registry_writability import probe_registry_write_path + + +REPO = Path(__file__).resolve().parents[2] +GOAL_ID = "goal-retirement-writer" +TODO_ID = "todo-retirement-writer" + +NODE = r""" +import {once} from "node:events"; +const input = JSON.parse(process.argv[1]); +const base = new URL(input.module_base); +const {openLocalAuthorityStore} = await import(new URL("local_authority_provider.ts", base)); +const {updateLocalCoordinationTodo} = await import(new URL("local_authority_runtime.ts", base)); +const store = await openLocalAuthorityStore(input.request.runtime_root, input.request.goal_id); +if (input.mode === "inspect") { + process.stdout.write(JSON.stringify(await store.loadAuthority()) + "\n"); +} else { + const actualCommit = store.commitAuthority.bind(store); + store.commitAuthority = async (commit) => { + process.stdout.write("BARRIER commit\n"); + await once(process.stdin, "data"); + return await actualCommit(commit); + }; + process.stdout.write(JSON.stringify( + await updateLocalCoordinationTodo(input.request, {createStore: () => store}), + ) + "\n"); +} +""" + + +def _node_command(mode: str, request: dict[str, object]) -> list[str]: + module_base = (REPO / "loopx/control_plane/coordination").as_uri() + "/" + return [ + "node", + "--no-warnings", + "--experimental-sqlite", + "--experimental-strip-types", + "--input-type=module", + "-e", + NODE, + json.dumps({"mode": mode, "request": request, "module_base": module_base}), + ] + + +def _stop(process: subprocess.Popen[str]) -> None: + if process.poll() is None: + process.terminate() + try: + process.communicate(timeout=5) + except subprocess.TimeoutExpired: + process.kill() + process.communicate(timeout=5) + + +@pytest.mark.parametrize("provider", ["file", "sqlite"]) +def test_global_goal_retirement_waits_for_an_admitted_canonical_update( + tmp_path: Path, + provider: str, + monkeypatch: pytest.MonkeyPatch, +) -> None: + project = tmp_path / "project" + registry = project / ".loopx" / "registry.json" + state = project / "ACTIVE_GOAL_STATE.md" + registry.parent.mkdir(parents=True) + state.write_text("# Goal\n", encoding="utf-8") + registry.write_text( + json.dumps( + { + "schema_version": 1, + "common_runtime_root": str(tmp_path / "runtime"), + "goals": [ + { + "id": GOAL_ID, + "repo": str(project), + "state_file": state.name, + "coordination": { + "registered_agents": ["agent-a", "agent-b"] + }, + } + ], + } + ), + encoding="utf-8", + ) + runtime = tmp_path / "runtime" + todo = { + "schema_version": TODO_DOMAIN_ITEM_SCHEMA_VERSION, + "todo_id": TODO_ID, + "text": "Before retirement", + "role": "agent", + "status": "open", + "done": False, + "archive_state": "active", + "claimed_by": "agent-a", + "note": None, + "required_capabilities": [], + "excluded_agents": [], + "evidence": None, + } + projection = { + "goal_id": GOAL_ID, + "handoff_mode": "soft_claim", + "todos": [todo], + "leases": [], + "todo_read_model": { + "schema_version": TODO_DOMAIN_READ_RECORD_SCHEMA_VERSION, + "todo_count": 1, + "records_sha256": hashlib.sha256(canonical_bytes([todo])).hexdigest(), + "contract_fields": list(TODO_DOMAIN_RECORD_FIELDS), + }, + } + initialize_canonical_authority( + runtime, + GOAL_ID, + projection, + state_path=state, + provider=provider, + ) + synced = sync_project_registry_to_global( + registry_path=registry, + runtime_root_override=str(runtime), + dry_run=False, + ) + assert synced["ok"] is True, synced + + request = { + "schema_version": "loopx_local_coordination_todo_update_request_v0", + "runtime_root": str(runtime), + "goal_id": GOAL_ID, + "todo_id": TODO_ID, + "role": "agent", + "actor_agent_id": "agent-a", + "registered_agents": ["agent-a", "agent-b"], + "operation_id": f"update-before-{provider}-retirement", + "patch": {"text": "Committed before retirement"}, + "clear_fields": [], + "dry_run": False, + "observed_at": "2026-10-07T06:00:00Z", + } + writer = subprocess.Popen( + _node_command("update", request), + cwd=REPO, + stdin=subprocess.PIPE, + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + text=True, + ) + retirement_finished = threading.Event() + retirement_preflight_finished = threading.Event() + retirement_results: list[dict[str, object]] = [] + retirement_errors: list[BaseException] = [] + actual_probe = probe_registry_write_path + + def observed_probe( + path: Path, + *, + create_parent: bool = True, + ) -> dict[str, Any]: + result = actual_probe(path, create_parent=create_parent) + retirement_preflight_finished.set() + return result + + monkeypatch.setattr( + global_registry_module, + "probe_registry_write_path", + observed_probe, + ) + + def retire() -> None: + try: + retirement_results.append( + retire_global_registry_goals( + runtime_root_override=str(runtime), + goal_ids=[GOAL_ID], + execute=True, + ) + ) + except BaseException as error: + retirement_errors.append(error) + finally: + retirement_finished.set() + + thread = threading.Thread(target=retire) + try: + assert writer.stdout is not None + ready, _, _ = select.select([writer.stdout], [], [], 10) + assert ready and writer.stdout.readline().strip() == "BARRIER commit" + registry.unlink() + state.unlink() + thread.start() + assert retirement_preflight_finished.wait(timeout=5) + assert not retirement_finished.wait(timeout=0.2), ( + "the global route retired while an admitted canonical update was " + "paused before provider commit" + ) + + stdout, stderr = writer.communicate("continue\n", timeout=20) + assert writer.returncode == 0, stdout + stderr + update = json.loads(stdout) + assert update["status"] == "applied", update + thread.join(timeout=10) + assert not thread.is_alive() + assert retirement_errors == [] + assert retirement_results and retirement_results[0]["ok"] is True + + global_goals = json.loads( + global_registry_path(runtime).read_text(encoding="utf-8") + )["goals"] + assert global_goals == [] + inspected = subprocess.run( + _node_command("inspect", request), + cwd=REPO, + capture_output=True, + text=True, + check=True, + timeout=20, + ) + head = json.loads(inspected.stdout) + assert head["head"]["todos"][0]["text"] == "Committed before retirement" + finally: + _stop(writer) + thread.join(timeout=10) diff --git a/tests/test_global_registry_write_serialization.py b/tests/test_global_registry_write_serialization.py index 0bfbf45c2c..a81c8e0e8b 100644 --- a/tests/test_global_registry_write_serialization.py +++ b/tests/test_global_registry_write_serialization.py @@ -244,6 +244,7 @@ def test_retire_reads_and_writes_inside_the_global_registry_lock( path.unlink() held: list[Path] = [] + acquired: list[Path] = [] events: list[str] = [] real_write = global_registry.write_json real_load = global_registry._load_global_registry @@ -251,6 +252,7 @@ def test_retire_reads_and_writes_inside_the_global_registry_lock( @contextmanager def recording_lock(path: Path, **kwargs: Any) -> Iterator[Path]: held.append(path) + acquired.append(path) try: yield path finally: @@ -281,6 +283,13 @@ def recording_write(path: Path, payload: dict[str, Any]) -> None: assert result["ok"] is True, result assert result["wrote"] is True, result assert held == [] + assert acquired == [ + shadow_maintenance_lock_target( + runtime_root, + "goal-alpha", + ), + global_path, + ] assert "read:locked" in events, events assert "backup:locked" in events, events assert "write:locked" in events, events