Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
37 changes: 27 additions & 10 deletions loopx/global_registry.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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"]

Expand Down
4 changes: 2 additions & 2 deletions loopx/semantics/project_registry_io_manifest_v1.json
Original file line number Diff line number Diff line change
Expand Up @@ -1975,15 +1975,15 @@
},
{
"site": "loopx/global_registry.py::<module>._sync_project_registry_to_global_once::codec_read:load_registry#1",
"line": 649,
"line": 666,
"column": 24,
"kind": "codec_read",
"api": "load_registry",
"classification": "codec_api"
},
{
"site": "loopx/global_registry.py::<module>.sync_project_registry_to_global::codec_read:load_registry#1",
"line": 857,
"line": 874,
"column": 22,
"kind": "codec_read",
"api": "load_registry",
Expand Down
249 changes: 249 additions & 0 deletions tests/control_plane/test_global_goal_retirement_writer_fence.py
Original file line number Diff line number Diff line change
@@ -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)
9 changes: 9 additions & 0 deletions tests/test_global_registry_write_serialization.py
Original file line number Diff line number Diff line change
Expand Up @@ -244,13 +244,15 @@ 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

@contextmanager
def recording_lock(path: Path, **kwargs: Any) -> Iterator[Path]:
held.append(path)
acquired.append(path)
try:
yield path
finally:
Expand Down Expand Up @@ -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
Expand Down
Loading