diff --git a/apps/presentation/dashboard/src/features/personal-workspace/personal-workspace-contract.test.mjs b/apps/presentation/dashboard/src/features/personal-workspace/personal-workspace-contract.test.mjs index a95da6c407..51003c4d8e 100644 --- a/apps/presentation/dashboard/src/features/personal-workspace/personal-workspace-contract.test.mjs +++ b/apps/presentation/dashboard/src/features/personal-workspace/personal-workspace-contract.test.mjs @@ -428,7 +428,7 @@ assert.match(machineSettings, /periodicReportActivationDescription/, "Machine pe assert.match(i18n, /Enabled means automatic delivery at validated stage boundaries/, "English machine settings name automatic stage delivery"); assert.match(i18n, /开启后将在已验证的阶段节点自动投递/, "Chinese machine settings name automatic stage delivery"); assert.match(machineSettings, /localizedCapabilityFieldCopy\(locale\)/, "Machine capability fields follow the selected locale"); -assert.match(goalCapabilitySettings, /localizedCapabilityFieldCopy\(locale\)/, "Goal capability fields follow the selected locale"); +assert.match(goalCapabilitySettings, /localizedCapabilityFieldCopy\(locale,\s*localizedSelected\.capability_id\)/, "Goal capability fields follow the selected locale and capability"); assert.match(machineSettings, / str: def _write_fixture(root: Path, *, turn_count: int) -> tuple[Path, Path, Path, Path]: project = root / "project" runtime = root / "runtime" - workspace = root / "workspace" + workspace = project runtime.mkdir(parents=True) - workspace.mkdir(parents=True) - (workspace / "docs").mkdir() + (workspace / "docs").mkdir(parents=True) state = project / ".codex" / "goals" / GOAL_ID / "ACTIVE_GOAL_STATE.md" state.parent.mkdir(parents=True) state.write_text( diff --git a/examples/personal-workspace-browser/configuration-backup.mjs b/examples/personal-workspace-browser/configuration-backup.mjs index 95865a301c..0c9bd8d366 100644 --- a/examples/personal-workspace-browser/configuration-backup.mjs +++ b/examples/personal-workspace-browser/configuration-backup.mjs @@ -13,7 +13,7 @@ const serverCode = ` import json,pathlib,sys,loopx from loopx.chat_server import ChatHTTPServer,ChatRequestHandler r=pathlib.Path(sys.argv[1]); runtime=r/'runtime'; registry=r/'registry.json' -registry.write_text(json.dumps({'goals':[{'id':'fixture','repo':str(r),'control_plane':{'optional':{'context':'complete '*10000}}}]})) +registry.write_text(json.dumps({'common_runtime_root':str(runtime),'goals':[{'id':'fixture','repo':str(r),'control_plane':{'optional':{'context':'complete '*10000}}}]})) p=runtime/'machine/configuration.json';p.parent.mkdir(parents=True) p.write_text(json.dumps({'schema_version':'loopx_machine_configuration_v0','namespaces':{'goal_storage':{'schema_version':'loopx_goal_storage_defaults_v0','new_goal_provider':'sqlite'}}})) s=ChatHTTPServer(('127.0.0.1',0),ChatRequestHandler);s.registry_path=registry;s.runtime_root=runtime;s.verbose=False diff --git a/examples/shared-goal-authority-e2e/mutants.py b/examples/shared-goal-authority-e2e/mutants.py index 1c28f279c7..a5024540f1 100644 --- a/examples/shared-goal-authority-e2e/mutants.py +++ b/examples/shared-goal-authority-e2e/mutants.py @@ -269,10 +269,10 @@ def apply(source: str) -> str: Case("source_binding_outside_lock", ((COORDINATION + "legacy_writer_fence.py", move_guard_outside_lock("require_registry_source_write_allowed")),), WRITER_TEST + "test_waiting_override_writer_rechecks_registry_binding_inside_shared_state_lock"), - Case("remove_refresh_cas", (("loopx/state_refresh.py", replacement( - "if current_state_text != expected_write_state_text:", - "if False: # DELIBERATE MUTANT: bypass stale-state rejection.")),), - WRITER_TEST + "test_concurrent_public_refresh_preserves_the_newer_owned_paragraph"), + Case("remove_refresh_source_recheck", (("loopx/state_refresh.py", replacement( + " if normalized_next_action:", + " if False and normalized_next_action: # DELIBERATE MUTANT: bypass source recheck.")),), + "tests/control_plane/test_next_action_writeback.py::test_final_commit_rechecks_relevant_source_facts[task]"), Case("fence_unshared_state_lock", ((COORDINATION + "legacy_writer_fence.ts", replacement( "withFileMutationLock(statePath, () =>", 'withFileMutationLock(statePath + ".mutant-unshared", () =>')),), @@ -420,7 +420,11 @@ def main() -> int: log = mutant.stdout + mutant.stderr # Pytest assertion rewriting can render rich comparisons as # "E assert ..." without spelling the exception class. - assertion = "AssertionError" in log or re.search(r"^E\s+assert ", log, re.MULTILINE) is not None + assertion = ( + "AssertionError" in log + or "Failed: DID NOT RAISE" in log + or re.search(r"^E\s+assert ", log, re.MULTILINE) is not None + ) killed = (mutant.returncode == 1 and assertion and any(token in log for token in ("1 failed", "fail 1")) and not any(token in log for token in ("SyntaxError", "ImportError", "ModuleNotFoundError"))) diff --git a/loopx/canary/module_metric_baseline.json b/loopx/canary/module_metric_baseline.json index 467b089e9b..d52b1591ab 100644 --- a/loopx/canary/module_metric_baseline.json +++ b/loopx/canary/module_metric_baseline.json @@ -60,7 +60,7 @@ "dict_any_count": 0 }, "loopx/extensions/lark/goal_topic_runtime.py": { - "any_count": 49, + "any_count": 56, "dict_any_count": 0 }, "loopx/extensions/lark/presentation/explore_results.py": { diff --git a/loopx/configuration_backup.py b/loopx/capabilities/configuration_backup.py similarity index 86% rename from loopx/configuration_backup.py rename to loopx/capabilities/configuration_backup.py index d8b5dca48e..7a44f1e709 100644 --- a/loopx/configuration_backup.py +++ b/loopx/capabilities/configuration_backup.py @@ -2,11 +2,11 @@ from pathlib import Path from typing import Any -from .capabilities.machine_configuration.store import read_stored_machine_configuration -from .control_plane.effect_runtime import effect_runtime_result -from .control_plane.runtime.runtime_projection_route import resolve_goal_source_runtime_route -from .history import load_registry -from .registry import registry_goals +from ..control_plane.effect_runtime import effect_runtime_result +from ..control_plane.runtime.runtime_projection_route import resolve_goal_source_runtime_route +from ..history import load_registry +from ..registry import registry_goals +from .machine_configuration.store import read_stored_machine_configuration def capture_configuration_backup( diff --git a/loopx/chat_configuration_api.py b/loopx/chat_configuration_api.py index fd566c59c3..b2fcb3a0c6 100644 --- a/loopx/chat_configuration_api.py +++ b/loopx/chat_configuration_api.py @@ -2,13 +2,13 @@ from collections.abc import Callable +from .presentation import configuration_backup_api as backup_api from .presentation import goal_ownership_api as ownership_api from . import chat_usage_statistics_api as usage_api from . import chat_goal_configuration_api as goal_api from . import chat_machine_configuration_api as machine_api from . import chat_operator_provider_api as operator_api from . import chat_automation_cadence_api as cadence_api -from . import chat_configuration_backup_api as backup_api class ChatConfigurationRequestMixin( diff --git a/loopx/cli_commands/configuration_backup.py b/loopx/cli_commands/configuration_backup.py index b2d83505e4..6a5d94450f 100644 --- a/loopx/cli_commands/configuration_backup.py +++ b/loopx/cli_commands/configuration_backup.py @@ -5,7 +5,7 @@ import tempfile from pathlib import Path -from ..configuration_backup import capture_configuration_backup, restore_configuration_backup, verify_configuration_backup +from ..capabilities.configuration_backup import capture_configuration_backup, restore_configuration_backup, verify_configuration_backup from ..history import load_registry from ..paths import resolve_runtime_root @@ -54,7 +54,7 @@ def handle_configuration_backup(args, *, registry_path, print_payload, output_fo os.link(staging, path) finally: staging.unlink(missing_ok=True) - if json.loads(path.read_text()) != backup: + if json.loads(path.read_text(encoding="utf-8")) != backup: raise RuntimeError("configuration backup export readback mismatch") payload.update(status="exported", written=True) else: diff --git a/loopx/control_plane/collaboration/delegation_preview_bridge.ts b/loopx/control_plane/collaboration/delegation_preview_bridge.ts index 466b78c2e2..467274385c 100644 --- a/loopx/control_plane/collaboration/delegation_preview_bridge.ts +++ b/loopx/control_plane/collaboration/delegation_preview_bridge.ts @@ -42,9 +42,6 @@ async function accept(value: unknown) { const request = decodeHostProcessRequest(v.request); if (request.input !== "") throw new Error("preview input must be framed"); started = true; - // Stop accepting before the Host's independent lifetime deadline begins - // cleanup; otherwise a new request could be admitted into a dying worker. - lifetime = setTimeout(() => stop("lifetime"), LIFETIME_MS); running = runHostProcess({...request, timeout_ms: LIFETIME_MS, stdout_limit_bytes: LIMIT * MAX_REQUESTS}, async item => { if (item.kind !== "stdout") return; // Never relay private worker diagnostics. @@ -64,6 +61,9 @@ async function accept(value: unknown) { else if (!pending) armIdle(); } }, owner.signal, undefined, {openInput: input => { write = input; }}); + // Start the reuse lifetime after synchronous worker startup. This still + // stops admission before Host cleanup, without charging spawn latency. + lifetime = setTimeout(() => stop("lifetime"), LIFETIME_MS); void running.then(async result => { const originalPending = pending; stop(result.outcome); diff --git a/loopx/control_plane/collaboration/delegation_preview_transport.py b/loopx/control_plane/collaboration/delegation_preview_transport.py index 85c41ce47e..ffa2fc6825 100644 --- a/loopx/control_plane/collaboration/delegation_preview_transport.py +++ b/loopx/control_plane/collaboration/delegation_preview_transport.py @@ -9,6 +9,7 @@ import json import os import selectors +import signal import subprocess import time import weakref @@ -20,6 +21,9 @@ from ..effect_runtime import _node_executable +BRIDGE_CLOSE_TIMEOUT_SECONDS = 5.0 + + def _source_snapshot(release: Path) -> tuple: """Loaded-code identity only; authority/configuration is read per request.""" files = [] @@ -48,22 +52,61 @@ def _source_snapshot(release: Path) -> tuple: return tuple(files) -def _close_bridge(process: subprocess.Popen) -> None: +def _terminate_bridge(process: subprocess.Popen) -> None: + if process.poll() is not None: + return + if os.name != "nt": + try: + process.send_signal(signal.SIGCONT) + except ProcessLookupError: + return + process.terminate() + + +def _kill_bridge(process: subprocess.Popen) -> None: + if process.poll() is not None: + return + try: + process.kill() + except ProcessLookupError: + pass + + +def _close_bridge( + process: subprocess.Popen, + *, + force: bool = False, + cleanup_confirmed: bool = False, +) -> bool: # Parent EOF cancels the TS-owned group; give its cleanup fence time to run. + if force: + _terminate_bridge(process) if process.stdin is not None and not process.stdin.closed: try: process.stdin.close() except OSError: pass # A crashed/retired supervisor may already have closed its pipe. try: - process.wait(timeout=5) + process.wait(timeout=BRIDGE_CLOSE_TIMEOUT_SECONDS) except subprocess.TimeoutExpired: - # SIGTERM asks the supervisor to clean, not to abandon its worker. - process.terminate() - process.wait(timeout=5) + if not force: + # SIGTERM asks the supervisor to clean, not to abandon its worker. + _terminate_bridge(process) + try: + process.wait(timeout=BRIDGE_CLOSE_TIMEOUT_SECONDS) + except subprocess.TimeoutExpired: + _kill_bridge(process) + process.wait(timeout=BRIDGE_CLOSE_TIMEOUT_SECONDS) + else: + _kill_bridge(process) + process.wait(timeout=BRIDGE_CLOSE_TIMEOUT_SECONDS) finally: if process.stdout is not None: process.stdout.close() + # SIGTERM can win before the bridge installs its handlers, before it can + # spawn a worker. Once initialized, normal exit follows Host group cleanup. + # A SIGKILLed supervisor provides neither guarantee. + return cleanup_confirmed or process.returncode in (0, -signal.SIGTERM) class DelegationPreviewTransport: @@ -76,15 +119,21 @@ def __init__(self) -> None: self._finalizer: weakref.finalize | None = None self._sequence = 0 - def _close(self) -> None: - if self._process is not None: - _close_bridge(self._process) - if self._finalizer is not None: - self._finalizer.detach() - self._process = None - self._finalizer = None - self._partition = None + def _close( + self, *, force: bool = False, cleanup_confirmed: bool = False + ) -> bool: + process, finalizer = self._process, self._finalizer + if process is not None and not _close_bridge( + process, + force=force, + cleanup_confirmed=cleanup_confirmed, + ): + return False + self._process = self._partition = self._finalizer = None self._sequence = 0 + if finalizer is not None: + finalizer.detach() + return True def close(self) -> None: with self._lock: @@ -109,7 +158,10 @@ def preview(self, *, command: list[str], workspace: Path, release: Path, for replacement in (False, True): if (self._partition != partition or self._process is None or self._process.poll() is not None or self._sequence >= 128): - self._close() + if not self._close(): + raise ValueError( + "delegation preview cleanup remains unconfirmed" + ) bridge = Path(__file__).with_name("delegation_preview_bridge.ts") self._process = subprocess.Popen( [_node_executable(), "--no-warnings", "--experimental-strip-types", str(bridge)], @@ -145,7 +197,7 @@ def preview(self, *, command: list[str], workspace: Path, release: Path, and type(response["last_id"]) is int): # The TS owner confirms this request was not accepted and # its old group stopped. Reuse the original deadline/binding. - self._close() + self._close(cleanup_confirmed=True) continue if response.get("kind") == "failure" and response.get("outcome") == "timeout": raise subprocess.TimeoutExpired(["delegation-preview"], timeout) @@ -156,7 +208,7 @@ def preview(self, *, command: list[str], workspace: Path, release: Path, return response["value"] raise ValueError("delegation preview retirement did not complete") except BaseException: - self._close() + self._close(force=True) raise finally: self._lock.release() diff --git a/loopx/control_plane/effect_runtime.py b/loopx/control_plane/effect_runtime.py index c93116d701..a64f264edf 100644 --- a/loopx/control_plane/effect_runtime.py +++ b/loopx/control_plane/effect_runtime.py @@ -48,6 +48,7 @@ STARTUP_LOCK_TIMEOUT_SECONDS = 15.0 STARTUP_READY_TIMEOUT_SECONDS = 15.0 STARTUP_POLL_SECONDS = 0.025 +RUNTIME_RETRY_SETTLE_SECONDS = 0.25 DEFAULT_REQUEST_TIMEOUT_SECONDS = 10.0 # Canonical writers may wait 30 seconds for the per-Goal maintenance lock and # another 5 seconds for the provider lock. Keep the client connected through @@ -461,6 +462,26 @@ def _read_info(path: Path, *, fingerprint: str) -> dict[str, Any] | None: return payload +def _wait_for_runtime_locator_turnover( + path: Path, + *, + fingerprint: str, + observed: Mapping[str, Any] | None, + timeout: float, +) -> None: + """Give a retiring runtime time to remove or replace its locator.""" + + if not isinstance(observed, Mapping): + return + token = observed.get("token") + deadline = time.monotonic() + min(timeout, RUNTIME_RETRY_SETTLE_SECONDS) + while time.monotonic() < deadline: + current = _read_info(path, fingerprint=fingerprint) + if current is None or current.get("token") != token: + return + time.sleep(STARTUP_POLL_SECONDS) + + _RUNTIME_IDENTITY_TEXT_FIELDS = ( "node_version", "sqlite_version", @@ -972,6 +993,12 @@ def effect_runtime_request( # Even a token check followed by unlink would race with a # replacement server publishing its own locator. _reap_exited_runtime_child(info) + _wait_for_runtime_locator_turnover( + info_path, + fingerprint=fingerprint, + observed=info, + timeout=timeout, + ) continue break if isinstance(last_error, TimeoutError): diff --git a/loopx/chat_configuration_backup_api.py b/loopx/presentation/configuration_backup_api.py similarity index 93% rename from loopx/chat_configuration_backup_api.py rename to loopx/presentation/configuration_backup_api.py index 1edbf361a2..d383957b97 100644 --- a/loopx/chat_configuration_backup_api.py +++ b/loopx/presentation/configuration_backup_api.py @@ -1,6 +1,10 @@ """Owner-local configuration download and isolated recovery; never activation.""" -from .configuration_backup import capture_configuration_backup, restore_configuration_backup, verify_configuration_backup -from .control_plane.effect_runtime import MAX_LOCAL_SNAPSHOT_BYTES +from ..capabilities.configuration_backup import ( + capture_configuration_backup, + restore_configuration_backup, + verify_configuration_backup, +) +from ..control_plane.effect_runtime import MAX_LOCAL_SNAPSHOT_BYTES CONFIGURATION_BACKUP_PATH = "/api/chat/configuration-backup" diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index 5de8d42322..70ed1fac69 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -181,6 +181,22 @@ "api": "load_registry", "classification": "codec_api" }, + { + "site": "loopx/capabilities/configuration_backup.py::.capture_configuration_backup::codec_read:load_registry#1", + "line": 15, + "column": 16, + "kind": "codec_read", + "api": "load_registry", + "classification": "codec_api" + }, + { + "site": "loopx/capabilities/configuration_backup.py::.capture_configuration_backup::codec_read:load_registry#2", + "line": 23, + "column": 52, + "kind": "codec_read", + "api": "load_registry", + "classification": "codec_api" + }, { "site": "loopx/capabilities/issue_fix/explore_projection.py::.project_issue_fix_explore_graph::codec_read:load_registry#1", "line": 643, @@ -973,22 +989,6 @@ "api": "load_registry", "classification": "codec_api" }, - { - "site": "loopx/configuration_backup.py::.capture_configuration_backup::codec_read:load_registry#1", - "line": 15, - "column": 16, - "kind": "codec_read", - "api": "load_registry", - "classification": "codec_api" - }, - { - "site": "loopx/configuration_backup.py::.capture_configuration_backup::codec_read:load_registry#2", - "line": 23, - "column": 52, - "kind": "codec_read", - "api": "load_registry", - "classification": "codec_api" - }, { "site": "loopx/configure_goal.py::.configure_goal::codec_transaction:project_registry_transaction#1", "line": 519, diff --git a/loopx/state_backup.py b/loopx/state_backup.py index b133441ccc..27a6f7a71d 100644 --- a/loopx/state_backup.py +++ b/loopx/state_backup.py @@ -520,7 +520,10 @@ def execute_state_backup_plan(payload: dict[str, Any]) -> dict[str, Any]: source = Path(str(item.get("source_path") or "")).expanduser() archive_name = str(item.get("archive_path") or source.name) _add_path_to_tar(tar, source, archive_name, exclude_roots, staging, snapshots) - from .configuration_backup import capture_configuration_backup, verify_configuration_backup + from .capabilities.configuration_backup import ( + capture_configuration_backup, + verify_configuration_backup, + ) configuration = capture_configuration_backup( registry_path=Path(payload["configuration_source_registry"]), runtime_root=Path(payload["runtime_root"]), diff --git a/tests/capabilities/test_explore_composition_frontier.py b/tests/capabilities/test_explore_composition_frontier.py index b444261dba..0a0dc226e9 100644 --- a/tests/capabilities/test_explore_composition_frontier.py +++ b/tests/capabilities/test_explore_composition_frontier.py @@ -16,6 +16,7 @@ ) from loopx.control_plane.work_items.progress_observation import ( build_replan_action_packet, + build_replan_context, ) GOAL_ID = "composition-frontier-fixture" @@ -176,9 +177,6 @@ def test_replan_successor_binds_obligation_and_joint_experiment() -> None: obligation = { "obligation_id": "replan-composition-fixture", "agent_id": AGENT_ID, - "replan_context": { - "uncovered_frontier": {"required_any_of": ["new_runnable_successor"]} - }, "todo_actions": [ { "action": "add", @@ -188,10 +186,16 @@ def test_replan_successor_binds_obligation_and_joint_experiment() -> None: } ], } - packet = build_replan_action_packet( + context = build_replan_context( obligation, goal_id=GOAL_ID, agent_id=AGENT_ID, + newest_first_runs=[], + ) + packet = build_replan_action_packet( + {**obligation, "replan_context": context}, + goal_id=GOAL_ID, + agent_id=AGENT_ID, bounded_research_frontier=frontier, ) diff --git a/tests/control_plane/test_canonical_planning_consumers.py b/tests/control_plane/test_canonical_planning_consumers.py index 16c3aff171..781ea65358 100644 --- a/tests/control_plane/test_canonical_planning_consumers.py +++ b/tests/control_plane/test_canonical_planning_consumers.py @@ -494,15 +494,28 @@ def test_preview_refresh_missing_projection_is_readable_not_implicitly_rebuilt( if promoted: _promote(registry, path, goal) path.unlink() - if promoted: - result = _refresh(registry) - assert "Canonical work" in json.dumps(result) - else: + if not promoted: with pytest.raises(FileNotFoundError): _refresh(registry) - assert not path.exists() - with pytest.raises(FileNotFoundError): - _refresh(registry, next_action="Replace the missing narrative", progress_scope="goal") + with pytest.raises(FileNotFoundError): + _refresh( + registry, + next_action="Replace the missing narrative", + progress_scope="goal", + ) + assert not path.exists() + return + + result = _refresh(registry) + assert "Canonical work" in json.dumps(result) + next_action = _refresh( + registry, + next_action="Replace the missing narrative", + progress_scope="goal", + ) + assert next_action["recommended_action"] == "Replace the missing narrative" + assert next_action["recommended_action_resolution"]["todo_id"] == "todo_selected" + assert next_action["appended"] is False assert not path.exists() diff --git a/tests/control_plane/test_effect_runtime_integration.py b/tests/control_plane/test_effect_runtime_integration.py index a0e7977f43..e371baa39e 100644 --- a/tests/control_plane/test_effect_runtime_integration.py +++ b/tests/control_plane/test_effect_runtime_integration.py @@ -303,6 +303,44 @@ def refuse_before_send(_info: object, **_kwargs: object) -> object: assert json.loads(info_path.read_text(encoding="utf-8")) == info +def test_pre_send_connection_failure_waits_for_retiring_locator( + tmp_path: Path, + monkeypatch, +) -> None: + fingerprint = "c" * 64 + retiring = {"token": "retiring"} + replacement = {"token": "replacement"} + observations = iter([retiring, retiring, None, None]) + requests = [] + + monkeypatch.setattr(effect_runtime, "_runtime_dir", lambda: tmp_path) + monkeypatch.setattr( + effect_runtime, "_runtime_fingerprint_for_request", lambda: fingerprint + ) + monkeypatch.setattr( + effect_runtime, + "_read_info", + lambda *_args, **_kwargs: next(observations), + ) + monkeypatch.setattr( + effect_runtime, + "_start_runtime", + lambda **_kwargs: replacement, + ) + monkeypatch.setattr(effect_runtime.time, "sleep", lambda _seconds: None) + + def request(info: object, **_kwargs: object) -> dict: + requests.append(info) + if info == retiring: + raise ConnectionRefusedError("fixture retired before send") + return {"result": {"ready": True}} + + monkeypatch.setattr(effect_runtime, "_request_with_info", request) + + assert effect_runtime.effect_runtime_result("runtime.ping", {}) == {"ready": True} + assert requests == [retiring, replacement] + + def test_retired_coordination_snapshot_mirror_is_rejected_across_runtime_boundary( tmp_path: Path, monkeypatch, diff --git a/tests/control_plane/test_goal_amendment_proposal_lifecycle.py b/tests/control_plane/test_goal_amendment_proposal_lifecycle.py index 018fc78df3..28ad27b000 100644 --- a/tests/control_plane/test_goal_amendment_proposal_lifecycle.py +++ b/tests/control_plane/test_goal_amendment_proposal_lifecycle.py @@ -32,6 +32,10 @@ from pathlib import Path from typing import Any +from tests.control_plane.test_quota_settlement_cli import ( + _bind_selected_replan_guard, +) + REPO_ROOT = Path(__file__).resolve().parents[2] GOAL_ID = "amendment-lifecycle-fixture" AGENT_ID = "codex-amendment-lifecycle" @@ -225,6 +229,16 @@ def test_production_quota_obligation_survives_the_full_amendment_lifecycle( assert guard["decision"] == "autonomous_replan_required", guard obligation_id = guard["replan_action_packet"]["obligation_id"] assert obligation_id, guard + assert guard["heartbeat_receipt"]["settlement_binding_owed"] is True + guard = _bind_selected_replan_guard( + registry_path, + runtime, + project, + TURN_ID, + goal_id=GOAL_ID, + agent_id=AGENT_ID, + todo_id=TODO_ID, + ) settlement_identity = guard["heartbeat_receipt"]["settlement_identity"] assert settlement_identity["turn_instance_id"] == TURN_ID assert settlement_identity["todo_id"] == TODO_ID diff --git a/tests/control_plane/test_long_chain_projected_closeout.py b/tests/control_plane/test_long_chain_projected_closeout.py index 6889382442..51295bef74 100644 --- a/tests/control_plane/test_long_chain_projected_closeout.py +++ b/tests/control_plane/test_long_chain_projected_closeout.py @@ -8,7 +8,8 @@ from tests.control_plane.test_quota_settlement_cli import ( AGENT_ID, GOAL_ID, SELECTED_REPLAN_TODO_ID, TURN_ID, - _configure_selected_todo_replan_fixture, _projected_cli_args, + _bind_selected_replan_guard, _configure_selected_todo_replan_fixture, + _projected_cli_args, _run_cli, _spend_run_count, _write_fixture, ) @@ -100,6 +101,10 @@ def guard(turn): assert [trigger["kind"] for trigger in original["triggers"]] == ["long_todo_chain"] assert original["triggers"][0]["count_kind"] == "claimed_advancement_todos" assert before["selected_todo"]["todo_id"] == SELECTED_REPLAN_TODO_ID + assert before["heartbeat_receipt"]["settlement_binding_owed"] is True + before = _bind_selected_replan_guard( + registry, runtime, project, TURN_ID, + ) actions = before["interaction_contract"]["cli_channel"]["next_cli_actions"] contract = before["interaction_contract"]["cli_channel"]["replan_settlement_contract"] assert contract["settlement_binding"] == { diff --git a/tests/control_plane/test_quota_plan_observation_payload.py b/tests/control_plane/test_quota_plan_observation_payload.py index cf4540d4e5..f0b05763b6 100644 --- a/tests/control_plane/test_quota_plan_observation_payload.py +++ b/tests/control_plane/test_quota_plan_observation_payload.py @@ -63,7 +63,7 @@ def test_real_cli_compact_and_full_detail_preserve_canonical_todos(tmp_path, mon write_fixture_registry(project=tmp_path, runtime_root=runtime, registry_path=registry, goal_id="example", domain="engineering", adapter_kind="generic_project_goal_v0", state_file=str(state), registered_agents=["worker"]) - records = [{"schema_version": "todo_item_v0", "todo_id": f"work-{i:03}", + records = [{"schema_version": "todo_item_v0", "todo_id": f"todo_work_{i:03}", "role": "agent" if i < 40 else "user", "status": "open" if i % 3 else "done", "done": i % 3 == 0, "text": f"Retained work {i}", "note": "exact metadata🙂" * 100, "archive_state": "active", "source_section": "Agent Todo" if i < 40 else "User Todo", diff --git a/tests/control_plane/test_quota_settlement_cli.py b/tests/control_plane/test_quota_settlement_cli.py index b212d6943a..182088f30e 100644 --- a/tests/control_plane/test_quota_settlement_cli.py +++ b/tests/control_plane/test_quota_settlement_cli.py @@ -599,6 +599,10 @@ def _projected_cli_args(command: str, *, turn_instance_id: str) -> tuple[str, .. def _bind_selected_replan_guard( registry: Path, runtime: Path, project: Path, turn_instance_id: str, + *, + goal_id: str = GOAL_ID, + agent_id: str = AGENT_ID, + todo_id: str = SELECTED_REPLAN_TODO_ID, ) -> dict[str, Any]: """Choose the fixture Todo explicitly, then consume the generated recovery. @@ -607,9 +611,9 @@ def _bind_selected_replan_guard( """ rc, deferred = _run_cli( registry, runtime, "quota", "should-run", "--codex-app", - "--goal-id", GOAL_ID, "--agent-id", AGENT_ID, + "--goal-id", goal_id, "--agent-id", agent_id, "--turn-instance-id", turn_instance_id, "--scan-path", str(project), - "--todo-id", SELECTED_REPLAN_TODO_ID, + "--todo-id", todo_id, ) if rc == 0: bound = deferred @@ -619,7 +623,7 @@ def _bind_selected_replan_guard( [command] = deferred["interaction_contract"]["cli_channel"]["next_cli_actions"] rc, bound = _run_generated_cli(command, registry_path=registry) assert rc == 0, bound - assert bound["heartbeat_receipt"]["settlement_identity"]["todo_id"] == SELECTED_REPLAN_TODO_ID + assert bound["heartbeat_receipt"]["settlement_identity"]["todo_id"] == todo_id return bound diff --git a/tests/control_plane/test_replan_successor_durable_ack.py b/tests/control_plane/test_replan_successor_durable_ack.py index 44cf5d076b..48643dd27f 100644 --- a/tests/control_plane/test_replan_successor_durable_ack.py +++ b/tests/control_plane/test_replan_successor_durable_ack.py @@ -17,6 +17,9 @@ ReplanWritebackRejected, enforce_open_replan_writeback, ) +from tests.control_plane.test_quota_settlement_cli import ( + _bind_selected_replan_guard, +) GOAL = "successor-review-fixture" AGENT = "fixture-agent" @@ -201,8 +204,17 @@ def call(*args: str, expected_error: str | None = None) -> dict: guard = call("quota", "should-run", "--codex-app", "--goal-id", GOAL, "--agent-id", AGENT, "--turn-instance-id", "turn-original-periodic-review") assert guard["selected_todo"]["todo_id"] == original_todo + assert guard["heartbeat_receipt"]["settlement_binding_owed"] is True # Selection is display until the caller binds the existing Todo explicitly. - guard = call("quota", "should-run", "--codex-app", *binding) + guard = _bind_selected_replan_guard( + registry, + runtime, + project, + "turn-original-periodic-review", + goal_id=GOAL, + agent_id=AGENT, + todo_id=original_todo, + ) obligation = guard["autonomous_replan_obligation"] added = call("todo", "add", "--goal-id", GOAL, "--role", "agent", "--claimed-by", AGENT, "--text", "Verify an independent source artifact", diff --git a/tests/test_claude_goal_release_qualification.py b/tests/test_claude_goal_release_qualification.py index 1d9b97a67a..ecfd873d44 100644 --- a/tests/test_claude_goal_release_qualification.py +++ b/tests/test_claude_goal_release_qualification.py @@ -346,10 +346,11 @@ def test_manifest_oracle_rejects_false_acceptance(tmp_path, defect): runner.verify_manifest_delivery(tmp_path) -def test_real_mcp_delivery_completes_and_settles_existing_plan(tmp_path): +def test_real_mcp_delivery_completes_and_settles_existing_plan(tmp_path, monkeypatch): from loopx.goal_mode_mcp import GoalModeMCPConfig, GoalModeMCPControlPlane project, runtime, launcher = runner.shared.setup(tmp_path) + monkeypatch.chdir(project) # Real delivery class: do not substitute same_agent_non_delivery to make # this acceptance test green. No live model or external side effect. (project / "delivery.txt").write_text("synthetic verified delivery\n") diff --git a/tests/test_configuration_backup.py b/tests/test_configuration_backup.py index 998932ab64..e8c2ba492c 100644 --- a/tests/test_configuration_backup.py +++ b/tests/test_configuration_backup.py @@ -10,9 +10,12 @@ import pytest -from loopx.configuration_backup import capture_configuration_backup, restore_configuration_backup from loopx.capabilities.machine_configuration.builtins import build_builtin_machine_configuration_registry from loopx.capabilities.machine_configuration.store import read_machine_configuration +from loopx.capabilities.configuration_backup import ( + capture_configuration_backup, + restore_configuration_backup, +) from loopx.control_plane.effect_runtime import restart_effect_runtime from loopx.state_backup import build_state_backup_plan, execute_state_backup_plan from tests.control_plane.canonical_authority_fixture import isolate_sqlite_runtime diff --git a/tests/test_delegation_preview_reuse.py b/tests/test_delegation_preview_reuse.py index 7cc582bb4a..aabb7c9eba 100644 --- a/tests/test_delegation_preview_reuse.py +++ b/tests/test_delegation_preview_reuse.py @@ -377,6 +377,72 @@ def test_partial_supervisor_frame_obeys_parent_deadline_and_eof_cleanup(): assert process.poll() == 0 +@pytest.mark.skipif(sys.platform == "win32", reason="POSIX forced cleanup signals") +def test_unconfirmed_supervisor_cleanup_cannot_start_a_second_worker( + tmp_path, monkeypatch +): + from loopx.control_plane.collaboration import delegation_preview_transport + + worker = ( + "import json,os,sys,time\nfrom pathlib import Path\n" + "marker=Path(sys.argv[1])\n" + "for line in sys.stdin:\n" + " json.loads(line);marker.write_text(str(os.getpid()));time.sleep(60)\n" + ) + transport = delegation_preview_transport.DelegationPreviewTransport() + monkeypatch.setattr( + delegation_preview_transport, + "BRIDGE_CLOSE_TIMEOUT_SECONDS", + 0.05, + ) + + def options(marker): + preload = ( + "import{existsSync}from'node:fs';" + f"const marker={json.dumps(str(marker))};" + "const timer=setInterval(()=>{if(existsSync(marker)){" + "clearInterval(timer);" + "Atomics.wait(new Int32Array(new SharedArrayBuffer(4)),0,0)}},1)" + ) + return { + "command": [sys.executable, "-c", worker, str(marker)], + "workspace": tmp_path, + "release": tmp_path, + "environment": { + **_pinned_release_environment(), + "NODE_OPTIONS": "--import=data:text/javascript," + + quote(preload, safe=""), + }, + "registry": tmp_path / "registry.json", + "runtime_root": tmp_path / "runtime", + "goal_id": "fixture-goal", + "agent_id": "fixture-agent", + "todo_id": "todo_fixture", + "argv": ("inspect",), + "timeout": 0.5, + } + + markers = [tmp_path / "worker-1.pid", tmp_path / "worker-2.pid"] + try: + with pytest.raises(subprocess.TimeoutExpired): + transport.preview(**options(markers[0])) + assert markers[0].exists() + os.killpg(int(markers[0].read_text()), 0) + assert transport._process is not None + assert transport._partition is not None + + with pytest.raises(ValueError, match="cleanup remains unconfirmed"): + transport.preview(**options(markers[1])) + assert not markers[1].exists() + finally: + for marker in markers: + if marker.exists(): + try: + os.killpg(int(marker.read_text()), signal.SIGKILL) + except ProcessLookupError: + pass + + @pytest.mark.skipif(sys.platform == "win32", reason="SIGSTOP fault injection requires POSIX") def test_backpressured_supervisor_input_uses_original_parent_deadline(tmp_path, monkeypatch): from loopx.control_plane.collaboration.delegation_preview_transport import DelegationPreviewTransport @@ -417,7 +483,7 @@ def measured_send(*args, **kwargs): transport.preview(**options, argv=("x" * 65536,), timeout=0.1) assert send_durations[-1] < 0.4, send_durations assert transport._process is None - assert process.poll() == 0 + assert process.poll() is not None finally: if timer: timer.cancel() diff --git a/tests/test_goal_mode_mcp_settlement.py b/tests/test_goal_mode_mcp_settlement.py index 3dfafc7fc9..5a76e5e292 100644 --- a/tests/test_goal_mode_mcp_settlement.py +++ b/tests/test_goal_mode_mcp_settlement.py @@ -200,7 +200,10 @@ def _write_fixture(tmp_path: Path) -> tuple[Path, Path]: "status": "active", "repo": str(project), "state_file": state_file.name, - "adapter": {"kind": "harness_self_improvement"}, + "adapter": { + "kind": "harness_self_improvement", + "status": "connected-read-only", + }, "quota": { "compute": 1.0, "window_hours": 24, @@ -464,7 +467,10 @@ def test_real_mcp_terminal_completion_closes_out_after_spend( @pytest.mark.parametrize("lost_after", ["lifecycle", "writeback", "spend"]) -def test_real_mcp_completion_recovers_a_lost_mutation_response(tmp_path, lost_after): +def test_real_mcp_completion_recovers_a_lost_mutation_response( + tmp_path: Path, + lost_after: str, +) -> None: """The real write commits, but its caller sees a failure: retry must not pay twice.""" registry, _ = _write_fixture(tmp_path) added = add_goal_todo( @@ -521,7 +527,9 @@ def lose_once(args, **kwargs): assert status["quota"]["spent_slots"] == 1 -def test_real_mcp_links_existing_successor_without_creating_another_todo(tmp_path): +def test_real_mcp_links_existing_successor_without_creating_another_todo( + tmp_path: Path, +) -> None: registry, state_file = _write_fixture(tmp_path) ids = [str(add_goal_todo( registry_path=registry, goal_id=GOAL_ID, role="agent", text=text, diff --git a/tests/test_host_vision_recovery.py b/tests/test_host_vision_recovery.py index deff2a2eb1..49aa8a1a2f 100644 --- a/tests/test_host_vision_recovery.py +++ b/tests/test_host_vision_recovery.py @@ -14,6 +14,17 @@ def vision(state="no_followup"): + path_outcome = { + "no_followup": "no_change", + "vision_closed": "replan", + }.get(state, "continue") + path_change = { + "no_followup": {"stopped": ["No further fixture delivery is claimed."]}, + "vision_closed": {"changed": ["Close the validated fixture stage."]}, + }.get( + state, + {"retained": ["Keep the remaining explicit Todo as the delivery frontier."]}, + ) return { "schema_version": "goal_vision_replan_contract_v0", "state": state, "vision_patch": { @@ -21,6 +32,14 @@ def vision(state="no_followup"): "acceptance_summary": "Replay, reversal and CLI atomic output verified.", "last_patch_summary": "Local acceptance tests passed; no external delivery is requested.", }, + "path_delta": { + "schema_version": "goal_path_delta_v0", + "outcome": path_outcome, + "prior_assumption": "The current milestone still required validation.", + "observed_reality": "The fixture acceptance passed with a verified Todo transition.", + "evidence_refs": ["fixture:synthetic-lifecycle-acceptance"], + **path_change, + }, }