Skip to content
Open
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
306 changes: 306 additions & 0 deletions packages/effect-codex-app-server/src/protocol.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<CodexError.CodexAppServerError>();
const rawLines: Array<unknown> = [];
const decoded: Array<unknown> = [];
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<CodexError.CodexAppServerError>();
const rawLines: Array<unknown> = [];
const decoded: Array<unknown> = [];
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) =>
Expand Down Expand Up @@ -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<Uint8Array>;
}> = [
{ 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<CodexProtocol.CodexAppServerIncomingNotification>();
const termination = yield* Deferred.make<CodexError.CodexAppServerError>();
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<CodexProtocol.CodexAppServerIncomingNotification>();
const termination = yield* Deferred.make<CodexError.CodexAppServerError>();
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<CodexProtocol.CodexAppServerIncomingNotification> = [];
const termination = yield* Deferred.make<CodexError.CodexAppServerError>();
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<CodexProtocol.CodexAppServerIncomingNotification> = [];
const termination = yield* Deferred.make<CodexError.CodexAppServerError>();
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<CodexError.CodexAppServerError>();
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();
Expand Down
Loading
Loading