Skip to content
Merged
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
60 changes: 60 additions & 0 deletions apps/server/src/project/AgentSessionJson.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,60 @@
import { describe, expect, it } from "@effect/vitest";

import { createTranscriptJsonReader, TranscriptJsonLimitError } from "./AgentSessionJson.ts";

function read(
json: string,
size: number,
select: (path: ReadonlyArray<string | number | null>) => boolean = () => true,
) {
const reader = createTranscriptJsonReader(() => {}, select);
for (let offset = 0; offset < json.length; offset += size)
reader.write(json.slice(offset, offset + size));
return reader.finish();
}

describe("transcript JSON projection", () => {
it.each([1, 2, 7, 64, 1024])("matches JSON.parse across %i-character boundaries", (size) => {
for (const text of [
'{"a":1,"a":2,"b":"before","b":"after"}',
'{"a":{"x":1},"a":{"y":2},"b":[],"c":{}}',
'{"a":[null,true,false,1,-2.3e4,"😀\\u0061\\\\\\\"",{},[],[1,2]]}',
'{"a":"s","a":null,"b":null,"b":"s","c":0,"c":false}',
'{"__proto__":{"polluted":true},"constructor":1,"__proto__":2}',
])
expect(read(text, size)).toEqual(JSON.parse(text));
});

it("projects siblings and array elements without merging repeated parent objects", () => {
const text =
'{"message":{"usage":{"input":100},"content":"large"},"message":{"usage":{"output":5},"content":[1,2]},"rows":[{"keep":1,"drop":2},{"keep":3}],"drop":{"keep":4}}';
const projected = read(text, 1, (path) => {
if (path[0] === "drop") return false;
return !path.includes("content") && !path.includes("drop");
});
expect(projected).toEqual({
message: { usage: { output: 5 } },
rows: [{ keep: 1 }, { keep: 3 }],
});
});

it.each(['{"a":', '{"a":1} trailing', '{"a":1}{"a":2}', '{"a":"bad\\x"}', '{"a":[1,]}'])(
"rejects malformed input %s",
(text) => {
expect(read(text, 1)).toBeUndefined();
},
);

it("retains the import allocation and depth limits", () => {
const limited = createTranscriptJsonReader(
() => {
throw new TranscriptJsonLimitError("budget");
},
() => true,
);
expect(() => limited.write('{"a":1}')).toThrow(TranscriptJsonLimitError);
expect(() => read("[".repeat(129) + "0" + "]".repeat(129), 10)).toThrow(
TranscriptJsonLimitError,
);
});
});
72 changes: 59 additions & 13 deletions apps/server/src/project/AgentSessionJson.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,6 @@
import * as SchemaAST from "effect/SchemaAST";
import { isMany, none, type Many } from "stream-chain/defs.js";
import { none, type Many } from "stream-chain/defs.js";
import { Assembler } from "stream-json/core/assembler.js";
import { filter } from "stream-json/core/filters/filter.js";
import * as StreamJson from "stream-json/core/parser.js";
import type { ParserOptions, Token } from "stream-json/core/parser.js";

