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
39 changes: 39 additions & 0 deletions apps/server/src/provider/Layers/CodexAdapter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
22 changes: 21 additions & 1 deletion apps/server/src/provider/Layers/CodexAdapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, unknown>;
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,
Expand Down Expand Up @@ -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,
Expand Down
233 changes: 233 additions & 0 deletions apps/server/src/provider/Layers/EventNdjsonLogger.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down Expand Up @@ -127,6 +128,164 @@ 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"');

// 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 });
}
}),
);

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: 64,
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<string, string>;
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<string, string>;
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",
() =>
Expand Down Expand Up @@ -453,6 +612,80 @@ 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("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-"));
Expand Down
Loading
Loading