diff --git a/packages/effect-codex-app-server/src/protocol.test.ts b/packages/effect-codex-app-server/src/protocol.test.ts index 7249afff1071..26757cb26015 100644 --- a/packages/effect-codex-app-server/src/protocol.test.ts +++ b/packages/effect-codex-app-server/src/protocol.test.ts @@ -325,6 +325,104 @@ it.layer(NodeServices.layer)("effect-codex-app-server protocol", (it) => { }), ); + it.effect("rejects an oversized fragmented message before join or parse", () => + Effect.gen(function* () { + const { stdio, input, output } = yield* makeInMemoryStdio(); + const termination = yield* Deferred.make(); + const rawLines: Array = []; + const decoded: Array = []; + let notificationCount = 0; + const transport = yield* CodexProtocol.makeCodexAppServerPatchedProtocol({ + stdio, + maxIncomingMessageBytes: 32, + logIncoming: true, + logger: (event) => + Effect.sync(() => { + if (event.stage === "raw") { + rawLines.push(event.payload); + } + if (event.stage === "decoded") { + decoded.push(event.payload); + } + }), + onNotification: () => Effect.sync(() => notificationCount++).pipe(Effect.asVoid), + onTermination: (error) => Deferred.succeed(termination, error).pipe(Effect.asVoid), + }); + const pending = yield* transport.request("thread/read", {}).pipe(Effect.forkScoped); + yield* Queue.take(output); + + yield* Queue.offer(input, encoder.encode('{"method":"x/huge"')); + yield* Queue.offer(input, encoder.encode(',"params":"xxxxxxxx')); + yield* Queue.offer(input, encoder.encode('xxxxxxxx"}\n')); + + const error = yield* Deferred.await(termination); + assert.instanceOf(error, CodexError.CodexAppServerTransportError); + assert.equal(error.operation, "read-input-stream"); + assert.equal(rawLines.length, 0); + assert.equal(decoded.length, 0); + assert.equal(notificationCount, 0); + + const pendingError = yield* Fiber.join(pending).pipe( + Effect.match({ + onFailure: (failure) => failure, + onSuccess: () => assert.fail("Expected the oversized message to fail the request"), + }), + ); + assert.strictEqual(pendingError, error); + }), + ); + + it.effect("rejects an oversized multibyte message counted in UTF-8 bytes", () => + Effect.gen(function* () { + const { stdio, input, output } = yield* makeInMemoryStdio(); + const termination = yield* Deferred.make(); + const rawLines: Array = []; + const decoded: Array = []; + let notificationCount = 0; + const line = '{"method":"x","params":"你好你好"}'; + assert.ok(line.length <= 32); + assert.ok(encoder.encode(line).byteLength > 32); + const transport = yield* CodexProtocol.makeCodexAppServerPatchedProtocol({ + stdio, + maxIncomingMessageBytes: 32, + logIncoming: true, + logger: (event) => + Effect.sync(() => { + if (event.stage === "raw") { + rawLines.push(event.payload); + } + if (event.stage === "decoded") { + decoded.push(event.payload); + } + }), + onNotification: () => Effect.sync(() => notificationCount++).pipe(Effect.asVoid), + onTermination: (error) => Deferred.succeed(termination, error).pipe(Effect.asVoid), + }); + const pending = yield* transport.request("thread/read", {}).pipe(Effect.forkScoped); + yield* Queue.take(output); + + yield* Queue.offer(input, encoder.encode('{"method":"x"')); + yield* Queue.offer(input, encoder.encode(',"params":"你好')); + yield* Queue.offer(input, encoder.encode('你好"}\n')); + + const error = yield* Deferred.await(termination); + assert.instanceOf(error, CodexError.CodexAppServerTransportError); + assert.equal(error.operation, "read-input-stream"); + assert.equal(rawLines.length, 0); + assert.equal(decoded.length, 0); + assert.equal(notificationCount, 0); + + const pendingError = yield* Fiber.join(pending).pipe( + Effect.match({ + onFailure: (failure) => failure, + onSuccess: () => + assert.fail("Expected the oversized multibyte message to fail the request"), + }), + ); + assert.strictEqual(pendingError, error); + }), + ); + it.effect.each([1, 7, 1024])( "preserves JSONL framing and UTF-8 across %i-byte input chunks", (chunkSize) => @@ -787,6 +885,214 @@ it.layer(NodeServices.layer)("effect-codex-app-server protocol", (it) => { }), ); + it.effect("accepts LF- and CRLF-framed messages that exactly fill maxIncomingMessageBytes", () => + Effect.gen(function* () { + const message = { method: "x" }; + const encoded = encodeUnknownJsonString(message); + const limit = encoder.encode(encoded).byteLength; + + const framings: ReadonlyArray<{ + readonly label: string; + readonly chunks: ReadonlyArray; + }> = [ + { label: "LF single chunk", chunks: [encoder.encode(`${encoded}\n`)] }, + { label: "CRLF single chunk", chunks: [encoder.encode(`${encoded}\r\n`)] }, + { + label: "CRLF split across chunks", + chunks: [encoder.encode(`${encoded}\r`), encoder.encode("\n")], + }, + { + label: "CRLF with empty chunk between CR and LF", + chunks: [encoder.encode(`${encoded}\r`), new Uint8Array(), encoder.encode("\n")], + }, + ]; + + for (const framing of framings) { + const { stdio, input } = yield* makeInMemoryStdio(); + const received = yield* Deferred.make(); + const termination = yield* Deferred.make(); + yield* CodexProtocol.makeCodexAppServerPatchedProtocol({ + stdio, + maxIncomingMessageBytes: limit, + onNotification: (notification) => + Deferred.succeed(received, notification).pipe(Effect.asVoid), + onTermination: (error) => Deferred.succeed(termination, error).pipe(Effect.asVoid), + }); + + for (const chunk of framing.chunks) { + yield* Queue.offer(input, chunk); + } + yield* Queue.end(input); + + assert.deepEqual( + yield* Deferred.await(received), + message, + `${framing.label} at exact limit was not routed`, + ); + assert.instanceOf( + yield* Deferred.await(termination), + CodexError.CodexAppServerInputStreamEndedError, + `${framing.label} should end cleanly, not by limit`, + ); + } + }), + ); + + it.effect("coalesces many tiny unterminated chunks without exploding remainder count", () => + Effect.gen(function* () { + const { stdio, input } = yield* makeInMemoryStdio(); + const received = yield* Deferred.make(); + const termination = yield* Deferred.make(); + const message = { method: "x", params: { n: 1 } }; + const encoded = encodeUnknownJsonString(message); + const bytes = encoder.encode(`${encoded}\n`); + yield* CodexProtocol.makeCodexAppServerPatchedProtocol({ + stdio, + maxIncomingMessageBytes: bytes.byteLength, + onNotification: (notification) => + Deferred.succeed(received, notification).pipe(Effect.asVoid), + onTermination: (error) => Deferred.succeed(termination, error).pipe(Effect.asVoid), + }); + + // One-byte chunks would previously push one remainder entry per byte. With + // coalescing, remainder stays O(1) entries while bytes stay under the limit. + for (let offset = 0; offset < bytes.length - 1; offset++) { + yield* Queue.offer(input, bytes.subarray(offset, offset + 1)); + } + yield* Queue.offer(input, bytes.subarray(bytes.length - 1)); + yield* Queue.end(input); + + assert.deepEqual(yield* Deferred.await(received), message); + assert.instanceOf( + yield* Deferred.await(termination), + CodexError.CodexAppServerInputStreamEndedError, + ); + }), + ); + + it.effect( + "dispatches a complete message before failing an oversized line in the same chunk", + () => + Effect.gen(function* () { + const message = { method: "x" }; + const encoded = encodeUnknownJsonString(message); + const limit = encoder.encode(encoded).byteLength + 8; + const oversized = `{"method":"x/huge","params":"${"y".repeat(limit)}"}`; + assert.ok(encoder.encode(oversized).byteLength > limit); + + const { stdio, input } = yield* makeInMemoryStdio(); + const notifications: Array = []; + const termination = yield* Deferred.make(); + yield* CodexProtocol.makeCodexAppServerPatchedProtocol({ + stdio, + maxIncomingMessageBytes: limit, + onNotification: (notification) => + Effect.sync(() => { + notifications.push(notification); + }), + onTermination: (error) => Deferred.succeed(termination, error).pipe(Effect.asVoid), + }); + + yield* Queue.offer(input, encoder.encode(`${encoded}\n${oversized}\n`)); + + const error = yield* Deferred.await(termination); + assert.instanceOf(error, CodexError.CodexAppServerTransportError); + assert.equal(error.operation, "read-input-stream"); + assert.deepEqual(notifications, [message]); + }), + ); + + it.effect("dispatches collected lines before failing an oversized unterminated remainder", () => + Effect.gen(function* () { + const message = { method: "x" }; + const encoded = encodeUnknownJsonString(message); + const limit = encoder.encode(encoded).byteLength + 8; + const oversized = `{"method":"x/huge","params":"${"y".repeat(limit)}"}`; + assert.ok(encoder.encode(oversized).byteLength > limit); + + const { stdio, input } = yield* makeInMemoryStdio(); + const notifications: Array = []; + const termination = yield* Deferred.make(); + yield* CodexProtocol.makeCodexAppServerPatchedProtocol({ + stdio, + maxIncomingMessageBytes: limit, + onNotification: (notification) => + Effect.sync(() => { + notifications.push(notification); + }), + onTermination: (error) => Deferred.succeed(termination, error).pipe(Effect.asVoid), + }); + + // No trailing newline: the overflow is the incomplete remainder appended + // after the complete line was already collected from this same chunk. + yield* Queue.offer(input, encoder.encode(`${encoded}\n${oversized}`)); + + const error = yield* Deferred.await(termination); + assert.instanceOf(error, CodexError.CodexAppServerTransportError); + assert.equal(error.operation, "read-input-stream"); + assert.deepEqual(notifications, [message]); + }), + ); + + it.effect("rejects CRLF-framed messages whose content exceeds maxIncomingMessageBytes", () => + Effect.gen(function* () { + const message = { method: "x", p: "y" }; + const encoded = encodeUnknownJsonString(message); + const limit = encoder.encode(encoded).byteLength - 1; + + const { stdio, input } = yield* makeInMemoryStdio(); + const termination = yield* Deferred.make(); + let notificationCount = 0; + yield* CodexProtocol.makeCodexAppServerPatchedProtocol({ + stdio, + maxIncomingMessageBytes: limit, + onNotification: () => Effect.sync(() => notificationCount++).pipe(Effect.asVoid), + onTermination: (error) => Deferred.succeed(termination, error).pipe(Effect.asVoid), + }); + + yield* Queue.offer(input, encoder.encode(`${encoded}\r\n`)); + + const error = yield* Deferred.await(termination); + assert.instanceOf(error, CodexError.CodexAppServerTransportError); + assert.equal(error.operation, "read-input-stream"); + assert.equal(notificationCount, 0); + }), + ); + + it.effect("rejects invalid maxIncomingMessageBytes at construction", () => + Effect.gen(function* () { + for (const invalid of [ + Number.NaN, + Number.POSITIVE_INFINITY, + Number.NEGATIVE_INFINITY, + -1, + 1.5, + ]) { + const { stdio } = yield* makeInMemoryStdio(); + const exit = yield* Effect.exit( + Effect.scoped( + CodexProtocol.makeCodexAppServerPatchedProtocol({ + stdio, + maxIncomingMessageBytes: invalid, + }), + ), + ); + assert.equal( + exit._tag, + "Failure", + `expected construction to fail for maxIncomingMessageBytes=${String(invalid)}`, + ); + if (exit._tag === "Failure") { + assert.include( + String(exit.cause), + "maxIncomingMessageBytes must be a finite non-negative integer", + `unexpected cause for maxIncomingMessageBytes=${String(invalid)}`, + ); + } + } + }), + ); + it.effect("classifies an input stream ending without inventing a cause", () => Effect.gen(function* () { const { stdio, input } = yield* makeInMemoryStdio(); diff --git a/packages/effect-codex-app-server/src/protocol.ts b/packages/effect-codex-app-server/src/protocol.ts index 4a32973a988b..39a29b0b2606 100644 --- a/packages/effect-codex-app-server/src/protocol.ts +++ b/packages/effect-codex-app-server/src/protocol.ts @@ -16,6 +16,39 @@ const isJsonRpcId = Schema.is(JsonRpcId); const isJsonRpcResponseEnvelope = Schema.is(JsonRpcResponseEnvelope); const isCodexAppServerError = Schema.is(CodexError.CodexAppServerError); const MAX_BUFFERED_RAW_MESSAGES = 32; +// UTF-8 byte size of decoded remainder before join/parse. 128 MiB sits above +// observed Codex diffs (~49M characters) and Effect ndjson's 16 MiB default, +// and well below a V8 heap-threatening line. Tests inject a smaller ceiling. +const MAX_INCOMING_MESSAGE_BYTES = 128 * 1024 * 1024; + +// Counts UTF-8 bytes for `chunk.slice(from, to)` without allocating the +// encoded copy, and stops early once the total exceeds `limit`. Unpaired +// surrogates match TextEncoder's replacement-character output (3 bytes), +// so the returned count matches `utf8.encode(fragment).byteLength` whenever +// it stays within `limit`. +const countUtf8Bytes = (chunk: string, from: number, to: number, limit: number): number => { + let bytes = 0; + for (let i = from; i < to; i++) { + const code = chunk.charCodeAt(i); + if (code < 0x80) { + bytes += 1; + } else if (code < 0x800) { + bytes += 2; + } else if (code >= 0xd800 && code <= 0xdbff && i + 1 < to) { + const next = chunk.charCodeAt(i + 1); + if (next >= 0xdc00 && next <= 0xdfff) { + bytes += 4; + i += 1; + } else { + bytes += 3; + } + } else { + bytes += 3; + } + if (bytes > limit) return bytes; + } + return bytes; +}; export interface CodexAppServerProtocolLogEvent { readonly direction: "incoming" | "outgoing"; @@ -39,6 +72,7 @@ export interface CodexAppServerPatchedProtocolOptions { readonly terminationError?: Effect.Effect; readonly logIncoming?: boolean; readonly logOutgoing?: boolean; + readonly maxIncomingMessageBytes?: number; readonly logger?: (event: CodexAppServerProtocolLogEvent) => Effect.Effect; readonly onNotification?: ( notification: CodexAppServerIncomingNotification, @@ -164,7 +198,21 @@ export const makeCodexAppServerPatchedProtocol = Effect.fn("makeCodexAppServerPa yield* Queue.sliding(MAX_BUFFERED_RAW_MESSAGES); const pending = yield* Ref.make(new Map()); const nextRequestId = yield* Ref.make(1); + const maxIncomingMessageBytes = options.maxIncomingMessageBytes ?? MAX_INCOMING_MESSAGE_BYTES; + if (!Number.isInteger(maxIncomingMessageBytes) || maxIncomingMessageBytes < 0) { + return yield* Effect.die( + new Error( + `Codex App Server maxIncomingMessageBytes must be a finite non-negative integer; received ${String(options.maxIncomingMessageBytes)}.`, + ), + ); + } const remainder: Array = []; + let remainderBytes = 0; + // Tracks a trailing carriage return whose partner \n may arrive in a later + // chunk. Kept out of `remainderBytes` so a CRLF-framed message that would + // exactly fill `maxIncomingMessageBytes` still fits after the terminator + // is stripped, matching LF-framed behavior at the limit. + let pendingCr = false; const terminationHandled = yield* Ref.make(false); const terminationFailure = yield* Ref.make(Option.none()); const terminationSignal = yield* Deferred.make(); @@ -395,29 +443,111 @@ export const makeCodexAppServerPatchedProtocol = Effect.fn("makeCodexAppServerPa ); }; + // Append into the existing remainder entry so fragmented input scales with + // payload bytes, not chunk count (one entry per byte would exhaust the heap + // at the 128 MiB ceiling before the size check could fail). + const appendRemainder = (fragment: string) => { + if (remainder.length === 0) { + remainder.push(fragment); + } else { + remainder[remainder.length - 1] += fragment; + } + }; + yield* options.stdio.stdin.pipe( Stream.interruptWhen(Deferred.await(terminationSignal)), Stream.decodeText(), Stream.runForEach((chunk) => - Effect.sync(() => { + Effect.suspend(() => { const lines: Array = []; let start = 0; - for ( - let newline = chunk.indexOf("\n"); - newline !== -1; - newline = chunk.indexOf("\n", start) - ) { - remainder.push(chunk.slice(start, newline)); - lines.push(remainder.join("").replace(/\r$/, "")); + const retainRange = (from: number, to: number) => { + if (from >= to) return true; + const limit = maxIncomingMessageBytes - remainderBytes; + const fragmentLength = countUtf8Bytes(chunk, from, to, limit); + if (fragmentLength > limit) { + remainder.length = 0; + remainderBytes = 0; + pendingCr = false; + return false; + } + appendRemainder(chunk.slice(from, to)); + remainderBytes += fragmentLength; + return true; + }; + // A chunk can already hold complete messages before the fragment that + // crosses the limit. Deliver those lines, then fail. Leave the + // oversized fragment itself unparsed. + const failAfterCollectedLines = () => + Effect.forEach(lines, handleLine, { discard: true }).pipe( + Effect.andThen( + Effect.fail( + new CodexError.CodexAppServerTransportError({ + operation: "read-input-stream", + cause: new Error( + `Incoming message exceeded ${String(maxIncomingMessageBytes)} bytes.`, + ), + }), + ), + ), + ); + if (pendingCr) { + // Empty decodeText chunks must not commit the deferred \r; the + // partner \n may still arrive in a later chunk. + if (chunk.length === 0) return Effect.void; + if (chunk.charCodeAt(0) === 0x0a) { + // Previous chunk's trailing \r plus this chunk's leading \n form + // a CRLF terminator; neither byte is charged against the limit. + pendingCr = false; + lines.push(remainder.join("")); + remainder.length = 0; + remainderBytes = 0; + start = 1; + } else { + // The deferred \r is line content after all; commit its byte now. + pendingCr = false; + if (remainderBytes + 1 > maxIncomingMessageBytes) { + remainder.length = 0; + remainderBytes = 0; + return Effect.fail( + new CodexError.CodexAppServerTransportError({ + operation: "read-input-stream", + cause: new Error( + `Incoming message exceeded ${String(maxIncomingMessageBytes)} bytes.`, + ), + }), + ); + } + appendRemainder("\r"); + remainderBytes += 1; + } + } + while (start < chunk.length) { + const newline = chunk.indexOf("\n", start); + if (newline === -1) break; + const hasCr = newline > start && chunk.charCodeAt(newline - 1) === 0x0d; + const rangeEnd = hasCr ? newline - 1 : newline; + if (!retainRange(start, rangeEnd)) { + return failAfterCollectedLines(); + } + lines.push(remainder.join("")); remainder.length = 0; + remainderBytes = 0; start = newline + 1; } // Keep unfinished lines in fragments so each chunk is scanned only once. if (start < chunk.length) { - remainder.push(chunk.slice(start)); + // Defer a trailing \r; only the next chunk (or stream end) reveals + // whether it is a CRLF terminator or literal content. + const endsWithCr = chunk.charCodeAt(chunk.length - 1) === 0x0d; + const rangeEnd = endsWithCr ? chunk.length - 1 : chunk.length; + if (!retainRange(start, rangeEnd)) { + return failAfterCollectedLines(); + } + if (endsWithCr) pendingCr = true; } - return lines; - }).pipe(Effect.flatMap((lines) => Effect.forEach(lines, handleLine, { discard: true }))), + return Effect.forEach(lines, handleLine, { discard: true }); + }), ), Effect.matchEffect({ onFailure: (error) => @@ -425,12 +555,31 @@ export const makeCodexAppServerPatchedProtocol = Effect.fn("makeCodexAppServerPa Effect.succeed(normalizeIncomingError(error, "read-input-stream")), ), onSuccess: () => - Effect.sync(() => { + Effect.suspend(() => { + if (pendingCr) { + // The stream ended before \n arrived, so the deferred \r is + // literal content on the final line and must be charged. + pendingCr = false; + appendRemainder("\r"); + remainderBytes += 1; + } + if (remainderBytes > maxIncomingMessageBytes) { + remainder.length = 0; + remainderBytes = 0; + return Effect.fail( + new CodexError.CodexAppServerTransportError({ + operation: "read-input-stream", + cause: new Error( + `Incoming message exceeded ${String(maxIncomingMessageBytes)} bytes.`, + ), + }), + ); + } const line = remainder.join(""); remainder.length = 0; - return line; + remainderBytes = 0; + return handleLine(line); }).pipe( - Effect.flatMap(handleLine), Effect.matchEffect({ onFailure: (error) => handleTermination(() => Effect.succeed(error)), onSuccess: () =>