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 a2d89d2dd4..884acf8543 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_priority", "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/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 bcaeb848a7..668864597d 100644 --- a/loopx/canary/module_metric_baseline.json +++ b/loopx/canary/module_metric_baseline.json @@ -59,7 +59,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/capabilities/manager_context/__init__.py b/loopx/capabilities/manager_context/__init__.py index 63626e0362..cf45ade48f 100644 --- a/loopx/capabilities/manager_context/__init__.py +++ b/loopx/capabilities/manager_context/__init__.py @@ -11,6 +11,7 @@ POLICY_SCHEMA as POLICY_SCHEMA, registered_context_recipients, source_context_authority, + source_context_target_authority, ) from ...control_plane.collaboration.goal_instance_scope import ( collaboration_goal_scope, @@ -86,6 +87,15 @@ def authority( grant = source_context_authority(runtime_root, registry_path, session, turn) return {**grant, "instruction": INSTRUCTION} if grant["mode"] == "context_only" else grant + +def target_authority( + runtime_root: Path, *, session: dict, turn: dict, target: dict +) -> dict: + """Authorize one target already validated by an exact Goal scope.""" + grant = source_context_target_authority(runtime_root, session, turn, target) + return {**grant, "instruction": INSTRUCTION} if grant["mode"] == "context_only" else grant + + def deliver( runtime_root: Path, registry_path: Path, *, session: dict, turn: dict, request: dict ) -> dict: @@ -101,7 +111,16 @@ def deliver( goal_scope, operation="request_create", ) - grant = authority(runtime_root, registry_path, session, turn) + grant = ( + target_authority( + runtime_root, + session=session, + turn=turn, + target=target, + ) + if goal_scope.exact + else authority(runtime_root, registry_path, session, turn) + ) if target not in grant["targets"]: raise ValueError("context recipient is not authorized or registered") content = str(turn.get("message") or "") diff --git a/loopx/capabilities/manager_context/roundtrip.py b/loopx/capabilities/manager_context/roundtrip.py index a801e34766..f4e56beb15 100644 --- a/loopx/capabilities/manager_context/roundtrip.py +++ b/loopx/capabilities/manager_context/roundtrip.py @@ -12,7 +12,7 @@ from datetime import datetime, timezone, timedelta from uuid import uuid4 -from . import _root, _read, _write, _hash, authority +from . import _root, _read, _write, _hash, authority, target_authority from .tracking import _entry, _now from ...file_lock import ( LockAcquisitionPolicy, @@ -380,7 +380,7 @@ def _exact_return_scope(registry, reply): return collaboration_goal_scope( registry, goal_id=reply["goal_id"], - agents=(), + agents=(reply["agent_id"],), caller_goal_ref=reply["goal_ref"], ) @@ -418,8 +418,13 @@ def _exact_return_context(root, registry, store, path, state_path, now): or not turn ): raise ValueError("original_conversation_unavailable") - grant = authority(root, registry, session, turn) target = {key: row[key] for key in ("goal_id", "agent_id")} + grant = target_authority( + root, + session=session, + turn=turn, + target=target, + ) if ( target not in grant["targets"] or grant.get("source_id") != row["source_id"] 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..cfc0b19d5d 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/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 bfb7da4475..238467f7af 100644 --- a/loopx/control_plane/collaboration/delegation_preview_transport.py +++ b/loopx/control_plane/collaboration/delegation_preview_transport.py @@ -8,6 +8,7 @@ import hashlib import json import os +import signal import subprocess import time import weakref @@ -19,6 +20,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 = [] @@ -47,22 +51,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: @@ -75,15 +118,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: @@ -108,7 +157,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)], @@ -144,7 +196,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) @@ -155,7 +207,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/collaboration/source_grant_observation.py b/loopx/control_plane/collaboration/source_grant_observation.py index fb88066738..a3f92fff48 100644 --- a/loopx/control_plane/collaboration/source_grant_observation.py +++ b/loopx/control_plane/collaboration/source_grant_observation.py @@ -33,27 +33,26 @@ def registered_context_recipients(registry: dict) -> dict: return {"active_goal_ids": active_goals, "available": available} -def source_context_authority( - runtime_root: Path, registry_path: Path, session: dict, turn: dict +def _source_context_grant( + runtime_root: Path, + session: dict, + turn: dict, + available_rows: list[dict], ) -> dict: - """Return only a write-only recipient catalog; no cross-audience Goal evidence.""" - if registry_path is None: - return {"mode": "unavailable", "targets": []} - try: - registry = load_project_registry(registry_path) - if not isinstance(registry, dict): - raise ValueError("invalid registry") - require_runtime_compatible_project_registry( - registry, operation="context source recipient observation" - ) - except (OSError, ValueError, TypeError): - return {"mode": "unavailable", "targets": []} - observed = registered_context_recipients(registry) - available = {(row["goal_id"], row["agent_id"]) for row in observed["available"]} + available = { + (row["goal_id"], row["agent_id"]) + for row in available_rows + if isinstance(row, dict) + and isinstance(row.get("goal_id"), str) + and isinstance(row.get("agent_id"), str) + } scope = conversation_scope(session, origin=turn.get("origin", "unknown")) if scope["private_conversation"] and turn.get("origin") == "web": - allowed = {target for target in available - if scope["goal_ids"] is None or target[0] in scope["goal_ids"]} + allowed = { + target + for target in available + if scope["goal_ids"] is None or target[0] in scope["goal_ids"] + } source_id = "web:" + _hash([session["session_id"], turn["client_turn_id"]]) else: if scope["kind"] != "external_audience": @@ -74,19 +73,68 @@ def source_context_authority( if policy.get("schema_version") != POLICY_SCHEMA: raise ValueError("invalid policy") grants = policy.get("sources", {}).get(ingress["channel"], {}) - selected = effect_runtime_result("collaboration.source.recipients", { - "source": grants, "sender_id": ingress["sender_id"], - "available": observed["available"], - }) - allowed = {(v["goal_id"], v["agent_id"]) for v in selected["targets"]} + selected = effect_runtime_result( + "collaboration.source.recipients", + { + "source": grants, + "sender_id": ingress["sender_id"], + "available": available_rows, + }, + ) + allowed = { + (value["goal_id"], value["agent_id"]) + for value in selected["targets"] + } source_id = ingress["source_id"] - except (OSError, ValueError, KeyError, TypeError, AttributeError, EffectRuntimeRejected): + except ( + OSError, + ValueError, + KeyError, + TypeError, + AttributeError, + EffectRuntimeRejected, + ): return {"mode": "unavailable", "targets": []} targets = [ - {"goal_id": g, "agent_id": a} for g, a in sorted(allowed & available) + {"goal_id": goal_id, "agent_id": agent_id} + for goal_id, agent_id in sorted(allowed & available) ] return { "mode": "context_only", "targets": targets, "source_id": source_id, } + + +def source_context_target_authority( + runtime_root: Path, + session: dict, + turn: dict, + target: dict, +) -> dict: + """Authorize one target whose exact Goal scope was already validated.""" + return _source_context_grant(runtime_root, session, turn, [target]) + + +def source_context_authority( + runtime_root: Path, registry_path: Path, session: dict, turn: dict +) -> dict: + """Return only a write-only recipient catalog; no cross-audience Goal evidence.""" + if registry_path is None: + return {"mode": "unavailable", "targets": []} + try: + registry = load_project_registry(registry_path) + if not isinstance(registry, dict): + raise ValueError("invalid registry") + require_runtime_compatible_project_registry( + registry, operation="context source recipient observation" + ) + except (OSError, ValueError, TypeError): + return {"mode": "unavailable", "targets": []} + observed = registered_context_recipients(registry) + return _source_context_grant( + runtime_root, + session, + turn, + observed["available"], + ) 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/control_plane/goals/acceptance_contract.ts b/loopx/control_plane/goals/acceptance_contract.ts index c4e8adc0ae..7d96b82465 100644 --- a/loopx/control_plane/goals/acceptance_contract.ts +++ b/loopx/control_plane/goals/acceptance_contract.ts @@ -187,7 +187,8 @@ export function normalizeGoalAcceptanceDocument(value: unknown): AcceptanceDocum const NON_WORK_FIELDS = new Set([ "schema_version", "source_section", "index", "title", "priority", "status", "done", "archive_state", "claimed_by", "created_by", "last_actor_agent_id", "updated_at", "completed_at", "completion_turn_key", - "completion_validation_sha256", "completion_recovery", "completion_continuation", "no_followup", "decision_outcome", + "completion_validation_sha256", "completion_recovery", "completion_continuation", "completion_receipt_id", + "no_followup", "decision_outcome", "completion_result", "decision_scope_outcomes", "note", "evidence", "reason", "handoff_note", "resume_ready", "resume_monitor_generation", "last_checked_at", "result_hash", "consecutive_no_change", diff --git a/loopx/control_plane/quota/settlement_precedence.py b/loopx/control_plane/quota/settlement_precedence.py index a486952222..1061297080 100644 --- a/loopx/control_plane/quota/settlement_precedence.py +++ b/loopx/control_plane/quota/settlement_precedence.py @@ -1,5 +1,6 @@ from __future__ import annotations from .effective_action import EffectiveAction +from .selected_todo_projection import selected_todo_projection from ..work_items.work_lane import work_lane_contract_is_receipt_bound_monitor_settled from typing import Any @@ -129,10 +130,24 @@ def apply_settled_monitor_precedence(payload: dict[str, Any]) -> None: recorded_action = payload.get("agent_lane_next_action") clear_quota_action_projections(payload) payload.update(settled_replay_fields()) - if ( - isinstance(recorded_action, dict) - and recorded_action.get("selection_binding") == "heartbeat_receipt" - and isinstance(lane, dict) - and recorded_action.get("todo_id") == lane.get("selected_todo_id") - ): - payload["agent_lane_next_action"] = recorded_action + bound_action = next( + ( + candidate + for candidate in ( + recorded_action, + lane.get("receipt_bound_monitor_item"), + ) + if isinstance(candidate, dict) + and candidate.get("selection_binding") == "heartbeat_receipt" + and candidate.get("todo_id") == lane.get("selected_todo_id") + ), + None, + ) + if bound_action is not None: + payload["agent_lane_next_action"] = bound_action + selected_todo = selected_todo_projection( + agent_lane_next_action=bound_action, + work_lane_contract=lane, + ) + if selected_todo is not None: + payload["selected_todo"] = selected_todo 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 91dec88735..3a64da3bb5 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, @@ -231,7 +247,7 @@ }, { "site": "loopx/capabilities/manager_context/__init__.py::.configure_delivery_target.update::codec_read:load_project_registry#1", - "line": 358, + "line": 377, "column": 51, "kind": "codec_read", "api": "load_project_registry", @@ -239,7 +255,7 @@ }, { "site": "loopx/capabilities/manager_context/__init__.py::.configure_evidence_scope::codec_read:load_project_registry#1", - "line": 313, + "line": 332, "column": 16, "kind": "codec_read", "api": "load_project_registry", @@ -263,7 +279,7 @@ }, { "site": "loopx/capabilities/manager_context/roundtrip.py::.drain::codec_read:load_project_registry#1", - "line": 879, + "line": 884, "column": 22, "kind": "codec_read", "api": "load_project_registry", @@ -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": 518, @@ -1039,7 +1039,7 @@ }, { "site": "loopx/control_plane/collaboration/source_grant_observation.py::.source_context_authority::codec_read:load_project_registry#1", - "line": 43, + "line": 126, "column": 20, "kind": "codec_read", "api": "load_project_registry", 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/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_acceptance_stale_replan.py b/tests/control_plane/test_goal_acceptance_stale_replan.py index 09e8150445..f45af1f5fc 100644 --- a/tests/control_plane/test_goal_acceptance_stale_replan.py +++ b/tests/control_plane/test_goal_acceptance_stale_replan.py @@ -81,7 +81,8 @@ def test_stale_binding_cannot_replan_another_agents_work_or_disabled_contract(): def test_stale_binding_reaches_quota_frontier_projection(): context = build_goal_frontier_projection_context_from_status( - goal_id="goal-a", agent_id="agent-a", status_payload={}, item={}, + goal_id="goal-a", agent_id="agent-a", + status_payload={"run_history": {"runs": []}}, item={}, project_asset=None, user_todo_summary={"open_count": 0}, agent_todo_summary=_summary(), agent_todo_source_items=_source(), work_lane_contract={"lane": "advancement_task", "must_attempt_work": True}, 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_monitor_followthrough_contract.py b/tests/control_plane/test_monitor_followthrough_contract.py index ddb3a98f6c..4170f37213 100644 --- a/tests/control_plane/test_monitor_followthrough_contract.py +++ b/tests/control_plane/test_monitor_followthrough_contract.py @@ -634,7 +634,8 @@ def test_same_turn_material_monitor_poll_is_no_spend_closeout_before_successor( runtime_root=runtime, ) assert replay["heartbeat_receipt"]["settlement_identity"]["todo_id"] == admitted["todo_id"] - assert replay.get("selected_todo") is None + assert replay["selected_todo"]["todo_id"] == admitted["todo_id"] + assert replay["agent_lane_next_action"]["receipt_bound_monitor_phase"] == "settled" assert replay["should_run"] is False assert replay["effective_action"] == "heartbeat_settled_skip" assert replay["execution_obligation"]["must_attempt_work"] is False @@ -800,7 +801,8 @@ def test_same_turn_unchanged_monitor_poll_is_already_settled(tmp_path: Path) -> runtime_root=runtime, ) assert replay["heartbeat_receipt"]["settlement_identity"]["todo_id"] == monitor["todo_id"] - assert replay.get("selected_todo") is None + assert replay["selected_todo"]["todo_id"] == monitor["todo_id"] + assert replay["agent_lane_next_action"]["receipt_bound_monitor_phase"] == "settled" assert replay["effective_action"] == "heartbeat_settled_skip" assert replay["execution_obligation"]["must_attempt_work"] is False assert replay["heartbeat_recommendation"]["agent_must_attempt"] is False 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 5034c35904..5a90508a2c 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,15 +611,15 @@ 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, ) assert rc == 1 and deferred["action_selection_qualification"]["state"] == "deferred", deferred [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_refresh_checkpoint_recovery.py b/tests/control_plane/test_refresh_checkpoint_recovery.py index 6ea9324edc..d3e9dfbfc9 100644 --- a/tests/control_plane/test_refresh_checkpoint_recovery.py +++ b/tests/control_plane/test_refresh_checkpoint_recovery.py @@ -18,6 +18,7 @@ SELECTED_REPLAN_TODO_ID, TODO_ID, TURN_ID, + _bind_selected_replan_guard, _configure_selected_todo_replan_fixture, _initialize_git_checkout, _run_cli, @@ -352,6 +353,13 @@ def test_checkpoint_only_recovery_bypasses_open_todo_completion_validation( assert rc == 0, guard assert guard["decision"] == "autonomous_replan_required" assert guard["selected_todo"]["todo_id"] == SELECTED_REPLAN_TODO_ID + assert guard["heartbeat_receipt"]["settlement_binding_owed"] is True + _bind_selected_replan_guard( + registry, + runtime, + project, + turn_id, + ) delivery = ( "refresh-state", diff --git a/tests/control_plane/test_replan_successor_durable_ack.py b/tests/control_plane/test_replan_successor_durable_ack.py index 75eed7adf9..476530aa3f 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,6 +204,16 @@ 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 + 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_runtime_source_read_batching.py b/tests/control_plane/test_runtime_source_read_batching.py index 1cd1f167c6..f593771637 100644 --- a/tests/control_plane/test_runtime_source_read_batching.py +++ b/tests/control_plane/test_runtime_source_read_batching.py @@ -1,6 +1,7 @@ from __future__ import annotations import hashlib +import os from pathlib import Path import pytest @@ -9,8 +10,11 @@ def serial_fingerprint(root: Path) -> str: - """Independent reference: names and raw bytes, not decoded source text.""" + """Independent reference: release identity, names, and raw source bytes.""" digest = hashlib.sha256() + digest.update(b"loopx_effect_runtime_source_instance_v1\0") + digest.update(os.fsencode(os.path.normcase(os.fspath(root.resolve())))) + digest.update(b"\0") for path in sorted(p for p in root.rglob("*") if p.suffix in {".ts", ".json"}): digest.update(path.relative_to(root).as_posix().encode("utf-8")) digest.update(path.read_bytes()) diff --git a/tests/control_plane/test_settled_monitor_user_gate.py b/tests/control_plane/test_settled_monitor_user_gate.py index 12182da3ca..1db31a85ca 100644 --- a/tests/control_plane/test_settled_monitor_user_gate.py +++ b/tests/control_plane/test_settled_monitor_user_gate.py @@ -83,7 +83,8 @@ def call(*args): claimed_by=AGENT_ID, agent_id=AGENT_ID) replay = call(*guard_args) assert replay["heartbeat_receipt"]["settlement_identity"] == admitted["heartbeat_receipt"]["settlement_identity"] - assert replay.get("selected_todo") is None + assert replay["selected_todo"]["todo_id"] == monitor["todo_id"] + assert replay["agent_lane_next_action"]["receipt_bound_monitor_phase"] == "settled" assert replay["should_run"] is False assert replay["safe_bypass_allowed"] is False assert replay["execution_obligation"]["must_attempt_work"] is False diff --git a/tests/control_plane_ts/goal_acceptance_authority.test.ts b/tests/control_plane_ts/goal_acceptance_authority.test.ts index 4ca3fafed6..ba0e81167e 100644 --- a/tests/control_plane_ts/goal_acceptance_authority.test.ts +++ b/tests/control_plane_ts/goal_acceptance_authority.test.ts @@ -52,7 +52,8 @@ function originalHead() { test("terminal continuation observations preserve work while changed requirements invalidate it", () => { const work = todo("todo_first"); const completed = {...work, status: "done", done: true, no_followup: true, - completion_continuation: "no_followup", note: "Bounded task completed"}; + completion_continuation: "no_followup", completion_receipt_id: `tcw_${"a".repeat(64)}`, + note: "Bounded task completed"}; assert.equal(goalAcceptanceTodoDigest(completed), goalAcceptanceTodoDigest(work)); assert.notEqual(goalAcceptanceTodoDigest({...completed, text: "Deliver different work"}), goalAcceptanceTodoDigest(work)); assert.notEqual(goalAcceptanceTodoDigest({...completed, completion_validation_required: true}), goalAcceptanceTodoDigest(work)); diff --git a/tests/test_collaboration_goal_instance.py b/tests/test_collaboration_goal_instance.py index 55b25232c5..4a89d07c82 100644 --- a/tests/test_collaboration_goal_instance.py +++ b/tests/test_collaboration_goal_instance.py @@ -1186,7 +1186,9 @@ def test_failed_external_turn_returns_only_through_its_current_sender_grant(tmp_ if revoke: policy_path = _root(tmp_path) / "policy.json" policy = json.loads(policy_path.read_text()) - policy["sources"][session["channel_id"]]["targets"] = [] + source = policy["sources"][session["channel_id"]] + source["local_delivery_scope"] = "selected" + source["targets"] = [] _write(policy_path, policy) calls = [] 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 d1a3c44968..d0a57f5e13 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_loopx_turn_driver.py b/tests/test_loopx_turn_driver.py index 92afb64b26..4fb36d8ace 100644 --- a/tests/test_loopx_turn_driver.py +++ b/tests/test_loopx_turn_driver.py @@ -1984,8 +1984,7 @@ def interrupt_checkpoint(path, journal): ] ) assert policy_code == 0, policy_output.getvalue() - host_project = tmp_path / "isolated-host-workspace" - host_project.mkdir() + host_project = project host_script = """ import json import pathlib @@ -2287,8 +2286,7 @@ def test_turn_run_once_cli_completes_selected_todo_after_validation( next_action: str, ) -> None: project, runtime, registry = _write_live_fixture(tmp_path) - host_project = tmp_path / "isolated-host-workspace" - host_project.mkdir() + host_project = project host_script = """ import json import pathlib @@ -2468,8 +2466,7 @@ def test_promoted_turn_completion_replays_after_commit_before_journal_crash( ) -> None: project, runtime, registry = _write_live_fixture(tmp_path) _promote_turn_fixture(project, runtime) - host_project = tmp_path / "isolated-host-workspace" - host_project.mkdir() + host_project = project host_script, validation_script = _completion_host_and_validation_scripts() argv = _turn_run_once_completion_argv( host_project, @@ -2632,8 +2629,7 @@ def test_turn_run_once_cli_repairs_committed_quota_spend_after_receipt_crash( monkeypatch: pytest.MonkeyPatch, ) -> None: project, runtime, registry = _write_live_fixture(tmp_path) - host_project = tmp_path / "isolated-host-workspace" - host_project.mkdir() + host_project = project host_script, validation_script = _completion_host_and_validation_scripts() argv = _turn_run_once_completion_argv( host_project, @@ -2739,8 +2735,7 @@ def test_turn_run_once_cli_projects_declared_successor_continuation( "task_class=advancement_task priority=P2 -->", ), ) - host_project = tmp_path / "isolated-host-workspace" - host_project.mkdir() + host_project = project host_script, validation_script = _completion_host_and_validation_scripts() output = io.StringIO() with contextlib.redirect_stdout(output): @@ -2771,8 +2766,7 @@ def test_turn_run_once_cli_projects_durable_no_followup_continuation( tmp_path, todo_metadata_extra="no_followup=true", ) - host_project = tmp_path / "isolated-host-workspace" - host_project.mkdir() + host_project = project host_script, validation_script = _completion_host_and_validation_scripts() output = io.StringIO() with contextlib.redirect_stdout(output): @@ -2808,8 +2802,7 @@ def test_turn_run_once_cli_terminal_recovery_rejects_unowned_completion( tmp_path, todo_metadata_extra="no_followup=true", ) - host_project = tmp_path / "isolated-host-workspace" - host_project.mkdir() + host_project = project host_script, validation_script = _completion_host_and_validation_scripts() argv = _turn_run_once_completion_argv( host_project, @@ -2973,8 +2966,7 @@ def test_turn_run_once_cli_replays_declared_successor_after_interruption( "task_class=advancement_task priority=P2 -->", ), ) - host_project = tmp_path / "isolated-host-workspace" - host_project.mkdir() + host_project = project host_script, validation_script = _completion_host_and_validation_scripts() argv = _turn_run_once_completion_argv( host_project, diff --git a/tests/test_loopx_turn_executor.py b/tests/test_loopx_turn_executor.py index 4fd2bae845..3020f85520 100644 --- a/tests/test_loopx_turn_executor.py +++ b/tests/test_loopx_turn_executor.py @@ -1729,8 +1729,8 @@ def host(_request: dict[str, object]) -> dict[str, object]: "session_binding_resolver": lambda _turn_envelope: { "schema_version": "loopx_turn_session_binding_v0", "goal_id": "fixture-goal", - "agent_id": "codex-fixture", - "todo_id": "todo_from_another_turn", + "agent_id": "codex-from-another-agent", + "todo_id": "todo_fixture0001", }, "project": tmp_path, "runtime_root": tmp_path / "runtime",