Expand Down Expand Up @@ -55,6 +54,7 @@ export class TranscriptJsonLimitError extends Error {}
export function createTranscriptJsonReader(
reserve: (bytes: number) => void,
selectPath: (path: JsonPath) => boolean,
options?: { readonly maxDepth?: number },
) {
// The synchronous tokenizer is exported at runtime in 3.6.0, but omitted
// from its bundled types. Unlike parser(), it does not wrap tokens in an
Expand All @@ -65,9 +65,6 @@ export function createTranscriptJsonReader(
) => (input: string | typeof none) => Many<Token> | typeof none;
};
const tokenize = jsonParser({ packValues: false });
const select = filter({ filter: selectPath, streamKeys: false }) as (
input: Token | typeof none,
) => Token | Many<Token> | typeof none;
const assembler = new Assembler();
let key: string | null = null;
let value = "";
Expand Down Expand Up @@ -100,13 +97,62 @@ export function createTranscriptJsonReader(
assembler.consume(token);
}
};
// Forward actual selected keys instead of reconstructing them from path
// changes: adjacent duplicate keys have the same path but JSON.parse keeps
// the last value. Reconstructing paths can silently retain the first value.
const stack: Array<{ path: JsonPath; key: string | number | null; selected: boolean }> = [];
let selectedValue = false;
const startValue = () => {
const parent = stack.at(-1);
const path = parent?.selected ? [...parent.path, parent.key] : [];
const selected = (parent?.selected ?? true) && selectPath(path);
if (selected && typeof parent?.key === "string") {
assemble({ name: "keyValue", value: parent.key });
}
return { path, selected };
};
const endValue = () => {
const parent = stack.at(-1);
if (parent && typeof parent.key === "number") parent.key++;
};
const selectToken = (token: Token | typeof none) => {
const selected = select(token);
if (selected === none) return;
if (isMany(selected)) {
for (const item of selected.values) assemble(item);
} else {
assemble(selected);
if (token === none) return;
switch (token.name) {
case "keyValue": {
const parent = stack.at(-1);
if (parent) parent.key = token.value;
return;
}
case "startObject":
case "startArray": {
const frame = startValue();
stack.push({ ...frame, key: token.name === "startArray" ? 0 : null });
if (frame.selected) assemble(token);
return;
}
case "endObject":
case "endArray":
if (stack.pop()?.selected) assemble(token);
endValue();
return;
case "startString":
case "startNumber":
selectedValue = startValue().selected;
if (selectedValue) assemble(token);
return;
case "endString":
case "endNumber":
if (selectedValue) assemble(token);
endValue();
return;
case "nullValue":
case "trueValue":
case "falseValue":
if (startValue().selected) assemble(token);
endValue();
return;
default:
if (selectedValue) assemble(token);
}
};
const consume = (input: string | typeof none) => {
Expand All @@ -116,8 +162,8 @@ export function createTranscriptJsonReader(
if (tokens === none) return;
for (const token of tokens.values) {
if (token.name === "startObject" || token.name === "startArray") {
if (++depth > 128)
throw new TranscriptJsonLimitError("Transcript JSON nesting exceeds 128 levels");
if (++depth > (options?.maxDepth ?? 128))
throw new TranscriptJsonLimitError("Transcript JSON nesting exceeds the depth limit");
} else if (token.name === "endObject" || token.name === "endArray") {
if (--depth === 0) complete = true;
}
Expand Down
53 changes: 53 additions & 0 deletions apps/server/src/usage/UsageService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -645,6 +645,59 @@ describe("UsageService", () => {
}).pipe(Effect.scoped),
);

it.live(
"keeps large-record totals and costs exact through append, dedupe, restart and cleanup",
() =>
Effect.gen(function* () {
const { transcript, settings, home } = yield* setup;
const large = claudeLine(1, 9900).replace(
'"message":',
'"padding":' + encodeUnknownJsonString("x".repeat(9 * 1024 * 1024)) + ',"message":',
);
yield* Effect.promise(() => NodeFSP.writeFile(transcript, large));
yield* Effect.gen(function* () {
const service = yield* UsageService.make;
const first = yield* service.readSummary(WINDOW);
assert.strictEqual(totalOutputTokens(first), 9900);
assert.closeTo(
first.buckets.reduce((sum, bucket) => sum + bucket.costUsd, 0),
0.4951,
1e-12,
);
const warm = yield* service.readSummary(WINDOW);
assert.deepStrictEqual(warm.buckets, first.buckets);
// The repeated content block has the same message/request identity.
yield* Effect.promise(() => NodeFSP.appendFile(transcript, large + claudeLine(2, 100)));
const appended = yield* service.readSummary(WINDOW);
assert.strictEqual(totalOutputTokens(appended), 10000);
assert.strictEqual(
appended.buckets.reduce((sum, bucket) => sum + bucket.totals.uncachedInputTokens, 0),
20,
);
const restarted = yield* UsageService.make;
const restored = yield* restarted.readSummary(WINDOW);
assert.deepStrictEqual(restored.buckets, appended.buckets);
yield* Effect.promise(() => NodeFSP.rm(transcript));
const afterCleanup = yield* UsageService.make;
assert.deepStrictEqual(
(yield* afterCleanup.readSummary(WINDOW)).buckets,
appended.buckets,
);
}).pipe(
Effect.provide(
serviceLayers({
prefix: "usage-service-large-record-test",
home,
settings,
ratesDocument: {
"claude-fable-5": { input_cost_per_token: 1e-5, output_cost_per_token: 5e-5 },
},
}),
),
);
}).pipe(Effect.scoped),
);

it.live("preserves saved tokens, costs and sessions after transcript cleanup and restart", () =>
Effect.gen(function* () {
const { transcript, settings, home } = yield* setup;
Expand Down
Loading
Loading