From 5400141b0759f5e1dcc3e806a02a2220eb3a6490 Mon Sep 17 00:00:00 2001 From: Koushik Dey Date: Wed, 9 Sep 2026 22:37:06 +0530 Subject: [PATCH 1/5] fix(server): stop an oversized provider diff from exhausting the backend heap A Codex `turn/diff/updated` notification carries the whole turn diff, which reaches hundreds of MiB on a large turn. `CodexAdapter` put that same string on both `raw.payload.diff` and `payload.unifiedDiff`, and `EventNdjsonLogger` JSON-encoded the event whole before writing it, so a single notification materialized a multi-hundred-MiB log line and could stop the backend with an out-of-memory error. The logger now bounds an event's string values before serializing it, under two limits: no single value keeps more than `maxStringLength` characters, and the values of one record together keep no more than `maxRecordLength`, so many merely large values cannot add up to a line the per-value cap would have allowed on its own. This guards the native, canonical, and orchestration streams for every provider. `CodexAdapter` no longer mirrors the diff into `raw.payload`, because the canonical `payload.unifiedDiff` beside it already carries the same string. With a 64 MiB turn diff, the canonical record drops from 128.00 MiB to 256 KiB. Fixes #10924 --- .../src/provider/Layers/CodexAdapter.test.ts | 39 +++++ .../src/provider/Layers/CodexAdapter.ts | 22 ++- .../provider/Layers/EventNdjsonLogger.test.ts | 151 ++++++++++++++++++ .../src/provider/Layers/EventNdjsonLogger.ts | 137 +++++++++++++++- 4 files changed, 347 insertions(+), 2 deletions(-) diff --git a/apps/server/src/provider/Layers/CodexAdapter.test.ts b/apps/server/src/provider/Layers/CodexAdapter.test.ts index f7c6036885d9..bdd29a53a633 100644 --- a/apps/server/src/provider/Layers/CodexAdapter.test.ts +++ b/apps/server/src/provider/Layers/CodexAdapter.test.ts @@ -1233,6 +1233,45 @@ lifecycleLayer("CodexAdapterLive lifecycle", (it) => { }), ); + it.effect("carries a turn diff once instead of mirroring it into the raw payload", () => + Effect.gen(function* () { + const { adapter, runtime } = yield* startLifecycleRuntime(); + const firstEventFiber = yield* Stream.runHead(adapter.streamEvents).pipe(Effect.forkChild); + + const diff = `diff --git a/file.ts b/file.ts ++${"x".repeat(4_096)} +`; + yield* runtime.emit({ + id: asEventId("evt-turn-diff"), + kind: "notification", + provider: ProviderDriverKind.make("codex"), + createdAt: "2026-01-01T00:00:00.000Z", + method: "turn/diff/updated", + threadId: asThreadId("thread-1"), + turnId: asTurnId("turn-1"), + payload: { threadId: "thread-1", turnId: "turn-1", diff }, + }); + + const firstEvent = yield* Fiber.join(firstEventFiber); + NodeAssert.equal(firstEvent._tag, "Some"); + if (firstEvent._tag !== "Some") { + return; + } + const mapped = firstEvent.value; + NodeAssert.equal(mapped.type, "turn.diff.updated"); + if (mapped.type !== "turn.diff.updated") { + return; + } + NodeAssert.equal(mapped.payload.unifiedDiff, diff); + + const rawPayload = mapped.raw?.payload as { readonly diff?: unknown } | undefined; + NodeAssert.equal( + rawPayload?.diff, + `[omitted by t3, ${diff.length} characters in payload.unifiedDiff]`, + ); + }), + ); + it.effect("labels MCP lifecycle entries with server and tool names", () => Effect.gen(function* () { const { adapter, runtime } = yield* startLifecycleRuntime(); diff --git a/apps/server/src/provider/Layers/CodexAdapter.ts b/apps/server/src/provider/Layers/CodexAdapter.ts index b43755736ca3..1fe496f05488 100644 --- a/apps/server/src/provider/Layers/CodexAdapter.ts +++ b/apps/server/src/provider/Layers/CodexAdapter.ts @@ -990,6 +990,23 @@ function runtimeEventBase( }; } +/** + * Replaces the diff in a mirrored `turn/diff/updated` payload with a short + * marker. `payload.unifiedDiff` on the canonical event already carries the same + * string, and a turn diff is large enough that serializing it twice has + * exhausted the backend heap while logging the event. + */ +function elideNativeTurnDiff(payload: unknown): unknown { + if (typeof payload !== "object" || payload === null) return payload; + const fields = payload as Record; + const diff = fields.diff; + if (typeof diff !== "string") return payload; + return { + ...fields, + diff: `[omitted by t3, ${diff.length} characters in payload.unifiedDiff]`, + }; +} + function mapItemLifecycle( event: ProviderEvent, canonicalThreadId: ThreadId, @@ -1648,9 +1665,12 @@ function mapToRuntimeEvents( if (!payload) { return []; } + const base = runtimeEventBase(event, canonicalThreadId); + const raw = base.raw; return [ { - ...runtimeEventBase(event, canonicalThreadId), + ...base, + ...(raw ? { raw: { ...raw, payload: elideNativeTurnDiff(raw.payload) } } : {}), type: "turn.diff.updated", payload: { unifiedDiff: payload.diff, diff --git a/apps/server/src/provider/Layers/EventNdjsonLogger.test.ts b/apps/server/src/provider/Layers/EventNdjsonLogger.test.ts index 500f2815c539..b77e68edd2de 100644 --- a/apps/server/src/provider/Layers/EventNdjsonLogger.test.ts +++ b/apps/server/src/provider/Layers/EventNdjsonLogger.test.ts @@ -21,6 +21,7 @@ import { } from "./EventNdjsonLogger.ts"; const encodeUnknownJson = Schema.encodeUnknownSync(Schema.fromJsonString(Schema.Unknown)); +const decodeUnknownJson = Schema.decodeUnknownSync(Schema.fromJsonString(Schema.Unknown)); function ownedLogPath(basePath: string, segment: string): string { const basename = NodePath.basename(basePath); @@ -127,6 +128,156 @@ describe("EventNdjsonLogger", () => { }), ); + it.effect("truncates oversized string values before serializing an event", () => + Effect.gen(function* () { + const tempDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-provider-log-")); + const basePath = NodePath.join(tempDir, "provider-canonical.ndjson"); + + try { + const logger = yield* makeEventNdjsonLogger(basePath, { + stream: "canonical", + maxStringLength: 64, + }); + assert.notEqual(logger, undefined); + if (!logger) { + return; + } + + // A real Codex turn diff reaches hundreds of MiB. Serializing it whole, + // twice over, is what exhausted the backend heap. + const diff = "d".repeat(200_000); + yield* logger.write( + { + type: "turn.diff.updated", + id: "evt-diff", + payload: { unifiedDiff: diff }, + raw: { method: "turn/diff/updated", payload: { diff } }, + }, + ThreadId.make("thread-diff"), + ); + yield* logger.close(); + + const line = NodeFS.readFileSync(ownedLogPath(basePath, "thread-diff"), "utf8").trim(); + assert.equal(line.length < 1_000, true); + assert.notInclude(line, "d".repeat(65)); + assert.include(line, "[truncated by t3, 200000 characters total]"); + assert.include(line, '"id":"evt-diff"'); + } finally { + NodeFS.rmSync(tempDir, { recursive: true, force: true }); + } + }), + ); + + it.effect("bounds the whole record, not only each value", () => + Effect.gen(function* () { + const tempDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-provider-log-")); + const basePath = NodePath.join(tempDir, "provider-canonical.ndjson"); + + try { + const logger = yield* makeEventNdjsonLogger(basePath, { + stream: "canonical", + maxStringLength: 64, + maxRecordLength: 128, + }); + assert.notEqual(logger, undefined); + if (!logger) { + return; + } + + // Every value here clears the per-value cap on its own; only the record + // budget stops ten of them from adding up to an oversized line. + const fields = Object.fromEntries( + Array.from({ length: 10 }, (_, index) => [`field${index}`, "v".repeat(64)]), + ); + yield* logger.write({ id: "evt-budget", ...fields }, ThreadId.make("thread-budget")); + yield* logger.close(); + + const line = NodeFS.readFileSync(ownedLogPath(basePath, "thread-budget"), "utf8").trim(); + assert.equal((line.match(/v/gu) ?? []).length <= 128, true); + assert.include(line, "[truncated by t3, 64 characters total]"); + assert.include(line, '"id":"evt-budget"'); + } finally { + NodeFS.rmSync(tempDir, { recursive: true, force: true }); + } + }), + ); + + it.effect("charges truncation markers against the record budget", () => + Effect.gen(function* () { + const tempDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-provider-log-")); + const basePath = NodePath.join(tempDir, "provider-canonical.ndjson"); + + try { + const logger = yield* makeEventNdjsonLogger(basePath, { + stream: "canonical", + maxStringLength: 32, + maxRecordLength: 2_048, + }); + assert.notEqual(logger, undefined); + if (!logger) { + return; + } + + // Two hundred oversized values. Each marker runs to about 40 characters, + // so leaving them unbilled would spend roughly 8,000 against a budget of + // 2,048 that already reads as exhausted. + const fields = Object.fromEntries( + Array.from({ length: 200 }, (_, index) => [`f${index}`, "w".repeat(5_000)]), + ); + yield* logger.write(fields, ThreadId.make("thread-markers")); + yield* logger.close(); + + const line = NodeFS.readFileSync(ownedLogPath(basePath, "thread-markers"), "utf8").trim(); + const record = decodeUnknownJson(parseLogLine(line).payload) as Record; + const retained = Object.values(record).reduce((total, value) => total + value.length, 0); + assert.equal(retained <= 2_048, true); + assert.include(line, "[truncated by t3, 5000 characters total]"); + } finally { + NodeFS.rmSync(tempDir, { recursive: true, force: true }); + } + }), + ); + + it.effect("keeps short identifying values behind an oversized one", () => + Effect.gen(function* () { + const tempDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-provider-log-")); + const basePath = NodePath.join(tempDir, "provider-canonical.ndjson"); + + try { + const logger = yield* makeEventNdjsonLogger(basePath, { + stream: "canonical", + maxStringLength: 4_096, + maxRecordLength: 512, + }); + assert.notEqual(logger, undefined); + if (!logger) { + return; + } + + // Key order puts a provider's bulky `raw` payload ahead of `type`, so + // without a reserve the identifiers behind it are emptied and the record + // no longer says what it is. + yield* logger.write( + { + raw: "r".repeat(100_000), + type: "turn.diff.updated", + eventId: "evt-legible", + }, + ThreadId.make("thread-legible"), + ); + yield* logger.close(); + + const line = NodeFS.readFileSync(ownedLogPath(basePath, "thread-legible"), "utf8").trim(); + const record = decodeUnknownJson(parseLogLine(line).payload) as Record; + assert.equal(record.type, "turn.diff.updated"); + assert.equal(record.eventId, "evt-legible"); + assert.include(record.raw ?? "", "[truncated by t3, 100000 characters total]"); + } finally { + NodeFS.rmSync(tempDir, { recursive: true, force: true }); + } + }), + ); + it.effect( "falls back to a global segment when orchestration thread id is missing or invalid", () => diff --git a/apps/server/src/provider/Layers/EventNdjsonLogger.ts b/apps/server/src/provider/Layers/EventNdjsonLogger.ts index de297ee020ea..e228c643010b 100644 --- a/apps/server/src/provider/Layers/EventNdjsonLogger.ts +++ b/apps/server/src/provider/Layers/EventNdjsonLogger.ts @@ -32,6 +32,17 @@ const DEFAULT_MAX_AGE_MS = 14 * DAY_MS; const DEFAULT_RETENTION_CHECK_INTERVAL_MS = 5 * 60 * 1_000; const DEFAULT_MAX_BUFFERED_BYTES = MEBIBYTE; const DEFAULT_MAX_BUFFERED_RECORDS = 512; +// Per string value, in UTF-16 code units rather than bytes: `String.length` is +// O(1), while measuring real bytes would scan every string on the write path. +const DEFAULT_MAX_STRING_LENGTH = 256 * 1024; +// Across all string values of one record, so that many merely large values +// cannot add up to a line the per-value cap would have allowed individually. +const DEFAULT_MAX_RECORD_LENGTH = 4 * MEBIBYTE; +// Values this short are the identifiers that make a record legible - `type`, +// `method`, ids. They draw from a reserve a long value may not touch, because +// key order puts a provider's bulky `raw` payload ahead of `type`, so without +// one a single large field would empty every identifier behind it. +const SHORT_VALUE_RESERVE = 256; const GLOBAL_THREAD_SEGMENT = "_global"; const LOG_SCOPE = "provider-observability"; const encodeUnknownJsonString = Schema.encodeUnknownEffect(Schema.fromJsonString(Schema.Unknown)); @@ -80,6 +91,8 @@ export interface EventNdjsonLogStoreOptions { readonly retentionCheckIntervalMs?: number; readonly maxBufferedBytes?: number; readonly maxBufferedRecords?: number; + readonly maxStringLength?: number; + readonly maxRecordLength?: number; readonly attribution?: ResourceAttribution["Service"]; } @@ -126,6 +139,8 @@ interface ResolvedOptions { readonly retentionCheckIntervalMs: number; readonly maxBufferedBytes: number; readonly maxBufferedRecords: number; + readonly maxStringLength: number; + readonly maxRecordLength: number; readonly attribution: ResourceAttribution["Service"] | undefined; } @@ -373,6 +388,8 @@ function resolveOptions( options.retentionCheckIntervalMs ?? DEFAULT_RETENTION_CHECK_INTERVAL_MS, maxBufferedBytes: options.maxBufferedBytes ?? DEFAULT_MAX_BUFFERED_BYTES, maxBufferedRecords: options.maxBufferedRecords ?? DEFAULT_MAX_BUFFERED_RECORDS, + maxStringLength: options.maxStringLength ?? DEFAULT_MAX_STRING_LENGTH, + maxRecordLength: options.maxRecordLength ?? DEFAULT_MAX_RECORD_LENGTH, attribution: options.attribution, } satisfies ResolvedOptions; @@ -385,6 +402,8 @@ function resolveOptions( ["retentionCheckIntervalMs", resolved.retentionCheckIntervalMs, 1], ["maxBufferedBytes", resolved.maxBufferedBytes, 1], ["maxBufferedRecords", resolved.maxBufferedRecords, 1], + ["maxStringLength", resolved.maxStringLength, 1], + ["maxRecordLength", resolved.maxRecordLength, 1], ] as const; for (const [option, value, minimum] of validations) { @@ -494,6 +513,115 @@ function drainPending(input: { ]; } +function truncationMarker(length: number): string { + return `[truncated by t3, ${length} characters total]`; +} + +/** + * Keeps the head of an oversized value and records how long the original was, so + * a truncated record still says which file a diff touched and how much was cut. + */ +function truncateString(value: string, headLength: number): string { + // Never cut between a surrogate pair; a lone surrogate would survive into the log line. + const lastCode = value.charCodeAt(headLength - 1); + const keep = lastCode >= 0xd800 && lastCode <= 0xdbff ? headLength - 1 : headLength; + return `${value.slice(0, Math.max(keep, 0))}${truncationMarker(value.length)}`; +} + +/** + * Answers whether a value can be rebuilt field by field. Anything else reaches the + * encoder untouched, because a class instance or a `toJSON` carrier such as `Date` + * would not survive being copied into a plain object. + */ +function isPlainContainer(value: object): boolean { + if (Array.isArray(value)) return true; + const prototype = Reflect.getPrototypeOf(value); + return prototype === Object.prototype || prototype === null; +} + +/** + * Bounds the string values of one event before it is serialized. A provider can + * hand us a payload far larger than any log record should be (a Codex + * `turn/diff/updated` diff reaches hundreds of MiB on a big turn), and encoding + * it whole materializes an equally large line that can exhaust the heap. + * + * Two bounds apply. No single value keeps more than `maxLength` characters, and + * the string values of one record, truncation markers included, together keep no + * more than the record budget, so an event carrying many merely large values + * cannot add up to a line the per-value cap would have allowed on its own. Once + * the budget cannot even fit a marker, the value is dropped. The budget is spent + * in encounter order, which is stable for a given event shape, and short values + * keep a reserve of it so the identifiers behind a bulky field survive. + * + * Values that already fit are returned by reference, so an event that needs no + * truncation copies no object or array. Anything that is not a plain object or + * array is left alone, because `toJSON` carriers such as `Date` must reach the + * encoder intact. A cycle is returned untouched at the point it closes, so + * serialization still fails the way it did before and `serializeEvent` reports it. + */ +function clampStrings( + value: unknown, + maxLength: number, + budget: { remaining: number }, + ancestors: Set, +): unknown { + if (typeof value === "string") { + const spendable = + value.length <= SHORT_VALUE_RESERVE + ? budget.remaining + : Math.max(budget.remaining - SHORT_VALUE_RESERVE, 0); + + if (value.length <= maxLength && value.length <= spendable) { + budget.remaining -= value.length; + return value; + } + // The marker is itself part of the record, so it is charged like any other + // retained text. Otherwise a run of oversized values would keep spending + // marker-sized bites of a budget that reads as exhausted. + const marker = truncationMarker(value.length); + if (marker.length > spendable) return ""; + + const replacement = truncateString(value, Math.min(maxLength, spendable - marker.length)); + budget.remaining -= replacement.length; + return replacement; + } + if (typeof value !== "object" || value === null) return value; + if (!isPlainContainer(value) || ancestors.has(value)) return value; + + ancestors.add(value); + try { + if (Array.isArray(value)) { + const entries = value as ReadonlyArray; + let clampedEntries: Array | undefined; + for (let index = 0; index < entries.length; index += 1) { + const entry = entries[index]; + const clamped = clampStrings(entry, maxLength, budget, ancestors); + if (clamped !== entry && clampedEntries === undefined) { + clampedEntries = entries.slice(0, index); + } + clampedEntries?.push(clamped); + } + return clampedEntries ?? value; + } + + const fields = value as Record; + let clampedFields: Record | undefined; + // `for...in` rather than `Object.entries`: this runs on every logged event, + // and the entries array would be allocated even when nothing is truncated. + for (const key in fields) { + if (!Object.hasOwn(fields, key)) continue; + const entry = fields[key]; + const clamped = clampStrings(entry, maxLength, budget, ancestors); + if (clamped === entry) continue; + clampedFields ??= { ...fields }; + clampedFields[key] = clamped; + } + return clampedFields ?? value; + } finally { + ancestors.delete(value); + } +} + const serializeEvent = Effect.fnUntraced(function* (event: unknown) { return yield* encodeUnknownJsonString(event).pipe( Effect.catch((error) => @@ -613,7 +741,14 @@ export const makeEventNdjsonLogStore = Effect.fnUntraced(function* ( const write = Effect.fnUntraced(function* (event: unknown, threadId: ThreadId | null) { if (!shouldPersist(stream, event)) return; - const payload = yield* serializeEvent(event); + const payload = yield* serializeEvent( + clampStrings( + event, + resolved.maxStringLength, + { remaining: resolved.maxRecordLength }, + new Set(), + ), + ); if (payload === undefined) return; const observedAt = yield* DateTime.now.pipe(Effect.map(DateTime.formatIso)); From 77942637e412dfd87288bca40af5576c666e2e16 Mon Sep 17 00:00:00 2001 From: Koushik Dey Date: Fri, 11 Sep 2026 02:20:44 +0530 Subject: [PATCH 2/5] fix(server): keep a truncated log value within its own cap A truncated value kept up to `maxStringLength` characters of the original and then appended the truncation marker, so the replacement ran past the cap it was supposed to honor. The marker now counts against the per-value cap as well as the record budget, and a value whose cap cannot fit even a marker is dropped, the same rule the record budget already applied. --- .../provider/Layers/EventNdjsonLogger.test.ts | 10 ++++++- .../src/provider/Layers/EventNdjsonLogger.ts | 27 ++++++++++--------- 2 files changed, 24 insertions(+), 13 deletions(-) diff --git a/apps/server/src/provider/Layers/EventNdjsonLogger.test.ts b/apps/server/src/provider/Layers/EventNdjsonLogger.test.ts index b77e68edd2de..d798de2958d4 100644 --- a/apps/server/src/provider/Layers/EventNdjsonLogger.test.ts +++ b/apps/server/src/provider/Layers/EventNdjsonLogger.test.ts @@ -162,6 +162,14 @@ describe("EventNdjsonLogger", () => { assert.notInclude(line, "d".repeat(65)); assert.include(line, "[truncated by t3, 200000 characters total]"); assert.include(line, '"id":"evt-diff"'); + + // The marker is part of the value, so the cap bounds the whole replacement. + const record = decodeUnknownJson(parseLogLine(line).payload) as { + readonly payload: { readonly unifiedDiff: string }; + readonly raw: { readonly payload: { readonly diff: string } }; + }; + assert.equal(record.payload.unifiedDiff.length <= 64, true); + assert.equal(record.raw.payload.diff.length <= 64, true); } finally { NodeFS.rmSync(tempDir, { recursive: true, force: true }); } @@ -210,7 +218,7 @@ describe("EventNdjsonLogger", () => { try { const logger = yield* makeEventNdjsonLogger(basePath, { stream: "canonical", - maxStringLength: 32, + maxStringLength: 64, maxRecordLength: 2_048, }); assert.notEqual(logger, undefined); diff --git a/apps/server/src/provider/Layers/EventNdjsonLogger.ts b/apps/server/src/provider/Layers/EventNdjsonLogger.ts index e228c643010b..6f71a035fb45 100644 --- a/apps/server/src/provider/Layers/EventNdjsonLogger.ts +++ b/apps/server/src/provider/Layers/EventNdjsonLogger.ts @@ -545,13 +545,14 @@ function isPlainContainer(value: object): boolean { * `turn/diff/updated` diff reaches hundreds of MiB on a big turn), and encoding * it whole materializes an equally large line that can exhaust the heap. * - * Two bounds apply. No single value keeps more than `maxLength` characters, and - * the string values of one record, truncation markers included, together keep no - * more than the record budget, so an event carrying many merely large values - * cannot add up to a line the per-value cap would have allowed on its own. Once - * the budget cannot even fit a marker, the value is dropped. The budget is spent - * in encounter order, which is stable for a given event shape, and short values - * keep a reserve of it so the identifiers behind a bulky field survive. + * Two bounds apply, and a truncation marker counts against both. No single value + * keeps more than `maxLength` characters, and the string values of one record + * together keep no more than the record budget, so an event carrying many merely + * large values cannot add up to a line the per-value cap would have allowed on + * its own. Once a limit cannot even fit a marker, the value is dropped. The + * budget is spent in encounter order, which is stable for a given event shape, + * and short values keep a reserve of it so the identifiers behind a bulky field + * survive. * * Values that already fit are returned by reference, so an event that needs no * truncation copies no object or array. Anything that is not a plain object or @@ -575,13 +576,15 @@ function clampStrings( budget.remaining -= value.length; return value; } - // The marker is itself part of the record, so it is charged like any other - // retained text. Otherwise a run of oversized values would keep spending - // marker-sized bites of a budget that reads as exhausted. + // The marker is itself part of the value, so it counts against both limits + // like any other retained text. Otherwise a truncated value would outgrow its + // cap, and a run of oversized values would keep spending marker-sized bites + // of a budget that reads as exhausted. + const room = Math.min(maxLength, spendable); const marker = truncationMarker(value.length); - if (marker.length > spendable) return ""; + if (marker.length > room) return ""; - const replacement = truncateString(value, Math.min(maxLength, spendable - marker.length)); + const replacement = truncateString(value, room - marker.length); budget.remaining -= replacement.length; return replacement; } From ce02736dd607dca507b380a4b434ba09f6a3bbb8 Mon Sep 17 00:00:00 2001 From: Koushik Dey Date: Fri, 11 Sep 2026 02:30:13 +0530 Subject: [PATCH 3/5] fix(server): contain a throwing accessor while bounding a log record Bounding a record reads every enumerable property before the encoder runs, and it ran outside the guard that turns a serialization failure into a dropped record. An event whose getter throws therefore failed the write outright, where `main` drops it with a warning. Bounding now sits behind the same guard. The record budget's comment also says what it measures: retained text, not the encoded line. --- .../provider/Layers/EventNdjsonLogger.test.ts | 36 +++++++++++++++++++ .../src/provider/Layers/EventNdjsonLogger.ts | 36 +++++++++++++------ 2 files changed, 62 insertions(+), 10 deletions(-) diff --git a/apps/server/src/provider/Layers/EventNdjsonLogger.test.ts b/apps/server/src/provider/Layers/EventNdjsonLogger.test.ts index d798de2958d4..f7bd902a7656 100644 --- a/apps/server/src/provider/Layers/EventNdjsonLogger.test.ts +++ b/apps/server/src/provider/Layers/EventNdjsonLogger.test.ts @@ -612,6 +612,42 @@ describe("EventNdjsonLogger", () => { }), ); + it.effect("drops a record whose enumerable accessor throws while it is bounded", () => + Effect.gen(function* () { + const tempDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-provider-log-")); + const basePath = NodePath.join(tempDir, "events.log"); + + try { + const logger = yield* makeEventNdjsonLogger(basePath, { + stream: "canonical", + batchWindowMs: 0, + }); + assert.exists(logger); + if (!logger) return; + + // Bounding reads every enumerable property before the encoder runs, so a + // throwing getter has to be contained there as well as in encoding. + const hostile = { + id: "evt-hostile", + get payload(): unknown { + throw new Error("blocked"); + }, + }; + yield* logger.write(hostile, ThreadId.make("thread-hostile-getter")); + yield* logger.write({ id: "evt-after" }, ThreadId.make("thread-hostile-getter")); + yield* logger.close(); + + const lines = NodeFS.readFileSync(ownedLogPath(basePath, "thread-hostile-getter"), "utf8") + .trim() + .split("\n"); + assert.equal(lines.length, 1); + assert.include(lines[0] ?? "", '"id":"evt-after"'); + } finally { + NodeFS.rmSync(tempDir, { recursive: true, force: true }); + } + }), + ); + it.effect("serializes concurrent first writes for the same segment", () => Effect.gen(function* () { const tempDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-provider-log-")); diff --git a/apps/server/src/provider/Layers/EventNdjsonLogger.ts b/apps/server/src/provider/Layers/EventNdjsonLogger.ts index 6f71a035fb45..63d11b05c785 100644 --- a/apps/server/src/provider/Layers/EventNdjsonLogger.ts +++ b/apps/server/src/provider/Layers/EventNdjsonLogger.ts @@ -37,6 +37,8 @@ const DEFAULT_MAX_BUFFERED_RECORDS = 512; const DEFAULT_MAX_STRING_LENGTH = 256 * 1024; // Across all string values of one record, so that many merely large values // cannot add up to a line the per-value cap would have allowed individually. +// It bounds retained text, not the encoded line: keys and JSON syntax are not +// counted, and escaping can grow the text at most sixfold (`\u0000`). const DEFAULT_MAX_RECORD_LENGTH = 4 * MEBIBYTE; // Values this short are the identifiers that make a record legible - `type`, // `method`, ids. They draw from a reserve a long value may not touch, because @@ -126,6 +128,17 @@ export class EventNdjsonLogDirectoryError extends Schema.TaggedError()( + "EventNdjsonLogRecordError", + { + cause: Schema.Defect(), + }, +) { + override get message(): string { + return "Failed to read a provider event while bounding its log record"; + } +} + export type EventNdjsonLogStoreError = | EventNdjsonLogConfigurationError | EventNdjsonLogDirectoryError; @@ -625,8 +638,18 @@ function clampStrings( } } -const serializeEvent = Effect.fnUntraced(function* (event: unknown) { - return yield* encodeUnknownJsonString(event).pipe( +const serializeEvent = Effect.fnUntraced(function* ( + event: unknown, + limits: Pick, +) { + // Bounding the record reads every enumerable property before the encoder does, + // so a hostile accessor can throw there too; both sit behind the same guard. + return yield* Effect.try({ + try: () => + clampStrings(event, limits.maxStringLength, { remaining: limits.maxRecordLength }, new Set()), + catch: (cause) => new EventNdjsonLogRecordError({ cause }), + }).pipe( + Effect.flatMap((bounded) => encodeUnknownJsonString(bounded)), Effect.catch((error) => logWarning("failed to serialize provider event log record", { errorTag: errorTag(error), @@ -744,14 +767,7 @@ export const makeEventNdjsonLogStore = Effect.fnUntraced(function* ( const write = Effect.fnUntraced(function* (event: unknown, threadId: ThreadId | null) { if (!shouldPersist(stream, event)) return; - const payload = yield* serializeEvent( - clampStrings( - event, - resolved.maxStringLength, - { remaining: resolved.maxRecordLength }, - new Set(), - ), - ); + const payload = yield* serializeEvent(event, resolved); if (payload === undefined) return; const observedAt = yield* DateTime.now.pipe(Effect.map(DateTime.formatIso)); From 6f1a55b628fa37ba560e30ea591681b9f71da746 Mon Sep 17 00:00:00 2001 From: Koushik Dey Date: Fri, 11 Sep 2026 02:45:50 +0530 Subject: [PATCH 4/5] fix(server): read each field once when rebuilding a bounded log record Rebuilding a record copied an object with a spread as soon as one field needed cutting, and the spread read every other field again. An accessor that returned a different value on that second read reached the encoder unbounded. Bounding now takes two passes. A read-only pass settles whether the event fits, spending the budget exactly as the rebuild would, so an ordinary event is still written as it came without copying anything. Only an event that does not fit is rebuilt, and that pass reads every field once and stores it as data. Copies have no prototype, so a `__proto__` key stays an ordinary field. --- .../provider/Layers/EventNdjsonLogger.test.ts | 38 ++++++ .../src/provider/Layers/EventNdjsonLogger.ts | 127 +++++++++++++----- 2 files changed, 133 insertions(+), 32 deletions(-) diff --git a/apps/server/src/provider/Layers/EventNdjsonLogger.test.ts b/apps/server/src/provider/Layers/EventNdjsonLogger.test.ts index f7bd902a7656..1549fa07be88 100644 --- a/apps/server/src/provider/Layers/EventNdjsonLogger.test.ts +++ b/apps/server/src/provider/Layers/EventNdjsonLogger.test.ts @@ -648,6 +648,44 @@ describe("EventNdjsonLogger", () => { }), ); + it.effect("reads each field once when rebuilding a record that needs bounding", () => + Effect.gen(function* () { + const tempDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-provider-log-")); + const basePath = NodePath.join(tempDir, "provider-canonical.ndjson"); + + try { + const logger = yield* makeEventNdjsonLogger(basePath, { + stream: "canonical", + maxStringLength: 64, + }); + assert.exists(logger); + if (!logger) return; + + // An accessor that grows between reads. If the rebuilt record re-read it + // instead of keeping the value it bounded, the later, oversized read would + // reach the encoder unbounded. + let reads = 0; + const event = { + get note(): string { + reads += 1; + return reads === 1 ? "short" : "n".repeat(100_000); + }, + raw: "r".repeat(100_000), + id: "evt-accessor", + }; + yield* logger.write(event, ThreadId.make("thread-accessor")); + yield* logger.close(); + + const line = NodeFS.readFileSync(ownedLogPath(basePath, "thread-accessor"), "utf8").trim(); + assert.equal(line.length < 1_000, true); + assert.notInclude(line, "n".repeat(65)); + assert.include(line, '"id":"evt-accessor"'); + } finally { + NodeFS.rmSync(tempDir, { recursive: true, force: true }); + } + }), + ); + it.effect("serializes concurrent first writes for the same segment", () => Effect.gen(function* () { const tempDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-provider-log-")); diff --git a/apps/server/src/provider/Layers/EventNdjsonLogger.ts b/apps/server/src/provider/Layers/EventNdjsonLogger.ts index 63d11b05c785..ec8449e6ec5d 100644 --- a/apps/server/src/provider/Layers/EventNdjsonLogger.ts +++ b/apps/server/src/provider/Layers/EventNdjsonLogger.ts @@ -553,10 +553,60 @@ function isPlainContainer(value: object): boolean { } /** - * Bounds the string values of one event before it is serialized. A provider can - * hand us a payload far larger than any log record should be (a Codex - * `turn/diff/updated` diff reaches hundreds of MiB on a big turn), and encoding - * it whole materializes an equally large line that can exhaust the heap. + * How much of the record budget a string may draw on. Values this short are the + * identifiers that make a record legible, so they may also spend the reserve. + */ +function spendableFor(value: string, budget: { remaining: number }): number { + return value.length <= SHORT_VALUE_RESERVE + ? budget.remaining + : Math.max(budget.remaining - SHORT_VALUE_RESERVE, 0); +} + +/** + * Answers whether an event already fits both limits, reading it without copying + * anything. It spends the budget exactly the way `clampStrings` does, so an event + * it accepts would come out of `clampStrings` unchanged. This is the only pass an + * ordinary event takes, which is why it walks objects with `for...in` rather + * than allocating an `Object.entries` array for each one. + */ +function fitsLimits( + value: unknown, + maxLength: number, + budget: { remaining: number }, + ancestors: Set, +): boolean { + if (typeof value === "string") { + if (value.length > maxLength || value.length > spendableFor(value, budget)) return false; + budget.remaining -= value.length; + return true; + } + if (typeof value !== "object" || value === null) return true; + if (!isPlainContainer(value) || ancestors.has(value)) return true; + + ancestors.add(value); + try { + if (Array.isArray(value)) { + const entries = value as ReadonlyArray; + for (let index = 0; index < entries.length; index += 1) { + if (!fitsLimits(entries[index], maxLength, budget, ancestors)) return false; + } + return true; + } + + const fields = value as Record; + for (const key in fields) { + if (Object.hasOwn(fields, key) && !fitsLimits(fields[key], maxLength, budget, ancestors)) { + return false; + } + } + return true; + } finally { + ancestors.delete(value); + } +} + +/** + * Rebuilds an event that does not fit, bounding its string values on the way. * * Two bounds apply, and a truncation marker counts against both. No single value * keeps more than `maxLength` characters, and the string values of one record @@ -567,10 +617,12 @@ function isPlainContainer(value: object): boolean { * and short values keep a reserve of it so the identifiers behind a bulky field * survive. * - * Values that already fit are returned by reference, so an event that needs no - * truncation copies no object or array. Anything that is not a plain object or - * array is left alone, because `toJSON` carriers such as `Date` must reach the - * encoder intact. A cycle is returned untouched at the point it closes, so + * Every plain object and array is copied, and each field is read exactly once + * and stored as data, so an accessor cannot hand the encoder a different, + * unbounded value afterwards. Objects are copied without a prototype so a + * `__proto__` key stays an ordinary field. Anything that is not a plain object + * or array is left alone, because `toJSON` carriers such as `Date` must reach + * the encoder intact. A cycle is returned untouched at the point it closes, so * serialization still fails the way it did before and `serializeEvent` reports it. */ function clampStrings( @@ -580,11 +632,7 @@ function clampStrings( ancestors: Set, ): unknown { if (typeof value === "string") { - const spendable = - value.length <= SHORT_VALUE_RESERVE - ? budget.remaining - : Math.max(budget.remaining - SHORT_VALUE_RESERVE, 0); - + const spendable = spendableFor(value, budget); if (value.length <= maxLength && value.length <= spendable) { budget.remaining -= value.length; return value; @@ -608,36 +656,52 @@ function clampStrings( try { if (Array.isArray(value)) { const entries = value as ReadonlyArray; - let clampedEntries: Array | undefined; + const bounded: Array = []; for (let index = 0; index < entries.length; index += 1) { - const entry = entries[index]; - const clamped = clampStrings(entry, maxLength, budget, ancestors); - if (clamped !== entry && clampedEntries === undefined) { - clampedEntries = entries.slice(0, index); - } - clampedEntries?.push(clamped); + bounded.push(clampStrings(entries[index], maxLength, budget, ancestors)); } - return clampedEntries ?? value; + return bounded; } const fields = value as Record; - let clampedFields: Record | undefined; - // `for...in` rather than `Object.entries`: this runs on every logged event, - // and the entries array would be allocated even when nothing is truncated. + const bounded: Record = Object.create(null); for (const key in fields) { if (!Object.hasOwn(fields, key)) continue; - const entry = fields[key]; - const clamped = clampStrings(entry, maxLength, budget, ancestors); - if (clamped === entry) continue; - clampedFields ??= { ...fields }; - clampedFields[key] = clamped; + bounded[key] = clampStrings(fields[key], maxLength, budget, ancestors); } - return clampedFields ?? value; + return bounded; } finally { ancestors.delete(value); } } +/** + * Bounds the string values of one event before it is serialized. A provider can + * hand us a payload far larger than any log record should be (a Codex + * `turn/diff/updated` diff reaches hundreds of MiB on a big turn), and encoding + * it whole materializes an equally large line that can exhaust the heap. An + * event that already fits is returned as it came, so the ordinary path copies + * no object or array; only one that does not fit is rebuilt. + */ +function boundEvent( + event: unknown, + limits: Pick, +): unknown { + const fits = fitsLimits( + event, + limits.maxStringLength, + { remaining: limits.maxRecordLength }, + new Set(), + ); + if (fits) return event; + return clampStrings( + event, + limits.maxStringLength, + { remaining: limits.maxRecordLength }, + new Set(), + ); +} + const serializeEvent = Effect.fnUntraced(function* ( event: unknown, limits: Pick, @@ -645,8 +709,7 @@ const serializeEvent = Effect.fnUntraced(function* ( // Bounding the record reads every enumerable property before the encoder does, // so a hostile accessor can throw there too; both sit behind the same guard. return yield* Effect.try({ - try: () => - clampStrings(event, limits.maxStringLength, { remaining: limits.maxRecordLength }, new Set()), + try: () => boundEvent(event, limits), catch: (cause) => new EventNdjsonLogRecordError({ cause }), }).pipe( Effect.flatMap((bounded) => encodeUnknownJsonString(bounded)), From b58e90611107b134104b5d7d1efb0db0d6853f23 Mon Sep 17 00:00:00 2001 From: Koushik Dey Date: Fri, 11 Sep 2026 02:54:58 +0530 Subject: [PATCH 5/5] docs(server): scope the log bounding read-once claim to the rebuild The rebuild reads each field once, but an event that already fits is passed to the encoder as it came, which reads its accessors again, and a rebuilt event has them read by both passes. Say so instead of implying the guarantee covers every event. --- .../src/provider/Layers/EventNdjsonLogger.ts | 21 ++++++++++++------- 1 file changed, 14 insertions(+), 7 deletions(-) diff --git a/apps/server/src/provider/Layers/EventNdjsonLogger.ts b/apps/server/src/provider/Layers/EventNdjsonLogger.ts index ec8449e6ec5d..8b707e92371f 100644 --- a/apps/server/src/provider/Layers/EventNdjsonLogger.ts +++ b/apps/server/src/provider/Layers/EventNdjsonLogger.ts @@ -617,13 +617,14 @@ function fitsLimits( * and short values keep a reserve of it so the identifiers behind a bulky field * survive. * - * Every plain object and array is copied, and each field is read exactly once - * and stored as data, so an accessor cannot hand the encoder a different, - * unbounded value afterwards. Objects are copied without a prototype so a - * `__proto__` key stays an ordinary field. Anything that is not a plain object - * or array is left alone, because `toJSON` carriers such as `Date` must reach - * the encoder intact. A cycle is returned untouched at the point it closes, so - * serialization still fails the way it did before and `serializeEvent` reports it. + * Every plain object and array is copied, and this pass reads each field once + * and stores it as data, so the encoder sees only values that were bounded + * rather than whatever an accessor returns on a later read. Objects are copied + * without a prototype so a `__proto__` key stays an ordinary field. Anything + * that is not a plain object or array is left alone, because `toJSON` carriers + * such as `Date` must reach the encoder intact. A cycle is returned untouched at + * the point it closes, so serialization still fails the way it did before and + * `serializeEvent` reports it. */ function clampStrings( value: unknown, @@ -682,6 +683,12 @@ function clampStrings( * it whole materializes an equally large line that can exhaust the heap. An * event that already fits is returned as it came, so the ordinary path copies * no object or array; only one that does not fit is rebuilt. + * + * The read-once guarantee therefore belongs to the rebuild. An event that fits + * reaches the encoder as it came, which reads its accessors again, exactly as it + * did before any bounding existed; an event that is rebuilt has its accessors + * read twice, once by each pass. Logged events are JSON-decoded payloads and + * object literals, so neither case arises in practice. */ function boundEvent( event: unknown,