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
23 changes: 23 additions & 0 deletions docs/architecture/rfcs/typescript-control-plane-migration-v0.md
Original file line number Diff line number Diff line change
Expand Up @@ -1190,6 +1190,29 @@ qualify distributed execution. Cursor/checkpoint reduction remains a measured
follow-up with complete-source parity, not a second Python policy. See the
[history decision evidence](ledger/typescript-control-plane-migration-v0/2026-09-22-replan-history-policy.md).

**Checkpoint read-context transport.** A long-lived Goal can exceed the Effect
runtime's 2 MiB socket limit even when the original Turn receipt is small:
complete shared Goal prose and archived Todo facts are part of the read basis,
and the basis may also exceed the response limit. Checkpoint source, evaluation,
replay inspection and commit explicitly opt into same-UID private request and
response files with byte counts and SHA-256 digests. The 2 MiB socket boundary
and default behavior of other effects remain intact. File/SQLite and legacy
Markdown still use the same TypeScript checkpoint reducer and exact receipt;
the transport neither truncates history nor creates a Python decision or new
authority. Missing, changed, non-private or over-64-MiB files fail closed.
After a handler may have committed, an unverifiable response stays ambiguous
and requires exact receipt readback, never automatic mutation retry.

This removes the immediate transport ceiling, not the cost of projecting a
complete multi-megabyte basis. The next measured T3 cut should combine the
canonical source read and checkpoint reduction inside one TypeScript call, then
offer a versioned manifest with bounded pages for human/Agent inspection.
Every page must bind to the same source head and disclose omitted components;
the receipt must still hash the complete relevant Todo/dependency, User Todo,
Goal prose, acceptance and vision basis. A display limit must never become a
settlement limit. Retain the current complete read until that parity and stale-
head recovery are qualified on legacy, File and SQLite backends.

