Skip to content
Closed
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
42 changes: 17 additions & 25 deletions apps/server/src/provider/Layers/CodexSessionRuntime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@ import * as Stream from "effect/Stream";
import { ChildProcess, ChildProcessSpawner } from "effect/unstable/process";
import * as CodexClient from "effect-codex-app-server/client";
import * as CodexErrors from "effect-codex-app-server/errors";
import { makeLineFramer } from "effect-codex-app-server/lineFramer";
import * as CodexRpc from "effect-codex-app-server/rpc";
import * as EffectCodexSchema from "effect-codex-app-server/schema";

Expand Down Expand Up @@ -1621,34 +1622,25 @@ export const makeCodexSessionRuntime = (
Effect.forkIn(runtimeScope),
);

const stderrRemainderRef = yield* Ref.make("");
const stderrLineFramer = makeLineFramer();
yield* child.stderr.pipe(
Stream.decodeText(),
Stream.runForEach((chunk) =>
Ref.modify(stderrRemainderRef, (current) => {
const combined = current + chunk;
const lines = combined.split("\n");
const remainder = lines.pop() ?? "";
return [lines.map((line) => line.replace(/\r$/, "")), remainder] as const;
}).pipe(
Effect.flatMap((lines) =>
Effect.forEach(
lines,
(line) => {
const classified = classifyCodexStderrLine(line);
if (!classified) {
return Effect.void;
}
return emitEvent({
kind: "notification",
threadId: options.threadId,
method: "process/stderr",
message: classified.message,
});
},
{ discard: true },
),
),
Effect.forEach(
stderrLineFramer.push(chunk),
(line) => {
const classified = classifyCodexStderrLine(line);
if (!classified) {
return Effect.void;
}
return emitEvent({
kind: "notification",
threadId: options.threadId,
method: "process/stderr",
message: classified.message,
});
},
{ discard: true },
),
),
Effect.forkIn(runtimeScope),
Expand Down
4 changes: 4 additions & 0 deletions packages/effect-codex-app-server/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,10 @@
"types": "./src/protocol.ts",
"import": "./src/protocol.ts"
},
"./lineFramer": {
"types": "./src/lineFramer.ts",
"import": "./src/lineFramer.ts"
},
"./errors": {
"types": "./src/errors.ts",
"import": "./src/errors.ts"
Expand Down
56 changes: 56 additions & 0 deletions packages/effect-codex-app-server/src/lineFramer.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,56 @@
import { assert, it } from "@effect/vitest";

import { makeLineFramer } from "./lineFramer.ts";

const collectLines = (chunks: ReadonlyArray<string>) => {
const framer = makeLineFramer();
const lines = chunks.flatMap((chunk) => framer.push(chunk));
const finalLine = framer.finish();
return finalLine === undefined ? lines : [...lines, finalLine];
};

it("preserves LF, CRLF, empty-line, and unterminated-line semantics across chunks", () => {
assert.deepEqual(collectLines(["alpha\r", "\n", "\nb", "eta\n", "gamma", "\r", "\ndelta"]), [
"alpha",
"",
"beta",
"gamma",
"delta",
]);
});

it("handles chunk boundaries immediately before, after, and within line endings", () => {
assert.deepEqual(collectLines(["one", "\n", "two\n", "\r", "\n", "three\r", "\nfour"]), [
"one",
"two",
"",
"three",
"four",
]);
});

it("retains a large fragmented line without emitting partial records", () => {
const framer = makeLineFramer();
const recordSize = 20 * 1024 * 1024;
const record = "0123456789abcdef".repeat(recordSize / 16);
const chunkSizes = [4_093, 8_191, 16_381];
let offset = 0;
let chunkIndex = 0;
const startedAt = performance.now();

while (offset < record.length) {
const nextOffset = Math.min(
record.length,
offset + chunkSizes[chunkIndex % chunkSizes.length]!,
);
assert.deepEqual(framer.push(record.slice(offset, nextOffset)), []);
offset = nextOffset;
chunkIndex++;
}

const completed = framer.push("\n");
assert.lengthOf(completed, 1);
assert.equal(completed[0], record);
assert.isUndefined(framer.finish());
assert.isBelow(performance.now() - startedAt, 2_000);
});
48 changes: 48 additions & 0 deletions packages/effect-codex-app-server/src/lineFramer.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,48 @@
export interface LineFramer {
readonly push: (chunk: string) => ReadonlyArray<string>;
readonly finish: () => string | undefined;
}

export function makeLineFramer(): LineFramer {
// Keep incomplete lines fragmented so every incoming chunk is copied at most once.
let fragments: Array<string> = [];

const completeLine = (segment: string) => {
if (fragments.length === 0) {
return segment.endsWith("\r") ? segment.slice(0, -1) : segment;
}
if (segment.length > 0) {
fragments.push(segment);
}
const line = fragments.join("");
fragments = [];
return line.endsWith("\r") ? line.slice(0, -1) : line;
};

return {
push: (chunk) => {
const lines: Array<string> = [];
let segmentStart = 0;
let newlineIndex = chunk.indexOf("\n");

while (newlineIndex !== -1) {
lines.push(completeLine(chunk.slice(segmentStart, newlineIndex)));
segmentStart = newlineIndex + 1;
newlineIndex = chunk.indexOf("\n", segmentStart);
}

if (segmentStart < chunk.length) {
fragments.push(chunk.slice(segmentStart));
}
return lines;
},
finish: () => {
if (fragments.length === 0) {
return undefined;
}
const line = fragments.join("");
fragments = [];
return line;
},
};
}
63 changes: 63 additions & 0 deletions packages/effect-codex-app-server/src/protocol.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -290,6 +290,69 @@ it.layer(NodeServices.layer)("effect-codex-app-server protocol", (it) => {
}),
);

it.effect("routes a CRLF-framed message split across arbitrary input chunks", () =>
Effect.gen(function* () {
const { stdio, input } = yield* makeInMemoryStdio();
const transport = yield* CodexProtocol.makeCodexAppServerPatchedProtocol({ stdio });
const notification = yield* transport.incomingNotifications.pipe(
Stream.take(1),
Stream.runCollect,
Effect.forkScoped,
);
const encoded = encoder.encode(
`${encodeUnknownJsonString({ method: "thread/started", params: { threadId: "thread-1" } })}\r\n`,
);

yield* Queue.offer(input, encoded.slice(0, 7));
yield* Queue.offer(input, encoded.slice(7, -1));
yield* Queue.offer(input, encoded.slice(-1));

assert.deepEqual(yield* Fiber.join(notification), [
{ method: "thread/started", params: { threadId: "thread-1" } },
]);
}),
);

it.effect("routes a final unterminated message before handling input stream completion", () =>
Effect.gen(function* () {
const { stdio, input } = yield* makeInMemoryStdio();
const termination = yield* Deferred.make<CodexError.CodexAppServerError>();
const lifecycle: Array<"notification" | "termination"> = [];
const transport = yield* CodexProtocol.makeCodexAppServerPatchedProtocol({
stdio,
onNotification: () =>
Effect.sync(() => {
lifecycle.push("notification");
}),
onTermination: (error) =>
Effect.sync(() => {
lifecycle.push("termination");
}).pipe(Effect.andThen(Deferred.succeed(termination, error)), Effect.asVoid),
});
const notification = yield* transport.incomingNotifications.pipe(
Stream.take(1),
Stream.runCollect,
Effect.forkScoped,
);
const encoded = encoder.encode(
encodeUnknownJsonString({ method: "thread/started", params: { threadId: "thread-1" } }),
);

yield* Queue.offer(input, encoded.slice(0, 11));
yield* Queue.offer(input, encoded.slice(11));
yield* Queue.end(input);

assert.deepEqual(yield* Fiber.join(notification), [
{ method: "thread/started", params: { threadId: "thread-1" } },
]);
assert.instanceOf(
yield* Deferred.await(termination),
CodexError.CodexAppServerInputStreamEndedError,
);
assert.deepEqual(lifecycle, ["notification", "termination"]);
}),
);

it.effect("surfaces JSON encoding failures as protocol parse errors", () =>
Effect.gen(function* () {
const { stdio } = yield* makeInMemoryStdio();
Expand Down
16 changes: 7 additions & 9 deletions packages/effect-codex-app-server/src/protocol.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import * as Stdio from "effect/Stdio";
import * as Stream from "effect/Stream";

import * as CodexError from "./errors.ts";
import { makeLineFramer } from "./lineFramer.ts";
import { JsonRpcId, JsonRpcResponseEnvelope } from "./_internal/shared.ts";
const isJsonRpcId = Schema.is(JsonRpcId);
const isJsonRpcResponseEnvelope = Schema.is(JsonRpcResponseEnvelope);
Expand Down Expand Up @@ -157,7 +158,7 @@ export const makeCodexAppServerPatchedProtocol = Effect.fn("makeCodexAppServerPa
const incomingRequests = yield* Queue.unbounded<CodexAppServerIncomingRequest>();
const pending = yield* Ref.make(new Map<string, CodexAppServerPendingRequest>());
const nextRequestId = yield* Ref.make(1);
const remainder = yield* Ref.make("");
const lineFramer = makeLineFramer();
const terminationHandled = yield* Ref.make(false);

const logProtocol = (event: CodexAppServerProtocolLogEvent) => {
Expand Down Expand Up @@ -354,21 +355,18 @@ export const makeCodexAppServerPatchedProtocol = Effect.fn("makeCodexAppServerPa
yield* options.stdio.stdin.pipe(
Stream.decodeText(),
Stream.runForEach((chunk) =>
Ref.modify(remainder, (current) => {
const combined = current + chunk;
const lines = combined.split("\n");
const nextRemainder = lines.pop() ?? "";
return [lines.map((line) => line.replace(/\r$/, "")), nextRemainder] as const;
}).pipe(Effect.flatMap((lines) => Effect.forEach(lines, handleLine, { discard: true }))),
Effect.forEach(lineFramer.push(chunk), handleLine, { discard: true }),
),
Effect.matchEffect({
onFailure: (error) =>
handleTermination(() =>
Effect.succeed(normalizeIncomingError(error, "read-input-stream")),
),
onSuccess: () =>
Ref.get(remainder).pipe(
Effect.flatMap((line) => (line.trim().length === 0 ? Effect.void : handleLine(line))),
Effect.succeed(lineFramer.finish()).pipe(
Effect.flatMap((line) =>
line === undefined || line.trim().length === 0 ? Effect.void : handleLine(line),
),
Effect.matchEffect({
onFailure: (error) => handleTermination(() => Effect.succeed(error)),
onSuccess: () =>
Expand Down
Loading