From 2566a3871ebe15024c2f045c0a321c61bb4782e3 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Fri, 25 Sep 2026 22:47:57 +0800 Subject: [PATCH 1/2] Bound canonical authority RPC waits and large file reads Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- loopx/cli_commands/todo_continuation.py | 12 ++- .../coordination/file_authority_store.ts | 92 +++++++++++++------ .../coordination/local_authority.py | 6 +- loopx/control_plane/effect_runtime.py | 6 ++ .../scheduler/provider_monitor_poll.py | 8 +- loopx/control_plane/todos/provider_create.py | 7 +- .../todos/provider_handoff_mode.py | 7 +- .../todos/provider_terminal_lifecycle.py | 17 ++-- loopx/control_plane/todos/provider_update.py | 15 ++- .../work_items/task_lease_acquire_adapter.py | 35 +++++-- .../test_shadow_native_todo_update_e2e.py | 59 +++++++++++- .../control_plane_ts/authority_store.test.ts | 32 +++++++ 12 files changed, 240 insertions(+), 56 deletions(-) diff --git a/loopx/cli_commands/todo_continuation.py b/loopx/cli_commands/todo_continuation.py index 3146d9417f..6591317dfa 100644 --- a/loopx/cli_commands/todo_continuation.py +++ b/loopx/cli_commands/todo_continuation.py @@ -5,7 +5,11 @@ from ..agent_registry import registered_agent_ids_from_registry from ..history import load_registry from ..paths import resolve_runtime_root -from ..control_plane.effect_runtime import effect_runtime_result +from ..control_plane.effect_runtime import ( + CANONICAL_AUTHORITY_READ_TIMEOUT_SECONDS, + CANONICAL_AUTHORITY_WRITE_TIMEOUT_SECONDS, + effect_runtime_result, +) def _render_digest(payload: dict) -> str: @@ -166,7 +170,11 @@ def handle_todo_continuation(args, *, registry_path, runtime_root_arg, output_fo version = args.task_lease_expected_version if key is not None or version is not None: request["lease_proof"] = {"idempotency_key": key, "expected_version": version} - payload = effect_runtime_result("coordination.local_authority.todo_continuation", request) + payload = effect_runtime_result( + "coordination.local_authority.todo_continuation", request, + timeout=(CANONICAL_AUTHORITY_READ_TIMEOUT_SECONDS if args.action == "inspect" + else CANONICAL_AUTHORITY_WRITE_TIMEOUT_SECONDS), + ) except (ValueError, RuntimeError, OSError, json.JSONDecodeError) as exc: payload = {"ok": False, "status": "failed", "reason_code": "invalid_continuation_request", "reason": str(exc)} diff --git a/loopx/control_plane/coordination/file_authority_store.ts b/loopx/control_plane/coordination/file_authority_store.ts index 969d24524e..d781c04ed3 100644 --- a/loopx/control_plane/coordination/file_authority_store.ts +++ b/loopx/control_plane/coordination/file_authority_store.ts @@ -31,27 +31,56 @@ 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. +// opens a new store handle for each request. Keep one verified read view across +// handles, keyed by exact bytes and store identity. Large journals retain only +// the head and receipt index in memory; commits and scans still load and verify +// the complete history. This is a bounded read optimization, not a new source +// of authority or a substitute for the SQLite long-goal profile. const MAX_CACHED_DOCUMENT_BYTES = 128 * 1024 * 1024; -let verifiedDocument: { +const MAX_CACHED_READ_VIEW_BYTES = 16 * 1024 * 1024; +interface VerifiedDocument { path: string; identity: string; digest: string; - document: FileAuthorityStoreDocument; -} | null = null; + head: JsonObject; + providerRevision: string; + cursor: string; + receipts: ReadonlyMap; + document?: FileAuthorityStoreDocument; +} +let verifiedDocument: VerifiedDocument | 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; + digest: string, document: FileAuthorityStoreDocument, + maxDocumentBytes: number): VerifiedDocument { + const receipts = new Map(); + let viewBytes = Buffer.byteLength(JSON.stringify(document.head)); + for (const entry of document.committed) { + receipts.set(entry.operation_id, {cursor: entry.cursor, + providerRevision: entry.provider_revision, receipts: entry.receipts}); + viewBytes += Buffer.byteLength(entry.operation_id) + Buffer.byteLength(entry.cursor) + + Buffer.byteLength(entry.provider_revision) + Buffer.byteLength(JSON.stringify(entry.receipts)); + } + const view: VerifiedDocument = {path, identity, digest, head: document.head, + providerRevision: document.provider_revision, cursor: document.cursor, + receipts, document}; + verifiedDocument = raw.byteLength <= maxDocumentBytes ? view + : viewBytes <= MAX_CACHED_READ_VIEW_BYTES + ? {...view, document: undefined} + : null; + return view; } interface FileAuthorityStoreDocument extends JsonObject, RetainedAuthorityJournal { @@ -195,6 +224,11 @@ export class FileAuthorityStore implements AuthorityStore { return decodeDocument(value, this.goalId, identity); } + /** Test seam for the compact-view path; production keeps the fixed memory cap. */ + protected fullDocumentCacheLimitBytes(): number { + return MAX_CACHED_DOCUMENT_BYTES; + } + private async readStoreIdentity(createIfMissing = !this.existingOnly): Promise { try { const identity = await readFile(this.identityPath, "utf8"); @@ -241,7 +275,7 @@ export class FileAuthorityStore implements AuthorityStore { } } - private async readDocument(knownIdentity?: string): Promise { + private async readVerified(knownIdentity?: string, requireHistory = false): Promise { let raw: Buffer; try { raw = await readFile(this.path); @@ -255,12 +289,13 @@ export class FileAuthorityStore implements AuthorityStore { try { const digest = documentDigest(raw); if (verifiedDocument?.path === this.path && - verifiedDocument.identity === identity && verifiedDocument.digest === digest) { - return verifiedDocument.document; + verifiedDocument.identity === identity && verifiedDocument.digest === digest && + (!requireHistory || verifiedDocument.document !== undefined)) { + return verifiedDocument; } const document = this.decodeStoredDocument(JSON.parse(raw.toString("utf8")), identity); - rememberVerifiedDocument(this.path, identity, raw, digest, document); - return document; + return rememberVerifiedDocument(this.path, identity, raw, digest, document, + this.fullDocumentCacheLimitBytes()); } catch (error) { if (error instanceof SyntaxError) { throw new AuthorityStoreProtocolError(`file authority store JSON is invalid: ${error.message}`); @@ -269,15 +304,19 @@ export class FileAuthorityStore implements AuthorityStore { } } + private async readDocument(knownIdentity?: string): Promise { + return (await this.readVerified(knownIdentity, true))?.document ?? null; + } + async loadAuthority(): Promise { try { - const document = await this.readDocument(); - return document + const verified = await this.readVerified(); + return verified ? { status: "loaded", - head: structuredClone(document.head), - provider_revision: document.provider_revision, - cursor: document.cursor, + head: structuredClone(verified.head), + provider_revision: verified.providerRevision, + cursor: verified.cursor, } : { status: "missing" }; } catch (error) { @@ -356,7 +395,8 @@ export class FileAuthorityStore implements AuthorityStore { try { const bytes = canonicalAuthorityBytes(document); await this.replaceDurably(this.path, bytes); - rememberVerifiedDocument(this.path, identity, bytes, documentDigest(bytes), document); + rememberVerifiedDocument(this.path, identity, bytes, documentDigest(bytes), document, + this.fullDocumentCacheLimitBytes()); } 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 @@ -393,15 +433,13 @@ export class FileAuthorityStore implements AuthorityStore { }; } try { - const transaction = (await this.readDocument())?.committed.find( - (entry) => entry.operation_id === normalized, - ); + const transaction = (await this.readVerified())?.receipts.get(normalized); return transaction ? { status: "found", cursor: transaction.cursor, - provider_revision: transaction.provider_revision, - receipts: structuredClone(transaction.receipts), + provider_revision: transaction.providerRevision, + receipts: structuredClone([...transaction.receipts]), } : { status: "missing" }; } catch (error) { diff --git a/loopx/control_plane/coordination/local_authority.py b/loopx/control_plane/coordination/local_authority.py index 6eb4197e7e..3d6c6587d6 100644 --- a/loopx/control_plane/coordination/local_authority.py +++ b/loopx/control_plane/coordination/local_authority.py @@ -15,7 +15,10 @@ from ...agent_registry import registered_agent_ids_from_registry from .authority_source_capture import authority_registry_source from ..runtime.time import now_local_iso as now_local -from ..effect_runtime import effect_runtime_result +from ..effect_runtime import ( + CANONICAL_AUTHORITY_WRITE_TIMEOUT_SECONDS, + effect_runtime_result, +) from .coordination_state_contract import ( TODO_CANONICAL_READ_RECORD_SCHEMA_VERSION, TODO_DOMAIN_READ_RECORD_SCHEMA_VERSION, @@ -133,6 +136,7 @@ def claim_canonical_todo_if_promoted( "observed_at": now_local(), "dry_run": dry_run, }, + timeout=CANONICAL_AUTHORITY_WRITE_TIMEOUT_SECONDS, ) if not isinstance(result, Mapping): raise LocalCoordinationAuthorityUnavailable( diff --git a/loopx/control_plane/effect_runtime.py b/loopx/control_plane/effect_runtime.py index ab727111b4..59b76efc84 100644 --- a/loopx/control_plane/effect_runtime.py +++ b/loopx/control_plane/effect_runtime.py @@ -38,6 +38,12 @@ STARTUP_READY_TIMEOUT_SECONDS = 15.0 STARTUP_POLL_SECONDS = 0.025 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 +# that declared critical section and a bounded readback; a shorter RPC budget +# turns an in-flight write into an avoidable ambiguous response. +CANONICAL_AUTHORITY_WRITE_TIMEOUT_SECONDS = 45.0 +CANONICAL_AUTHORITY_READ_TIMEOUT_SECONDS = 15.0 _NODE_VERSION_RE = re.compile(r"^v?(\d+)\.(\d+)\.(\d+)(?:[-+].*)?$") _RUNTIME_SOURCE_SUFFIXES = frozenset({".json", ".ts"}) _RuntimeSourceSnapshot = tuple[tuple[str, int, int, int], ...] diff --git a/loopx/control_plane/scheduler/provider_monitor_poll.py b/loopx/control_plane/scheduler/provider_monitor_poll.py index 68b159bd87..a392de9e9a 100644 --- a/loopx/control_plane/scheduler/provider_monitor_poll.py +++ b/loopx/control_plane/scheduler/provider_monitor_poll.py @@ -13,11 +13,13 @@ local_authority_is_promoted, read_canonical_todos_if_promoted, ) -from ..effect_runtime import effect_runtime_result +from ..effect_runtime import ( + CANONICAL_AUTHORITY_WRITE_TIMEOUT_SECONDS, + effect_runtime_result, +) 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: @@ -44,7 +46,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) + }, timeout=CANONICAL_AUTHORITY_WRITE_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_create.py b/loopx/control_plane/todos/provider_create.py index e4c4c383c3..192370f6d9 100644 --- a/loopx/control_plane/todos/provider_create.py +++ b/loopx/control_plane/todos/provider_create.py @@ -16,7 +16,11 @@ LocalCoordinationAuthorityUnavailable, read_canonical_todos_if_promoted, ) -from ..effect_runtime import EffectRuntimeResponseAmbiguous, effect_runtime_result +from ..effect_runtime import ( + CANONICAL_AUTHORITY_WRITE_TIMEOUT_SECONDS, + EffectRuntimeResponseAmbiguous, + effect_runtime_result, +) from .contract import ( normalize_todo_metadata_for_write, normalize_todo_task_class, @@ -107,6 +111,7 @@ def create_canonical_todo_if_promoted( "dry_run": dry_run, "observed_at": now_local(), }, + timeout=CANONICAL_AUTHORITY_WRITE_TIMEOUT_SECONDS, ) except EffectRuntimeResponseAmbiguous as error: raise LocalCoordinationAuthorityUnavailable( diff --git a/loopx/control_plane/todos/provider_handoff_mode.py b/loopx/control_plane/todos/provider_handoff_mode.py index 1ae1c2b0fa..7328ae282f 100644 --- a/loopx/control_plane/todos/provider_handoff_mode.py +++ b/loopx/control_plane/todos/provider_handoff_mode.py @@ -8,7 +8,10 @@ LocalCoordinationAuthorityUnavailable, local_authority_is_promoted, read_canonical_todos_if_promoted, ) -from ..effect_runtime import effect_runtime_result +from ..effect_runtime import ( + CANONICAL_AUTHORITY_WRITE_TIMEOUT_SECONDS, + effect_runtime_result, +) from ..runtime.time import now_local_iso @@ -29,7 +32,7 @@ def set_canonical_handoff_mode(*, runtime_root: Path, goal_id: str, mode: str, "runtime_root": str(runtime_root.expanduser().resolve()), "goal_id": goal_id, "operation_id": (operation_id if operation_id is not None else f"handoff-mode:{goal_id}:{uuid4().hex}"), "requested_mode": mode, "observed_at": now_local_iso(), "dry_run": dry_run, - }) + }, timeout=CANONICAL_AUTHORITY_WRITE_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 or result.get("decision_read_from_provider") is not True or result.get("legacy_fallback_used") is not False): diff --git a/loopx/control_plane/todos/provider_terminal_lifecycle.py b/loopx/control_plane/todos/provider_terminal_lifecycle.py index cba643db08..502fdadb9c 100644 --- a/loopx/control_plane/todos/provider_terminal_lifecycle.py +++ b/loopx/control_plane/todos/provider_terminal_lifecycle.py @@ -22,7 +22,10 @@ read_canonical_todos_if_promoted, ) from ..coordination.local_authority_shadow_adapter import effective_runtime_root -from ..effect_runtime import effect_runtime_result +from ..effect_runtime import ( + CANONICAL_AUTHORITY_WRITE_TIMEOUT_SECONDS, + effect_runtime_result, +) from .completion_policy import ( build_completion_policy_request, linked_successor_from_todo, @@ -45,10 +48,6 @@ _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]] @@ -414,7 +413,7 @@ def terminal_canonical_todo_if_promoted( } result = effect_runtime_result( "coordination.local_authority.todo_terminal", request, - timeout=_TERMINAL_RUNTIME_TIMEOUT_SECONDS, + timeout=CANONICAL_AUTHORITY_WRITE_TIMEOUT_SECONDS, ) if isinstance(result, Mapping) and result.get("status") == "resolve_validation": # Admission and receipt recovery precede host-local declaration IO. @@ -430,7 +429,7 @@ def terminal_canonical_todo_if_promoted( ) result = effect_runtime_result( "coordination.local_authority.todo_terminal", request, - timeout=_TERMINAL_RUNTIME_TIMEOUT_SECONDS, + timeout=CANONICAL_AUTHORITY_WRITE_TIMEOUT_SECONDS, ) completion_validation_executed = False if isinstance(result, Mapping) and result.get("status") == "execute_validation": @@ -458,7 +457,7 @@ def terminal_canonical_todo_if_promoted( request["observed_at"] = now_local() result = effect_runtime_result( "coordination.local_authority.todo_terminal", request, - timeout=_TERMINAL_RUNTIME_TIMEOUT_SECONDS, + timeout=CANONICAL_AUTHORITY_WRITE_TIMEOUT_SECONDS, ) if not isinstance(result, Mapping): raise LocalCoordinationAuthorityUnavailable( @@ -574,6 +573,7 @@ def archive_canonical_todos_if_promoted( "dry_run": dry_run, "observed_at": now_local(), }, + timeout=CANONICAL_AUTHORITY_WRITE_TIMEOUT_SECONDS, ) if not isinstance(result, Mapping) or result.get("status") not in _ACCEPTED: payload = dict(result) if isinstance(result, Mapping) else {} @@ -612,6 +612,7 @@ def archive_canonical_todos_if_promoted( "role": role, "operation_id": response.get("operation_id"), }, + timeout=CANONICAL_AUTHORITY_WRITE_TIMEOUT_SECONDS, ) response["archive_delivery_ack"] = ( dict(acknowledgement) diff --git a/loopx/control_plane/todos/provider_update.py b/loopx/control_plane/todos/provider_update.py index 78b3b9d5a6..b136513fbb 100644 --- a/loopx/control_plane/todos/provider_update.py +++ b/loopx/control_plane/todos/provider_update.py @@ -18,7 +18,10 @@ local_authority_is_promoted, read_canonical_todos_if_promoted, ) -from ..effect_runtime import effect_runtime_result +from ..effect_runtime import ( + CANONICAL_AUTHORITY_WRITE_TIMEOUT_SECONDS, + effect_runtime_result, +) from .contract import compact_todo_text from .provider_projection import settle_canonical_todo_projection from .text import normalize_new_todo @@ -221,7 +224,10 @@ def update_canonical_todo_if_promoted( prepare_completion_validation_declaration( runtime_root=runtime_root, goal_id=goal_id, declaration=completion_validation_revision ) - result = effect_runtime_result("coordination.local_authority.todo_update", request) + result = effect_runtime_result( + "coordination.local_authority.todo_update", request, + timeout=CANONICAL_AUTHORITY_WRITE_TIMEOUT_SECONDS, + ) completion_validation_executed = False if isinstance(result, dict) and result.get("status") == "execute_validation": if completion is None: @@ -231,7 +237,10 @@ def update_canonical_todo_if_promoted( result, registry_path=registry_path, goal_id=goal_id)) completion_validation_executed = True request["observed_at"] = now_local() - result = effect_runtime_result("coordination.local_authority.todo_update", request) + result = effect_runtime_result( + "coordination.local_authority.todo_update", request, + timeout=CANONICAL_AUTHORITY_WRITE_TIMEOUT_SECONDS, + ) if isinstance(result, dict): failure = completion_validation_failure(result, goal_id=goal_id, todo_id=todo_id, dry_run=dry_run) if failure is not None: diff --git a/loopx/control_plane/work_items/task_lease_acquire_adapter.py b/loopx/control_plane/work_items/task_lease_acquire_adapter.py index 162fd1e0ff..7b577aeb8b 100644 --- a/loopx/control_plane/work_items/task_lease_acquire_adapter.py +++ b/loopx/control_plane/work_items/task_lease_acquire_adapter.py @@ -371,7 +371,10 @@ def execute_native_task_lease_acquire( ) -> dict[str, Any]: """Transport one compact acquire request to the native TypeScript owner.""" - from ..effect_runtime import effect_runtime_result + from ..effect_runtime import ( + CANONICAL_AUTHORITY_WRITE_TIMEOUT_SECONDS, + effect_runtime_result, + ) from ..coordination.local_authority import local_authority_is_promoted @@ -402,7 +405,10 @@ def execute_native_task_lease_acquire( "provider": "file_v0", } payload = _require_native_acquire_shape( - effect_runtime_result("task_lease.acquire.native", request, timeout=15.0) + effect_runtime_result( + "task_lease.acquire.native", request, + timeout=(CANONICAL_AUTHORITY_WRITE_TIMEOUT_SECONDS if canonical else 15.0), + ) ) if ( payload.get("error_code") == "authority_source_changed" @@ -724,12 +730,16 @@ def execute_native_task_lease_lifecycle( compacted_todo = _compact_lifecycle_todo(todo, todo_id=str(todo_id)) if compacted_todo is not None and not canonical_lifecycle: request["todo"] = compacted_todo - from ..effect_runtime import effect_runtime_result + from ..effect_runtime import ( + CANONICAL_AUTHORITY_WRITE_TIMEOUT_SECONDS, + effect_runtime_result, + ) payload = effect_runtime_result( "task_lease.lifecycle.native", request, - timeout=15.0, + timeout=(CANONICAL_AUTHORITY_WRITE_TIMEOUT_SECONDS + if canonical_lifecycle else 15.0), ) if not isinstance(payload, dict): raise RuntimeError("native task-lease lifecycle result shape mismatch") @@ -780,7 +790,11 @@ def inspect_native_task_lease( LOCAL_AUTHORITY_SOURCES, LocalCoordinationAuthorityUnavailable, local_authority_is_promoted, ) - from ..effect_runtime import effect_runtime_result + from ..effect_runtime import ( + CANONICAL_AUTHORITY_READ_TIMEOUT_SECONDS, + DEFAULT_REQUEST_TIMEOUT_SECONDS, + effect_runtime_result, + ) for attempt in range(TASK_LEASE_AUTHORITY_SNAPSHOT_ATTEMPTS): canonical = local_authority_is_promoted(runtime_root=runtime_root, goal_id=goal_id) @@ -796,7 +810,11 @@ def inspect_native_task_lease( "runtime_root": str(runtime_root.resolve()), "goal_id": goal_id, "todo_id": todo_id, "authority": authority, } - result = effect_runtime_result("task_lease.inspect.native", request) + result = effect_runtime_result( + "task_lease.inspect.native", request, + timeout=(CANONICAL_AUTHORITY_READ_TIMEOUT_SECONDS + if canonical else DEFAULT_REQUEST_TIMEOUT_SECONDS), + ) if isinstance(result, dict) and result.get("todo_projection_required") is True: if canonical or result.get("schema_version") != TASK_LEASE_SCHEMA_VERSION or result.get("ok") is not True or result.get("action") != "inspect": raise RuntimeError("native lease inspection requested an invalid source projection") @@ -806,7 +824,10 @@ def inspect_native_task_lease( ) # Re-read the lease and fence: neither the record nor its expiry is # assumed unchanged while Python prepares the Todo projection. - result = effect_runtime_result("task_lease.inspect.native", request) + result = effect_runtime_result( + "task_lease.inspect.native", request, + timeout=DEFAULT_REQUEST_TIMEOUT_SECONDS, + ) if not isinstance(result, dict) or result.get("schema_version") != TASK_LEASE_SCHEMA_VERSION or result.get("action") != "inspect" or not isinstance(result.get("ok"), bool): raise RuntimeError("native lease inspection result shape mismatch") if result.get("error_code") == "authority_source_changed" and attempt + 1 < TASK_LEASE_AUTHORITY_SNAPSHOT_ATTEMPTS: diff --git a/tests/control_plane/test_shadow_native_todo_update_e2e.py b/tests/control_plane/test_shadow_native_todo_update_e2e.py index f21ce584d8..bad3f1b8c3 100644 --- a/tests/control_plane/test_shadow_native_todo_update_e2e.py +++ b/tests/control_plane/test_shadow_native_todo_update_e2e.py @@ -14,6 +14,7 @@ import select import subprocess import sys +import time import pytest @@ -30,6 +31,7 @@ shadow_management_state_path, ) from loopx.control_plane.effect_runtime import effect_runtime_result +from loopx.control_plane.todos import provider_update REPO = Path(__file__).resolve().parents[2] GOAL, TODO = "goal-update", "todo_update_probe" @@ -88,7 +90,7 @@ def workspace(tmp_path: Path, request: pytest.FixtureRequest) -> Workspace: import fs from 'node:fs'; import {syncBuiltinESMExports} from 'node:module'; import {once} from 'node:events'; -import {join} from 'node:path'; +import {dirname, join} from 'node:path'; const input = JSON.parse(process.argv[1]); const base = new URL(input.module_base); const {FileAuthorityStore} = await import(new URL('file_authority_store.ts', base)); @@ -109,7 +111,16 @@ def workspace(tmp_path: Path, request: pytest.FixtureRequest) -> Workspace: syncBuiltinESMExports(); } let result; -if (input.mode === 'inspect') { +if (input.mode === 'hold_writer_lock') { + const {withFileMutationLock} = await import(new URL('../effect_runtime_io.ts', base)); + const lockPath = management.shadowMaintenanceLockPath(input.request.runtime_root, input.request.goal_id); + await fs.promises.mkdir(dirname(lockPath), {recursive:true}); + result = await withFileMutationLock(lockPath, async () => { + process.stdout.write('BARRIER lock-held\n'); + await new Promise(resolve => setTimeout(resolve, 14_000)); + return {status:'released'}; + }); +} else if (input.mode === 'inspect') { result = {head:await store.loadAuthority(), receipt:input.operation_id ? await store.readReceipt(input.operation_id) : null, document_path:store.path}; } else if (input.mode.startsWith('update')) { @@ -228,6 +239,50 @@ def test_native_update_preserves_complete_records_and_exact_receipts(workspace: assert inspect(workspace)["head"] == after["head"] +@pytest.mark.parametrize("workspace", ["native"], indirect=True) +def test_canonical_update_rpc_budget_covers_the_writer_lock_and_receipt( + workspace: Workspace, monkeypatch: pytest.MonkeyPatch, +) -> None: + observed: list[float] = [] + original = provider_update.effect_runtime_result + + def timed_call(method: str, request: dict, **kwargs: object) -> object: + assert method == "coordination.local_authority.todo_update" + observed.append(float(kwargs["timeout"])) + return original(method, request, **kwargs) + + monkeypatch.setattr(provider_update, "effect_runtime_result", timed_call) + result = provider_update.update_canonical_todo_if_promoted( + registry_path=workspace.registry, runtime_root=workspace.runtime, + goal_id=GOAL, todo_id=TODO, actor_agent_id="agent-a", role="agent", + text="Updated with the bounded canonical writer budget", note=None, + dry_run=False, operation_id="budgeted-update", + ) + assert result is not None and result["status"] == "applied", result + assert observed and all(timeout > 35 for timeout in observed) + receipt = inspect(workspace, "budgeted-update")["receipt"] + assert receipt["status"] == "found" + + +@pytest.mark.parametrize("workspace", ["native"], indirect=True) +def test_cli_update_waits_past_default_rpc_budget_for_real_writer_lock(workspace: Workspace) -> None: + holder = start(node_command("hold_writer_lock", workspace.request())) + try: + expect_barrier(holder, "lock-held") + started = time.monotonic() + result = subprocess.run(workspace.command(), cwd=REPO, + capture_output=True, text=True, timeout=30) + elapsed = time.monotonic() - started + assert result.returncode == 0, result.stdout + result.stderr + payload = json.loads(result.stdout) + assert elapsed > 10, elapsed + assert payload["status"] == "applied", payload + receipt = payload["original_receipt"] + assert inspect(workspace, receipt["operation_id"])["receipt"]["status"] == "found" + finally: + stop(holder) + + @pytest.mark.parametrize("transport", ["cli", "rpc"]) @pytest.mark.parametrize("management", ["corrupt", "bootstrapping", "rolling_back"]) def test_native_update_holds_before_primary_for_management(workspace: Workspace, transport: str, management: str) -> None: diff --git a/tests/control_plane_ts/authority_store.test.ts b/tests/control_plane_ts/authority_store.test.ts index 0bb641e29f..694b0558e0 100644 --- a/tests/control_plane_ts/authority_store.test.ts +++ b/tests/control_plane_ts/authority_store.test.ts @@ -143,6 +143,38 @@ test("file verification is reused only for exact bytes and store identity", asyn assert.equal(CountingStore.validations, 3); }); +test("large file read view reuses verified head and receipts without retaining history", async (t) => { + const {root, store} = await fixture(t); + assert.equal((await store.commitAuthority(commit(null, "operation-a", 1, 1))).status, "applied"); + const original = await readFile(store.path, "utf8"); + await writeFile(store.path, `${original}\n`, "utf8"); + + class CompactStore extends FileAuthorityStore { + static validations = 0; + protected override fullDocumentCacheLimitBytes() { return 1; } + protected override decodeStoredDocument(value: unknown, identity: string) { + CompactStore.validations += 1; + return super.decodeStoredDocument(value, identity); + } + } + const first = new CompactStore(root, "goal-a"); + const second = new CompactStore(root, "goal-a"); + const head = await first.loadAuthority(); + assert.equal(head.status, "loaded"); + assert.deepEqual(await second.loadAuthority(), head); + assert.equal((await second.readReceipt("operation-a")).status, "found"); + assert.equal(CompactStore.validations, 1); + assert.equal((await second.scanCommitted(null, 1)).status, "page"); + assert.equal(CompactStore.validations, 2, "history scans still verify the complete journal"); + + const changed = JSON.parse(original); + changed.committed[0].projection.authority_revision = 99; + changed.head.authority_revision = 99; + await writeFile(store.path, JSON.stringify(changed), "utf8"); + assert.equal((await second.loadAuthority()).status, "failed"); + assert.equal((await second.readReceipt("operation-a")).status, "failed"); +}); + 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 df9d5fcd0f1de4fbed473d0d549f1be6ea1783a4 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Fri, 25 Sep 2026 22:48:12 +0800 Subject: [PATCH 2/2] Record file authority timeout boundary in shared RFC Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../rfcs/shared-goal-authority-state-provider-v0.md | 10 ++++++++++ .../shared-goal-authority-state-provider-v0.zh-CN.md | 8 ++++++++ 2 files changed, 18 insertions(+) diff --git a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md index 3276631f87..b6db7a1371 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md @@ -55,6 +55,16 @@ no durable authority and relaxes no writer fence or D2 gate. This runtime repair does not mechanically subtract one of the three remaining implementation packages. [Evidence and boundary](ledger/shared-goal-authority-state-provider-v0/2026-09-24-default-cutover-reconciliation.md#long-history-closeout-this-repair-and-its-remaining-boundary). +For existing file-v0 Goals with a journal above the 128 MiB full-document +cache limit, the Effect server now reuses a bounded, byte-and-identity-verified +head and receipt read view. Commits and history scans still verify the complete +journal. Canonical write callers allow the declared 30-second maintenance-lock +wait, 5-second provider-lock wait and a bounded readback before declaring an +ambiguous RPC; read-only lease inspection has its own budget. This is an +interim L2/L5 reliability repair for existing Goals, not bounded file-v0 write +amplification, SQLite D2 qualification or authorization to migrate a live Goal. +Section 7.2's capacity, recovery, soak and fenced migration gates remain. + Promotion admission now binds complete sources to a current registry witness and rechecks it inside the TS lock scope. Saved execution retains the reviewed handoff policy, and failures report durable fence presence. Recovery of a diff --git a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md index a51de6138d..167b5f17df 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md @@ -47,6 +47,14 @@ writer fence/D2。剩余三个实现边界不因该运行缺陷修复而机械减一。 [证据与边界](ledger/shared-goal-authority-state-provider-v0/2026-09-24-default-cutover-reconciliation.zh-CN.md#长历史收尾检查本次修复与剩余边界)。 +对 journal 超过 128 MiB 完整文档缓存上限的既有 file-v0 Goal,Effect server +现在复用按原始字节和 store identity 校验、容量有界的 head 与历史回执读取视图; +提交及历史扫描仍校验完整 journal。canonical 写入 caller 的 RPC 预算覆盖既有 +30 秒 maintenance lock 等待、5 秒 provider lock 等待和有界读回,避免过早报告 +结果不明;只读的 lease inspect 有独立预算。这只是面向既有 Goal 的 L2/L5 +可靠性修复,不解决 file-v0 写入放大,不代表 SQLite D2 资格通过,也不授权迁移 +活跃 Goal。第 7.2 节的容量、恢复、soak 和 fenced migration 门槛仍然有效。 + 晋升准入现将完整来源绑定到当前 registry witness,并在 TS 持锁范围内重新校验; 保存计划执行保留已审核的 handoff 策略,失败结果如实报告持久 fence。 已提交事务的恢复仍按原 fence/receipt,不要求失去权威的旧来源重新有效。