diff --git a/.github/workflows/node-server-tests.yml b/.github/workflows/node-server-tests.yml index 0b121c5bb2e..90c3fc68473 100644 --- a/.github/workflows/node-server-tests.yml +++ b/.github/workflows/node-server-tests.yml @@ -164,7 +164,7 @@ jobs: cache-pnpm-store-lookup-only: 'true' # Design D7 upgrade and rollback: the last Bun orcad, built from a main commit that shipped # it, beside this checkout's Node slot; the live-terminal hand-over skips once PROTOCOL_VERSION - # moves past the Bun daemon's. Its build uses this checkout's installed dependencies. + # moves past the Bun daemon's. Install its pinned dependencies independently of this checkout. - uses: oven-sh/setup-bun@0c5077e51419868618aeaa5fe8019c62421857d6 # v2.2.0 if: runner.os == 'Linux' with: @@ -179,7 +179,7 @@ jobs: if [ "$RUNNER_OS" != Linux ]; then exit 0; fi git fetch --no-tags --depth=1 origin "$BUN_ORCAD_COMMIT" git worktree add --detach "$RUNNER_TEMP/bun-orcad-source" "$BUN_ORCAD_COMMIT" - ln -s "$GITHUB_WORKSPACE/node_modules" "$RUNNER_TEMP/bun-orcad-source/node_modules" + pnpm --dir "$RUNNER_TEMP/bun-orcad-source" install --frozen-lockfile --ignore-scripts node "$RUNNER_TEMP/bun-orcad-source/config/scripts/build-orcad-bun.mjs" --out-dir "$RUNNER_TEMP/bun-orcad" echo "slot=$RUNNER_TEMP/bun-orcad" >> "$GITHUB_OUTPUT" echo "executable=$(command -v bun)" >> "$GITHUB_OUTPUT" diff --git a/config/patches/@streamparser__json@0.0.26.patch b/config/patches/@streamparser__json@0.0.26.patch deleted file mode 100644 index c93998f85fe..00000000000 --- a/config/patches/@streamparser__json@0.0.26.patch +++ /dev/null @@ -1,66 +0,0 @@ -diff --git a/dist/cjs/utils/bufferedString.js b/dist/cjs/utils/bufferedString.js -index 82f710a018f9771fe10335e2dcacd75d707d9062..f3cfa64cefa967d7a83c328e1abb9246135b067b 100644 ---- a/dist/cjs/utils/bufferedString.js -+++ b/dist/cjs/utils/bufferedString.js -@@ -16,7 +16,7 @@ class NonBufferedString { - constructor() { - // fatal: true makes invalid byte sequences (e.g. a lead byte followed by a - // non-continuation byte) throw instead of silently decoding to U+FFFD. -- this.decoder = new TextDecoder("utf-8", { fatal: true }); -+ this.decoder = new TextDecoder("utf-8", { fatal: true, ignoreBOM: true }); - // Pieces appended since the last toString(), not yet folded into `string`. - this.pending = []; - this.string = ""; -@@ -66,7 +66,7 @@ class BufferedString { - constructor(bufferSize) { - // fatal: true makes invalid byte sequences (e.g. a lead byte followed by a - // non-continuation byte) throw instead of silently decoding to U+FFFD. -- this.decoder = new TextDecoder("utf-8", { fatal: true }); -+ this.decoder = new TextDecoder("utf-8", { fatal: true, ignoreBOM: true }); - this.bufferOffset = 0; - this.string = ""; - this.byteLength = 0; -diff --git a/dist/mjs/utils/bufferedString.js b/dist/mjs/utils/bufferedString.js -index 0fb208d8615f5e20a086f75a37bef928b155c47c..0d8ca407d594d61798eca7f7253e1dfc1770d291 100644 ---- a/dist/mjs/utils/bufferedString.js -+++ b/dist/mjs/utils/bufferedString.js -@@ -13,7 +13,7 @@ export class NonBufferedString { - constructor() { - // fatal: true makes invalid byte sequences (e.g. a lead byte followed by a - // non-continuation byte) throw instead of silently decoding to U+FFFD. -- this.decoder = new TextDecoder("utf-8", { fatal: true }); -+ this.decoder = new TextDecoder("utf-8", { fatal: true, ignoreBOM: true }); - // Pieces appended since the last toString(), not yet folded into `string`. - this.pending = []; - this.string = ""; -@@ -62,7 +62,7 @@ export class BufferedString { - constructor(bufferSize) { - // fatal: true makes invalid byte sequences (e.g. a lead byte followed by a - // non-continuation byte) throw instead of silently decoding to U+FFFD. -- this.decoder = new TextDecoder("utf-8", { fatal: true }); -+ this.decoder = new TextDecoder("utf-8", { fatal: true, ignoreBOM: true }); - this.bufferOffset = 0; - this.string = ""; - this.byteLength = 0; -diff --git a/src/utils/bufferedString.ts b/src/utils/bufferedString.ts -index 482c7402899bb157249cfb7882d327b7d9477912..4e45ef5d5bd7ec66776ece54e917282441cededb 100644 ---- a/src/utils/bufferedString.ts -+++ b/src/utils/bufferedString.ts -@@ -40,7 +40,7 @@ export interface StringBuilder { - export class NonBufferedString implements StringBuilder { - // fatal: true makes invalid byte sequences (e.g. a lead byte followed by a - // non-continuation byte) throw instead of silently decoding to U+FFFD. -- private decoder = new TextDecoder("utf-8", { fatal: true }); -+ private decoder = new TextDecoder("utf-8", { fatal: true, ignoreBOM: true }); - // Pieces appended since the last toString(), not yet folded into `string`. - private pending: string[] = []; - private string = ""; -@@ -90,7 +90,7 @@ export class NonBufferedString implements StringBuilder { - export class BufferedString implements StringBuilder { - // fatal: true makes invalid byte sequences (e.g. a lead byte followed by a - // non-continuation byte) throw instead of silently decoding to U+FFFD. -- private decoder = new TextDecoder("utf-8", { fatal: true }); -+ private decoder = new TextDecoder("utf-8", { fatal: true, ignoreBOM: true }); - private buffer: Uint8Array; - private bufferOffset = 0; - private string = ""; diff --git a/config/scripts/json-parser-benchmark-fixtures.mjs b/config/scripts/json-parser-benchmark-fixtures.mjs new file mode 100644 index 00000000000..385b4010895 --- /dev/null +++ b/config/scripts/json-parser-benchmark-fixtures.mjs @@ -0,0 +1,65 @@ +import { writeFileSync } from 'node:fs' +import { join } from 'node:path' + +export const JSON_PARSER_CASES = [ + { name: 'rg-10k', kind: 'rg', count: 10_000, cap: 2000 }, + { name: 'rg-100k', kind: 'rg', count: 100_000, cap: 2000 }, + { name: 'rg-100k-cap1', kind: 'rg', count: 100_000, cap: 1 }, + { name: 'rg-unicode', kind: 'rg', count: 10_000, cap: 2000, unicode: true }, + { name: 'rg-large-dense', kind: 'rg', count: 900_000, cap: 2000, large: true }, + { name: 'rg-fast-path', kind: 'rg', count: 1, cap: 2000, large: true }, + { name: 'session-messages', kind: 'session' }, + { name: 'session-skipped-objects', kind: 'session' }, + { name: 'session-skipped-string', kind: 'session' }, + { name: 'session-selected-string', kind: 'session' } +] + +export function writeJsonParserFixtures(directory) { + for (const fixture of JSON_PARSER_CASES) { + let value + if (fixture.kind === 'rg') { + const text = fixture.unicode + ? '\ufeff日本語😀x' + : fixture.large && fixture.count > 1 + ? 'xxxx' + : 'x' + const matchBytes = Buffer.byteLength(text) + value = { + type: 'match', + data: { + path: { text: `${text}.ts` }, + lines: { + text: text.repeat(fixture.count === 1 ? 4 * 1024 * 1024 : fixture.count) + }, + line_number: 1, + submatches: Array.from({ length: fixture.count }, (_, index) => ({ + match: { text }, + start: index * matchBytes, + end: (index + 1) * matchBytes + })) + } + } + } else { + value = { id: 'synthetic', messages: [{ text: 'one' }] } + if (fixture.name === 'session-messages') { + value.agent = { model: 'model', other: 'ignored' } + value.messages = Array.from({ length: 20_000 }, (_, index) => ({ + role: index % 2 ? 'assistant' : 'user', + text: '\ufeff日本語😀 hello world '.repeat(16), + timestamp: index, + metadata: { model: 'synthetic', tokens: 512 } + })) + } else if (fixture.name === 'session-skipped-objects') { + value.ignored = Array.from({ length: 200_000 }, (_, index) => ({ + id: index, + data: { text: 'x'.repeat(64), values: [1, 2, 3] } + })) + } else if (fixture.name === 'session-skipped-string') { + value.ignored = '日本語😀x'.repeat(1_500_000) + } else { + value.messages = [{ text: '日本語😀x'.repeat(1_500_000) }] + } + } + writeFileSync(join(directory, `${fixture.name}.json`), JSON.stringify(value)) + } +} diff --git a/config/scripts/json-parser-migration-benchmark.mjs b/config/scripts/json-parser-migration-benchmark.mjs new file mode 100644 index 00000000000..0aca1734e1b --- /dev/null +++ b/config/scripts/json-parser-migration-benchmark.mjs @@ -0,0 +1,166 @@ +import assert from 'node:assert/strict' +import { createHash } from 'node:crypto' +import { spawnSync } from 'node:child_process' +import { createReadStream, readFileSync, statSync, mkdtempSync, rmSync } from 'node:fs' +import { tmpdir } from 'node:os' +import { join, resolve } from 'node:path' +import { pathToFileURL } from 'node:url' +import { buildCounterbalancedSchedule } from './counterbalanced-benchmark-schedule.mjs' +import { summarizeBenchmarkSamples } from './benchmark-sample-summary.mjs' +import { JSON_PARSER_CASES, writeJsonParserFixtures } from './json-parser-benchmark-fixtures.mjs' + +// Bundle each revision's two consumers as {baseline,candidate}-{rg,session}.mjs first. +const [bundleDirectory, workerFixture, workerArm] = process.argv.slice(2) +if (!bundleDirectory || !global.gc) { + throw new Error('Usage: node --expose-gc json-parser-migration-benchmark.mjs BUNDLE_DIRECTORY') +} + +async function load(arm, kind) { + return import(pathToFileURL(resolve(bundleDirectory, `${arm}-${kind}.mjs`)).href) +} + +function prepareRun(module, fixture, file) { + if (fixture.kind === 'rg') { + const text = readFileSync(file, 'utf8') + return () => + module.parseRipgrepMatchJson(text, fixture.cap, { + structuralTokens: 32 * 1024, + nestingDepth: 16 + }) + } + return () => + module.readStreamedSessionDocument({ + bytes: createReadStream(file, { highWaterMark: 64 * 1024 }), + arrayKey: 'messages', + fields: ['id'], + objectFields: { agent: ['model'] }, + create: () => ({ count: 0, textLength: 0 }), + consume(state, value) { + state.count++ + state.textLength += typeof value?.text === 'string' ? value.text.length : 0 + } + }) +} + +function digest(value) { + return createHash('sha256').update(JSON.stringify(value)).digest('hex') +} + +async function consumedContentDigest(module, file) { + const result = await module.readStreamedSessionDocument({ + bytes: createReadStream(file, { highWaterMark: 64 * 1024 }), + arrayKey: 'messages', + fields: ['id'], + objectFields: { agent: ['model'] }, + create: () => createHash('sha256'), + consume(hash, value) { + hash.update(JSON.stringify(value)).update('\n') + } + }) + return { record: result.record, consumedSha256: result.state.digest('hex') } +} + +if (workerFixture) { + const fixture = JSON_PARSER_CASES.find((item) => workerFixture.endsWith(`${item.name}.json`)) + assert(fixture) + const module = await load(workerArm, fixture.kind) + const run = prepareRun(module, fixture, workerFixture) + global.gc() + const before = process.memoryUsage() + let running = true + let maxLoopGapMs = 0 + let previous = performance.now() + const observe = () => { + const now = performance.now() + maxLoopGapMs = Math.max(maxLoopGapMs, now - previous) + previous = now + if (running) { + setImmediate(observe) + } + } + setImmediate(observe) + const started = performance.now() + const result = await run() + const elapsedMs = performance.now() - started + await new Promise((done) => setImmediate(done)) + running = false + const peakRssMiB = process.resourceUsage().maxRSS / 1024 + global.gc() + const retainedHeapDeltaMiB = (process.memoryUsage().heapUsed - before.heapUsed) / 1024 ** 2 + console.log( + JSON.stringify({ + elapsedMs, + peakRssMiB, + retainedHeapDeltaMiB, + maxLoopGapMs, + digest: digest(result) + }) + ) +} else { + const directory = mkdtempSync(join(tmpdir(), 'orca-json-parser-benchmark-')) + try { + writeJsonParserFixtures(directory) + global.gc() + const results = [] + for (const fixture of JSON_PARSER_CASES) { + const file = join(directory, `${fixture.name}.json`) + const runs = {} + const contents = {} + for (const arm of ['baseline', 'candidate']) { + const module = await load(arm, fixture.kind) + runs[arm] = prepareRun(module, fixture, file) + if (fixture.kind === 'session') { + contents[arm] = await consumedContentDigest(module, file) + } + } + assert.deepEqual(contents.candidate, contents.baseline) + for (let warmup = 0; warmup < 3; warmup++) { + assert.deepEqual(await runs.candidate(), await runs.baseline()) + } + const samples = { baseline: [], candidate: [] } + for (const pair of buildCounterbalancedSchedule(12, 'baseline', 'candidate')) { + for (const arm of pair) { + const started = performance.now() + await runs[arm]() + samples[arm].push(performance.now() - started) + } + } + const memory = { baseline: [], candidate: [] } + for (const pair of buildCounterbalancedSchedule(2, 'baseline', 'candidate')) { + for (const arm of pair) { + const child = spawnSync( + process.execPath, + ['--expose-gc', import.meta.filename, bundleDirectory, file, arm], + { + encoding: 'utf8', + env: { ...process.env, ORCA_BACKGROUND_LAUNCH: '1' }, + windowsHide: true + } + ) + assert.equal(child.status, 0, child.stderr) + memory[arm].push(JSON.parse(child.stdout)) + } + } + for (const sample of [...memory.baseline, ...memory.candidate]) { + assert.equal(sample.digest, memory.baseline[0].digest) + } + results.push({ + name: fixture.name, + bytes: statSync(file).size, + baseline: summarizeBenchmarkSamples(samples.baseline), + candidate: summarizeBenchmarkSamples(samples.candidate), + samples, + memory + }) + } + console.log( + JSON.stringify( + { node: process.version, platform: process.platform, arch: process.arch, results }, + null, + 2 + ) + ) + } finally { + rmSync(directory, { recursive: true, force: true }) + } +} diff --git a/config/scripts/node-server-change-scope.test.mjs b/config/scripts/node-server-change-scope.test.mjs index 4216beb2556..38b56f0d600 100644 --- a/config/scripts/node-server-change-scope.test.mjs +++ b/config/scripts/node-server-change-scope.test.mjs @@ -310,6 +310,10 @@ it('runs the Bun and Node cross-runtime tests on Linux against pinned inputs', ( expect(build.if).toBeUndefined() expect(build['continue-on-error']).toBeUndefined() expect(build.run).toMatch(/^if \[ "\$RUNNER_OS" != Linux \]; then exit 0; fi\n/) + expect(build.run).toContain( + 'pnpm --dir "$RUNNER_TEMP/bun-orcad-source" install --frozen-lockfile --ignore-scripts' + ) + expect(build.run).not.toContain('"$GITHUB_WORKSPACE/node_modules"') expect(build.run).toContain('echo "slot=$RUNNER_TEMP/bun-orcad" >> "$GITHUB_OUTPUT"') expect(build.run).toContain('echo "executable=$(command -v bun)" >> "$GITHUB_OUTPUT"') expect(build.run).not.toContain('GITHUB_ENV') diff --git a/config/tsconfig.cli.json b/config/tsconfig.cli.json index 44edcfd90b7..6afcf4f142f 100644 --- a/config/tsconfig.cli.json +++ b/config/tsconfig.cli.json @@ -258,8 +258,8 @@ ], "compilerOptions": { "composite": true, - // TypeScript 7 removed node10 resolution; Node16 preserves CommonJS emit for this package. - "module": "Node16", + // The CLI runs on Node 24; Node20 models synchronous ESM imports from CommonJS. + "module": "Node20", "moduleResolution": "Node16", "rootDir": "../src", "outDir": "../out" diff --git a/config/tsconfig.tc.cli.json b/config/tsconfig.tc.cli.json index bd59f47e7be..c613dd688d9 100644 --- a/config/tsconfig.tc.cli.json +++ b/config/tsconfig.tc.cli.json @@ -1,7 +1,7 @@ { "extends": "./tsconfig.cli.json", "compilerOptions": { - "module": "node16", + "module": "Node20", "moduleResolution": "node16" } } diff --git a/electron.vite.config.ts b/electron.vite.config.ts index a899fd97f9f..ac23e12afba 100644 --- a/electron.vite.config.ts +++ b/electron.vite.config.ts @@ -13,7 +13,8 @@ import { import packageJson from './package.json' with { type: 'json' } const BUNDLED_MAIN_DEPENDENCIES = new Set([ - '@streamparser/json', + 'stream-json', + 'stream-chain', '@xterm/headless', '@xterm/addon-serialize', 'tldts', diff --git a/package.json b/package.json index 3e67dfb4813..3c1979d5ff0 100644 --- a/package.json +++ b/package.json @@ -182,7 +182,6 @@ "@floating-ui/dom": "1.8.0", "@linear/sdk": "^97.0.0", "@parcel/watcher": "^2.5.6", - "@streamparser/json": "0.0.26", "@xterm/addon-serialize": "0.15.0-beta.300", "@xterm/headless": "6.1.0-beta.302", "agent-browser": "~0.27.0", @@ -198,6 +197,8 @@ "sherpa-onnx": "1.12.37", "smol-toml": "1.8.0", "ssh2": "^1.17.0", + "stream-chain": "4.2.6", + "stream-json": "3.7.0", "tldts": "7.4.16", "tweetnacl": "^1.0.3", "ws": "^8.22.0", diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 2c58ae7a17a..fbf12ea1ae1 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -166,7 +166,6 @@ overrides: monaco-editor>dompurify: 3.4.16 patchedDependencies: - '@streamparser/json@0.0.26': b2cf43861e5b4e485e97ffa7449ea65acab4d65dd983ab8c0508f1c68f7d80c9 '@vscode/windows-process-tree@0.8.0': 9da74aa3d17243aa53dcdc95c9f06e97437e7fbccf098aeb017579e2d24cbac2 '@xterm/addon-image@0.10.0-beta.300': e5254a46d6f57bef4a8a19683bfa685afa0ca0127545aea53b54a48104ca3562 '@xterm/addon-ligatures@0.11.0-beta.300': 47405b9994b5acf1b4e90b49250358c1ca03649854d59560e7732b72fe336920 @@ -197,9 +196,6 @@ importers: '@parcel/watcher': specifier: ^2.5.6 version: 2.5.6 - '@streamparser/json': - specifier: 0.0.26 - version: 0.0.26(patch_hash=b2cf43861e5b4e485e97ffa7449ea65acab4d65dd983ab8c0508f1c68f7d80c9) '@xterm/addon-serialize': specifier: 0.15.0-beta.300 version: 0.15.0-beta.300(patch_hash=b35533fe252e7e45433150170348889f4e08a6c17f7017ac34ea694d831fec7f)(@xterm/xterm@6.1.0-beta.303(patch_hash=dd0ccc59cd1ccf99f4d76e5aa2456da165fa0804dce19a833d7638bd07ffa393)) @@ -245,6 +241,12 @@ importers: ssh2: specifier: ^1.17.0 version: 1.17.0 + stream-chain: + specifier: 4.2.6 + version: 4.2.6 + stream-json: + specifier: 3.7.0 + version: 3.7.0 tldts: specifier: 7.4.16 version: 7.4.16 @@ -3186,9 +3188,6 @@ packages: '@standard-schema/spec@1.1.0': resolution: {integrity: sha512-l2aFy5jALhniG5HgqrD6jXLi/rUWrKvqN/qJx6yoJsgKhblVd+iqqU4RCXavm/jPityDo5TCvKMnpjKnOriy0w==} - '@streamparser/json@0.0.26': - resolution: {integrity: sha512-46597LNFI+MFdUnzX2QJWwmdTRdq0XVD+vVNJTtGVzIrnCuhG9pFo1OAzbNBqci8UJgk/X5KJZ6LcV+y7PTuDQ==} - '@swc/core-darwin-arm64@1.15.46': resolution: {integrity: sha512-IsISIT22EfktVJrlvIpnAxG2u/A9aob9l99HMlx80x72WlFmFPk1V3UhkEzx86eJP8hw049KTFv/RISho2cq2Q==} engines: {node: '>=10'} @@ -7355,6 +7354,12 @@ packages: resolution: {integrity: sha512-eCPu1qRxPVkl5605OTWF8Wz40b4Mf45NY5LQmVPQ599knfs5QhASUm9GbJ5BDMDOXgrnh0wyEdvzmL//YMlw0A==} engines: {node: '>=18'} + stream-chain@4.2.6: + resolution: {integrity: sha512-1zeJ8CrtJfmiba26ui8jXkq/xLRFvhzkdH02D5QLO9Cnovgeb28IJxg88DdraWnjnkbOol/++uIAybJbJhk7ig==} + + stream-json@3.7.0: + resolution: {integrity: sha512-rCSBdcBP/bPk6T8QFcxAj1MSzAuc5i49cYW6IE7sYObNEPccBJIUiL6fU9c9BWt40aJK0siVynLX7lWaW8alXw==} + strict-event-emitter@0.5.1: resolution: {integrity: sha512-vMgjE/GGEPEFnhFub6pa4FmJBRBVOLpIII2hvCZ8Kzb7K0hlHo7mQv6xYrBvCL2LtAIBwFUK8wvuJgTVSQ5MFQ==} @@ -10209,8 +10214,6 @@ snapshots: '@standard-schema/spec@1.1.0': {} - '@streamparser/json@0.0.26(patch_hash=b2cf43861e5b4e485e97ffa7449ea65acab4d65dd983ab8c0508f1c68f7d80c9)': {} - '@swc/core-darwin-arm64@1.15.46': optional: true @@ -14966,6 +14969,12 @@ snapshots: stdin-discarder@0.3.2: {} + stream-chain@4.2.6: {} + + stream-json@3.7.0: + dependencies: + stream-chain: 4.2.6 + strict-event-emitter@0.5.1: optional: true diff --git a/pnpm-workspace.yaml b/pnpm-workspace.yaml index f00e23d7f2f..e3c133b29ee 100644 --- a/pnpm-workspace.yaml +++ b/pnpm-workspace.yaml @@ -63,4 +63,3 @@ patchedDependencies: lint-staged@16.4.0: config/patches/lint-staged@16.4.0.patch '@vscode/windows-process-tree@0.8.0': config/patches/@vscode__windows-process-tree@0.8.0.patch i18next-cli@1.74.2: config/patches/i18next-cli@1.74.2.patch - '@streamparser/json@0.0.26': config/patches/@streamparser__json@0.0.26.patch diff --git a/src/main/ai-vault/session-document-stream-contract.test.ts b/src/main/ai-vault/session-document-stream-contract.test.ts new file mode 100644 index 00000000000..3ab4053b303 --- /dev/null +++ b/src/main/ai-vault/session-document-stream-contract.test.ts @@ -0,0 +1,166 @@ +import { describe, expect, it } from 'vitest' +import { readStreamedSessionDocument } from './session-document-stream' + +async function* bytes(content: string, chunkSize = 7): AsyncGenerator { + const buffer = Buffer.from(content) + for (let offset = 0; offset < buffer.length; offset += chunkSize) { + yield buffer.subarray(offset, offset + chunkSize) + } +} + +function read( + content: string, + overrides: Partial>[0]> = {} +) { + return readStreamedSessionDocument({ + bytes: bytes(content), + arrayKey: 'messages', + fields: ['id'], + objectFields: { agent: ['model'] }, + create: (): unknown[] => [], + consume: (state, value) => { + state.push(value) + }, + ...overrides + }) +} + +describe('streamed session document contracts', () => { + it.each(['[]', 'null', 'false', '0', '"not an array"', '{}'])( + 'a later messages value %s clears an earlier fold', + async (last) => { + expect(await read(`{"messages":[1,2],"messages":${last}}`)).toEqual({ record: {}, state: [] }) + } + ) + + it('replaces a failed earlier array and continues with its successor', async () => { + const consume = (state: unknown[], value: unknown): void => { + if (value === 'bad') { + throw new Error('discarded failure') + } + state.push(value) + } + expect(await read('{"messages":["bad",1],"messages":[2,3]}', { consume })).toEqual({ + record: {}, + state: [2, 3] + }) + expect(await read('{"messages":["bad"],"messages":[]}', { consume })).toEqual({ + record: {}, + state: [] + }) + }) + + it('keeps the first consumer failure unless a later array replaces it', async () => { + const failure = new Error('consumer failed') + let calls = 0 + await expect( + read('{"messages":[1,2]}', { + consume: () => { + calls++ + throw failure + } + }) + ).rejects.toBe(failure) + expect(calls).toBe(1) + }) + + it('gives malformed trailing JSON precedence over a consumer failure', async () => { + await expect( + read('{"messages":[1]} trailing', { + consume: () => { + throw new Error('consumer') + } + }) + ).rejects.toBeInstanceOf(SyntaxError) + }) + + it('keeps projected object resets and full-field overlap precedence', async () => { + expect(await read('{"agent":{"model":"old"},"agent":null}')).toEqual({ + record: { agent: {} }, + state: [] + }) + expect(await read('{"agent":{"model":"m","extra":[1,2]}}', { fields: ['agent'] })).toEqual({ + record: { agent: { model: 'm', extra: [1, 2] } }, + state: [] + }) + expect(await read('{"messages":[1,2]}', { objectFields: { messages: ['model'] } })).toEqual({ + record: { messages: {} }, + state: [1, 2] + }) + expect( + await read('{"messages":{"model":"m","ignored":1}}', { + objectFields: { messages: ['model'] } + }) + ).toEqual({ record: { messages: { model: 'm' } }, state: [] }) + }) + + it('preserves selected prototype keys as own properties', async () => { + const content = + '{"id":{"__proto__":{"polluted":true}},"agent":{"model":{"constructor":"value","__proto__":{"x":1}}}}' + const parsed = await read(content) + expect(parsed?.record).toEqual(JSON.parse(content)) + expect(Object.getPrototypeOf(parsed?.record)).toBeNull() + expect(Object.prototype).not.toHaveProperty('polluted') + }) + + it.each(['', ' ', '{', '{"messages":[1,]}', '{"messages":[]}{"messages":[]}'])( + 'rejects malformed input %j', + async (content) => { + await expect(read(content)).rejects.toBeInstanceOf(SyntaxError) + } + ) + + it.each(['null', '[]', '123', '"text"'])( + 'does not publish a nonobject document %s', + async (content) => { + expect(await read(content)).toBeNull() + } + ) + + it('validates discarded subtrees and closes their source on malformed input', async () => { + let closed = false + async function* malformed() { + try { + yield Buffer.from('{"ignored":[{"deep":1},]}') + throw new Error('must not request another chunk') + } finally { + closed = true + } + } + await expect(read('', { bytes: malformed() })).rejects.toBeInstanceOf(SyntaxError) + expect(closed).toBe(true) + }) + + it('preserves source failure and closes the source even after a consumer failure', async () => { + const failure = new Error('disk failure') + let closed = false + async function* failing() { + try { + yield Buffer.from('{"messages":[1]') + throw failure + } finally { + closed = true + } + } + await expect( + read('', { + bytes: failing(), + consume: () => { + throw new Error('consumer') + } + }) + ).rejects.toBe(failure) + expect(closed).toBe(true) + }) + + it('preserves escaped astral text across large and byte-sized chunks', async () => { + const value = `${'x'.repeat(65530)}日本語😀\\literal` + const content = JSON.stringify({ messages: [value, '\ud800', '\udc00'] }) + for (const size of [1, 65536, content.length * 4]) { + expect(await read(content, { bytes: bytes(content, size) })).toEqual({ + record: {}, + state: [value, '\ud800', '\udc00'] + }) + } + }) +}) diff --git a/src/main/ai-vault/session-document-stream.ts b/src/main/ai-vault/session-document-stream.ts index 4d12ab0fcc2..4fcf24f7a55 100644 --- a/src/main/ai-vault/session-document-stream.ts +++ b/src/main/ai-vault/session-document-stream.ts @@ -1,6 +1,7 @@ import { StringDecoder } from 'node:string_decoder' -import { JSONParser, TokenizerError, TokenParserError, TokenType } from '@streamparser/json' +import { FlexAssembler, arrayRule, objectRule } from 'stream-json/core/utils/flex-assembler.js' import { setImmediate as yieldToEventLoop } from 'node:timers/promises' +import { createJsonTokenReader } from '../../shared/json-token-reader' import { throwIfAiVaultScanCancelled } from './ai-vault-scan-cancellation' /** Fold one root array while retaining only the root fields the agent parser uses. */ @@ -13,92 +14,108 @@ export async function readStreamedSessionDocument(args: { consume: (state: T, value: unknown) => void signal?: AbortSignal }): Promise<{ record: Record; state: T } | null> { - const parser = new JSONParser({ - paths: [ - ...args.fields.map((field) => `$.${field}`), - ...Object.entries(args.objectFields ?? {}).flatMap(([root, fields]) => - fields.map((field) => `$.${root}.${field}`) - ), - ...(args.arrayKey ? [`$.${args.arrayKey}`, `$.${args.arrayKey}.*`] : []) - ], - keepStack: false, - stringBufferSize: 64 * 1024 - }) const record: Record = Object.create(null) const fields = new Set(args.fields) - let depth = 0 - let expectingRootKey = false - parser.onToken = ({ token, value }) => { - if (depth === 1 && expectingRootKey && token === TokenType.STRING) { - if (typeof value === 'string' && Object.hasOwn(args.objectFields ?? {}, value)) { - record[value] = Object.create(null) - } - expectingRootKey = false - } - if (token === TokenType.LEFT_BRACE || token === TokenType.LEFT_BRACKET) { - if (depth === 0 && token === TokenType.LEFT_BRACE) { - expectingRootKey = true - } - depth++ - } else if (token === TokenType.RIGHT_BRACE || token === TokenType.RIGHT_BRACKET) { - depth-- - } else if (token === TokenType.COMMA && depth === 1) { - expectingRootKey = true - } - } + const objectFields = new Map( + Object.entries(args.objectFields ?? {}).map(([root, keys]) => [root, new Set(keys)]) + ) + const foldedArray = Symbol('folded session array') let state = args.create() - let currentArray: unknown = null let consumeFailure: { error: unknown } | undefined + let inFoldedArray = false + const reset = (): void => { + state = args.create() + consumeFailure = undefined + } + const retain = (path: (string | number)[]): boolean => { + const [root, child] = path + return ( + typeof root === 'string' && + ((root !== args.arrayKey && fields.has(root)) || + (inFoldedArray && root === args.arrayKey && path.length >= 2) || + (typeof child === 'string' && objectFields.get(root)?.has(child) === true)) + ) + } + // Every discarded container needs a rule; dropping only its parent still builds large children. + const discard = { + filter: (path: (string | number)[]) => !retain(path), + create: () => null, + add: () => {} + } + const assembler = new FlexAssembler({ + maxDepth: Infinity, + objectRules: [ + objectRule>({ + filter: (path) => path.length === 0, + create: () => record, + add: (target, key, value) => { + if (key === args.arrayKey) { + if (value !== foldedArray) { + reset() + } + } else if (fields.has(key)) { + target[key] = value + } + } + }), + objectRule>({ + filter: (path) => + path.length === 1 && + typeof path[0] === 'string' && + objectFields.has(path[0]) && + (!fields.has(path[0]) || path[0] === args.arrayKey), + create: (path) => { + const projected: Record = Object.create(null) + record[String(path[0])] = projected + return projected + }, + add: (target, key, value) => { + const root = assembler.path[0] + if (typeof root === 'string' && objectFields.get(root)?.has(key)) { + target[key] = value + } + } + }), + discard + ], + arrayRules: [ + arrayRule({ + filter: (path) => Boolean(args.arrayKey) && path.length === 1 && path[0] === args.arrayKey, + create: () => { + reset() + inFoldedArray = true + return null + }, + add: (_target, value) => { + if (!consumeFailure) { + try { + args.consume(state, value) + } catch (error) { + consumeFailure = { error } + } + } + }, + finalize: () => { + inFoldedArray = false + return foldedArray + } + }), + discard + ] + }) + const parser = createJsonTokenReader((token) => { + if ( + token.name === 'keyValue' && + assembler.depth === 1 && + !assembler.isArray && + objectFields.has(token.value) + ) { + record[token.value] = Object.create(null) + } + assembler.consume(token) + }) const decoder = new StringDecoder('utf8') let objectRoot: boolean | undefined - parser.onValue = ({ key, value, parent, stack }) => { - if (stack.length === 2 && stack[1].key === args.arrayKey && Array.isArray(parent)) { - if (parent !== currentArray) { - state = args.create() - consumeFailure = undefined - currentArray = parent - } - if (!consumeFailure) { - try { - args.consume(state, value) - } catch (error) { - consumeFailure = { error } - } - } - // The parser's array cursor is independent of retained array slots. - parent.pop() - } else if ( - stack.length === 2 && - typeof stack[1].key === 'string' && - typeof key === 'string' && - parent && - !Array.isArray(parent) - ) { - const root = stack[1].key - const projected = record[root] - if ( - Object.hasOwn(args.objectFields ?? {}, root) && - args.objectFields?.[root]?.includes(key) && - projected && - typeof projected === 'object' - ) { - Reflect.set(projected, key, value) - } - } else if (stack.length === 1 && typeof key === 'string') { - if (key === args.arrayKey) { - if (value !== currentArray || !Array.isArray(value)) { - state = args.create() - consumeFailure = undefined - } - currentArray = null - } else if (fields.has(key)) { - record[key] = value - } - if (parent && typeof parent === 'object') { - Reflect.deleteProperty(parent, key) - } - } - } for await (const chunk of args.bytes) { throwIfAiVaultScanCancelled(args.signal) if (objectRoot === undefined) { @@ -107,37 +124,17 @@ export async function readStreamedSessionDocument(args: { objectRoot = first === 123 } } - parseJson(() => parser.write(decoder.write(chunk))) + parser.write(decoder.write(chunk)) await yieldToEventLoop() } const tail = decoder.end() if (tail) { - parseJson(() => parser.write(tail)) - } - if (objectRoot === undefined) { - throw new SyntaxError('Unexpected end of JSON input') - } - if (!parser.isEnded) { - parseJson(() => parser.end(), true) + parser.write(tail) } + parser.end() throwIfAiVaultScanCancelled(args.signal) if (consumeFailure) { throw consumeFailure.error } return objectRoot ? { record, state } : null } - -function parseJson(run: () => void, ending = false): void { - try { - run() - } catch (error) { - if ( - error instanceof TokenizerError || - error instanceof TokenParserError || - (ending && error instanceof Error) - ) { - throw new SyntaxError(error.message) - } - throw error - } -} diff --git a/src/shared/json-token-reader.test.ts b/src/shared/json-token-reader.test.ts new file mode 100644 index 00000000000..9cba3e8a869 --- /dev/null +++ b/src/shared/json-token-reader.test.ts @@ -0,0 +1,47 @@ +import { expect, it } from 'vitest' +import { FlexAssembler } from 'stream-json/core/utils/flex-assembler.js' +import { createJsonTokenReader } from './json-token-reader' + +it.each([1, 7, 64 * 1024, Infinity])('preserves Unicode at chunk size %s', (size) => { + const value = { '\ufeffkey': `\ufeffstart${'\n'.repeat(64 * 1024)}\ufeff😀end` } + const literal = JSON.stringify(value) + for (const text of [literal, literal.replaceAll('\ufeff', '\\uFEFF')]) { + const assembler = new FlexAssembler() + const reader = createJsonTokenReader((token) => { + assembler.consume(token) + }) + for (let offset = 0; offset < text.length; offset += size) { + reader.write(text.slice(offset, offset + size)) + } + reader.end() + expect(assembler.current).toEqual(value) + } +}) + +it.each(['', ' ', '{', '{"a":1,}', '{}{}', '[01]', '\ufeff{}'])( + 'rejects malformed JSON %j as a syntax error', + (text) => { + const reader = createJsonTokenReader(() => {}) + expect(() => { + reader.write(text) + reader.end() + }).toThrow(SyntaxError) + } +) + +it('preserves consumer errors and completes numbers at EOF synchronously', () => { + const values: string[] = [] + const reader = createJsonTokenReader((token) => { + if (token.name === 'numberValue') { + values.push(token.value) + } + }) + reader.write('12') + reader.end() + expect(values).toEqual(['12']) + const failure = new RangeError('consumer failure') + const throwing = createJsonTokenReader(() => { + throw failure + }) + expect(() => throwing.write('{}')).toThrow(failure) +}) diff --git a/src/shared/json-token-reader.ts b/src/shared/json-token-reader.ts new file mode 100644 index 00000000000..c0081ee016f --- /dev/null +++ b/src/shared/json-token-reader.ts @@ -0,0 +1,48 @@ +import parser, { type Token } from 'stream-json/core/parser.js' +import exec from 'stream-chain/exec.js' +import { none } from 'stream-chain/defs.js' + +/** Bounds token batches while preserving caller-owned UTF-8 decoding. */ +export function createJsonTokenReader( + consume: (token: Token) => void, + batchCodeUnits = 64 * 1024 +): { + write: (text: string) => void + end: () => void +} { + if (!Number.isSafeInteger(batchCodeUnits) || batchCodeUnits < 1) { + throw new RangeError('JSON token batch size must be a positive safe integer') + } + const parse = exec(parser({ streamValues: false })) + function run(input: string | typeof none): void { + let consumerFailed = false + try { + const pending = parse(input, (token: Token) => { + try { + consume(token) + } catch (error) { + consumerFailed = true + throw error + } + }) + if (pending) { + throw new Error('JSON token reader requires a synchronous parser') + } + } catch (error) { + if (!consumerFailed && error instanceof Error) { + throw new SyntaxError(error.message) + } + throw error + } + } + return { + write(text) { + for (let offset = 0; offset < text.length; offset += batchCodeUnits) { + run(text.slice(offset, offset + batchCodeUnits)) + } + }, + end() { + run(none) + } + } +} diff --git a/src/shared/ripgrep-dense-match-json.test.ts b/src/shared/ripgrep-dense-match-json.test.ts index 5a4cdd41e29..ec62caddb2f 100644 --- a/src/shared/ripgrep-dense-match-json.test.ts +++ b/src/shared/ripgrep-dense-match-json.test.ts @@ -1,23 +1,6 @@ import { expect, it } from 'vitest' -import { JSONParser } from '@streamparser/json' import { parseDenseRipgrepMatchJson } from './ripgrep-dense-match-json' -it.each([0, 64 * 1024])('preserves literal and escaped BOMs with buffer size %i', (size) => { - const source = { '\ufeffkey': `\ufeffstart${'\n'.repeat(64 * 1024)}\ufeffend` } - const literal = JSON.stringify(source) - for (const record of [literal, literal.replaceAll('\ufeff', '\\uFEFF')]) { - const parser = new JSONParser({ stringBufferSize: size }) - let parsed: unknown - parser.onValue = ({ value, stack }) => { - if (stack.length === 0) { - parsed = value - } - } - parser.write(record) - expect(parsed).toEqual(JSON.parse(record)) - } -}) - it.each(['', '\\', '\n', '\n'.repeat(64 * 1024), `${'x'.repeat(64 * 1024)}\n`])( 'preserves U+FEFF in text and filenames across string-buffer boundaries (%#)', (prefix) => { @@ -32,7 +15,9 @@ it.each(['', '\\', '\n', '\n'.repeat(64 * 1024), `${'x'.repeat(64 * 1024)}\n`])( } } const record = JSON.stringify(source) - expect(parseDenseRipgrepMatchJson(record, 1, 16)).toEqual(JSON.parse(record)) + for (const encoded of [record, record.replaceAll('\ufeff', '\\uFEFF')]) { + expect(parseDenseRipgrepMatchJson(encoded, 1, 16)).toEqual(JSON.parse(record)) + } } ) @@ -76,3 +61,33 @@ it.each(['{}', '{"type":"begin","data":{}}', '{"data":null}', '{"data":[]}'])( expect(projected.data?.submatches).toEqual([]) } ) + +it('validates submatches after reaching the requested result cap', () => { + const source = JSON.stringify({ + type: 'match', + data: { + submatches: [ + { start: 0, end: 1 }, + { start: null, end: 2 } + ] + } + }) + expect(() => parseDenseRipgrepMatchJson(source, 1, 16)).toThrow('Invalid rg submatch') +}) + +it('rejects duplicate coordinate fields whose final value is not numeric', () => { + const source = '{"data":{"submatches":[{"start":0,"start":{},"end":1}]}}' + expect(() => parseDenseRipgrepMatchJson(source, 1, 16)).toThrow('Invalid rg submatch') +}) + +it.each([16_378, 16_379])('retains the existing per-element token budget at %i values', (count) => { + const source = JSON.stringify({ + data: { submatches: [{ start: 0, end: 1, other: Array(count).fill(0) }] } + }) + const parse = (): unknown => parseDenseRipgrepMatchJson(source, 1, 16) + if (count === 16_378) { + expect(parse()).toMatchObject({ data: { submatches: [{ start: 0, end: 1 }] } }) + } else { + expect(parse).toThrow('rg submatch structure exceeds limit') + } +}) diff --git a/src/shared/ripgrep-dense-match-json.ts b/src/shared/ripgrep-dense-match-json.ts index 821415a0748..5b54069fb10 100644 --- a/src/shared/ripgrep-dense-match-json.ts +++ b/src/shared/ripgrep-dense-match-json.ts @@ -3,7 +3,7 @@ import { JsonTextStructureCapacityError, type JsonTextStructureLimits } from './json-text-structure-limit' -import { JSONParser, TokenType } from '@streamparser/json' +import { createJsonTokenReader } from './json-token-reader' export type RipgrepMatchMessage = { type?: string @@ -34,102 +34,132 @@ export function parseRipgrepMatchJson( } } +type RipgrepJsonFrame = { + kind: 'object' | 'array' + context: 'root' | 'data' | 'path' | 'lines' | 'matches' | 'match' | null + key?: string + values: number + keys?: Set + start?: number + end?: number +} + /** Dense rg records retain only the requested ranges while validating the entire record. */ export function parseDenseRipgrepMatchJson( line: string, maxMatches: number, nestingDepth: number ): RipgrepMatchMessage { - const parser = new JSONParser({ - paths: [ - '$.type', - '$.data.path.text', - '$.data.lines.text', - '$.data.lines.bytes', - '$.data.line_number', - '$.data.submatches.*' - ], - keepStack: false, - stringBufferSize: 64 * 1024 - }) const data: NonNullable = { submatches: [] } const result: RipgrepMatchMessage = { data } - const containers: { object: boolean; expectingKey: boolean; keys?: Set }[] = [] + const frames: RipgrepJsonFrame[] = [] let elementTokens = 0 - parser.onToken = ({ token, value }) => { - const current = containers.at(-1) - if (current?.object && current.expectingKey && token === TokenType.STRING) { - if (typeof value !== 'string') { - throw new SyntaxError('Invalid rg object key') - } - if (current.keys?.has(value)) { - throw new SyntaxError('Duplicate rg object key') - } - current.keys?.add(value) - if ((current.keys?.size ?? 0) > 128) { - throw new Error('Too many rg record fields') - } - current.expectingKey = false - } - if (token === TokenType.LEFT_BRACE || token === TokenType.LEFT_BRACKET) { - containers.push({ - object: token === TokenType.LEFT_BRACE, - expectingKey: token === TokenType.LEFT_BRACE, - // rg's envelope keys are unique; reject duplicates instead of mixing projections. - keys: containers.length < 2 ? new Set() : undefined - }) - if (containers.length > nestingDepth) { - throw new Error('rg record nesting exceeds limit') - } - if (containers.length === 4) { - elementTokens = 0 - } - } else if (token === TokenType.RIGHT_BRACE || token === TokenType.RIGHT_BRACKET) { - containers.pop() - } else if (token === TokenType.COMMA && current?.object) { - current.expectingKey = true - } - if (containers.length >= 4 && ++elementTokens > 32 * 1024) { + const countTokens = (amount = 1): void => { + if (frames.length >= 4 && (elementTokens += amount) > 32 * 1024) { throw new Error('rg submatch structure exceeds limit') } } - parser.onValue = ({ key, value, parent, stack }) => { - if (stack.length === 1 && key === 'type' && typeof value === 'string') { - result.type = value - } else if (stack.length === 2 && stack[1].key === 'data' && key === 'line_number') { - if (typeof value === 'number') { - data.line_number = value + const parser = createJsonTokenReader((token) => { + const current = frames.at(-1) + if (token.name === 'keyValue') { + if (!current || current.kind !== 'object') { + throw new SyntaxError('Unexpected rg object key') } - } else if (stack.length === 3 && stack[1].key === 'data') { - if (stack[2].key === 'submatches' && Array.isArray(parent)) { - if ( - !value || - typeof value !== 'object' || - Array.isArray(value) || - typeof value.start !== 'number' || - typeof value.end !== 'number' - ) { + if (current.values++ > 0) { + countTokens() + } + // Packed tokens omit commas and colons; include them in the existing structure budget. + countTokens(2) + current.key = token.value + if (current.keys?.has(token.value)) { + throw new SyntaxError('Duplicate rg object key') + } + current.keys?.add(token.value) + if ((current.keys?.size ?? 0) > 128) { + throw new Error('Too many rg record fields') + } + return + } + if (token.name === 'endObject' || token.name === 'endArray') { + const frame = frames.pop() + countTokens() + if (frame?.context === 'match') { + if (typeof frame.start !== 'number' || typeof frame.end !== 'number') { throw new SyntaxError('Invalid rg submatch') } if (data.submatches && data.submatches.length < maxMatches) { - data.submatches.push({ start: value.start, end: value.end }) - } - parent.pop() - } else if (stack[2].key === 'path' && key === 'text' && typeof value === 'string') { - data.path = { text: value } - } else if (stack[2].key === 'lines' && typeof value === 'string') { - if (key === 'text') { - data.lines = { ...data.lines, text: value } - } - if (key === 'bytes') { - data.lines = { ...data.lines, bytes: value } + data.submatches.push({ start: frame.start, end: frame.end }) } } + return } - } + if (current?.kind === 'array' && current.values++ > 0) { + countTokens() + } + if (current?.context === 'match' && (current.key === 'start' || current.key === 'end')) { + current[current.key] = undefined + } + if (token.name === 'startObject' || token.name === 'startArray') { + const kind = token.name === 'startObject' ? 'object' : 'array' + let context: RipgrepJsonFrame['context'] = null + if (!current && kind === 'object') { + context = 'root' + } else if (current?.context === 'root' && current.key === 'data' && kind === 'object') { + context = 'data' + } else if (current?.context === 'data') { + if ((current.key === 'lines' || current.key === 'path') && kind === 'object') { + context = current.key + } + if (current.key === 'submatches' && kind === 'array') { + context = 'matches' + } + } else if (current?.context === 'matches') { + if (kind !== 'object') { + throw new SyntaxError('Invalid rg submatch') + } + context = 'match' + } + frames.push({ + kind, + context, + values: 0, + keys: kind === 'object' && frames.length < 2 ? new Set() : undefined + }) + if (frames.length > nestingDepth) { + throw new Error('rg record nesting exceeds limit') + } + if (frames.length === 4) { + elementTokens = 0 + } + countTokens() + return + } + countTokens() + if (current?.context === 'matches') { + throw new SyntaxError('Invalid rg submatch') + } + if (token.name === 'numberValue') { + const value = Number(token.value) + if (current?.context === 'match' && (current.key === 'start' || current.key === 'end')) { + current[current.key] = value + } + if (current?.context === 'data' && current.key === 'line_number') { + data.line_number = value + } + } else if (token.name === 'stringValue') { + if (current?.context === 'root' && current.key === 'type') { + result.type = token.value + } + if (current?.context === 'path' && current.key === 'text') { + data.path = { text: token.value } + } + if (current?.context === 'lines' && (current.key === 'text' || current.key === 'bytes')) { + data.lines ??= {} + data.lines[current.key] = token.value + } + } + }, 8 * 1024) parser.write(line) - if (!parser.isEnded) { - parser.end() - } + parser.end() return result }