ouroboros/tests/test_extension_reconcile_queue.py
Ouroboros d32c99798f v7next F1: domain D14 - extension_loader six-leaf and skill_review four-leaf splits from tip bytes; upstream cycles re-home honored
Module side (41 D14 owners at tip): 15 byte-identical across tip/reference/
base (event_bus, extension_isolated_deps, extension_reconcile_queue,
marketplace/{__init__,adapter,clawhub,fetcher,install,install_specs,
isolated_deps}, skill_owner_attestation, skill_readiness,
skill_repair_admission, skill_review_status, skill_token); 16 pure upstream
drift - upstream bytes stand (extension_companion, extension_health,
extension_process_runner, extension_ui_validation,
marketplace/{ouroboroshub,provenance}, skill_dependencies,
skill_lifecycle_queue, skill_loader, skill_publish_eligibility,
skill_review_history, skill_review_passes, skill_review_runner,
tools/skill_{exec,preflight,publish}); 8 NEW post-cutoff upstream modules
with no ledger rows - untouched (betterleaks_runtime, skill_payload_binding,
skill_publish_{github,result,scanner,snapshot} of 8cc2ac69,
skill_review_cycles of 386e9417, skill_review_usage of f18da8c3). The
reference's unrowed failure_kind delta on extension_process_runner is NOT
replayed (typed-dispatch family, F3 territory); the supervised-future leak
in PluginAPIImpl is preserved per plan (F3-acceptance carries the direct
regression test). PluginAPI semantics untouched - mechanical byte-preserving
splits only.

The splits: extension_loader.py 2194 -> 960 into six leaves per ledger rows
2467-2519 (registry_state 87, surface_names 132, child_catalog 109,
import_staging 173, liveness 242, plugin_api 744); 49/53 spans
byte-identical to the oracle leaves, 4 byte-falsified by pure upstream drift
(widget-geometry promotion, durable companion-health overlay) and re-emitted
from tip bytes; two unrowed tip riders travel with their readers
(_widget_geometry_from_render -> surface_names beside its rowed sibling,
_apply_durable_extension_health -> liveness). skill_review.py 1600 -> 841
into four leaves (packs 166, rebuttals 127, prompt 376, output 245); 25/31
ledger rows executed from tip bytes, 6 SUPERSEDED by upstream's own re-home
into skill_review_cycles (386e9417) - upstream ownership stands, the facade
keeps the historical aliases and the identity suite pins them.
_ws_broadcaster deliberately not aliased on the facade (rebindable global
keeps one binding at its owner - the reference contract, pinned). Proof
green: transplant-tool --check per leaf - ast=tokens=bytes=True on all 80
spans, undeclared_top_level=[], leaf_invariants=[], exit 0; facade audit:
17+14 retained top-level statements byte-identical to tip, every tip name
accounted, every moved name re-imported (minus the pinned _ws_broadcaster);
import smoke + shared-registry object-identity checks green.

