Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,不要求失去权威的旧来源重新有效。
Expand Down
12 changes: 10 additions & 2 deletions loopx/cli_commands/todo_continuation.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -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)}
Expand Down
92 changes: 65 additions & 27 deletions loopx/control_plane/coordination/file_authority_store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, {
cursor: string;
providerRevision: string;
receipts: readonly JsonObject[];
}>;
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<string, {
cursor: string;
providerRevision: string;
receipts: readonly JsonObject[];
}>();
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 {
Expand Down Expand Up @@ -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<string> {
try {
const identity = await readFile(this.identityPath, "utf8");
Expand Down Expand Up @@ -241,7 +275,7 @@ export class FileAuthorityStore implements AuthorityStore {
}
}

private async readDocument(knownIdentity?: string): Promise<FileAuthorityStoreDocument | null> {
private async readVerified(knownIdentity?: string, requireHistory = false): Promise<VerifiedDocument | null> {
let raw: Buffer;
try {
raw = await readFile(this.path);
Expand All @@ -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}`);
Expand All @@ -269,15 +304,19 @@ export class FileAuthorityStore implements AuthorityStore {
}
}

private async readDocument(knownIdentity?: string): Promise<FileAuthorityStoreDocument | null> {
return (await this.readVerified(knownIdentity, true))?.document ?? null;
}

async loadAuthority(): Promise<AuthorityStoreLoadResult> {
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) {
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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) {
Expand Down
6 changes: 5 additions & 1 deletion loopx/control_plane/coordination/local_authority.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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(
Expand Down
6 changes: 6 additions & 0 deletions loopx/control_plane/effect_runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -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], ...]
Expand Down
8 changes: 5 additions & 3 deletions loopx/control_plane/scheduler/provider_monitor_poll.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand All @@ -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
Expand Down
7 changes: 6 additions & 1 deletion loopx/control_plane/todos/provider_create.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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(
Expand Down
7 changes: 5 additions & 2 deletions loopx/control_plane/todos/provider_handoff_mode.py
Original file line number Diff line number Diff line change
Expand Up @@ -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


Expand All @@ -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):
Expand Down
17 changes: 9 additions & 8 deletions loopx/control_plane/todos/provider_terminal_lifecycle.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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]]
Expand Down Expand Up @@ -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.
Expand All @@ -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":
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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 {}
Expand Down Expand Up @@ -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)
Expand Down
Loading
Loading