diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index 2ce3ed6223..1975e0bab3 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -189,6 +189,8 @@ touch, mention it in the PR instead of fixing it there. | `package-smoke.yml` | extension package paths | no | extension packages install, entrypoints, and example schemas | | `release-artifacts.yml` | `loopx/`, packaging, and lockfile paths | no | release identity and a release build from this source | | `ark-turn.yml` | Turn driver and collaboration paths | no | optional Ark Turn package, stdio MCP, and DSH parity | +| `dsh-plugin.yml` | DSH plugin and distribution workflow paths | no | Linux/Windows package install and uninstall, typed host contracts, and Linux real LoopX/web runtime admission | +| `dsh-plugin-publish.yml` | no (manual, published plugin tag) | no | merged release identity, immutable npm package publication, and distribution readback | | `frontstage-pages.yml` | README, dashboard, and chat bundle paths | no | public Pages build | | `desktop-release-artifacts.yml` | desktop app and dashboard paths | no | macOS and Windows desktop builds | | `desktop-updater.yml` | desktop app paths | no | desktop app build and updater feed | 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 6a275ee803..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 @@ -459,7 +459,7 @@ for (const capabilityId of [ const matches = capabilityLocalization.match(new RegExp(`${capabilityId}:`, "g")) ?? []; assert.equal(matches.length, 2, `${capabilityId} has English and Simplified Chinese metadata`); } -for (const fieldKey of ["allowed_domains", "coordinator_agent_id", "eligible_endpoints", "enabled", "executor_endpoint", "executor_model", "executor_reasoning_effort", "max_children", "profile", "profile_preset", "review_order", "route_ref", "safe_fix", "selection_policy", "strict_receipt", "timezone"]) { +for (const fieldKey of ["agent_orders", "allowed_domains", "coordinator_agent_id", "eligible_endpoints", "enabled", "executor_endpoint", "executor_model", "executor_reasoning_effort", "max_children", "profile", "profile_preset", "review_order", "route_ref", "safe_fix", "selection_policy", "strict_receipt", "timezone"]) { const matches = capabilityLocalization.match(new RegExp(`^\\s+${fieldKey}:`, "gm")) ?? []; assert.equal(matches.length, 2, `${fieldKey} has English and Simplified Chinese field copy`); } diff --git a/docs/reference/protocols/active-state-structured-projection-v0.md b/docs/reference/protocols/active-state-structured-projection-v0.md index 4a61afad75..a1a00b2f8c 100644 --- a/docs/reference/protocols/active-state-structured-projection-v0.md +++ b/docs/reference/protocols/active-state-structured-projection-v0.md @@ -226,11 +226,14 @@ title still fails parity, even when its first 500 characters match. This is not permission to accept malformed records or to shorten an already committed Todo to repair its display. -Status, `todo list` (including an exact ID), dashboard and chat attention views -keep their existing bounded summaries. The canonical provider and regenerated -active state retain the complete source; a display summary is not an input to -source serialization. No new frontend setting, Lark command or parallel state -store is introduced. This change stays in the permanent Python Markdown/legacy +Status, unfiltered `todo list`, `todo list --thin` (including an exact ID), +dashboard and chat attention views keep their existing bounded summaries. +An exact `todo list --todo-id ID` without `--thin` returns the complete current +or retained request text through the existing Task detail reader. The canonical +provider and regenerated active state retain the complete source; a display +summary is not an input to source serialization. No new frontend setting, Lark +command or parallel state store is introduced. This change stays in the +permanent Python Markdown/legacy I/O adapter; the TypeScript authority, admission and delivery-confirmation owners are unchanged. diff --git a/examples/control_plane/heartbeat-prompt-smoke.py b/examples/control_plane/heartbeat-prompt-smoke.py index 3cbb929a0d..245e0ffa29 100644 --- a/examples/control_plane/heartbeat-prompt-smoke.py +++ b/examples/control_plane/heartbeat-prompt-smoke.py @@ -587,7 +587,7 @@ def main() -> int: "current quota claim/lease and workspace contract plus repository rules", "Within authority/budget, deliver verifiable results", "Task-scoped coordination grants no authority over other agents", - "Keep scope in the heartbeat prompt, not todo metadata", + "Keep scope in this prompt, not todo metadata", "Normal turns use CLI `interaction_contract`; use `loopx-project` for " "lifecycle/registry and `loopx-self-repair` for runtime/projection drift", "use selection_command when required", diff --git a/examples/loopx-turn-codex-cli-e2e-smoke.py b/examples/loopx-turn-codex-cli-e2e-smoke.py index 338c62e78f..ae5d6d5d23 100755 --- a/examples/loopx-turn-codex-cli-e2e-smoke.py +++ b/examples/loopx-turn-codex-cli-e2e-smoke.py @@ -45,10 +45,9 @@ def _marker_value(turn_number: int) -> 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 d8f444289c..c1b86f7a74 100644 --- a/examples/shared-goal-authority-e2e/mutants.py +++ b/examples/shared-goal-authority-e2e/mutants.py @@ -269,6 +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_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("public_refresh_writes_owned_paragraph", (("loopx/state_refresh.py", replacement( "if not dry_run:\n runs_dir.mkdir(parents=True, exist_ok=True)\n json_path, markdown_path = reserve_unique_run_paths(runs_dir, generated_at)", "if not dry_run:\n resolved_state_file.write_text(\n current_text + \"\\n- \" + normalized_next_action, encoding=\"utf-8\"\n )\n runs_dir.mkdir(parents=True, exist_ok=True)\n json_path, markdown_path = reserve_unique_run_paths(runs_dir, generated_at)")),), @@ -420,7 +424,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 3d2171e08b..bfa0caf5c3 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_server.py b/loopx/chat_server.py index ed71205816..77f3fccd7e 100644 --- a/loopx/chat_server.py +++ b/loopx/chat_server.py @@ -11,7 +11,7 @@ from typing import Any, Literal from urllib.parse import parse_qs, urlparse -from . import chat_configuration_api as config_api +from .presentation import configuration_api as config_api from .attached_session_api import AttachedSessionRequestMixin from .chat import ( TodoReviewPreviewConflict, diff --git a/loopx/cli_commands/configuration_backup.py b/loopx/cli_commands/configuration_backup.py index 1996477a39..fa6f5719dc 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 diff --git a/loopx/cli_commands/quota.py b/loopx/cli_commands/quota.py index 672c73cce0..b21619f0c9 100644 --- a/loopx/cli_commands/quota.py +++ b/loopx/cli_commands/quota.py @@ -423,6 +423,45 @@ def _project_quota_cli_payload( degraded["turn_envelope_skipped"] = str(envelope_error)[:200] return degraded + +def _emit_quota_result( + payload: dict[str, object], + args: argparse.Namespace, + *, + context: QuotaCommandContext | None, + registry_path: Path, + heartbeat_turn_id: str | None, + usage_quota_started: int, + capture_directory: Path | None, + detail_sections: frozenset[str], + print_payload: PrintPayload, +) -> int: + """Capture the full decision before projecting and printing the CLI view.""" + if context is not None: + observe_quota_result( + args, payload, registry_path=registry_path, runtime_root=context.runtime_root, + turn_id=_effective_spend_turn_instance_id(payload, heartbeat_turn_id=heartbeat_turn_id), + started_at=usage_quota_started, + ) + capture_decision(capture_directory, payload) + payload = _project_quota_cli_payload( + payload, args, detail_sections, + context.scheduler_context if context is not None else None, + captured_decision_path=( + str(capture_directory / "decision.json") if capture_directory is not None else None + ), + ) + if args.quota_command == "should-run" and context is not None: + attach_host_poll_receipt( + context.status_payload, + args, + payload, + registry_path=registry_path, + ) + print_payload(payload, args.format, _quota_renderer(args)) + return 0 if payload.get("ok") else 1 + + def handle_quota_command( args: argparse.Namespace, *, @@ -894,26 +933,14 @@ def handle_quota_command( goal_id=args.goal_id, agent_id=args.agent_id, ) - if context is not None: - observe_quota_result( - args, payload, registry_path=registry_path, runtime_root=context.runtime_root, - turn_id=_effective_spend_turn_instance_id(payload, heartbeat_turn_id=heartbeat_turn_id), - started_at=usage_quota_started, - ) - capture_decision(capture_directory, payload) - payload = _project_quota_cli_payload( - payload, args, detail_sections, - context.scheduler_context if context is not None else None, - captured_decision_path=( - str(capture_directory / "decision.json") if capture_directory is not None else None - ), + return _emit_quota_result( + payload, + args, + context=context, + registry_path=registry_path, + heartbeat_turn_id=heartbeat_turn_id, + usage_quota_started=usage_quota_started, + capture_directory=capture_directory, + detail_sections=detail_sections, + print_payload=print_payload, ) - if args.quota_command == "should-run" and context is not None: - attach_host_poll_receipt( - context.status_payload, - args, - payload, - registry_path=registry_path, - ) - print_payload(payload, args.format, _quota_renderer(args)) - return 0 if payload.get("ok") else 1 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 05dc39ff6b..313c65876d 100644 --- a/loopx/control_plane/effect_runtime.py +++ b/loopx/control_plane/effect_runtime.py @@ -47,6 +47,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 @@ -479,6 +480,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", @@ -1008,6 +1029,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/control_plane/heartbeat/agent.py b/loopx/control_plane/heartbeat/agent.py index 2a8cadb416..db5361e29c 100644 --- a/loopx/control_plane/heartbeat/agent.py +++ b/loopx/control_plane/heartbeat/agent.py @@ -152,7 +152,7 @@ def render_peer_agent_scope_instruction( f"Equal peer `{identity}` (peer_v1); scope: {scope_text}. Follow the current quota " "claim/lease and workspace contract plus repository rules. Follow todo " "continuation policy. Task-scoped coordination grants no authority over " - "other agents. Keep scope in the heartbeat prompt, not todo metadata." + "other agents. Keep scope in this prompt, not todo metadata." ) if compact: return ( diff --git a/loopx/control_plane/quota/goal_boundary.py b/loopx/control_plane/quota/goal_boundary.py index c4e0877e3d..7335868065 100644 --- a/loopx/control_plane/quota/goal_boundary.py +++ b/loopx/control_plane/quota/goal_boundary.py @@ -206,7 +206,13 @@ def _reward_memory_enablement_projection( "writability_verified", "exact_readback_verified", ) - return {field: status[field] for field in fields if field in status} + projection = {field: status[field] for field in fields if field in status} + if status.get("config_schema_version"): + projection["config_schema_version"] = str(status["config_schema_version"]) + config_runtime_route = status.get("config_runtime_route") + if isinstance(config_runtime_route, Mapping): + projection["config_runtime_route"] = dict(config_runtime_route) + return projection def _reward_memory_automation_projection( @@ -378,15 +384,6 @@ def goal_boundary( reward_capability.update( _reward_memory_enablement_projection(reward_memory_experiment_status) ) - if reward_memory_experiment_status.get("config_schema_version"): - reward_capability["config_schema_version"] = str( - reward_memory_experiment_status["config_schema_version"] - ) - config_runtime_route = reward_memory_experiment_status.get( - "config_runtime_route" - ) - if isinstance(config_runtime_route, Mapping): - reward_capability["config_runtime_route"] = dict(config_runtime_route) if agent_id is not None: reward_capability.update( { @@ -442,10 +439,6 @@ def goal_boundary( if repository_identity and repository_identity.startswith("git:"): boundary["task_repository"] = repository_identity project_asset_source = item if item is not None else goal - for policy_source in (goal, project_asset_source): - if not isinstance(policy_source, dict): - continue - break if isinstance(project_asset_source, dict) and project_asset_source.get( "project_asset" ): diff --git a/loopx/control_plane/testing/cli_output_budget.py b/loopx/control_plane/testing/cli_output_budget.py index 314c80e2d0..4de48ad2a7 100644 --- a/loopx/control_plane/testing/cli_output_budget.py +++ b/loopx/control_plane/testing/cli_output_budget.py @@ -170,16 +170,12 @@ class CliOutputCommandClassification: markdown_anchor="# LoopX Turn Plan", max_chars={ "small": {"json": 12_000, "markdown": 300}, - # The crowded fixture exercises the required-vision route. Its - # TurnEnvelope intentionally carries the complete authoring schema - # that the validator accepts, plus typed executor/selection facts - # and explicit registry routing. Completing its evidence-linked - # example while removing duplicate prose changes the same fixture - # from 14,360 to 14,571 characters; retain the 14,600 ceiling. - # The over-target TurnEnvelope diagnostic remains visible instead - # of hiding authority overflow; latest main renders it in 542 - # characters, leaving a narrow 58-character presentation margin. - "crowded": {"json": 14_600, "markdown": 600}, + # Required vision carries the validator's complete authoring schema, + # executable registry-bound commands and the overflow diagnostic. + # The same fixed-path fixture emits 14,647 chars on base and head; + # 15,000 leaves 353 chars without dropping these decision inputs. + # Keep the line, per-Todo and fixed semantic-growth guards below. + "crowded": {"json": 15_000, "markdown": 600}, "multi_agent": {"json": 12_000, "markdown": 300}, }, max_lines={ @@ -236,7 +232,11 @@ class CliOutputCommandClassification: markdown_anchor="# LoopX Diagnosis Packet", max_chars={ "small": {"json": 21_000, "markdown": 4_300}, - "crowded": {"json": 44_000, "markdown": 4_500}, + # The unchanged selected/Goal-array diagnostic emits 44,126 chars + # on the same fixed-path base/head fixture. Preserve both consumers + # and their required-replan evidence; 45,000 leaves 874 chars. + # Line and per-Todo/fixed-growth limits remain independently active. + "crowded": {"json": 45_000, "markdown": 4_500}, "multi_agent": {"json": 21_000, "markdown": 4_300}, }, max_lines={ diff --git a/loopx/control_plane/testing/cli_output_semantics.py b/loopx/control_plane/testing/cli_output_semantics.py index e8921c06b5..a020021b47 100644 --- a/loopx/control_plane/testing/cli_output_semantics.py +++ b/loopx/control_plane/testing/cli_output_semantics.py @@ -32,7 +32,7 @@ def heartbeat_peer_admission_prompt_revision(text: str) -> str | None: block = ( "Follow the current quota claim/lease and workspace contract plus repository rules. " "Follow todo continuation policy. Task-scoped coordination grants no authority over " - "other agents. Keep scope in the heartbeat prompt, not todo metadata." + "other agents. Keep scope in this prompt, not todo metadata." ) return "heartbeat_peer_admission_v1" if block in text else None diff --git a/loopx/control_plane/testing/replan_semantic_action_behavior.py b/loopx/control_plane/testing/replan_semantic_action_behavior.py index 8f740de088..f332e7d811 100644 --- a/loopx/control_plane/testing/replan_semantic_action_behavior.py +++ b/loopx/control_plane/testing/replan_semantic_action_behavior.py @@ -534,7 +534,34 @@ def _is_replan_successor_create(command: str) -> bool: ) +def _required_explore_read_command(packet: Mapping[str, Any]) -> str | None: + """Accept only this fixture's current, scoped turn-start read contract.""" + reads = packet.get("required_reads") or [] + if not reads: + return None + if not isinstance(reads, list) or len(reads) != 1: + raise ValueError("unexpected_command") + read = reads[0] + if not isinstance(read, Mapping) or ( + read.get("kind") != "explore_turn_context" + or read.get("source") != "turn_start_capability_hook" + or read.get("ordering") != "before_work" + ): + raise ValueError("unexpected_command") + command = str(read.get("command") or "") + tokens = loopx_command_tokens(command) or [] + if "explore" not in tokens or ( + tokens[tokens.index("explore"):tokens.index("explore") + 2] + != ["explore", "turn-context"] + or argument_value(tokens, "--goal-id") != _FIXTURE_GOAL_ID + or argument_value(tokens, "--agent-id") != _FIXTURE_AGENT_ID + ): + raise ValueError("unexpected_command") + return command + + def _quota_behavior_observation(packet: Mapping[str, Any]) -> dict[str, Any]: + _required_explore_read_command(packet) signature = quota_action_signature_document(packet) action = dict(signature.get("action") or {}) user = dict(signature.get("user") or {}) @@ -552,7 +579,6 @@ def _quota_behavior_observation(packet: Mapping[str, Any]) -> dict[str, Any]: or context.get("delivery") != "host_projected" or context.get("evidence_source") != "compact_run_history" or dict(context.get("delivery_receipt") or {}).get("status") != "delivered" - or packet.get("required_reads") not in (None, []) ): raise ValueError("quota does not expose the semantic replan contract") if not ( @@ -1119,6 +1145,22 @@ def _dispatch_behavior_command( ) -> tuple[str, str, bool]: if _is_quota_guard(command): return _handle_quota_command(command, state), "quota_should_run", False + if state.quota_packet is not None and ( + required := _required_explore_read_command(state.quota_packet) + ): + observed = any(step["kind"] == "explore_turn_context" for step in state.steps) + if shlex.split(command) == shlex.split(required): + if observed: + raise ValueError("repeated_explore_turn_context") + output = _execute_loopx( + command, fixture=state.fixture, turn_instance_id=state.turn_instance_id + ) + context = json.loads(output) + if context.get("goal_id") != _FIXTURE_GOAL_ID or context.get("agent_id") != _FIXTURE_AGENT_ID: + raise ValueError("unexpected_command") + return output, "explore_turn_context", False + if not observed: + raise ValueError("required_explore_turn_context_missing") if (clock_output := _clock_output(command)) is not None: return _handle_clock_command(clock_output, state), "clock", False if _bounded_workspace_read_plan(command, fixture=state.fixture) is not None: @@ -1150,6 +1192,10 @@ def _behavior_command_kind( if _is_quota_guard(command): return "quota_should_run" + if state.quota_packet is not None and ( + required := _required_explore_read_command(state.quota_packet) + ) and shlex.split(command) == shlex.split(required): + return "explore_turn_context" if _clock_output(command) is not None: return "clock" if _bounded_workspace_read_plan(command, fixture=state.fixture) is not None: @@ -1166,6 +1212,8 @@ def _behavior_command_kind( _EXPECTED_BEHAVIOR_FAILURES = frozenset( { + "repeated_explore_turn_context", + "required_explore_turn_context_missing", "repeated_quota_should_run", "repeated_clock", "semantic_action_before_quota", diff --git a/loopx/chat_configuration_api.py b/loopx/presentation/configuration_api.py similarity index 87% rename from loopx/chat_configuration_api.py rename to loopx/presentation/configuration_api.py index fd566c59c3..6d3a10c4de 100644 --- a/loopx/chat_configuration_api.py +++ b/loopx/presentation/configuration_api.py @@ -2,13 +2,13 @@ from collections.abc import Callable -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 +from . import configuration_backup_api as backup_api +from . 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 class ChatConfigurationRequestMixin( 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 64fd4c241f..adfc346128 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, @@ -989,22 +1005,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/architecture/test_contributor_task_board.py b/tests/architecture/test_contributor_task_board.py index 9ba19eb717..45d2dcf078 100644 --- a/tests/architecture/test_contributor_task_board.py +++ b/tests/architecture/test_contributor_task_board.py @@ -47,7 +47,16 @@ def run(updated: str = board) -> None: ) check() - return run, board, tmp_path + # Negative controls target a disposable row rather than a real task that + # legitimately disappears when delivered. The default run still checks + # the untouched production board against the real RFC generator. + first_row = smoke["contributor_board_claimable_rows"](board)[0] + fixture_row = ( + "| GH-C9999 | [RFC](../architecture/rfcs/goal-direction-baseline-v0.md) " + "| Synthetic admission gap. Exit: preserve canonical admission. " + "| Focused admission negative controls | Available |" + ) + return run, board.replace(first_row, fixture_row, 1), tmp_path def edit_row( @@ -58,7 +67,7 @@ def edit_row( task_id: str | None = None, gap: str | None = None, ) -> str: - line = next(line for line in board.splitlines() if line.startswith("| GH-C89b |")) + line = next(line for line in board.splitlines() if line.startswith("| GH-C9999 |")) cells = [cell.strip() for cell in line.strip("|").split("|")] for index, value in ((0, task_id), (1, anchor), (2, gap), (4, status)): if value is not None: diff --git a/tests/capabilities/test_capability_configuration_ui.py b/tests/capabilities/test_capability_configuration_ui.py index 7fadf12723..39787436c7 100644 --- a/tests/capabilities/test_capability_configuration_ui.py +++ b/tests/capabilities/test_capability_configuration_ui.py @@ -71,7 +71,13 @@ def test_pull_request_review_editor_supports_machine_and_goal_ci_policy() -> Non assert editor["writable_scopes"] == ["machine", "goal"] assert editor["fields"][0]["key"] == "wait_for_ci" assert editor["fields"][0]["input_kind"] == "boolean" - assert editor["fields"][1:] == [ + assert [field["key"] for field in editor["fields"]] == [ + "wait_for_ci", "owner_logins", "review_order", + ] + owner = editor["fields"][1] + assert owner["input_kind"] == "string_list" + assert "does not infer membership or grant review/merge authority" in owner["description"] + assert editor["fields"][2:] == [ { "key": "review_order", "label": "Review direction", 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_cli_output_differential.py b/tests/control_plane/test_cli_output_differential.py index a53fcfb8d6..72c15e44c3 100644 --- a/tests/control_plane/test_cli_output_differential.py +++ b/tests/control_plane/test_cli_output_differential.py @@ -91,7 +91,7 @@ def test_readable_peer_admission_growth_is_one_time_and_thin_only(output_format) block = ( "Follow the current quota claim/lease and workspace contract plus repository rules. " "Follow todo continuation policy. Task-scoped coordination grants no authority over " - "other agents. Keep scope in the heartbeat prompt, not todo metadata." + "other agents. Keep scope in this prompt, not todo metadata." ) marker = heartbeat_peer_admission_prompt_revision(block) assert marker == "heartbeat_peer_admission_v1" diff --git a/tests/control_plane/test_effect_runtime_integration.py b/tests/control_plane/test_effect_runtime_integration.py index 7dcc857555..72c5f96eca 100644 --- a/tests/control_plane/test_effect_runtime_integration.py +++ b/tests/control_plane/test_effect_runtime_integration.py @@ -302,6 +302,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_heartbeat_agent_input.py b/tests/control_plane/test_heartbeat_agent_input.py index bcd6ea98fb..65080ae64c 100644 --- a/tests/control_plane/test_heartbeat_agent_input.py +++ b/tests/control_plane/test_heartbeat_agent_input.py @@ -61,7 +61,7 @@ def test_thin_cli_preserves_readable_admission_and_authority_with_host_scope(tmp assert "quota claim/lease and workspace contract plus repository rules" in body assert "Follow todo continuation policy" in body assert "Task-scoped coordination grants no authority over other agents" in body - assert "Keep scope in the heartbeat prompt, not todo metadata" in body + assert "Keep scope in this prompt, not todo metadata" in body for capability in ("filesystem_read", "filesystem_write", "shell"): assert capability in body 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_monitor_observation_admission.py b/tests/control_plane/test_monitor_observation_admission.py index 99e6d39dde..39bd4ecfdc 100644 --- a/tests/control_plane/test_monitor_observation_admission.py +++ b/tests/control_plane/test_monitor_observation_admission.py @@ -1,6 +1,7 @@ from __future__ import annotations import json +from collections.abc import Iterator from pathlib import Path import pytest @@ -17,6 +18,14 @@ _spend_run_count, _write_fixture, ) +from tests.control_plane.test_cli_output_budget import _stable_budget_fixture_root + + +@pytest.fixture +def settlement_budget_root(tmp_path: Path) -> Iterator[Path]: + # Measure the actual output with stable paths, as other CLI budgets do. + with _stable_budget_fixture_root(tmp_path) as root: + yield root def test_settled_advancement_hold_preserves_independent_monitor( @@ -204,12 +213,13 @@ def test_receipt_bound_advancement_allows_one_auxiliary_due_monitor_receipt( ], ) def test_completed_advancement_retains_auxiliary_monitor_admission( - tmp_path: Path, + settlement_budget_root: Path, monkeypatch: pytest.MonkeyPatch, provider: str, writeback: bool, spend_first: bool, ) -> None: + tmp_path = settlement_budget_root from canonical_authority_fixture import ( initialize_canonical_authority, isolate_sqlite_runtime, 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..7a7ab0341c 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 @@ -5086,7 +5090,7 @@ def test_pending_action_selection_does_not_commit_after_new_user_gate( assert all(not event["details"].get("settlement_effect_id") for event in events) -def test_todoless_autonomous_replan_settles_quota_refresh_spend_chain( +def test_todoless_replan_settles_once_and_rearms_missing_vision( tmp_path: Path, ) -> None: project, runtime, registry_path = _write_fixture(tmp_path) @@ -5279,11 +5283,17 @@ def test_todoless_autonomous_replan_settles_quota_refresh_spend_chain( ) assert fresh_rc == 0, fresh - assert fresh["decision"] == "skip", fresh - assert fresh["effective_action"] == "monitor_quiet_skip" - assert fresh["execution_obligation"]["must_attempt_work"] is False - assert fresh.get("autonomous_replan_obligation") is None - assert fresh.get("replan_action_packet") is None + # The original Turn is settled, but its material write did not establish + # a per-agent Vision baseline. The next Turn must address that independent + # gap without replaying the previous obligation or spending a second slot. + assert fresh["decision"] == "autonomous_replan_required", fresh + assert fresh["execution_obligation"]["must_attempt_work"] is True + assert fresh["autonomous_replan_obligation"]["rearmed_after_obligation_id"] == obligation_id + assert fresh["autonomous_replan_obligation"]["obligation_id"] != obligation_id + assert any( + gap["kind"] == "vision_checkpoint_missing" and gap["missing_baseline"] + for gap in fresh["goal_frontier_projection"]["acceptance_gaps"] + ) assert fresh["heartbeat_receipt"]["turn_instance_id"] == fresh_turn_id assert _spend_run_count(runtime) == 1 diff --git a/tests/control_plane/test_replan_semantic_action_behavior.py b/tests/control_plane/test_replan_semantic_action_behavior.py index d0c9b3d01e..cb1e72044d 100644 --- a/tests/control_plane/test_replan_semantic_action_behavior.py +++ b/tests/control_plane/test_replan_semantic_action_behavior.py @@ -89,6 +89,15 @@ def _exhausted_action( ) +def _explore_context_action( + request: Mapping[str, object], +) -> ScriptedExecToolAction: + read = _latest_quota_packet(request)["required_reads"][0] + assert read["kind"] == "explore_turn_context" + assert read["ordering"] == "before_work" + return ScriptedExecToolAction(command=read["command"]) + + def _composition_successor_action( request: Mapping[str, object], ) -> ScriptedExecToolAction: @@ -426,6 +435,7 @@ def test_real_tool_loop_selects_composition_gap_and_creates_bound_successor( transport = ScriptedDoubaoExecTransport( [ ScriptedExecToolAction(command=fixture.quota_guard_command), + _explore_context_action, ScriptedExecToolAction(command="cat replan-frontier.json"), ScriptedExecToolAction(command="cat fixture/permission-config.json"), _composition_successor_action, @@ -443,6 +453,7 @@ def test_real_tool_loop_selects_composition_gap_and_creates_bound_successor( assert receipt["qualification_passed"] is True, receipt assert receipt["observed_tool_sequence"] == [ "quota_should_run", + "explore_turn_context", "workspace_read", "workspace_read", "replan_successor_create", @@ -468,6 +479,40 @@ def test_real_tool_loop_selects_composition_gap_and_creates_bound_successor( ] == selected_gap["experiment_node_ref"] +@pytest.mark.parametrize("attempt", ["omit", "repeat", "wrong_scope"]) +def test_composition_read_obligation_cannot_be_skipped_replayed_or_retargeted( + tmp_path: Path, attempt: str, +) -> None: + fixture = _build_fixture(tmp_path / "oracle", composition_frontier=True) + actions = [ScriptedExecToolAction(command=fixture.quota_guard_command)] + if attempt == "omit": + actions.append(ScriptedExecToolAction(command="cat replan-frontier.json")) + elif attempt == "repeat": + actions.extend([_explore_context_action, _explore_context_action]) + else: + def wrong_scope(request: Mapping[str, object]) -> ScriptedExecToolAction: + action = _explore_context_action(request) + return ScriptedExecToolAction( + command=action.command.replace( + "--goal-id replan-semantic-action-fixture", "--goal-id another-goal" + ) + ) + actions.append(wrong_scope) + receipt = DoubaoReplanSemanticActionBehaviorActor( + api_key="test-only-placeholder", transport=ScriptedDoubaoExecTransport(actions), + ).qualify( + qualification_id=f"composition-read-{attempt}", + fixture_root=tmp_path / "actor", composition_frontier=True, + ) + assert receipt["qualification_passed"] is False + assert receipt["failure_code"] == ( + "repeated_explore_turn_context" + if attempt == "repeat" else "required_explore_turn_context_missing" + ) + assert receipt["semantic_action_accepted"] is False + assert "replan_successor_create" not in receipt["observed_tool_sequence"] + + def test_action_outside_observed_frontier_is_rejected(tmp_path: Path) -> None: fixture = _build_fixture(tmp_path / "oracle") transport = ScriptedDoubaoExecTransport( 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/control_plane/test_todo_projection_recovery.py b/tests/control_plane/test_todo_projection_recovery.py index b522352fa2..b60f444e9f 100644 --- a/tests/control_plane/test_todo_projection_recovery.py +++ b/tests/control_plane/test_todo_projection_recovery.py @@ -122,11 +122,18 @@ def test_long_committed_todo_rebuilds_from_the_fresh_head_without_a_second_creat assert state.read_text() == rendered and _read(runtime) == before code, listed = _cli(registry, "list", "--goal-id", "goal-a", "--todo-id", record["todo_id"]) - # The CLI remains a bounded attention view, not source serialization. + # Exact cold reads preserve the current request; attention views remain bounded. summary_text = normalize_todo_text(text) - assert code == 0 and listed["todo"]["text"] == summary_text, listed + assert code == 0 and listed["todo"]["text"] == text, listed assert listed["todo"]["title"] == todo_priority_parts(summary_text)[1] assert listed["authority_read"]["provider_revision"] == before["provider_revision"] + code, thin = _cli(registry, "list", "--goal-id", "goal-a", "--todo-id", record["todo_id"], "--thin") + assert code == 0 and thin["todo"]["todo_id"] == record["todo_id"], thin + assert len(thin["todo"]["text"]) <= 500 + assert "retain the final obligation" not in thin["todo"]["text"] + code, hot = _cli(registry, "list", "--goal-id", "goal-a") + assert code == 0 and hot["todos"][0]["text"] == summary_text, hot + assert _read(runtime) == before manager = read_manager_goal_details(registry, runtime, "goal-a", owner_scope=True) assert manager["status"] == "read" and manager["coverage"]["active"] == 1 assert manager["authority_revision"] == before["provider_revision"] diff --git a/tests/host_surface_cli_probes.py b/tests/host_surface_cli_probes.py index bd71587d9a..24bf399790 100644 --- a/tests/host_surface_cli_probes.py +++ b/tests/host_surface_cli_probes.py @@ -12,6 +12,7 @@ import json import os import shlex +import shutil import subprocess import sys from pathlib import Path @@ -124,6 +125,9 @@ def onboarding_setup_command_installs( """agent-onboard hands back a setup command; executing it must provision the host it named, from any cwd. The surface's skills come from the LoopX installer, not from a host that manages skills itself.""" + node = shutil.which("node") + assert node is not None, "The real CLI onboarding probe requires Node.js" + env = {**env, "PATH": os.pathsep.join((str(Path(node).parent), env["PATH"]))} onboard = run_cli( "agent-onboard", "--agent-type", diff --git a/tests/test_chat_machine_configuration_api.py b/tests/test_chat_machine_configuration_api.py index 9188faffe9..4805cfe360 100644 --- a/tests/test_chat_machine_configuration_api.py +++ b/tests/test_chat_machine_configuration_api.py @@ -444,7 +444,7 @@ def test_machine_catalog_discovers_goal_features_without_granting_machine_writes assert [ field["key"] for field in machine["pull_request_review"]["configuration_editor"]["fields"] - ] == ["wait_for_ci", "review_order"] + ] == ["wait_for_ci", "owner_logins", "review_order"] assert "multi_subagent" in machine for capability_id, item in machine.items(): assert "current" not in item 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 10688890d7..f255ed1fda 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_completion_validation.py b/tests/test_goal_mode_mcp_completion_validation.py index e0573167e4..2da1e2e47c 100644 --- a/tests/test_goal_mode_mcp_completion_validation.py +++ b/tests/test_goal_mode_mcp_completion_validation.py @@ -48,7 +48,10 @@ def _write_fixture(tmp_path: Path) -> tuple[Path, Path]: "status": "active", "repo": str(repo), "state_file": state.name, - "adapter": {"kind": "harness_self_improvement"}, + "adapter": { + "kind": "harness_self_improvement", + "status": "connected-read-only", + }, "coordination": { "agent_model": "peer_v1", "registered_agents": [AGENT], 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, + }, } diff --git a/tests/test_kiro_cli_mcp.py b/tests/test_kiro_cli_mcp.py index 519a3d9b6e..93b297b6fe 100644 --- a/tests/test_kiro_cli_mcp.py +++ b/tests/test_kiro_cli_mcp.py @@ -66,7 +66,10 @@ def _write_project( "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, @@ -181,10 +184,8 @@ def test_completion_through_the_kiro_server_settles_once_under_its_profile( gate = json.loads(control.should_run()) assert gate["goal_id"] == GOAL_ID assert gate["agent_identity"]["agent_id"] == AGENT_ID - assert any( - f"--agent-id {AGENT_ID} --runtime-profile kiro_cli" in action - for action in gate["interaction_contract"]["cli_channel"]["next_cli_actions"] - ), gate["interaction_contract"]["cli_channel"] + assert gate["should_run"] is True + assert gate["scheduler_hint"]["execution_context"]["host_surface"] == "kiro_cli" # The server refuses to act for an agent the session is not bound to. foreign = json.loads(control.claim_task(todo_id, OTHER_AGENT_ID)) @@ -222,6 +223,35 @@ def test_completion_through_the_kiro_server_settles_once_under_its_profile( assert by_id[todo_id]["status"] == "done" +def test_completion_requires_a_connected_delivery_frontier(tmp_path: Path) -> None: + project, registry, state_file = _write_project(tmp_path) + document = json.loads(registry.read_text(encoding="utf-8")) + document["goals"][0]["adapter"].pop("status") + registry.write_text(json.dumps(document), encoding="utf-8") + added = add_goal_todo( + registry_path=registry, + goal_id=GOAL_ID, + role="agent", + text="Retain the unadmitted Kiro task.", + task_class="advancement_task", + claimed_by=AGENT_ID, + ) + before = state_file.read_bytes() + control = GoalModeMCPControlPlane( + CONFIG, lambda: goal_context(project, {KIRO_CLI_SESSION_ID_ENV: SESSION_ID}), + ) + control.command_prefix = lambda: [sys.executable, "-m", "loopx.cli"] + result = json.loads(control.complete_task(str(added["todo_id"]), AGENT_ID, "not admitted")) + assert result["ok"] is False and result["settlement_blocked_completion"] is True + assert result["settlement"]["failed_stage"] == "guard" + assert state_file.read_bytes() == before + code, status = _run_cli( + registry, "quota", "should-run", "--goal-id", GOAL_ID, + "--agent-id", AGENT_ID, "--runtime-profile", CONFIG.runtime_profile, + ) + assert code == 0 and status["quota"]["spent_slots"] == 0 + + def test_server_entrypoint_starts_as_a_standalone_script() -> None: """Kiro launches the registered script with the provisioned interpreter, not as an installed module, so it must import cleanly by path."""