Replace patched JSON parser with stream-json (#25202)

* Replace patched JSON parser with stream-json

* Isolate dependencies for historical server compatibility builds
This commit is contained in:
Neil
2026-10-04 13:05:30 -07:00
committed by GitHub
parent 41cc77509f
commit b32462f246
17 changed files with 765 additions and 283 deletions
+2 -2
View File
@@ -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"
@@ -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 = "";
@@ -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))
}
}
@@ -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 })
}
}
@@ -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')
+2 -2
View File
@@ -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"
+1 -1
View File
@@ -1,7 +1,7 @@
{
"extends": "./tsconfig.cli.json",
"compilerOptions": {
"module": "node16",
"module": "Node20",
"moduleResolution": "node16"
}
}
+2 -1
View File
@@ -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',
+2 -1
View File
@@ -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",
+18 -9
View File
@@ -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
-1
View File
@@ -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
@@ -0,0 +1,166 @@
import { describe, expect, it } from 'vitest'
import { readStreamedSessionDocument } from './session-document-stream'
async function* bytes(content: string, chunkSize = 7): AsyncGenerator<Buffer> {
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<Parameters<typeof readStreamedSessionDocument<unknown[]>>[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']
})
}
})
})
+101 -104
View File
@@ -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<T>(args: {
consume: (state: T, value: unknown) => void
signal?: AbortSignal
}): Promise<{ record: Record<string, unknown>; 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<string, unknown> = 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<Record<string, unknown>>({
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<Record<string, unknown>>({
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<string, unknown> = 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<null>({
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<T>(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
}
}
+47
View File
@@ -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)
})
+48
View File
@@ -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)
}
}
}
+33 -18
View File
@@ -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')
}
})
+108 -78
View File
@@ -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<string>
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<RipgrepMatchMessage['data']> = { submatches: [] }
const result: RipgrepMatchMessage = { data }
const containers: { object: boolean; expectingKey: boolean; keys?: Set<string> }[] = []
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
}