From 726bc18e51f454c7a17528efdda075803698fbb9 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Thu, 24 Sep 2026 07:47:12 +0800 Subject: [PATCH 1/8] fix(control-plane): preserve ambiguous Effect responses and expose exact receipts Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- loopx/cli_commands/todo.py | 44 +++++++++- .../cli_commands/todo_argument_validation.py | 12 +++ loopx/cli_commands/todo_event.py | 3 +- loopx/cli_commands/todo_registration.py | 5 ++ .../coordination/local_authority_read.ts | 33 ++++++++ .../coordination/local_authority_write.ts | 7 +- loopx/control_plane/effect_runtime.py | 66 ++++++++++----- .../control_plane/effect_runtime_handlers.ts | 4 +- loopx/control_plane/effect_runtime_server.ts | 43 ++++++---- .../scheduler/provider_monitor_poll.py | 5 +- .../todos/provider_terminal_lifecycle.py | 15 +++- .../test_effect_runtime_integration.py | 81 +++++++++++++++++++ .../test_local_coordination_authority.py | 12 ++- .../test_todo_operation_receipt_cli.py | 47 +++++++++++ .../local_authority_operation_receipt.test.ts | 32 ++++++++ .../local_authority_provider.test.ts | 6 ++ .../shadow_native_writer_boundary.test.ts | 31 +++++++ 17 files changed, 399 insertions(+), 47 deletions(-) create mode 100644 tests/control_plane/test_todo_operation_receipt_cli.py create mode 100644 tests/control_plane_ts/local_authority_operation_receipt.test.ts diff --git a/loopx/cli_commands/todo.py b/loopx/cli_commands/todo.py index 8acac619da..08cf6d94b4 100644 --- a/loopx/cli_commands/todo.py +++ b/loopx/cli_commands/todo.py @@ -4,7 +4,11 @@ from collections.abc import Callable, Sequence from pathlib import Path -from ..control_plane.coordination.local_authority import read_canonical_todo_fields_if_promoted +from ..control_plane.coordination.local_authority import ( + local_authority_is_promoted, + read_canonical_todo_fields_if_promoted, +) +from ..control_plane.effect_runtime import effect_runtime_result from ..control_plane.agents.workspace_guard import capture_delivery_workspace from ..control_plane.todos.contract import ( replan_successor_semantic_binding, @@ -47,6 +51,7 @@ validate_todo_claim_options, validate_todo_complete_options, validate_todo_list_options, + validate_todo_receipt_options, validate_todo_project_markdown_options, validate_todo_plan_options, validate_todo_supersede_options, @@ -182,6 +187,23 @@ def _todo_path_args(args: argparse.Namespace) -> dict[str, Path | None]: } +def _render_todo_receipt(payload: dict[str, object]) -> str: + lines = [ + "# LoopX Canonical Operation Receipt", + "", + f"- status: `{payload.get('status')}`", + f"- goal_id: `{payload.get('goal_id')}`", + f"- operation_id: `{payload.get('operation_id')}`", + f"- source_authority: `{payload.get('source_authority')}`", + f"- provider_revision: `{payload.get('provider_revision')}`", + f"- cursor: `{payload.get('cursor')}`", + "- note: Historical readback only; it does not grant a current lease or a retry.", + ] + if payload.get("error") or payload.get("reason"): + lines.append(f"- error: `{payload.get('error') or payload.get('reason')}`") + return "\n".join(lines) + + def handle_todo_command( args: argparse.Namespace, *, @@ -194,8 +216,8 @@ def handle_todo_command( post_writeback_projection_builder: PostWritebackProjectionBuilder | None = None, ) -> int: renderer = ( - render_task_planning_packet - if args.todo_command == "plan" + render_task_planning_packet if args.todo_command == "plan" + else _render_todo_receipt if args.todo_command == "receipt" else render_todo_markdown ) try: @@ -228,6 +250,22 @@ def handle_todo_command( **_todo_path_args(args), runtime_root_arg=runtime_root_arg, ) + elif args.todo_command == "receipt": + validate_todo_receipt_options(args) + runtime_root = resolve_runtime_root(load_registry(registry_path), runtime_root_arg) + if not local_authority_is_promoted(runtime_root=runtime_root, goal_id=args.goal_id): + raise ValueError("todo receipt requires promoted canonical authority; no legacy fallback") + result = effect_runtime_result( + "coordination.local_authority.operation_receipt", + {"schema_version": "loopx_local_coordination_operation_receipt_request_v0", + "runtime_root": str(runtime_root.expanduser().resolve(strict=False)), + "goal_id": args.goal_id, "operation_id": args.operation_id}, + timeout=15.0, + ) + if not isinstance(result, dict): + raise RuntimeError("canonical operation receipt returned an invalid result") + payload = {"ok": result.get("status") in {"found", "missing"}, + "command": "receipt", **result} elif args.todo_command == "project-markdown": validate_todo_project_markdown_options(args) registry = load_registry(registry_path) diff --git a/loopx/cli_commands/todo_argument_validation.py b/loopx/cli_commands/todo_argument_validation.py index e88721ba39..76b836ce5d 100644 --- a/loopx/cli_commands/todo_argument_validation.py +++ b/loopx/cli_commands/todo_argument_validation.py @@ -11,6 +11,7 @@ ("--priority", "priority"), ("--clear-priority", "clear_priority"), ("--todo-id", "todo_id"), + ("--operation-id", "operation_id"), ("--claim-operation-id", "claim_operation_id"), ("--update-operation-id", "update_operation_id"), ("--update-expected-provider-revision", "update_expected_provider_revision"), @@ -300,6 +301,15 @@ def validate_todo_list_options(args: argparse.Namespace) -> None: ) +def validate_todo_receipt_options(args: argparse.Namespace) -> None: + _validate_todo_option_subset( + args, {"operation_id"}, + "todo receipt only accepts --goal-id, --operation-id, and --format; unsupported: ", + ) + if not args.operation_id: + raise ValueError("todo receipt requires --operation-id") + + def validate_todo_plan_options(args: argparse.Namespace) -> None: _validate_todo_option_subset( args, {"text", "agent_id"}, @@ -493,6 +503,8 @@ def validate_todo_archive_completed_options(args: argparse.Namespace) -> None: def validate_shared_todo_options(args: argparse.Namespace) -> None: + if getattr(args, "operation_id", None) and args.todo_command != "receipt": + raise ValueError("--operation-id is supported only by todo receipt") agent_id_allowed_for_user_authoring = ( args.todo_command == "add" and args.role == "user" diff --git a/loopx/cli_commands/todo_event.py b/loopx/cli_commands/todo_event.py index 392a6ad397..ea5edcc170 100644 --- a/loopx/cli_commands/todo_event.py +++ b/loopx/cli_commands/todo_event.py @@ -81,7 +81,8 @@ def append_todo_rollout_event( and getattr(args, "no_follow_up", False) ) if ( - not payload.get("ok") + args.todo_command == "receipt" + or not payload.get("ok") or payload.get("dry_run") or (payload.get("idempotent_replay") and not turn_instance_id) ): diff --git a/loopx/cli_commands/todo_registration.py b/loopx/cli_commands/todo_registration.py index 4c95047f78..237f284e4a 100644 --- a/loopx/cli_commands/todo_registration.py +++ b/loopx/cli_commands/todo_registration.py @@ -31,6 +31,7 @@ def register_todo_command( choices=[ "add", "list", + "receipt", "claim", "update", "complete", @@ -53,6 +54,10 @@ def register_todo_command( todo_parser.add_argument("--priority", choices=["P0", "P1", "P2", "P3", "P4"], help="For add/update, declare Todo priority independently of text; omission retains the current value.") todo_parser.add_argument("--clear-priority", action="store_true", help="For update, explicitly remove priority; cannot be combined with --priority.") todo_parser.add_argument("--todo-id", help="Structured todo id from status/quota, such as todo_ab12cd34ef56.") + todo_parser.add_argument( + "--operation-id", + help="For todo receipt, read the exact historical canonical operation after an ambiguous response; this does not grant a retry or lease.", + ) todo_parser.add_argument( "--update-operation-id", help=("For promoted text/note, planning, validator revision or User completion update, reuse this operation id after a lost response; " diff --git a/loopx/control_plane/coordination/local_authority_read.ts b/loopx/control_plane/coordination/local_authority_read.ts index 78c9f9e9ae..71d99e847b 100644 --- a/loopx/control_plane/coordination/local_authority_read.ts +++ b/loopx/control_plane/coordination/local_authority_read.ts @@ -12,6 +12,39 @@ import {LOCAL_COORDINATION_TODO_LIST_REQUEST_SCHEMA, LOCAL_COORDINATION_TODO_LIS LOCAL_COORDINATION_TODO_READ_REQUEST_SCHEMA, LOCAL_COORDINATION_TODO_READ_RESULT_SCHEMA} from "./coordination_state_contract.generated.ts"; import {decodeProjectionReadback, confirmProjectionReadback} from "../todos/projection_delivery.ts"; +const LOCAL_COORDINATION_OPERATION_RECEIPT_REQUEST_SCHEMA = "loopx_local_coordination_operation_receipt_request_v0"; +const LOCAL_COORDINATION_OPERATION_RECEIPT_RESULT_SCHEMA = "loopx_local_coordination_operation_receipt_result_v0"; + +/** Exact historical operation readback. It grants no current lease or retry authority. */ +export async function readLocalCoordinationOperationReceipt( + value: unknown, + dependencies: LocalAuthorityProviderDependencies = {}, +): Promise { + let sourceAuthority = "canonical_unavailable"; + try { + const input = requireJsonObject(value, "local coordination operation receipt request"); + if (input.schema_version !== LOCAL_COORDINATION_OPERATION_RECEIPT_REQUEST_SCHEMA) { + throw new Error("local coordination operation receipt request schema mismatch"); + } + const root = runtimeRoot(input.runtime_root); + const goalId = requireAuthorityStoreId(input.goal_id, "goal id"); + const operationId = requireAuthorityStoreId(input.operation_id, "operation id"); + const store = await openRuntimeStore(root, goalId, dependencies); + sourceAuthority = sourceAuthorityFor(store); + const receipt = await store.readReceipt(operationId); + return {schema_version: LOCAL_COORDINATION_OPERATION_RECEIPT_RESULT_SCHEMA, + goal_id: goalId, operation_id: operationId, ...receipt, + source_authority: sourceAuthority, decision_read_from_provider: true, + legacy_fallback_used: false}; + } catch (error) { + return {schema_version: LOCAL_COORDINATION_OPERATION_RECEIPT_RESULT_SCHEMA, + status: "failed", reason_code: "invalid_local_coordination_operation_receipt_request", + reason: error instanceof Error ? error.message : "invalid operation receipt request", + source_authority: sourceAuthority, decision_read_from_provider: false, + legacy_fallback_used: false, ...localAuthorityOpenFailure(error)}; + } +} + /** Provider-first exact Todo read. Missing/unavailable state never falls back. */ export async function readLocalCoordinationTodo( value: unknown, diff --git a/loopx/control_plane/coordination/local_authority_write.ts b/loopx/control_plane/coordination/local_authority_write.ts index 24826035d8..864da1ba9b 100644 --- a/loopx/control_plane/coordination/local_authority_write.ts +++ b/loopx/control_plane/coordination/local_authority_write.ts @@ -1,11 +1,16 @@ import {withFileMutationLock} from "../effect_runtime_io.ts"; import {requireShadowPrimaryWriteAllowed, shadowMaintenanceLockPath} from "./shadow_management.ts"; +// File-v0 still verifies and republishes its retained history on each write. +// Give a competing canonical writer a bounded wait for that real critical +// section; the provider CAS and hard lease remain the decision authority. +const CANONICAL_WRITER_LOCK_TIMEOUT_MS = 30_000; + /** Canonical command writers share the maintenance guard; provider CAS owns state. */ export async function withCanonicalWriter(root: string, goalId: string, dryRun: boolean, write: () => Promise): Promise { if (dryRun) return await write(); return await withFileMutationLock(shadowMaintenanceLockPath(root, goalId), async () => { await requireShadowPrimaryWriteAllowed(root, goalId); return await write(); - }); + }, CANONICAL_WRITER_LOCK_TIMEOUT_MS); } diff --git a/loopx/control_plane/effect_runtime.py b/loopx/control_plane/effect_runtime.py index fd8783e928..2e7a8550d1 100644 --- a/loopx/control_plane/effect_runtime.py +++ b/loopx/control_plane/effect_runtime.py @@ -197,6 +197,18 @@ def __init__(self, message: str, *, diagnostic_code: str) -> None: self.diagnostic_code = diagnostic_code +class EffectRuntimeResponseAmbiguous(EffectRuntimeStartupError): + """The request may have executed even though its response was lost.""" + + def __init__(self, method: str, *, timeout: float) -> None: + super().__init__( + f"TypeScript Effect runtime returned no verifiable response for {method} " + f"(request budget {timeout:g}s); the operation may have committed. Read its exact " + "durable receipt before any retry", + diagnostic_code="runtime_response_ambiguous", + ) + + def _control_plane_root() -> Path: return Path(__file__).resolve().parent @@ -530,30 +542,33 @@ def _request_with_info( with socket.create_connection( (str(info["host"]), int(info["port"])), timeout=timeout ) as connection: - connection.settimeout(timeout) - connection.sendall(encoded) - while True: - chunk = connection.recv(64 * 1024) - if not chunk: - break - chunks.append(chunk) - size += len(chunk) - if size > MAX_RESPONSE_BYTES: - raise RuntimeError("TypeScript Effect runtime response is oversized") - if b"\n" in chunk: - break + try: + connection.settimeout(timeout) + # sendall may have delivered a prefix before it raises. From this + # point onward the caller cannot prove that no effect ran. + connection.sendall(encoded) + while True: + chunk = connection.recv(64 * 1024) + if not chunk: + break + chunks.append(chunk) + size += len(chunk) + if size > MAX_RESPONSE_BYTES: + raise RuntimeError("TypeScript Effect runtime response is oversized") + if b"\n" in chunk: + break + except (OSError, RuntimeError) as exc: + raise EffectRuntimeResponseAmbiguous(method, timeout=timeout) from exc try: response = json.loads(b"".join(chunks).split(b"\n", 1)[0]) except (json.JSONDecodeError, IndexError): - raise RuntimeError( - "TypeScript Effect runtime returned malformed JSON" - ) from None + raise EffectRuntimeResponseAmbiguous(method, timeout=timeout) from None if ( not isinstance(response, dict) or response.get("schema_version") != EFFECT_RUNTIME_RESPONSE_SCHEMA_VERSION or response.get("request_id") != request_id ): - raise RuntimeError("TypeScript Effect runtime response shape mismatch") + raise EffectRuntimeResponseAmbiguous(method, timeout=timeout) if response.get("ok") is not True: raise _remote_runtime_error(response.get("error")) return response @@ -773,6 +788,7 @@ def effect_runtime_request( request_id = str(uuid.uuid4()) last_error: OSError | RuntimeError | None = None for attempt in range(2 if retry_safe else 1): + info: dict[str, Any] | None = None try: info = _read_info(info_path, fingerprint=fingerprint) if info is None: @@ -784,18 +800,30 @@ def effect_runtime_request( params=params, timeout=timeout, ) - except EffectRuntimeRemoteError: + except (EffectRuntimeRemoteError, EffectRuntimeResponseAmbiguous): raise except EffectRuntimeStartupError as exc: last_error = exc if attempt == 0 and retry_safe: - info_path.unlink(missing_ok=True) continue raise + except TimeoutError as exc: + # A connect timeout is not evidence that an existing runtime died. + # In particular it must not replace a live server which may still + # be completing an earlier mutation under the per-Goal lock. + raise EffectRuntimeStartupError( + f"TypeScript Effect runtime did not connect for {method} " + f"within {timeout:g}s", + diagnostic_code="runtime_request_timeout", + ) from exc except (OSError, RuntimeError) as exc: last_error = exc if attempt == 0 and retry_safe: - info_path.unlink(missing_ok=True) + # Only pre-send connection failures reach this branch. Do not + # remove a replacement runtime published by another caller. + current = _read_info(info_path, fingerprint=fingerprint) + if info is not None and current is not None and current.get("token") == info.get("token"): + info_path.unlink(missing_ok=True) continue break if isinstance(last_error, TimeoutError): diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index 558bcdb980..262b8177ad 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -173,7 +173,8 @@ import { executeReviewedCoordinationPromotion, terminalLifecycleLocalCoordinationTodo, } from "./coordination/local_authority_runtime.ts"; -import {listLocalCoordinationTodos, readLocalCoordinationTodo} from "./coordination/local_authority_read.ts"; +import {listLocalCoordinationTodos, readLocalCoordinationTodo, + readLocalCoordinationOperationReceipt} from "./coordination/local_authority_read.ts"; import { evaluateCoordinationTodoClaimDecision } from "./coordination/todo_claim.ts"; import { evaluateCoordinationTodoTerminalDecision, @@ -566,6 +567,7 @@ export function createEffectRuntimeHandlers( ["coordination.local_authority.todo_archive", archiveLocalCoordinationTodos], ["coordination.local_authority.todo_archive_ack", acknowledgeLocalCoordinationTodoArchive], ["coordination.local_authority.todo_read", readLocalCoordinationTodo], + ["coordination.local_authority.operation_receipt", readLocalCoordinationOperationReceipt], ["coordination.ownership_observation", projectOwnershipObservation], ["coordination.local_authority.ownership_observation", observeLocalCoordinationOwnership], ["coordination.local_authority.todo_list", listLocalCoordinationTodos], diff --git a/loopx/control_plane/effect_runtime_server.ts b/loopx/control_plane/effect_runtime_server.ts index 73c718fb71..87bf6b0ad9 100644 --- a/loopx/control_plane/effect_runtime_server.ts +++ b/loopx/control_plane/effect_runtime_server.ts @@ -1,6 +1,6 @@ import { writeSync } from "node:fs"; import { createServer, type Socket } from "node:net"; -import { chmod, rm } from "node:fs/promises"; +import { chmod, readFile, rm } from "node:fs/promises"; import type { JsonObject } from "./effect_program.ts"; import { @@ -11,7 +11,7 @@ import { EffectRuntimeRequestError, effectRuntimeErrorPayload, } from "./effect_runtime_errors.ts"; -import { atomicWriteJson } from "./effect_runtime_io.ts"; +import { atomicWriteJson, withFileMutationLock } from "./effect_runtime_io.ts"; import { sqliteRuntimeIdentity } from "./coordination/sqlite_runtime.ts"; import { requireJsonObject as requiredObject, @@ -184,23 +184,38 @@ const server = createServer((socket) => { }); server.on("close", () => { - void rm(infoPath, { force: true }).finally(() => process.exit(0)); + void withFileMutationLock(infoPath, async () => { + let published: Record; + try { + published = JSON.parse(await readFile(infoPath, "utf8")); + } catch { + return; + } + // A timed-out client may have published a replacement server. The old + // server must never erase that server's locator when it finally exits. + if (published.token === token && published.pid === process.pid && + published.fingerprint === fingerprint) { + await rm(infoPath, { force: true }); + } + }).finally(() => process.exit(0)); }); server.listen(0, "127.0.0.1", async () => { const address = server.address(); if (!address || typeof address === "string") throw new Error("invalid address"); - await atomicWriteJson(infoPath, { - schema_version: INFO_SCHEMA, - fingerprint, - pid: process.pid, - host: "127.0.0.1", - port: address.port, - token, - // A managed runtime is reused per source revision, so the Node/SQLite pair - // serving a goal is not necessarily the one the caller resolves from PATH. - runtime_identity: sqliteRuntimeIdentity(), + await withFileMutationLock(infoPath, async () => { + await atomicWriteJson(infoPath, { + schema_version: INFO_SCHEMA, + fingerprint, + pid: process.pid, + host: "127.0.0.1", + port: address.port, + token, + // A managed runtime is reused per source revision, so the Node/SQLite pair + // serving a goal is not necessarily the one the caller resolves from PATH. + runtime_identity: sqliteRuntimeIdentity(), + }); + await chmod(infoPath, 0o600); }); - await chmod(infoPath, 0o600); resetIdleTimer(server); }); diff --git a/loopx/control_plane/scheduler/provider_monitor_poll.py b/loopx/control_plane/scheduler/provider_monitor_poll.py index 6dd5374648..68b159bd87 100644 --- a/loopx/control_plane/scheduler/provider_monitor_poll.py +++ b/loopx/control_plane/scheduler/provider_monitor_poll.py @@ -17,6 +17,9 @@ from ..todos.provider_projection import settle_canonical_todo_projection +_MONITOR_POLL_RUNTIME_TIMEOUT_SECONDS = 45.0 + + def require_monitor_poll_source_available(*, runtime_root: Path, goal_id: str) -> None: """Fail closed on unavailable promoted authority before unrelated quota work.""" read_canonical_todos_if_promoted(runtime_root=runtime_root, goal_id=goal_id) @@ -41,7 +44,7 @@ def poll_canonical_monitor_if_promoted( "registered_agents": registered, "registry_source": registry_source, "dry_run": not execute, "observation": observation, "intent": intent, - }) + }, timeout=_MONITOR_POLL_RUNTIME_TIMEOUT_SECONDS) if (not isinstance(result, dict) or result.get("status") not in {"applied", "replayed", "recovered", "planned"} or result.get("source_authority") not in LOCAL_AUTHORITY_SOURCES diff --git a/loopx/control_plane/todos/provider_terminal_lifecycle.py b/loopx/control_plane/todos/provider_terminal_lifecycle.py index 5923dba006..d1f0ba0116 100644 --- a/loopx/control_plane/todos/provider_terminal_lifecycle.py +++ b/loopx/control_plane/todos/provider_terminal_lifecycle.py @@ -44,6 +44,10 @@ _ARCHIVE_REQUEST_SCHEMA = "loopx_local_coordination_todo_archive_request_v0" _ARCHIVE_ACK_REQUEST_SCHEMA = "loopx_local_coordination_todo_archive_ack_request_v0" _ACCEPTED = {"applied", "recovered", "replayed", "no_change", "planned"} +# A promoted File authority can verify and publish a retained journal while +# holding the canonical writer fence. This is an explicit terminal-command +# budget, not a global relaxation of the Effect runtime request deadline. +_TERMINAL_RUNTIME_TIMEOUT_SECONDS = 45.0 TodoMutation = Callable[..., dict[str, Any]] @@ -387,7 +391,8 @@ def terminal_canonical_todo_if_promoted( "observed_at": now_local(), } result = effect_runtime_result( - "coordination.local_authority.todo_terminal", request + "coordination.local_authority.todo_terminal", request, + timeout=_TERMINAL_RUNTIME_TIMEOUT_SECONDS, ) if isinstance(result, Mapping) and result.get("status") == "resolve_validation": # Admission and receipt recovery precede host-local declaration IO. @@ -401,7 +406,10 @@ def terminal_canonical_todo_if_promoted( registry_path=registry_path, goal_id=goal_id, todo_id=todo_id, role=role, persist_if_resolved=not dry_run, ) - result = effect_runtime_result("coordination.local_authority.todo_terminal", request) + result = effect_runtime_result( + "coordination.local_authority.todo_terminal", request, + timeout=_TERMINAL_RUNTIME_TIMEOUT_SECONDS, + ) completion_validation_executed = False if isinstance(result, Mapping) and result.get("status") == "execute_validation": request["validation_source_provider_revision"] = result["provider_revision"] @@ -413,7 +421,8 @@ def terminal_canonical_todo_if_promoted( completion_validation_executed = True request["observed_at"] = now_local() result = effect_runtime_result( - "coordination.local_authority.todo_terminal", request + "coordination.local_authority.todo_terminal", request, + timeout=_TERMINAL_RUNTIME_TIMEOUT_SECONDS, ) if not isinstance(result, Mapping): raise LocalCoordinationAuthorityUnavailable( diff --git a/tests/control_plane/test_effect_runtime_integration.py b/tests/control_plane/test_effect_runtime_integration.py index 93651583e8..3021d3bc38 100644 --- a/tests/control_plane/test_effect_runtime_integration.py +++ b/tests/control_plane/test_effect_runtime_integration.py @@ -156,6 +156,87 @@ def test_managed_runtime_is_reused_and_restart_safe_for_typed_write( ) +def test_response_timeout_keeps_live_runtime_locator_and_never_replays_write( + tmp_path: Path, + monkeypatch, +) -> None: + runtime_dir = tmp_path / "runtime" + runtime_dir.mkdir() + fingerprint = "a" * 64 + monkeypatch.setattr(effect_runtime, "_runtime_dir", lambda: runtime_dir) + monkeypatch.setattr( + effect_runtime, "_runtime_fingerprint_for_request", lambda: fingerprint + ) + monkeypatch.setattr( + effect_runtime, "_start_runtime", + lambda **_kwargs: pytest.fail("a response timeout must not start another server"), + ) + with socket.socket() as listener: + listener.bind(("127.0.0.1", 0)) + listener.listen() + listener.settimeout(2) + info_path = effect_runtime._runtime_info_path(fingerprint) + info = { + "schema_version": effect_runtime.EFFECT_RUNTIME_INFO_SCHEMA_VERSION, + "fingerprint": fingerprint, + "pid": os.getpid(), + "host": "127.0.0.1", + "port": listener.getsockname()[1], + "token": "fixture-token", + } + info_path.write_text(json.dumps(info), encoding="utf-8") + + def serve_one() -> int: + with listener.accept()[0] as connection: + request = b"" + while b"\n" not in request: + request += connection.recv(4096) + time.sleep(0.1) + try: + connection.sendall(b'{}\n') + except OSError: + pass + return 1 + + with ThreadPoolExecutor(max_workers=1) as executor: + served = executor.submit(serve_one) + with pytest.raises(effect_runtime.EffectRuntimeResponseAmbiguous) as error: + effect_runtime.effect_runtime_request( + "coordination.local_authority.todo_terminal", + {"operation_id": "fixture-operation"}, + timeout=0.02, + ) + assert error.value.diagnostic_code == "runtime_response_ambiguous" + assert served.result(timeout=2) == 1 + assert json.loads(info_path.read_text(encoding="utf-8")) == info + + +def test_old_server_close_cannot_remove_replacement_runtime_locator( + tmp_path: Path, + monkeypatch, +) -> None: + runtime_dir = tmp_path / "runtime" + monkeypatch.setattr(effect_runtime, "_runtime_dir", lambda: runtime_dir) + monkeypatch.setenv("LOOPX_EFFECT_RUNTIME_IDLE_MS", "60000") + effect_runtime.effect_runtime_result("runtime.ping", {}) + info_path = effect_runtime._runtime_info_path(effect_runtime._runtime_fingerprint()) + original = json.loads(info_path.read_text(encoding="utf-8")) + replacement = {**original, "pid": os.getpid(), "token": "replacement-token"} + info_path.write_text(json.dumps(replacement), encoding="utf-8") + + effect_runtime._request_with_info( + original, + request_id="shutdown-old-server", + method="runtime.shutdown", + params={}, + timeout=2, + ) + deadline = time.monotonic() + 2 + while effect_runtime._pid_is_alive(original["pid"]) and time.monotonic() < deadline: + time.sleep(0.025) + assert json.loads(info_path.read_text(encoding="utf-8")) == replacement + + def test_retired_coordination_snapshot_mirror_is_rejected_across_runtime_boundary( tmp_path: Path, monkeypatch, diff --git a/tests/control_plane/test_local_coordination_authority.py b/tests/control_plane/test_local_coordination_authority.py index 890d096d42..b5d4174906 100644 --- a/tests/control_plane/test_local_coordination_authority.py +++ b/tests/control_plane/test_local_coordination_authority.py @@ -1459,7 +1459,9 @@ def test_promoted_terminal_lifecycle_commits_successors_and_archive_natively( original_effect_runtime_result = provider_terminal_lifecycle.effect_runtime_result original_authority_runtime_result = local_authority_module.effect_runtime_result - def count_runtime_call(method: str, params: dict[str, object]) -> object: + def count_runtime_call( + method: str, params: dict[str, object], **kwargs: object + ) -> object: runtime_calls.append(method) if method == "coordination.local_authority.todo_archive": archive_operation_ids.append(str(params["operation_id"])) @@ -1467,7 +1469,7 @@ def count_runtime_call(method: str, params: dict[str, object]) -> object: assert params["schema_version"] == "loopx_local_coordination_todo_terminal_lifecycle_request_v3" assert "operation_id" not in params assert params["operation_identity"]["kind"] == "explicit" - result = original_effect_runtime_result(method, params) + result = original_effect_runtime_result(method, params, **kwargs) if method == "coordination.local_authority.todo_terminal": terminal_phases.append(str(result["status"])) return result @@ -1997,9 +1999,11 @@ def test_promoted_terminal_retry_reuses_receipt_after_projection_crash( original_effect_runtime_result = provider_terminal_lifecycle.effect_runtime_result original_authority_runtime_result = local_authority_module.effect_runtime_result - def count_runtime_call(method: str, params: dict[str, object]) -> object: + def count_runtime_call( + method: str, params: dict[str, object], **kwargs: object + ) -> object: runtime_calls.append(method) - return original_effect_runtime_result(method, params) + return original_effect_runtime_result(method, params, **kwargs) def count_authority_runtime_call( method: str, params: dict[str, object], **kwargs: object diff --git a/tests/control_plane/test_todo_operation_receipt_cli.py b/tests/control_plane/test_todo_operation_receipt_cli.py new file mode 100644 index 0000000000..d29f26a50d --- /dev/null +++ b/tests/control_plane/test_todo_operation_receipt_cli.py @@ -0,0 +1,47 @@ +"""The public CLI can inspect an exact canonical operation without retrying it.""" +from __future__ import annotations + +import json +import subprocess +import sys +from pathlib import Path + +from canonical_authority_fixture import initialize_canonical_authority +from loopx.control_plane.coordination.runtime_shadow import ( + build_todo_runtime_shadow_projection, +) + + +def test_todo_receipt_cli_reads_history_without_writing_an_event(tmp_path: Path) -> None: + runtime_root = tmp_path / "runtime" + project = tmp_path / "project" + project.mkdir() + state_path = project / "ACTIVE_GOAL_STATE.md" + state_path.write_text("# Goal\n\n## Agent Todo\n", encoding="utf-8") + registry_path = tmp_path / "registry.json" + registry_path.write_text(json.dumps({"schema_version": 1, + "common_runtime_root": str(runtime_root), "goals": [{"id": "goal-a", + "repo": str(project), "state_file": state_path.name, + "coordination": {"registered_agents": ["agent-a"]}}]}), encoding="utf-8") + projection = build_todo_runtime_shadow_projection( + goal_id="goal-a", todos=[], handoff_mode="soft_claim", + ) + initialize_canonical_authority(runtime_root, "goal-a", projection, state_path=state_path) + + base = [sys.executable, "-m", "loopx.cli", "--format", "json", + "--registry", str(registry_path), "todo", "receipt", "--goal-id", "goal-a"] + found = subprocess.run([*base, "--operation-id", "canonical-fixture"], + capture_output=True, text=True, timeout=30) + assert found.returncode == 0, found.stdout + found.stderr + payload = json.loads(found.stdout) + assert payload["status"] == "found" + assert payload["operation_id"] == "canonical-fixture" + assert payload["source_authority"] == "file_v0" + assert payload["decision_read_from_provider"] is True + assert payload["receipts"] == [] # Bootstrap had no business receipt. + + missing = subprocess.run([*base, "--operation-id", "other-operation"], + capture_output=True, text=True, timeout=30) + assert missing.returncode == 0, missing.stdout + missing.stderr + assert json.loads(missing.stdout)["status"] == "missing" + assert not (runtime_root / "goals" / "goal-a" / "runs" / "index.jsonl").exists() diff --git a/tests/control_plane_ts/local_authority_operation_receipt.test.ts b/tests/control_plane_ts/local_authority_operation_receipt.test.ts new file mode 100644 index 0000000000..6c38d3efe8 --- /dev/null +++ b/tests/control_plane_ts/local_authority_operation_receipt.test.ts @@ -0,0 +1,32 @@ +import assert from "node:assert/strict"; +import {mkdtemp, rm} from "node:fs/promises"; +import {tmpdir} from "node:os"; +import {join} from "node:path"; +import test from "node:test"; + +import {FileAuthorityStore} from "../../loopx/control_plane/coordination/file_authority_store.ts"; +import {readLocalCoordinationOperationReceipt} from "../../loopx/control_plane/coordination/local_authority_read.ts"; + +test("exact local receipt readback uses the selected provider without replaying an operation", async (t) => { + const root = await mkdtemp(join(tmpdir(), "loopx-operation-receipt-")); + t.after(() => rm(root, {recursive: true, force: true})); + const store = new FileAuthorityStore(join(root, "authority", "file-v0"), "goal-a"); + const committed = await store.commitAuthority({expected_provider_revision: null, + operation_id: "terminal:fixture", events: [], next_projection: {goal_id: "goal-a"}, + receipts: [{schema_version: "fixture_terminal_receipt_v0", operation_id: "terminal:fixture"}]}); + assert.equal(committed.status, "applied"); + + const request = {schema_version: "loopx_local_coordination_operation_receipt_request_v0", + runtime_root: root, goal_id: "goal-a", operation_id: "terminal:fixture"}; + const found = await readLocalCoordinationOperationReceipt(request); + assert.equal(found.status, "found"); + assert.equal(found.source_authority, "file_v0"); + assert.deepEqual(found.receipts, [{schema_version: "fixture_terminal_receipt_v0", + operation_id: "terminal:fixture"}]); + assert.equal(found.provider_revision, committed.provider_revision); + const missing = await readLocalCoordinationOperationReceipt({...request, operation_id: "terminal:other"}); + assert.equal(missing.status, "missing"); + const after = await store.loadAuthority(); + assert.equal(after.status, "loaded"); + if (after.status === "loaded") assert.equal(after.provider_revision, committed.provider_revision); +}); diff --git a/tests/control_plane_ts/local_authority_provider.test.ts b/tests/control_plane_ts/local_authority_provider.test.ts index d0d841219e..41b7c63645 100644 --- a/tests/control_plane_ts/local_authority_provider.test.ts +++ b/tests/control_plane_ts/local_authority_provider.test.ts @@ -116,6 +116,12 @@ function providerCalls(directory: string, revision: string, dryRun: boolean) { observeLocalCoordinationOwnership: [{...input, schema_version: "loopx_local_ownership_observation_request_v0"}], listLocalCoordinationTodos: [{...input, schema_version: runtime.LOCAL_COORDINATION_TODO_LIST_REQUEST_SCHEMA}], readLocalCoordinationTodo: [{...input, schema_version: runtime.LOCAL_COORDINATION_TODO_READ_REQUEST_SCHEMA}], + readLocalCoordinationOperationReceipt: [{ + runtime_root: directory, + goal_id: "goal-a", + schema_version: "loopx_local_coordination_operation_receipt_request_v0", + operation_id: "open-failure", + }], createLocalCoordinationTodo: [ {...input, schema_version: runtime.LOCAL_COORDINATION_TODO_CREATE_REQUEST_SCHEMA, todo: {}}, {...witnessed, schema_version: runtime.LOCAL_COORDINATION_TODO_CREATE_WITNESSED_REQUEST_SCHEMA, todo: {}}], diff --git a/tests/control_plane_ts/shadow_native_writer_boundary.test.ts b/tests/control_plane_ts/shadow_native_writer_boundary.test.ts index cb19d92bb3..a19304c787 100644 --- a/tests/control_plane_ts/shadow_native_writer_boundary.test.ts +++ b/tests/control_plane_ts/shadow_native_writer_boundary.test.ts @@ -5,6 +5,7 @@ import { tmpdir } from "node:os"; import { join } from "node:path"; import test from "node:test"; import { atomicWriteJson } from "../../loopx/control_plane/effect_runtime_io.ts"; +import { withCanonicalWriter } from "../../loopx/control_plane/coordination/local_authority_write.ts"; import { shadowMaintenanceLockPath, shadowManagementStatePath, @@ -105,6 +106,36 @@ import { taskLeaseLockPath } from "../../loopx/control_plane/work_items/task_lea import { withFileMutationLock } from "../../loopx/control_plane/effect_runtime_io.ts"; import { access } from "node:fs/promises"; +test("a canonical writer waits for a slow file-v0 critical section without losing the fence", async (t) => { + const root = await mkdtemp(join(tmpdir(), "loopx-canonical-long-writer-")); + t.after(() => rm(root, {recursive: true, force: true})); + let entered!: () => void; + let release!: () => void; + const firstEntered = new Promise((resolve) => { entered = resolve; }); + const firstRelease = new Promise((resolve) => { release = resolve; }); + const first = withCanonicalWriter(root, "goal-a", false, async () => { + entered(); + await firstRelease; + return "first"; + }); + await firstEntered; + let secondState = "pending"; + const second = withCanonicalWriter(root, "goal-a", false, async () => "second") + .then((value) => { secondState = "applied"; return value; }, (error: unknown) => { + secondState = "failed"; + return error; + }); + try { + await new Promise((resolve) => setTimeout(resolve, 5_200)); + assert.equal(secondState, "pending", "the old five-second lock deadline must not reject a live writer"); + } finally { + release(); + } + assert.equal(await first, "first"); + assert.equal(await second, "second"); + assert.equal(secondState, "applied"); +}); + function fenceRequest(root: string) { return {schema_version: "loopx_legacy_coordination_writer_fence_engage_request_v0", runtime_root: root, goal_id: "goal-a", state_path: join(root, "ACTIVE_GOAL_STATE.md"), fence: { From b2ab91123194918cb8e193d049c8243b2910f617 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Thu, 24 Sep 2026 07:52:47 +0800 Subject: [PATCH 2/8] fix(runtime): retain live locator on pre-send connection failure Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- loopx/control_plane/effect_runtime.py | 31 +++++++++++--- .../test_effect_runtime_integration.py | 40 +++++++++++++++++++ 2 files changed, 66 insertions(+), 5 deletions(-) diff --git a/loopx/control_plane/effect_runtime.py b/loopx/control_plane/effect_runtime.py index 2e7a8550d1..7888f5c5cc 100644 --- a/loopx/control_plane/effect_runtime.py +++ b/loopx/control_plane/effect_runtime.py @@ -371,6 +371,25 @@ def _pid_is_alive(value: object) -> bool: return process_is_alive(value) +def _reap_exited_runtime_child(info: Mapping[str, Any] | None) -> None: + """Let a dead directly spawned child fail the next non-signaling probe. + + A stopped child can remain a zombie until its Python parent reaps it; + ``kill(pid, 0)`` still reports that zombie as present. This helper does + nothing for a live child or a runtime owned by another process. + """ + + if os.name == "nt" or not isinstance(info, Mapping): + return + pid = info.get("pid") + if isinstance(pid, bool) or not isinstance(pid, int) or pid <= 0: + return + try: + os.waitpid(pid, os.WNOHANG) + except (ChildProcessError, OSError): + pass + + def _start_lock_holder_pid(path: Path) -> int | None: try: value = int(path.read_text(encoding="utf-8").strip()) @@ -819,11 +838,13 @@ def effect_runtime_request( except (OSError, RuntimeError) as exc: last_error = exc if attempt == 0 and retry_safe: - # Only pre-send connection failures reach this branch. Do not - # remove a replacement runtime published by another caller. - current = _read_info(info_path, fingerprint=fingerprint) - if info is not None and current is not None and current.get("token") == info.get("token"): - info_path.unlink(missing_ok=True) + # Only pre-send connection failures reach this branch. Re-read + # the locator on retry: reap our own exited child first so a + # zombie is rejected by _read_info, while a live server may + # simply be draining. + # Even a token check followed by unlink would race with a + # replacement server publishing its own locator. + _reap_exited_runtime_child(info) continue break if isinstance(last_error, TimeoutError): diff --git a/tests/control_plane/test_effect_runtime_integration.py b/tests/control_plane/test_effect_runtime_integration.py index 3021d3bc38..9d96d7fb36 100644 --- a/tests/control_plane/test_effect_runtime_integration.py +++ b/tests/control_plane/test_effect_runtime_integration.py @@ -237,6 +237,46 @@ def test_old_server_close_cannot_remove_replacement_runtime_locator( assert json.loads(info_path.read_text(encoding="utf-8")) == replacement +def test_pre_send_connection_failure_does_not_remove_live_locator( + tmp_path: Path, + monkeypatch, +) -> None: + runtime_dir = tmp_path / "runtime" + runtime_dir.mkdir() + fingerprint = "b" * 64 + monkeypatch.setattr(effect_runtime, "_runtime_dir", lambda: runtime_dir) + monkeypatch.setattr( + effect_runtime, "_runtime_fingerprint_for_request", lambda: fingerprint + ) + monkeypatch.setattr( + effect_runtime, "_start_runtime", + lambda **_kwargs: pytest.fail("a live locator must not start a replacement"), + ) + info_path = effect_runtime._runtime_info_path(fingerprint) + info = { + "schema_version": effect_runtime.EFFECT_RUNTIME_INFO_SCHEMA_VERSION, + "fingerprint": fingerprint, + "pid": os.getpid(), + "host": "127.0.0.1", + "port": 1, + "token": "still-live", + } + info_path.write_text(json.dumps(info), encoding="utf-8") + calls = 0 + + def refuse_before_send(_info: object, **_kwargs: object) -> object: + nonlocal calls + calls += 1 + raise ConnectionRefusedError("fixture refused before send") + + monkeypatch.setattr(effect_runtime, "_request_with_info", refuse_before_send) + with pytest.raises(effect_runtime.EffectRuntimeStartupError) as error: + effect_runtime.effect_runtime_request("runtime.ping", {}, retry_safe=True) + assert error.value.diagnostic_code == "runtime_request_failed" + assert calls == 2 + assert json.loads(info_path.read_text(encoding="utf-8")) == info + + def test_retired_coordination_snapshot_mirror_is_rejected_across_runtime_boundary( tmp_path: Path, monkeypatch, From 0fae92a602b9676998805e2cc34016a92ff33978 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Thu, 24 Sep 2026 08:06:52 +0800 Subject: [PATCH 3/8] fix(quota): resume pending monitor after append-only index growth Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../quota/monitor_poll_commit.ts | 40 +++++++++++----- .../quota_monitor_poll_commit.test.ts | 48 ++++++++++++++++++- 2 files changed, 76 insertions(+), 12 deletions(-) diff --git a/loopx/control_plane/quota/monitor_poll_commit.ts b/loopx/control_plane/quota/monitor_poll_commit.ts index f908e0bbdb..ed8c5d17cd 100644 --- a/loopx/control_plane/quota/monitor_poll_commit.ts +++ b/loopx/control_plane/quota/monitor_poll_commit.ts @@ -1326,6 +1326,20 @@ function matchingIndexRecord( return null; } +function pendingIndexHistoryIntact( + pending: PendingMonitorReceipt, + current: Buffer | null, +): boolean { + const expectedBytes = pending.expected_index_bytes; + if (!Number.isSafeInteger(expectedBytes) || expectedBytes < 0 || + (current?.length ?? 0) < expectedBytes) return false; + if (pending.expected_index_digest === null) { + return expectedBytes === 0; + } + return sha256Bytes((current ?? Buffer.alloc(0)).subarray(0, expectedBytes)) === + pending.expected_index_digest; +} + function indexRecordFor( request: MonitorRequest, record: JsonObject, @@ -2041,6 +2055,13 @@ export async function evaluateQuotaMonitorPollCommit( ); } const records = indexRecords(indexContent); + if (existing?.status === "provider_pending" && matchingIndexRecord(records, request.effect_id)) { + return result( + request, fingerprint, "conflict", null, {ok: false, appended: false}, + "quota monitor-poll effect identity exists while its provider receipt is pending", + currentDigest, null, null, {reason_code: "effect_id_conflict"}, + ); + } if (!existing && matchingIndexRecord(records, request.effect_id)) { return result( request, @@ -2057,11 +2078,7 @@ export async function evaluateQuotaMonitorPollCommit( } if (request.phase === "preflight") { - if ( - existing && - (currentDigest !== existing.expected_index_digest || - (indexBytes?.length ?? 0) !== existing.expected_index_bytes) - ) { + if (existing && !pendingIndexHistoryIntact(existing, indexBytes)) { return result( request, fingerprint, @@ -2139,10 +2156,7 @@ export async function evaluateQuotaMonitorPollCommit( "malformed_transaction_receipt", ); } - if ( - currentDigest !== existing.expected_index_digest || - (indexBytes?.length ?? 0) !== existing.expected_index_bytes - ) { + if (!pendingIndexHistoryIntact(existing, indexBytes)) { return result( request, fingerprint, @@ -2161,8 +2175,12 @@ export async function evaluateQuotaMonitorPollCommit( plan, validatedProviderReceipt(request.provider_receipt, plan), ); - expectedDigest = existing.expected_index_digest; - expectedBytes = existing.expected_index_bytes; + // The pending receipt freezes the admitted observation, not the entire + // append-only run index. An unrelated run may have committed while the + // provider wrote this exact effect. The new prepared WAL fences the + // current index under the same lock after its historical prefix passes. + expectedDigest = currentDigest; + expectedBytes = indexBytes?.length ?? 0; } else if (existing?.status === "provider_pending") { return effectConflict(request, fingerprint, currentDigest, existing); } diff --git a/tests/control_plane_ts/quota_monitor_poll_commit.test.ts b/tests/control_plane_ts/quota_monitor_poll_commit.test.ts index 5aa87ff20a..ea681c8f1e 100644 --- a/tests/control_plane_ts/quota_monitor_poll_commit.test.ts +++ b/tests/control_plane_ts/quota_monitor_poll_commit.test.ts @@ -810,7 +810,7 @@ test("retry repairs a truncated owned index tail and rejects artifact drift", as ); }); -test("provider retry fails before writeback when its preflight index fence is stale", async (t) => { +test("pending provider effect settles after an unrelated append-only index commit", async (t) => { const runtimeRoot = await tempRuntime(t); const pending = request({ phase: "preflight", @@ -836,6 +836,52 @@ test("provider retry fails before writeback when its preflight index fence is st "written", ); + const retry = await evaluateQuotaMonitorPollCommit(pending); + assert.equal(retry.status, "provider_required"); + const receipt = { + schema_version: "monitor_poll_todo_writeback_v0", + monitor_effect_id: "quota-monitor-poll:pending-provider", + dry_run: false, + goal_id: goalId, + todo_id: "todo_public_monitor", + target_key: null, + result_hash: "unchanged-42", + material_change: false, + material_change_generation: 0, + consecutive_no_change: 1, + last_checked_at: pending.generated_at, + next_due_at: null, + cadence: null, + todo_update: {ok: true}, + next_todos: [], + successor_receipts: [], + }; + const settled = await evaluateQuotaMonitorPollCommit({ + ...pending, phase: "commit", provider_receipt: receipt, + }); + assert.equal(settled.status, "written"); + assert.equal((await evaluateQuotaMonitorPollCommit({ + ...pending, phase: "commit", provider_receipt: receipt, + })).status, "replayed"); + const indexPath = join(runtimeRoot, "goals", goalId, "runs", "index.jsonl"); + assert.equal((await readFile(indexPath, "utf8")).trim().split("\n").length, 2); +}); + +test("pending provider effect still rejects changed preflight index history", async (t) => { + const runtimeRoot = await tempRuntime(t); + const prior = request({phase: "commit", runtime_root: runtimeRoot, execute: true, + effect_id: "quota-monitor-poll:prior-history"}); + const initial = await evaluateQuotaMonitorPollCommit(prior); + assert.equal(initial.status, "written"); + const pending = request({phase: "preflight", runtime_root: runtimeRoot, execute: true, + expected_index_digest: initial.index_digest, + effect_id: "quota-monitor-poll:pending-history", + observation: observation({todo_id: "todo_public_monitor", result_hash: "unchanged-42"})}); + assert.equal((await evaluateQuotaMonitorPollCommit(pending)).status, "provider_required"); + const indexPath = join(runtimeRoot, "goals", goalId, "runs", "index.jsonl"); + const original = await readFile(indexPath, "utf8"); + assert.match(original, /quota-monitor-poll:prior-history/); + await writeFile(indexPath, original.replaceAll("quota-monitor-poll:prior-history", "quota-monitor-poll:alter-history")); const retry = await evaluateQuotaMonitorPollCommit(pending); assert.equal(retry.status, "conflict"); assert.equal(retry.reason_code, "index_digest_conflict"); From 9862bd411d2bad55bf23e795177c7650580d9cc4 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Thu, 24 Sep 2026 08:45:33 +0800 Subject: [PATCH 4/8] test(control-plane): align fence probes with writer wait budget Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- tests/control_plane/checkpoint_process.py | 4 +++- tests/control_plane/test_checkpoint_provider_fence.py | 5 ++++- tests/control_plane/test_reviewed_terminal_actions.py | 4 ++-- 3 files changed, 9 insertions(+), 4 deletions(-) diff --git a/tests/control_plane/checkpoint_process.py b/tests/control_plane/checkpoint_process.py index 32f07f2759..9e208c2b64 100644 --- a/tests/control_plane/checkpoint_process.py +++ b/tests/control_plane/checkpoint_process.py @@ -88,7 +88,9 @@ def native(method, params): # must keep excluding a new index reader and receipt replacement. os._exit(0) child = start_probe(envelope) - stdout, stderr = child.communicate(timeout=30) + # The canonical writer may now wait 30s for the real provider fence; + # let that typed lock result win over the probe harness deadline. + stdout, stderr = child.communicate(timeout=45) assert child.returncode == 0, stdout + stderr return json.loads(stdout) diff --git a/tests/control_plane/test_checkpoint_provider_fence.py b/tests/control_plane/test_checkpoint_provider_fence.py index 9a099afae4..16cf52632c 100644 --- a/tests/control_plane/test_checkpoint_provider_fence.py +++ b/tests/control_plane/test_checkpoint_provider_fence.py @@ -95,7 +95,10 @@ def test_existing_cli_maintenance_guard_precedes_provider_commit(tmp_path, monke before = read_canonical_todos_if_promoted(runtime_root=runtime, goal_id=GOAL_ID)["provider_revision"] with context_io._source_guard(runtime, GOAL_ID, state): writer = public_writer(registry, runtime, barrier, provider) - stdout, stderr = writer.communicate(timeout=15) + # The canonical writer now waits up to 30s for a real retained File + # history commit. Give the subprocess enough time to return its typed + # lock error rather than timing out the test harness first. + stdout, stderr = writer.communicate(timeout=60) assert writer.returncode == 1, stdout + stderr assert "lock timed out" in json.loads(stdout)["error"] assert not (barrier / "writer-entered").exists() diff --git a/tests/control_plane/test_reviewed_terminal_actions.py b/tests/control_plane/test_reviewed_terminal_actions.py index a42c6252b3..a7794c492f 100644 --- a/tests/control_plane/test_reviewed_terminal_actions.py +++ b/tests/control_plane/test_reviewed_terminal_actions.py @@ -210,12 +210,12 @@ def test_monitor_completion_uses_explicit_identity_modes(tmp_path, monkeypatch, execute = provider_terminal_lifecycle.effect_runtime_result identities = [] - def capture(method, params): + def capture(method, params, **kwargs): if method == "coordination.local_authority.todo_terminal": assert params["schema_version"] == "loopx_local_coordination_todo_terminal_lifecycle_request_v3" assert "operation_id" not in params identities.append(params["operation_identity"]) - return execute(method, params) + return execute(method, params, **kwargs) monkeypatch.setattr(provider_terminal_lifecycle, "effect_runtime_result", capture) kwargs = dict(registry_path=registry, goal_id="goal-a", todo_id=todo_id, From d846e640f9fc0cd418299c5a56efb818ef96db57 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Thu, 24 Sep 2026 09:09:14 +0800 Subject: [PATCH 5/8] test(architecture): refresh Todo registry IO census Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../project_registry_io_manifest_v1.json | 24 ++++++++++++------- 1 file changed, 16 insertions(+), 8 deletions(-) diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index 3d09ed9372..a561a353ff 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -727,7 +727,7 @@ }, { "site": "loopx/cli_commands/todo.py::._validated_replan_successor_obligation::codec_read:load_registry#1", - "line": 121, + "line": 126, "column": 16, "kind": "codec_read", "api": "load_registry", @@ -735,31 +735,39 @@ }, { "site": "loopx/cli_commands/todo.py::.handle_todo_command::codec_read:load_registry#1", - "line": 233, - "column": 24, + "line": 255, + "column": 49, "kind": "codec_read", "api": "load_registry", "classification": "codec_api" }, { "site": "loopx/cli_commands/todo.py::.handle_todo_command::codec_read:load_registry#2", - "line": 413, - "column": 21, + "line": 271, + "column": 24, "kind": "codec_read", "api": "load_registry", "classification": "codec_api" }, { "site": "loopx/cli_commands/todo.py::.handle_todo_command::codec_read:load_registry#3", - "line": 595, - "column": 13, + "line": 451, + "column": 21, "kind": "codec_read", "api": "load_registry", "classification": "codec_api" }, { "site": "loopx/cli_commands/todo.py::.handle_todo_command::codec_read:load_registry#4", - "line": 637, + "line": 633, + "column": 13, + "kind": "codec_read", + "api": "load_registry", + "classification": "codec_api" + }, + { + "site": "loopx/cli_commands/todo.py::.handle_todo_command::codec_read:load_registry#5", + "line": 675, "column": 38, "kind": "codec_read", "api": "load_registry", From f8489b184906a5185c6cec7b190d4ad0e4e0a13e Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Thu, 24 Sep 2026 09:30:43 +0800 Subject: [PATCH 6/8] test(sqlite): gate deferred lifecycle on qualified runtime Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../deferred_hard_lease_lifecycle.test.ts | 19 +++++++++++++------ 1 file changed, 13 insertions(+), 6 deletions(-) diff --git a/tests/control_plane_ts/deferred_hard_lease_lifecycle.test.ts b/tests/control_plane_ts/deferred_hard_lease_lifecycle.test.ts index 73e1f4c9e5..445800918b 100644 --- a/tests/control_plane_ts/deferred_hard_lease_lifecycle.test.ts +++ b/tests/control_plane_ts/deferred_hard_lease_lifecycle.test.ts @@ -7,6 +7,7 @@ import test from "node:test"; import type {JsonObject} from "../../loopx/control_plane/effect_program.ts"; import {FileAuthorityStore} from "../../loopx/control_plane/coordination/file_authority_store.ts"; import {SqliteAuthorityStore} from "../../loopx/control_plane/coordination/sqlite_authority_store.ts"; +import {sqliteRuntimeIdentity} from "../../loopx/control_plane/coordination/sqlite_runtime.ts"; import {canonicalAuthoritySha256} from "../../loopx/control_plane/coordination/authority_store_codec.ts"; import {TODO_DOMAIN_ITEM_SCHEMA, TODO_DOMAIN_READ_RECORD_SCHEMA, TODO_DOMAIN_RECORD_CONTRACT} from "../../loopx/control_plane/coordination/coordination_state_contract.ts"; @@ -32,11 +33,12 @@ async function seeded(provider: "file" | "sqlite", overrides: JsonObject = {}, l resume_when: "resume_at:2026-09-05T22:00:00Z", ...overrides}; if (todo.claimed_by === null) delete todo.claimed_by; const todos = [todo]; - assert.equal((await store.commitAuthority({operation_id: "seed", expected_provider_revision: null, + const seededResult = await store.commitAuthority({operation_id: "seed", expected_provider_revision: null, events: [], receipts: [], next_projection: {goal_id: GOAL, handoff_mode: "hard_lease", todos, leases: lease ? [lease] : [], todo_read_model: {schema_version: TODO_DOMAIN_READ_RECORD_SCHEMA, todo_count: 1, records_sha256: canonicalAuthoritySha256(todos), - contract_fields: [...TODO_DOMAIN_RECORD_CONTRACT.fields]}}})).status, "applied"); + contract_fields: [...TODO_DOMAIN_RECORD_CONTRACT.fields]}}}); + assert.equal(seededResult.status, "applied", JSON.stringify(seededResult)); return store; } @@ -71,8 +73,13 @@ function oldLease(status: "active" | "released", expires_at: string): JsonObject version: 3, lease_epoch: 2, write_scopes: [], acquire_ttl_seconds: 900}; } +const sqliteSkip = sqliteRuntimeIdentity().sqlite_authority_qualified + ? false : "requires the qualified SQLite runtime"; + for (const provider of ["file", "sqlite"] as const) { - test(`${provider}: deferred resume is one CAS transition and never grants execution`, async () => { + const providerTest = (name: string, run: () => Promise) => + test(name, {skip: provider === "sqlite" && sqliteSkip}, run); + providerTest(`${provider}: deferred resume is one CAS transition and never grants execution`, async () => { const store = await seeded(provider); const before = await read(store); assert.equal(leaseOwnerRejection({status: "deferred", claimed_by: OWNER, excluded_agents: []}, OWNER, AGENTS), @@ -99,7 +106,7 @@ for (const provider of ["file", "sqlite"] as const) { "coordination_operation_identity_mismatch"); }); - test(`${provider}: expired lease is retired atomically, while live execution and foreign edits fail closed`, async () => { + providerTest(`${provider}: expired lease is retired atomically, while live execution and foreign edits fail closed`, async () => { const expired = await seeded(provider, {}, oldLease("active", "2026-09-05T22:30:00Z")); const resumed = await executeCoordinationTodoUpdate(expired, resume("retire-expired")); assert.equal(resumed.status, "applied", JSON.stringify(resumed)); @@ -123,7 +130,7 @@ for (const provider of ["file", "sqlite"] as const) { } }); - test(`${provider}: an unclaimed deferred Todo stays unclaimed; excluded actors cannot reopen it`, async () => { + providerTest(`${provider}: an unclaimed deferred Todo stays unclaimed; excluded actors cannot reopen it`, async () => { const unclaimed = await seeded(provider, {claimed_by: null}); const applied = await executeCoordinationTodoUpdate(unclaimed, resume("unclaimed-resume")); assert.equal(applied.status, "applied", JSON.stringify(applied)); @@ -138,7 +145,7 @@ for (const provider of ["file", "sqlite"] as const) { assert.deepEqual(await read(excluded), before); }); - test(`${provider}: deferred supersede closes one Todo and retires expired lease lineage`, async () => { + providerTest(`${provider}: deferred supersede closes one Todo and retires expired lease lineage`, async () => { const store = await seeded(provider, {}, oldLease("active", "2026-09-05T22:30:00Z")); const input = supersede("supersede-once"); const before = await read(store); From ac011885df2d975e945ca3939c73c9a26370a5b0 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Thu, 24 Sep 2026 11:24:51 +0800 Subject: [PATCH 7/8] perf(authority): reuse verified File journal across runtime requests Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../coordination/file_authority_store.ts | 55 ++++++++++++++++--- .../control_plane_ts/authority_store.test.ts | 36 ++++++++++++ 2 files changed, 84 insertions(+), 7 deletions(-) diff --git a/loopx/control_plane/coordination/file_authority_store.ts b/loopx/control_plane/coordination/file_authority_store.ts index 0a75a6012f..969d24524e 100644 --- a/loopx/control_plane/coordination/file_authority_store.ts +++ b/loopx/control_plane/coordination/file_authority_store.ts @@ -30,6 +30,29 @@ import {AuthorityJournalScan} from "./authority_journal_scan.ts"; const FILE_AUTHORITY_STORE_SCHEMA = "loopx_file_authority_store_v0"; const STORE_IDENTITY_PATTERN = /^file:[0-9a-f]{32}$/; +// File-v0 retains every projection in one envelope. A managed Effect server +// opens a new store handle for each request, so revalidating an unchanged +// journal on every read makes one Goal's history dominate the RPC budget. +// Keep only one verified document process-wide; bytes and store identity must +// both match before a later handle may reuse that validation. +const MAX_CACHED_DOCUMENT_BYTES = 128 * 1024 * 1024; +let verifiedDocument: { + path: string; + identity: string; + digest: string; + document: FileAuthorityStoreDocument; +} | null = null; + +function documentDigest(raw: Uint8Array): string { + return createHash("sha256").update(raw).digest("hex"); +} + +function rememberVerifiedDocument(path: string, identity: string, raw: Uint8Array, + digest: string, document: FileAuthorityStoreDocument): void { + verifiedDocument = raw.byteLength <= MAX_CACHED_DOCUMENT_BYTES + ? {path, identity, digest, document} + : null; +} interface FileAuthorityStoreDocument extends JsonObject, RetainedAuthorityJournal { schema_version: typeof FILE_AUTHORITY_STORE_SCHEMA; @@ -167,6 +190,11 @@ export class FileAuthorityStore implements AuthorityStore { /** Filesystem-only crash seam; the archive owner must still fsync both parents. */ protected async archiveRenamed(): Promise {} + /** Full-history verification seam; unchanged byte-identical reads may reuse it. */ + protected decodeStoredDocument(value: unknown, identity: string): FileAuthorityStoreDocument { + return decodeDocument(value, this.goalId, identity); + } + private async readStoreIdentity(createIfMissing = !this.existingOnly): Promise { try { const identity = await readFile(this.identityPath, "utf8"); @@ -213,19 +241,26 @@ export class FileAuthorityStore implements AuthorityStore { } } - private async readDocument(): Promise { - let raw: string; + private async readDocument(knownIdentity?: string): Promise { + let raw: Buffer; try { - raw = await readFile(this.path, "utf8"); + raw = await readFile(this.path); } catch (error) { if ((error as NodeJS.ErrnoException).code === "ENOENT") return null; throw new FileStoreUnavailableError( error instanceof Error ? error.message : "authority document unavailable", ); } - const identity = await this.readStoreIdentity(); + const identity = knownIdentity ?? await this.readStoreIdentity(); try { - return decodeDocument(JSON.parse(raw), this.goalId, identity); + const digest = documentDigest(raw); + if (verifiedDocument?.path === this.path && + verifiedDocument.identity === identity && verifiedDocument.digest === digest) { + return verifiedDocument.document; + } + const document = this.decodeStoredDocument(JSON.parse(raw.toString("utf8")), identity); + rememberVerifiedDocument(this.path, identity, raw, digest, document); + return document; } catch (error) { if (error instanceof SyntaxError) { throw new AuthorityStoreProtocolError(`file authority store JSON is invalid: ${error.message}`); @@ -284,7 +319,7 @@ export class FileAuthorityStore implements AuthorityStore { // A restored directory must not race a missing-head bootstrap and // bind new authority bytes to an identity observed before the lock. identity = await this.readStoreIdentity(); - current = await this.readDocument(); + current = await this.readDocument(identity); } catch (error) { return { status: "failed", @@ -319,8 +354,14 @@ export class FileAuthorityStore implements AuthorityStore { const document: FileAuthorityStoreDocument = {schema_version: FILE_AUTHORITY_STORE_SCHEMA, goal_id: this.goalId, store_identity: identity, ...journal}; try { - await this.replaceDurably(this.path, canonicalAuthorityBytes(document)); + const bytes = canonicalAuthorityBytes(document); + await this.replaceDurably(this.path, bytes); + rememberVerifiedDocument(this.path, identity, bytes, documentDigest(bytes), document); } catch (error) { + // A failure after rename may already have published the new bytes. + // The next read must prove the actual file rather than reuse either + // the previous or attempted document. + verifiedDocument = null; return { status: "ambiguous", reason_code: "commit_outcome_unknown", diff --git a/tests/control_plane_ts/authority_store.test.ts b/tests/control_plane_ts/authority_store.test.ts index ce494c346a..0bb641e29f 100644 --- a/tests/control_plane_ts/authority_store.test.ts +++ b/tests/control_plane_ts/authority_store.test.ts @@ -90,8 +90,12 @@ test("corrupt, cross-goal, or revision-divergent documents fail closed", async ( const { store } = await fixture(t); const applied = await store.commitAuthority(commit(null, "operation-a", 1, 1)); assert.equal(applied.status, "applied"); + const otherHandle = new FileAuthorityStore(store.directory, "goal-a"); + assert.deepEqual(await otherHandle.loadAuthority(), await store.loadAuthority()); const original = JSON.parse(await readFile(store.path, "utf8")); + // A same-length replacement must not inherit the verified journal simply + // because its path or filesystem size is unchanged. await writeFile(store.path, JSON.stringify({ ...original, goal_id: "goal-b" }), "utf8"); assert.equal((await store.loadAuthority()).status, "failed"); @@ -107,6 +111,38 @@ test("corrupt, cross-goal, or revision-divergent documents fail closed", async ( if (divergent.status === "failed") assert.match(divergent.reason, /revision lineage/); }); +test("file verification is reused only for exact bytes and store identity", async (t) => { + const { root, store } = await fixture(t); + const applied = await store.commitAuthority(commit(null, "operation-a", 1, 1)); + assert.equal(applied.status, "applied"); + const validBytes = await readFile(store.path, "utf8"); + // A benign external rewrite forces one full validation; a second handle + // reads the same proven bytes without validating the whole journal again. + await writeFile(store.path, `${validBytes}\n`, "utf8"); + class CountingStore extends FileAuthorityStore { + static validations = 0; + protected override decodeStoredDocument(value: unknown, identity: string) { + CountingStore.validations += 1; + return super.decodeStoredDocument(value, identity); + } + } + const first = new CountingStore(root, "goal-a"); + const second = new CountingStore(root, "goal-a"); + assert.equal((await first.loadAuthority()).status, "loaded"); + assert.equal((await second.loadAuthority()).status, "loaded"); + assert.equal(CountingStore.validations, 1); + + const originalIdentity = await readFile(store.identityPath, "utf8"); + await writeFile(store.identityPath, `file:${"f".repeat(32)}`, "utf8"); + assert.equal((await second.loadAuthority()).status, "failed"); + assert.equal(CountingStore.validations, 2); + await writeFile(store.identityPath, originalIdentity, "utf8"); + + await writeFile(store.path, validBytes.replace('"goal-a"', '"goal-b"'), "utf8"); + assert.equal((await second.loadAuthority()).status, "failed"); + assert.equal(CountingStore.validations, 3); +}); + test("store identity is one durable directory lineage and restored bytes are fenced", async (t) => { const { root, store } = await fixture(t); const handles = Array.from({ length: 8 }, () => new FileAuthorityStore(root, "goal-a")); From d2132d6db2d32262867c632a9b70cddfd44796c9 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Thu, 24 Sep 2026 11:58:05 +0800 Subject: [PATCH 8/8] fix(runtime): observe ambiguous shutdown before reporting restart Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- loopx/control_plane/effect_runtime.py | 7 +- .../test_effect_runtime_restart.py | 107 ++++++++++++++++++ 2 files changed, 113 insertions(+), 1 deletion(-) diff --git a/loopx/control_plane/effect_runtime.py b/loopx/control_plane/effect_runtime.py index 2c2d167f42..f0e75891eb 100644 --- a/loopx/control_plane/effect_runtime.py +++ b/loopx/control_plane/effect_runtime.py @@ -511,7 +511,12 @@ def restart_effect_runtime(*, timeout: float = 5.0) -> dict[str, Any]: params={}, timeout=timeout, ) - except (EffectRuntimeRejected, EffectRuntimeRemoteError, OSError): + except ( + EffectRuntimeRejected, + EffectRuntimeRemoteError, + EffectRuntimeResponseAmbiguous, + OSError, + ): # A runtime that is already closing must still be reported as pending # rather than as a failed restart. pass diff --git a/tests/control_plane/test_effect_runtime_restart.py b/tests/control_plane/test_effect_runtime_restart.py index b40769133a..06b8de6be4 100644 --- a/tests/control_plane/test_effect_runtime_restart.py +++ b/tests/control_plane/test_effect_runtime_restart.py @@ -1,12 +1,19 @@ """A reused managed runtime keeps its own Node; operators need a restart path.""" from __future__ import annotations +from concurrent.futures import ThreadPoolExecutor import json +import os +import socket import subprocess from pathlib import Path from types import SimpleNamespace +import time + +import pytest from loopx.cli import build_parser +from loopx.cli_commands import doctor as doctor_command from loopx.cli_commands.doctor import handle_doctor_command from loopx.control_plane import effect_runtime from loopx.doctor import render_doctor_markdown @@ -35,6 +42,106 @@ def test_restart_reports_not_running_without_a_serving_runtime( assert result["previous_runtime_identity"] is None +@pytest.mark.parametrize( + ("observed_token", "expected_status"), + [("replacement", "stopped"), ("serving", "shutdown_pending")], +) +def test_ambiguous_shutdown_still_observes_the_serving_runtime( + tmp_path: Path, + monkeypatch, + observed_token: str, + expected_status: str, +) -> None: + """A lost shutdown response is not proof of failure or permission to retry.""" + + info = {"pid": 12345, "token": "serving"} + monkeypatch.setattr(effect_runtime, "_runtime_fingerprint", lambda: "fixture") + monkeypatch.setattr( + effect_runtime, "_runtime_info_path", lambda _: tmp_path / "runtime.json" + ) + monkeypatch.setattr(effect_runtime, "_read_info", lambda *_args, **_kwargs: info) + requests = [] + + def lost_shutdown_response(*_args, **kwargs): + requests.append(kwargs["method"]) + raise effect_runtime.EffectRuntimeResponseAmbiguous( + "runtime.shutdown", timeout=0.01 + ) + + monkeypatch.setattr( + effect_runtime, + "_request_with_info", + lost_shutdown_response, + ) + monkeypatch.setattr( + effect_runtime, "_serving_token", lambda _: (True, observed_token) + ) + monkeypatch.setattr(effect_runtime, "_pid_is_alive", lambda _: True) + + result = effect_runtime.restart_effect_runtime(timeout=0.01) + assert result["status"] == expected_status + assert result["stopped"] is (expected_status == "stopped") + assert requests == ["runtime.shutdown"] + + if observed_token == "replacement": + monkeypatch.setattr(doctor_command, "collect_doctor", lambda **_: {"ok": True}) + captured: dict[str, object] = {} + args = SimpleNamespace( + deep=False, + agent_type=None, + installation_only=True, + restart_runtime=True, + subcommand_format="json", + format="json", + ) + assert handle_doctor_command( + args, lambda payload, _format, _render: captured.update(payload) + ) == 0 + assert captured["effect_runtime_restart"]["status"] == "stopped" + assert requests == ["runtime.shutdown", "runtime.shutdown"] + + +@pytest.mark.parametrize("lost_response", ["eof", "timeout"]) +def test_restart_observes_real_shutdown_transport_loss_without_retry( + tmp_path: Path, + monkeypatch, + lost_response: str, +) -> None: + with socket.socket() as listener: + listener.bind(("127.0.0.1", 0)) + listener.listen() + info = { + "host": "127.0.0.1", + "port": listener.getsockname()[1], + "pid": os.getpid(), + "token": "serving", + } + monkeypatch.setattr(effect_runtime, "_runtime_fingerprint", lambda: "fixture") + monkeypatch.setattr( + effect_runtime, "_runtime_info_path", lambda _: tmp_path / "runtime.json" + ) + monkeypatch.setattr(effect_runtime, "_read_info", lambda *_a, **_k: info) + monkeypatch.setattr( + effect_runtime, "_serving_token", lambda _: (True, "replacement") + ) + + def serve_one() -> str: + with listener.accept()[0] as connection: + request = b"" + while b"\n" not in request: + request += connection.recv(4096) + if lost_response == "timeout": + time.sleep(0.1) + return json.loads(request)["method"] + + with ThreadPoolExecutor(max_workers=1) as executor: + received = executor.submit(serve_one) + result = effect_runtime.restart_effect_runtime(timeout=0.02) + assert received.result(timeout=2) == "runtime.shutdown" + assert result["status"] == "stopped" + assert result["stopped"] is True + + def test_restart_stops_the_runtime_that_carries_the_serving_identity( tmp_path: Path, monkeypatch,