diff --git a/apps/server/src/provider/Layers/EventNdjsonLogger.test.ts b/apps/server/src/provider/Layers/EventNdjsonLogger.test.ts index 500f2815c539..e87372e03381 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); @@ -49,7 +50,7 @@ function parseLogLine(line: string) { } describe("EventNdjsonLogger", () => { - it.effect("logs bounded diagnostics when an event cannot be serialized", () => { + it.effect("summarizes circular events without exposing their contents in diagnostics", () => { const messages: Array = []; const logCapture = Logger.make(({ message }) => { if (Array.isArray(message)) { @@ -71,10 +72,14 @@ describe("EventNdjsonLogger", () => { assert.exists(logger); if (!logger) return; yield* logger.write(circular, ThreadId.make("thread-1")); + yield* logger.close(); const serialized = encodeUnknownJson(messages); assert.notInclude(serialized, secret); - assert.include(serialized, '"errorTag":"SchemaError"'); + const line = parseLogLine( + NodeFS.readFileSync(ownedLogPath(basePath, "thread-1"), "utf8").trim(), + ); + assert.equal(line.payload, '{"truncated":true}'); } finally { NodeFS.rmSync(tempDir, { recursive: true, force: true }); } @@ -314,6 +319,22 @@ describe("EventNdjsonLogger", () => { { method: "thread/realtime/transcript/delta", payload: circularDelta }, threadId, ); + yield* native.write({ method: "turn/diff/updated", payload: circularDelta }, threadId); + yield* native.write( + { + event: { + direction: "incoming", + stage: "decoded", + payload: { + method: "turn/diff/updated", + get params() { + throw new Error("unused diff snapshots must not be traversed"); + }, + }, + }, + }, + threadId, + ); yield* native.write( { event: { @@ -332,6 +353,41 @@ describe("EventNdjsonLogger", () => { }, threadId, ); + yield* native.write( + { + event: { + direction: "incoming", + stage: "decoded", + payload: { method: "item/commandExecution/outputDelta", params: circularDelta }, + }, + }, + threadId, + ); + yield* native.write( + { + event: { + direction: "incoming", + stage: "decoded", + payload: { + type: "stream_event", + event: { type: "content_block_delta", delta: circularDelta }, + }, + }, + }, + threadId, + ); + yield* native.write( + { + event: { + direction: "incoming", + stage: "raw", + get payload() { + throw new Error("raw frames must not be inspected"); + }, + }, + }, + threadId, + ); yield* native.write( { event: { @@ -375,6 +431,119 @@ describe("EventNdjsonLogger", () => { }), ); + it.effect("summarizes large histories without reading their items", () => + Effect.gen(function* () { + const tempDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-provider-log-")); + const basePath = NodePath.join(tempDir, "events.log"); + const turns = Array.from({ length: 10_000 }); + Object.defineProperty(turns, 0, { + get: () => { + throw new Error("history must not be serialized"); + }, + }); + + try { + const store = yield* makeEventNdjsonLogStore(basePath); + yield* store.logger("native").write( + { + provider: "codex", + event: { + direction: "incoming", + stage: "decoded", + payload: { id: 42, result: { thread: { id: "native-thread", turns } } }, + }, + }, + ThreadId.make("large-history"), + ); + yield* store.close(); + + const contents = NodeFS.readFileSync(ownedLogPath(basePath, "large-history"), "utf8"); + assert.isBelow(Buffer.byteLength(contents), 2_048); + const record = decodeUnknownJson(parseLogLine(contents.trim()).payload); + assert.nestedPropertyVal(record, "event.payload.id", 42); + assert.nestedPropertyVal(record, "event.payload.result.thread.id", "native-thread"); + assert.nestedPropertyVal(record, "event.payload.result.thread.turns.itemCount", 10_000); + } finally { + NodeFS.rmSync(tempDir, { recursive: true, force: true }); + } + }), + ); + + it.effect("bounds oversized records while retaining failure details", () => + Effect.gen(function* () { + const tempDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-provider-log-")); + const basePath = NodePath.join(tempDir, "events.log"); + + try { + const store = yield* makeEventNdjsonLogStore(basePath); + const logger = store.logger("native"); + const threadId = ThreadId.make("large-error"); + const failure = { + method: "error", + params: { + threadId: "native-thread", + turnId: "native-turn", + error: { message: "The provider is unavailable.", code: "overloaded" }, + output: "x".repeat(128 * 1_024), + }, + }; + yield* logger.write(failure, threadId); + yield* logger.write({ id: "escaped", output: "\u0000".repeat(20_000) }, threadId); + yield* store.close(); + + const contents = NodeFS.readFileSync(ownedLogPath(basePath, "large-error"), "utf8"); + const records = contents + .trim() + .split("\n") + .map((line) => decodeUnknownJson(parseLogLine(line).payload)); + assert.isBelow(Buffer.byteLength(contents), 64 * 1_024); + assert.equal(records.length, 2); + assert.nestedPropertyVal(records[0], "params.threadId", "native-thread"); + assert.nestedPropertyVal(records[0], "params.turnId", "native-turn"); + assert.nestedPropertyVal( + records[0], + "params.error.message", + "The provider is unavailable.", + ); + assert.nestedPropertyVal(records[0], "params.error.code", "overloaded"); + assert.propertyVal(records[1], "id", "escaped"); + } finally { + NodeFS.rmSync(tempDir, { recursive: true, force: true }); + } + }), + ); + + it.effect("bounds canonical diff snapshots before serializing their duplicate payloads", () => + Effect.gen(function* () { + const tempDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-provider-log-")); + const basePath = NodePath.join(tempDir, "events.log"); + const threadId = ThreadId.make("large-diff"); + const diff = "diff-payload".repeat(128 * 1_024); + try { + const store = yield* makeEventNdjsonLogStore(basePath, { batchWindowMs: 0 }); + yield* store.logger("canonical").write( + { + type: "turn.diff.updated", + threadId, + turnId: "native-turn", + raw: { method: "turn/diff/updated", payload: { diff } }, + payload: { unifiedDiff: diff }, + }, + threadId, + ); + yield* store.close(); + const contents = NodeFS.readFileSync(ownedLogPath(basePath, "large-diff"), "utf8"); + assert.isBelow(Buffer.byteLength(contents), 2_048); + const record = decodeUnknownJson(parseLogLine(contents.trim()).payload); + assert.propertyVal(record, "type", "turn.diff.updated"); + assert.propertyVal(record, "threadId", threadId); + assert.propertyVal(record, "turnId", "native-turn"); + } finally { + NodeFS.rmSync(tempDir, { recursive: true, force: true }); + } + }), + ); + it.effect("keeps OpenCode tool input, final output, and errors in native logs", () => 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 de297ee020ea..dda83fe4a007 100644 --- a/apps/server/src/provider/Layers/EventNdjsonLogger.ts +++ b/apps/server/src/provider/Layers/EventNdjsonLogger.ts @@ -32,6 +32,9 @@ 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; +const MAX_RECORD_CHARACTERS = 64 * 1024; +const MAX_RECORD_FIELDS = 1_024; +const MAX_RECORD_DEPTH = 16; const GLOBAL_THREAD_SEGMENT = "_global"; const LOG_SCOPE = "provider-observability"; const encodeUnknownJsonString = Schema.encodeUnknownEffect(Schema.fromJsonString(Schema.Unknown)); @@ -54,6 +57,7 @@ const transientNativeMethods = new Set([ "item/reasoning/textDelta", "thread/realtime/outputAudio/delta", "thread/realtime/transcript/delta", + "turn/diff/updated", ]); const transientAcpUpdates = new Set(["agent_message_chunk", "agent_thought_chunk"]); @@ -188,7 +192,7 @@ function providerLogPath(directory: string, prefix: string, threadSegment: strin return NodePath.join(directory, `${prefix}${threadSegment}.log`); } -function shouldPersist(stream: EventNdjsonStream, event: unknown): boolean { +function shouldPersistProviderEvent(stream: EventNdjsonStream, event: unknown): boolean { if (stream === "orchestration" || typeof event !== "object" || event === null) { return true; } @@ -200,7 +204,17 @@ function shouldPersist(stream: EventNdjsonStream, event: unknown): boolean { if (stream !== "native") return true; const nested = Reflect.get(event, "event"); - const nativeEvent = typeof nested === "object" && nested !== null ? nested : event; + const envelope = typeof nested === "object" && nested !== null ? nested : event; + // Decoded frames carry the same information as raw frames without another + // copy of every token delta. Decode failures have their own diagnostic frame. + if (Reflect.get(envelope, "stage") === "raw") return false; + const decodedPayload = Reflect.get(envelope, "payload"); + const nativeEvent = + Reflect.get(envelope, "stage") === "decoded" && + typeof decodedPayload === "object" && + decodedPayload !== null + ? decodedPayload + : envelope; const method = Reflect.get(nativeEvent, "method"); if ( typeof method === "string" && @@ -212,6 +226,16 @@ function shouldPersist(stream: EventNdjsonStream, event: unknown): boolean { const nativeType = Reflect.get(nativeEvent, "type"); if (nativeType === "message.part.delta") return false; + if (nativeType === "stream_event") { + const streamEvent = Reflect.get(nativeEvent, "event"); + if ( + typeof streamEvent === "object" && + streamEvent !== null && + Reflect.get(streamEvent, "type") === "content_block_delta" + ) { + return false; + } + } const payload = Reflect.get(nativeEvent, "payload"); if (typeof payload !== "object" || payload === null) return true; @@ -246,6 +270,110 @@ function shouldPersist(stream: EventNdjsonStream, event: unknown): boolean { } } +const summaryFields = [ + "provider", + "protocol", + "kind", + "providerSessionId", + "direction", + "stage", + "type", + "subtype", + "method", + "id", + "threadId", + "turnId", + "requestId", + "session_id", + "status", + "is_error", + "api_error_status", + "terminal_reason", + "stop_reason", + "operation", + "code", + "willRetry", + "message", + "event", + "payload", + "params", + "result", + "thread", + "turn", + "error", + "turns", + "items", + "content", +] as const; + +function summarizeProviderEvent(event: unknown): unknown { + let remainingFields = 128; + let remainingCharacters = 8 * 1024; + const summarize = (value: unknown, depth: number): unknown => { + if (typeof value === "string") { + if (value.length > Math.min(1_024, remainingCharacters)) { + return { omittedCharacters: value.length }; + } + remainingCharacters -= value.length; + return value; + } + if (value === null || typeof value === "number" || typeof value === "boolean") return value; + if (typeof value !== "object") return undefined; + if (Array.isArray(value)) return { itemCount: value.length }; + const summary: Record = { truncated: true }; + if (depth >= 6) return summary; + for (const key of summaryFields) { + if (remainingFields <= 0) break; + const nested = Reflect.get(value, key); + if (nested === undefined) continue; + remainingFields -= 1; + summary[key] = summarize(nested, depth + 1); + } + return summary; + }; + try { + return summarize(event, 0); + } catch { + return { truncated: true }; + } +} + +/** Bounds traversal before the logger encodes payloads. */ +function boundProviderEventForLogging(event: unknown): unknown { + let remainingCharacters = MAX_RECORD_CHARACTERS; + let remainingFields = MAX_RECORD_FIELDS; + const ancestors = new WeakSet(); + const fits = (value: unknown, depth: number): boolean => { + if (typeof value === "string") { + remainingCharacters -= value.length; + return remainingCharacters >= 0; + } + if (typeof value !== "object" || value === null) return true; + if (depth > MAX_RECORD_DEPTH || ancestors.has(value)) return false; + if (Array.isArray(value) && value.length > remainingFields) return false; + ancestors.add(value); + for (const key in value) { + if (!Object.hasOwn(value, key)) continue; + remainingFields -= 1; + remainingCharacters -= key.length; + if ( + remainingFields < 0 || + remainingCharacters < 0 || + !fits(Reflect.get(value, key), depth + 1) + ) + return false; + } + ancestors.delete(value); + return true; + }; + try { + if (fits(event, 0)) return event; + } catch { + // A failing accessor must not escape into provider processing. + } + return summarizeProviderEvent(event); +} + export function writeBatchedMessages( sink: Pick, records: ReadonlyArray, @@ -612,9 +740,15 @@ export const makeEventNdjsonLogStore = Effect.fnUntraced(function* ( if (existing) return existing; const write = Effect.fnUntraced(function* (event: unknown, threadId: ThreadId | null) { - if (!shouldPersist(stream, event)) return; - const payload = yield* serializeEvent(event); + if (!shouldPersistProviderEvent(stream, event)) return; + let payload = yield* serializeEvent(boundProviderEventForLogging(event)); if (payload === undefined) return; + // Escaping can expand strings beyond their input size. Keep that bounded + // serialization out of the file too, while retaining routing/error fields. + if (Buffer.byteLength(payload) > MAX_RECORD_CHARACTERS) { + payload = yield* serializeEvent(summarizeProviderEvent(event)); + if (payload === undefined) return; + } const observedAt = yield* DateTime.now.pipe(Effect.map(DateTime.formatIso)); const line = `[${observedAt}] ${resolveStreamLabel(stream)}: ${payload}\n`;