**Recovery boundary (2026-09-22).** The
[authority archive command](../../reference/authority-archive.md) places retained
history validation, delta reconstruction and resumable restore in the existing
Expand Down
6 changes: 6 additions & 0 deletions docs/reference/protocols/goal-vision-replan-contract-v0.md
Original file line number Diff line number Diff line change
Expand Up @@ -306,6 +306,12 @@ The basis covers the selected Todo, its dependency closure and recorded results,
shared Goal prose and User Todos, the owner acceptance document/revision when
configured, the current agent vision, and the local source binding. A replan
obligation covers the full Todo frontier. Archived dependencies remain inputs.
Large local bases use digest-checked private files across the Python/TypeScript
runtime boundary, including the response; the CLI still returns the complete
basis. The 2 MiB default RPC guard remains for other effects. File size is
bounded and an unverifiable response after a possible commit is ambiguous,
so the caller reads the exact receipt before retrying any mutation. Neither
transport nor a future paged presentation may silently omit a basis component.
Todo display positions, source headings, and the Goal's global `updated_at` are
excluded; an unrelated Agent Todo or run-history append does not invalidate an
otherwise unchanged Todo-bound basis. Shared prose is deliberately conservative:
Expand Down
145 changes: 113 additions & 32 deletions loopx/control_plane/effect_runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@
import time
import uuid
from collections.abc import Iterator, Mapping
from contextlib import contextmanager
from contextlib import ExitStack, contextmanager
from contextvars import ContextVar
from dataclasses import dataclass
from functools import lru_cache
Expand All @@ -33,6 +33,13 @@
MINIMUM_NODE_VERSION_TEXT = ".".join(str(part) for part in MINIMUM_NODE_VERSION)
MAX_RESPONSE_BYTES = 2 * 1024 * 1024
MAX_REQUEST_BYTES = 2 * 1024 * 1024
MAX_LOCAL_SNAPSHOT_BYTES = 64 * 1024 * 1024
LOCAL_SNAPSHOT_METHODS = frozenset({
"goal.checkpoint_read_context.source",
"goal.checkpoint_read_context.evaluate",
"goal.checkpoint_read_context.commit",
"goal.checkpoint_read_context.inspect_replay",
})
MAX_STARTUP_DIAGNOSTIC_BYTES = 8 * 1024
STARTUP_LOCK_TIMEOUT_SECONDS = 15.0
STARTUP_READY_TIMEOUT_SECONDS = 15.0
Expand Down Expand Up @@ -548,6 +555,7 @@ def _request_with_info(
method: str,
params: Mapping[str, Any],
timeout: float,
large_local_snapshot: bool = False,
) -> dict[str, Any]:
request = {
"schema_version": EFFECT_RUNTIME_REQUEST_SCHEMA_VERSION,
Expand All @@ -556,38 +564,74 @@ def _request_with_info(
"method": method,
"params": dict(params),
}
encoded = (json.dumps(request, separators=(",", ":")) + "\n").encode()
if len(encoded) > MAX_REQUEST_BYTES:
raise EffectRuntimeRejected(
"TypeScript Effect runtime request is oversized",
diagnostic_code="request_too_large",
)
chunks: list[bytes] = []
size = 0
with socket.create_connection(
(str(info["host"]), int(info["port"])), timeout=timeout
) as connection:
with ExitStack() as stack:
response_sink: Path | None = None
encoded = (json.dumps(request, separators=(",", ":")) + "\n").encode()
if large_local_snapshot:
if method not in LOCAL_SNAPSHOT_METHODS:
raise EffectRuntimeRejected(
"local snapshot transport is unavailable for this method",
diagnostic_code="invalid_request",
)
directory = Path(stack.enter_context(tempfile.TemporaryDirectory(prefix="loopx-effect-")))
response_sink = directory / "response.json"
descriptor = os.open(response_sink, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
os.close(descriptor)
request["response_sink"] = str(response_sink)
encoded = (json.dumps(request, separators=(",", ":")) + "\n").encode()
if len(encoded) > MAX_REQUEST_BYTES:
params_bytes = json.dumps(dict(params), separators=(",", ":")).encode()
if len(params_bytes) > MAX_LOCAL_SNAPSHOT_BYTES:
raise EffectRuntimeRejected(
"TypeScript Effect runtime local snapshot is oversized",
diagnostic_code="request_too_large",
)
params_path = directory / "params.json"
descriptor = os.open(params_path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
with os.fdopen(descriptor, "wb") as file:
file.write(params_bytes)
request.pop("params")
request["params_ref"] = {
"schema_version": "loopx_effect_runtime_snapshot_v0",
"path": str(params_path), "byte_count": len(params_bytes),
"sha256": hashlib.sha256(params_bytes).hexdigest(),
}
encoded = (json.dumps(request, separators=(",", ":")) + "\n").encode()
if len(encoded) > MAX_REQUEST_BYTES:
raise EffectRuntimeRejected(
"TypeScript Effect runtime request is oversized",
diagnostic_code="request_too_large",
)
chunks: list[bytes] = []
size = 0
with socket.create_connection(
(str(info["host"]), int(info["port"])), timeout=timeout
) as connection:
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:
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 EffectRuntimeResponseAmbiguous(method, timeout=timeout) from None
response = json.loads(b"".join(chunks).split(b"\n", 1)[0])
except (json.JSONDecodeError, IndexError):
raise EffectRuntimeResponseAmbiguous(method, timeout=timeout) from None
if isinstance(response, dict) and "result_ref" in response:
response = _read_local_snapshot_response(
response, response_sink, method=method, request_id=request_id, timeout=timeout,
)
if (
not isinstance(response, dict)
or response.get("schema_version") != EFFECT_RUNTIME_RESPONSE_SCHEMA_VERSION
Expand All @@ -599,6 +643,39 @@ def _request_with_info(
return response


def _read_local_snapshot_response(
envelope: dict[str, Any], sink: Path | None, *, method: str, request_id: str, timeout: float,
) -> dict[str, Any]:
"""Read an exact private response; unverifiable post-dispatch bytes are ambiguous."""
try:
ref = envelope["result_ref"]
if (sink is None or envelope.get("schema_version") != EFFECT_RUNTIME_RESPONSE_SCHEMA_VERSION
or envelope.get("request_id") != request_id or envelope.get("ok") is not True
or not isinstance(ref, dict)):
raise ValueError("invalid local snapshot envelope")
size, digest = ref.get("byte_count"), ref.get("sha256")
if (not isinstance(size, int) or isinstance(size, bool) or size <= 0
or size > MAX_LOCAL_SNAPSHOT_BYTES or not isinstance(digest, str)
or re.fullmatch(r"[a-f0-9]{64}", digest) is None):
raise ValueError("invalid local snapshot reference")
descriptor = os.open(sink, os.O_RDONLY | getattr(os, "O_NOFOLLOW", 0))
with os.fdopen(descriptor, "rb") as file:
metadata = os.fstat(file.fileno())
if (not os.path.isfile(sink) or metadata.st_size != size
or (hasattr(os, "getuid") and
(metadata.st_uid != os.getuid() or metadata.st_mode & 0o077))):
raise ValueError("invalid local snapshot file")
data = file.read(MAX_LOCAL_SNAPSHOT_BYTES + 1)
if len(data) != size or hashlib.sha256(data).hexdigest() != digest:
raise ValueError("local snapshot digest mismatch")
result = json.loads(data)
if not isinstance(result, dict):
raise ValueError("invalid local snapshot response")
return result
except (OSError, ValueError, json.JSONDecodeError) as exc:
raise EffectRuntimeResponseAmbiguous(method, timeout=timeout) from exc


def _remote_runtime_error(value: object) -> EffectRuntimeRemoteError:
if not isinstance(value, Mapping):
return EffectRuntimeInternalError(
Expand Down Expand Up @@ -805,6 +882,7 @@ def effect_runtime_request(
*,
timeout: float = DEFAULT_REQUEST_TIMEOUT_SECONDS,
retry_safe: bool = True,
large_local_snapshot: bool = False,
) -> dict[str, Any]:
"""Call the managed TS runtime, retrying only idempotent typed effects."""

Expand All @@ -824,6 +902,7 @@ def effect_runtime_request(
method=method,
params=params,
timeout=timeout,
large_local_snapshot=large_local_snapshot,
)
except (EffectRuntimeRemoteError, EffectRuntimeResponseAmbiguous):
raise
Expand Down Expand Up @@ -873,12 +952,14 @@ def effect_runtime_result(
*,
timeout: float = DEFAULT_REQUEST_TIMEOUT_SECONDS,
retry_safe: bool = True,
large_local_snapshot: bool = False,
) -> Any:
return effect_runtime_request(
method,
params,
timeout=timeout,
retry_safe=retry_safe,
large_local_snapshot=large_local_snapshot,
).get("result")


Expand Down
80 changes: 65 additions & 15 deletions loopx/control_plane/effect_runtime_server.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
import { writeSync } from "node:fs";
import { createServer, type Socket } from "node:net";
import { chmod, readFile, rm } from "node:fs/promises";
import { chmod, readFile, rm, type FileHandle } from "node:fs/promises";

import type { JsonObject } from "./effect_program.ts";
import {
Expand All @@ -13,6 +13,7 @@ import {
} from "./effect_runtime_errors.ts";
import { atomicWriteJson, withFileMutationLock } from "./effect_runtime_io.ts";
import { sqliteRuntimeIdentity } from "./coordination/sqlite_runtime.ts";
import {openPrivateResponseSink, readPrivateJsonSnapshot, writePrivateResponse} from "./effect_runtime_snapshot.ts";
import {
requireJsonObject as requiredObject,
requireNonEmptyString as requiredString,
Expand All @@ -23,6 +24,14 @@ const RESPONSE_SCHEMA = "loopx_effect_runtime_response_v1";
const INFO_SCHEMA = "loopx_effect_runtime_info_v0";
const STARTUP_ERROR_SCHEMA = "loopx_effect_runtime_startup_error_v0";
const MAX_REQUEST_BYTES = 2 * 1024 * 1024;
const MAX_INLINE_RESPONSE_BYTES = 2 * 1024 * 1024;
// Explicit opt-in: ordinary effects retain the 2 MiB request/response wire.
const LOCAL_SNAPSHOT_METHODS = new Set([
"goal.checkpoint_read_context.source",
"goal.checkpoint_read_context.evaluate",
"goal.checkpoint_read_context.commit",
"goal.checkpoint_read_context.inspect_replay",
]);
const DEFAULT_IDLE_MS = 5 * 60 * 1_000;
// Bounds for LOOPX_EFFECT_RUNTIME_IDLE_MS. The upper bound is the largest
// delay `setTimeout` accepts: a larger delay overflows and fires immediately,
Expand Down Expand Up @@ -108,8 +117,21 @@ function resetIdleTimer(server: ReturnType<typeof createServer>): void {
idleTimer.unref();
}

function writeResponse(socket: Socket, response: JsonObject): void {
socket.end(`${JSON.stringify(response)}\n`);
async function writeResponse(socket: Socket, response: JsonObject, sink: FileHandle | null): Promise<void> {
const encoded = Buffer.from(`${JSON.stringify(response)}\n`);
if (encoded.length <= MAX_INLINE_RESPONSE_BYTES) {
socket.end(encoded);
return;
}
if (sink === null) {
// The operation may already have committed. Never downgrade a lost large
// response to a safe request rejection or retry it automatically.
socket.destroy();
return;
}
const ref = await writePrivateResponse(sink, encoded);
socket.end(`${JSON.stringify({schema_version: RESPONSE_SCHEMA,
request_id: response.request_id, ok: true, result_ref: ref})}\n`);
}

const server = createServer((socket) => {
Expand All @@ -128,22 +150,24 @@ const server = createServer((socket) => {
raw = "";
socket.pause();
socket.removeAllListeners("data");
writeResponse(socket, {
socket.end(`${JSON.stringify({
schema_version: RESPONSE_SCHEMA,
request_id: "unknown",
ok: false,
error: effectRuntimeErrorPayload(new EffectRuntimeRequestError(
"Effect runtime request exceeds the 2 MiB limit",
"request_too_large",
)),
});
})}\n`);
return;
}
raw += chunk;
if (!raw.includes("\n")) return;
socket.pause();
void (async () => {
let requestId = "unknown";
let sink: FileHandle | null = null;
let dispatched = false;
try {
let parsed: unknown;
try {
Expand All @@ -164,25 +188,51 @@ const server = createServer((socket) => {
"authentication_failed",
);
}
const method = requiredString(request.method, "method");
if (request.params_ref !== undefined || request.response_sink !== undefined) {
if (!LOCAL_SNAPSHOT_METHODS.has(method) || request.response_sink === undefined ||
(request.params_ref !== undefined && request.params !== undefined)) {
throw new EffectRuntimeRequestError("local snapshot transport is unavailable for this request");
}
sink = await openPrivateResponseSink(request.response_sink);
}
let params: JsonObject;
if (request.params_ref !== undefined) {
const ref = requiredObject(request.params_ref, "params snapshot reference");
if (ref.schema_version !== "loopx_effect_runtime_snapshot_v0") {
throw new EffectRuntimeRequestError("invalid local snapshot schema");
}
params = requiredObject(await readPrivateJsonSnapshot(ref), "snapshot params");
} else {
params = asObject(request.params);
}
const result = await dispatchEffectRuntimeMethod(
handlers,
requiredString(request.method, "method"),
asObject(request.params),
method,
params,
);
writeResponse(socket, {
dispatched = true;
await writeResponse(socket, {
schema_version: RESPONSE_SCHEMA,
request_id: requestId,
ok: true,
result,
});
}, sink);
if (shutdownRequested) setImmediate(() => server.close());
} catch (error) {
writeResponse(socket, {
schema_version: RESPONSE_SCHEMA,
request_id: requestId,
ok: false,
error: effectRuntimeErrorPayload(error),
});
if (dispatched) socket.destroy();
else {
try {
await writeResponse(socket, {
schema_version: RESPONSE_SCHEMA,
request_id: requestId,
ok: false,
error: effectRuntimeErrorPayload(error),
}, sink);
} catch { socket.destroy(); }
}
} finally {
await sink?.close();
}
})();
});
Expand Down
Loading
Loading