Test side: identity suites carried - test_extension_loader_extraction.py
(+2 rider rows) and test_skill_review_extraction.py (tool_module_inventory
clauses reverse-mapped out: v7-only mechanism absent at this tip;
cycles-alias identity clause added; facade size bound 800 -> 900 for the
tip-retained lifecycle members, which join the patchable-seams pin).
test_extension_loader.py 1717 -> 359 split into 4 siblings +
_extension_loader_shared.py per 45 ledger rows (zero upstream drift on the
file; 43/45 moved bodies byte-identical, 1 oracle dual-supervisor-patch
adaptation kept, 1 reverse-mapped to the tip spelling supervisor/workers.py
- the reference's worker_process.py is the still-pending D08 split).
test_skill_review.py 1943 -> 579 into 5 siblings + _skill_review_shared.py
per 65 rows (58 byte-identical, 3 re-emitted from tip test drift, 3
patch-retarget adaptations kept; the upstream-deleted advisory test is not
resurrected - its f8d87c69 successors stay in the remainder on tip bytes,
theme re-home deferred to F5). Lossless both families: 52==52 and 74==74
test names, zero duplicate names, no new ast-identical bodies (10
pre-existing review_cycles duplicates at base noted, untouched). Dead-patch
class closed: the remainder's advisory pre-review patch -> prompt owner,
extension_companion supervisor patch doubled onto the plugin_api owner
(single-module patch proven dead by a red run); every other facade patch
site of moved names verified live (production consumers do call-time facade
imports).

size-ratchet manifest regenerated with the official tool (extension_loader,
test_extension_loader, test_skill_review leave GIANT_PATHS; no new band
entries); --check exact; ratchet 5 passed; ruff F clean tree-wide.
Receipts: 12 identity + 126 split-family + 1751 targeted/-n16-loadscope +
full CI-shape battery 11950 parallel + 618 serial, all rc=0; HEAD held and
worktree status unchanged through every pytest run. Ledger corrections:
docs/v7next/LEDGER_CORRECTIONS.md D14 section, entries 1-10.

(cherry picked from commit c7e3d04df95eb9470f5ca3ad824485083a201fe4)
2026-08-30 20:40:34 +00:00

221 lines
8.1 KiB
Python

"""The extension reconcile marker queue and companion pickup.
Split out of ``tests/test_extension_loader.py`` when that module was divided by
theme; every moved block is verbatim. Covers the worker-side reconcile writing
server markers for enable and disable, the server pickup spawning, stopping and
redriving companion processes, a newer marker surviving a processing race, the
failed-marker overflow queue, the supervisor's redrive surface and the server
lifespan wiring of the pickup loop.
"""
from __future__ import annotations
import json
import pathlib
from typing import Any, Dict
from ouroboros import extension_loader
from ouroboros import extension_plugin_api
from ouroboros.extension_companion import CompanionSupervisor, init_server_process_pid
from ouroboros.extension_reconcile_queue import (
MAX_ATTEMPTS,
list_extension_reconcile_requests,
process_extension_reconcile_requests,
request_extension_reconcile,
)
from ouroboros.skill_loader import save_enabled
from tests._extension_loader_shared import (
_prepare_extension,
)
from tests._extension_loader_shared import ( # noqa: F401 (autouse fixture applies on import)
_clear_loader_state,
)
def _prepare_companion_extension(tmp_path: pathlib.Path, name: str = "compskill"):
return _prepare_extension(
tmp_path,
name,
"def register(api):\n api.register_companion_process('daemon')\n",
permissions=["companion_process"],
extra_frontmatter=(
"companion_processes:\n"
" - name: daemon\n"
" runtime: python3\n"
" command: [\"python3\", \"scripts/daemon.py\"]\n"
),
)
def test_worker_reconcile_writes_server_marker_for_enable_and_disable(tmp_path: pathlib.Path) -> None:
init_server_process_pid(999999)
loaded, repo_root, drive_root = _prepare_companion_extension(tmp_path)
state = extension_loader.reconcile_extension(
loaded.name,
drive_root,
lambda: {},
repo_path=str(repo_root),
)
assert state["action"] == "extension_loaded"
requests = list_extension_reconcile_requests(drive_root)
assert [item["skill"] for item in requests] == [loaded.name]
save_enabled(drive_root, loaded.name, False)
state = extension_loader.reconcile_extension(
loaded.name,
drive_root,
lambda: {},
repo_path=str(repo_root),
)
assert state["action"] == "extension_unloaded"
requests = list_extension_reconcile_requests(drive_root)
assert {item["skill"] for item in requests} == {loaded.name}
assert "desired_disabled" in {item["reason"] for item in requests}
def test_server_pickup_spawns_stops_and_redrives_missing_companion(
tmp_path: pathlib.Path,
monkeypatch,
) -> None:
init_server_process_pid()
loaded, repo_root, drive_root = _prepare_companion_extension(tmp_path)
class FakeSupervisor:
def __init__(self):
self.runtimes: Dict[str, Dict[str, Any]] = {}
self.started: list[str] = []
self.stopped: list[str] = []
def start(self, descriptor):
key = f"{descriptor.skill_name}:{descriptor.name}"
self.runtimes[key] = {"skill_name": descriptor.skill_name, "name": descriptor.name}
self.started.append(key)
return True
def snapshot(self):
return dict(self.runtimes)
def stop(self, skill_name: str, name: str):
self.stopped.append(f"{skill_name}:{name}")
self.runtimes.pop(f"{skill_name}:{name}", None)
def stop_skill(self, skill_name: str):
self.stopped.append(skill_name)
self.runtimes = {
key: value
for key, value in self.runtimes.items()
if value.get("skill_name") != skill_name
}
fake = FakeSupervisor()
# The companion supervisor is read by PluginAPIImpl (its owner) at register
# time and by the loader's ensure_companions_running/unload paths.
monkeypatch.setattr(extension_plugin_api, "get_global_supervisor", lambda: fake)
monkeypatch.setattr(extension_loader, "get_global_supervisor", lambda: fake)
request_extension_reconcile(drive_root, loaded.name, reason="test")
processed = process_extension_reconcile_requests(drive_root, lambda: {}, repo_path=str(repo_root))
assert processed[0]["skill"] == loaded.name
assert fake.started == [f"{loaded.name}:daemon"]
assert list_extension_reconcile_requests(drive_root) == []
request_extension_reconcile(drive_root, loaded.name, reason="idempotent")
process_extension_reconcile_requests(drive_root, lambda: {}, repo_path=str(repo_root))
assert fake.started == [f"{loaded.name}:daemon"]
fake.runtimes.clear()
state = extension_loader.reconcile_extension(
loaded.name,
drive_root,
lambda: {},
repo_path=str(repo_root),
)
assert state["action"] == "extension_already_live"
assert fake.started == [f"{loaded.name}:daemon", f"{loaded.name}:daemon"]
fake.runtimes.clear()
request_extension_reconcile(drive_root, loaded.name, reason="redrive")
process_extension_reconcile_requests(drive_root, lambda: {}, repo_path=str(repo_root))
assert fake.started == [
f"{loaded.name}:daemon",
f"{loaded.name}:daemon",
f"{loaded.name}:daemon",
]
save_enabled(drive_root, loaded.name, False)
request_extension_reconcile(drive_root, loaded.name, reason="disable")
process_extension_reconcile_requests(drive_root, lambda: {}, repo_path=str(repo_root))
assert fake.stopped == [f"{loaded.name}:daemon", loaded.name]
assert fake.snapshot() == {}
assert list_extension_reconcile_requests(drive_root) == []
def test_pickup_keeps_newer_marker_written_during_processing(
tmp_path: pathlib.Path,
monkeypatch,
) -> None:
drive_root = tmp_path / "drive"
drive_root.mkdir()
request_extension_reconcile(drive_root, "race_skill", reason="old")
def fake_reconcile(skill_name, drive_root_arg, settings_reader, **kwargs):
request_extension_reconcile(drive_root_arg, skill_name, reason="newer")
return {"action": "extension_loaded"}
monkeypatch.setattr(extension_loader, "reconcile_extension", fake_reconcile)
monkeypatch.setattr(
extension_loader,
"ensure_companions_running",
lambda *args, **kwargs: {"action": "noop"},
)
processed = process_extension_reconcile_requests(drive_root, lambda: {})
assert processed[0]["marker_removed"] is True
requests = list_extension_reconcile_requests(drive_root)
assert len(requests) == 1
assert requests[0]["reason"] == "newer"
def test_repeatedly_failed_marker_moves_out_of_active_queue(
tmp_path: pathlib.Path,
monkeypatch,
) -> None:
drive_root = tmp_path / "drive"
drive_root.mkdir()
request_extension_reconcile(drive_root, "broken_skill", reason="test")
def fake_reconcile(*args, **kwargs):
raise RuntimeError("boom")
monkeypatch.setattr(extension_loader, "reconcile_extension", fake_reconcile)
for _ in range(MAX_ATTEMPTS):
process_extension_reconcile_requests(drive_root, lambda: {})
assert list_extension_reconcile_requests(drive_root) == []
failed = list((drive_root / "state" / "extension_reconcile" / "failed").glob("*.json"))
assert len(failed) == 1
assert json.loads(failed[0].read_text(encoding="utf-8"))["status"] == "failed"
def test_companion_supervisor_exposes_server_redrive_methods() -> None:
assert callable(getattr(CompanionSupervisor, "snapshot"))
assert callable(getattr(CompanionSupervisor, "stop_skill"))
def test_server_lifespan_wires_extension_reconcile_pickup() -> None:
server_py = pathlib.Path(__file__).resolve().parents[1] / "server.py"
text = server_py.read_text(encoding="utf-8")
assert "from ouroboros.extension_reconcile_queue import extension_reconcile_pickup_loop" in text
assert "extension_reconcile_task = asyncio.create_task" in text
assert "name=\"extension-reconcile-pickup\"" in text
assert "extension_reconcile_task.cancel()" in text
assert "await asyncio.wait_for(extension_reconcile_task, timeout=30)" in text