diff --git a/bun.lock b/bun.lock index 92876a58..262be0e6 100644 --- a/bun.lock +++ b/bun.lock @@ -59,6 +59,7 @@ "shell-quote": "^1.8.3", "signal-exit": "^4.1.0", "stack-utils": "^2.0.6", + "stream-json": "^3.7.0", "strip-ansi": "^7.2.0", "supports-hyperlinks": "^4.4.0", "tree-kill": "^1.2.2", @@ -949,6 +950,10 @@ "statuses": ["statuses@2.0.2", "https://registry.npmmirror.com/statuses/-/statuses-2.0.2.tgz", {}, "sha512-DvEy55V3DB7uknRo+4iOGT5fP1slR8wQohVdknigZPMpMstaKJQWhwiYBACJE3Ul2pTnATihhBYnRhZQHGBiRw=="], + "stream-chain": ["stream-chain@4.2.6", "https://registry.npmmirror.com/stream-chain/-/stream-chain-4.2.6.tgz", {}, "sha512-1zeJ8CrtJfmiba26ui8jXkq/xLRFvhzkdH02D5QLO9Cnovgeb28IJxg88DdraWnjnkbOol/++uIAybJbJhk7ig=="], + + "stream-json": ["stream-json@3.7.0", "https://registry.npmmirror.com/stream-json/-/stream-json-3.7.0.tgz", { "dependencies": { "stream-chain": "^4.2.5" } }, "sha512-rCSBdcBP/bPk6T8QFcxAj1MSzAuc5i49cYW6IE7sYObNEPccBJIUiL6fU9c9BWt40aJK0siVynLX7lWaW8alXw=="], + "string-width": ["string-width@8.2.0", "https://registry.npmmirror.com/string-width/-/string-width-8.2.0.tgz", { "dependencies": { "get-east-asian-width": "^1.5.0", "strip-ansi": "^7.1.2" } }, "sha512-6hJPQ8N0V0P3SNmP6h2J99RLuzrWz2gvT7VnK5tKvrNqJoyS9W4/Fb8mo31UiPvy00z7DQXkP2hnKBVav76thw=="], "strip-ansi": ["strip-ansi@7.2.0", "https://registry.npmmirror.com/strip-ansi/-/strip-ansi-7.2.0.tgz", { "dependencies": { "ansi-regex": "^6.2.2" } }, "sha512-yDPMNjp4WyfYBkHnjIRLfca1i6KMyGCtsVgoKe/z1+6vukgaENdgGBZt+ZmKPc4gavvEZ5OgHfHdrazhgNyG7w=="], diff --git a/package.json b/package.json index 080ea5a8..9558c7f0 100644 --- a/package.json +++ b/package.json @@ -103,6 +103,7 @@ "shell-quote": "^1.8.3", "signal-exit": "^4.1.0", "stack-utils": "^2.0.6", + "stream-json": "^3.7.0", "strip-ansi": "^7.2.0", "supports-hyperlinks": "^4.4.0", "tree-kill": "^1.2.2", diff --git a/src/server/__tests__/e2e/oversized-session.test.ts b/src/server/__tests__/e2e/oversized-session.test.ts new file mode 100644 index 00000000..514c43da --- /dev/null +++ b/src/server/__tests__/e2e/oversized-session.test.ts @@ -0,0 +1,90 @@ +import { expect, test } from 'bun:test' +import { appendFile, mkdir, mkdtemp, readFile, rm } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { fileURLToPath } from 'node:url' +import { createSandboxedTestEnvironment } from '../../../../scripts/pr/test-environment.js' + +test('an existing session with a large image opens checkpoint metadata and resumes without losing its history', async () => { + const home = await mkdtemp(join(tmpdir(), 'oversized-session-e2e-')) + const original = { ...process.env } + const env = createSandboxedTestEnvironment(home, { + CLAUDE_CLI_PATH: fileURLToPath(new URL('../fixtures/mock-sdk-cli.ts', import.meta.url)), + DISABLE_TELEMETRY: '1', DISABLE_ERROR_REPORTING: '1', DISABLE_AUTOUPDATER: '1', + }) + for (const key of Object.keys(process.env)) delete process.env[key] + Object.assign(process.env, env) + let server: ReturnType | undefined + let shutdown: (() => Promise) | undefined + let socket: WebSocket | undefined + try { + const workDir = join(home, 'project') + await mkdir(workDir) + const runtime = await import('../../index.js') + const { sessionService } = await import('../../services/sessionService.js') + shutdown = runtime.stopServerRuntimeForShutdown + server = runtime.startServer(0, '127.0.0.1') + const base = `http://127.0.0.1:${server.port}` + const created = await fetch(`${base}/api/sessions`, { + method: 'POST', headers: { 'content-type': 'application/json' }, body: JSON.stringify({ workDir }), + }) + expect(created.status).toBe(201) + const { sessionId } = await created.json() as { sessionId: string } + const initial = await sessionService.getSessionLaunchInfo(sessionId) + const image = { + uuid: crypto.randomUUID(), type: 'user', cwd: workDir, sessionId, + timestamp: '2026-09-27T01:00:00Z', + message: { role: 'user', content: [ + { type: 'text', text: 'Inspect this large image' }, + { type: 'image', source: { type: 'base64', media_type: 'image/png', data: 'A'.repeat(9 * 1024 * 1024) } }, + ] }, + } + await appendFile(initial!.filePath, [ + { uuid: crypto.randomUUID(), type: 'user', cwd: workDir, message: { role: 'user', content: 'Earlier request' } }, + { uuid: crypto.randomUUID(), type: 'assistant', message: { role: 'assistant', content: [{ type: 'text', text: 'Earlier completed reply' }] } }, + image, + ].map(entry => JSON.stringify(entry)).join('\n') + '\n') + const before = await readFile(initial!.filePath, 'utf8') + + // This is fetched automatically by the desktop when opening a completed + // conversation. It used to surface the screenshot's red metadata error. + const checkpoint = await fetch(`${base}/api/sessions/${sessionId}/turn-checkpoints`) + expect(checkpoint.status).toBe(200) + expect((await checkpoint.json() as { checkpoints: unknown[] }).checkpoints).toHaveLength(1) + expect((await sessionService.getSessionLaunchInfo(sessionId))?.transcriptMessageCount).toBe(3) + + const events: Record[] = [] + socket = new WebSocket(`ws://127.0.0.1:${server.port}/ws/${sessionId}`) + socket.onmessage = event => events.push(JSON.parse(String(event.data))) + async function until(predicate: () => boolean) { + const deadline = Date.now() + 10_000 + while (!predicate()) { + if (Date.now() > deadline) throw new Error(`Session did not recover: ${JSON.stringify(events)}`) + await Bun.sleep(20) + } + } + await until(() => events.some(event => event.type === 'connected')) + socket.send(JSON.stringify({ type: 'user_message', content: 'continue after the large image' })) + await until(() => events.some(event => event.type === 'message_complete')) + expect(events.filter(event => event.type === 'error')).toEqual([]) + expect(JSON.stringify(events)).toContain('Echo: continue after the large image') + // A large first turn must never be mistaken for a metadata-only placeholder + // and cleared by startup. Only new metadata may be appended. + expect((await readFile(initial!.filePath, 'utf8')).startsWith(before)).toBe(true) + + await appendFile(initial!.filePath, JSON.stringify({ + uuid: crypto.randomUUID(), type: 'assistant', timestamp: '2026-09-27T01:01:00Z', + message: { role: 'assistant', content: [{ type: 'text', text: 'Visible reply after the large image' }] }, + }) + '\n') + const history = await fetch(`${base}/api/sessions/${sessionId}/messages`) + expect(history.status).toBe(200) + expect(JSON.stringify(await history.json())).toContain('Visible reply after the large image') + } finally { + socket?.close() + await shutdown?.() + server?.stop(true) + for (const key of Object.keys(process.env)) delete process.env[key] + Object.assign(process.env, original) + await rm(home, { recursive: true, force: true }) + } +}, 30_000) diff --git a/src/server/services/sessionHistoryContext.test.ts b/src/server/services/sessionHistoryContext.test.ts index d7c8fe60..e9da132c 100644 --- a/src/server/services/sessionHistoryContext.test.ts +++ b/src/server/services/sessionHistoryContext.test.ts @@ -90,3 +90,77 @@ test('rebuilds scalar context after an in-place rewrite grows beyond the cached expect(rebuilt.contexts.get(0)?.suppressed).toBe(false) expect(rebuilt.scannedBytes).toBe(Number((await stat(file)).size)) }) + +const classifyMessage = (entry: Record) => { + const message = entry.message as { role?: string; content?: unknown } + const content = message?.content + const blocks = Array.isArray(content) ? content : [] + const texts = typeof content === 'string' ? [content] : blocks.filter(block => block.type === 'text').map(block => block.text) + const user = message?.role === 'user' && !entry.isMeta + return { + notification: user && texts.length > 0 && texts.every(text => //.test(text)), + reset: user && !blocks.some(block => block.type === 'tool_result'), + agentToolId: blocks.find(block => block.type === 'tool_use' && ['Agent', 'Task'].includes(block.name))?.id, + } +} + +test('an oversized ordinary image resets notification suppression and preserves following replies', async () => { + const notice = row('notice', { type: 'user', message: { role: 'user', content: 'done' } }) + '\n' + const image = row('image', { type: 'user', message: { role: 'user', content: [{ type: 'image', source: { type: 'base64', data: 'A'.repeat(9 * 1024 * 1024) } }] } }) + '\n' + const answer = row('answer', { type: 'assistant', message: { role: 'assistant', content: 'Visible reply' } }) + '\n' + await writeFile(file, notice + image + answer) + const answerOffset = Buffer.byteLength(notice + image) + const contexts = await readHistoryContexts({ filePath: file, sourceVersion: await version(), offsets: [0, Buffer.byteLength(notice), answerOffset], classify: classifyMessage }) + expect(contexts.contexts.get(0)?.suppressed).toBe(true) + expect(contexts.contexts.get(Buffer.byteLength(notice))?.suppressed).toBe(false) + expect(contexts.contexts.get(answerOffset)?.suppressed).toBe(false) + expect((await readHistoryContexts({ filePath: file, sourceVersion: await version(), offsets: [answerOffset], classify: classifyMessage })).scannedBytes).toBe(0) +}) + +test('oversized Agent input preserves the sidechain owner for later small records', async () => { + const parent = row('parent', { type: 'assistant', message: { role: 'assistant', content: [{ type: 'tool_use', id: 'agent-tool', name: 'Agent', input: { prompt: 'x'.repeat(9 * 1024 * 1024) } }] } }) + '\n' + const child = row('child', { type: 'assistant', parentUuid: 'parent', isSidechain: true, message: { role: 'assistant', content: 'Child result' } }) + '\n' + await writeFile(file, parent + child) + const offset = Buffer.byteLength(parent) + const contexts = await readHistoryContexts({ filePath: file, sourceVersion: await version(), offsets: [offset], classify: classifyMessage }) + expect(contexts.contexts.get(offset)).toEqual({ owner: 'agent-tool', suppressed: false }) +}) + +test('an oversized tool result does not reset preceding notification suppression', async () => { + const notice = row('notice', { type: 'user', message: { role: 'user', content: 'done' } }) + '\n' + const result = row('result', { type: 'user', message: { role: 'user', content: [{ type: 'tool_result', tool_use_id: 'tool', content: 'x'.repeat(9 * 1024 * 1024) }] } }) + '\n' + const answer = row('answer', { type: 'assistant', message: { role: 'assistant', content: 'Hidden notification response' } }) + '\n' + await writeFile(file, notice + result + answer) + const offset = Buffer.byteLength(notice + result) + expect((await readHistoryContexts({ filePath: file, sourceVersion: await version(), offsets: [offset], classify: classifyMessage })).contexts.get(offset)?.suppressed).toBe(true) +}) + +test('a notification beyond a projected text prefix remains suppressed until an explicit user reset', async () => { + const longText = row('long', { type: 'user', message: { role: 'user', content: 'x'.repeat(9 * 1024 * 1024) + 'done' } }) + '\n' + const hidden = row('hidden', { type: 'assistant', message: { role: 'assistant', content: 'Notification response' } }) + '\n' + const reset = row('reset', { type: 'user', message: { role: 'user', content: 'New request' } }) + '\n' + const visible = row('visible', { type: 'assistant', message: { role: 'assistant', content: 'Normal reply' } }) + '\n' + await writeFile(file, longText + hidden + reset + visible) + const hiddenOffset = Buffer.byteLength(longText) + const visibleOffset = Buffer.byteLength(longText + hidden + reset) + const contexts = await readHistoryContexts({ filePath: file, sourceVersion: await version(), offsets: [hiddenOffset, visibleOffset], classify: classifyMessage }) + expect(contexts.contexts.get(hiddenOffset)?.suppressed).toBe(true) + expect(contexts.contexts.get(visibleOffset)?.suppressed).toBe(false) +}) + +test('malformed oversized records preserve uncertainty rather than exposing subsequent notification replies', async () => { + const malformed = '{"message":{"content":"' + 'x'.repeat(9 * 1024 * 1024) + '\n' + const reply = row('reply', { type: 'assistant', message: { role: 'assistant', content: 'Unknown context' } }) + '\n' + await writeFile(file, malformed + reply) + const contexts = await readHistoryContexts({ filePath: file, sourceVersion: await version(), offsets: [Buffer.byteLength(malformed)], classify: classifyMessage }) + expect(contexts.contexts.get(Buffer.byteLength(malformed))?.suppressed).toBe(true) +}) + +test('ordinary long text inside the existing semantic limit remains a complete reset', async () => { + const notice = row('notice', { type: 'user', message: { role: 'user', content: 'done' } }) + '\n' + const request = row('request', { type: 'user', message: { role: 'user', content: 'x'.repeat(60 * 1024) } }) + '\n' + const answer = row('answer', { type: 'assistant', message: { role: 'assistant', content: 'Visible reply' } }) + '\n' + await writeFile(file, notice + request + answer) + const offset = Buffer.byteLength(notice + request) + expect((await readHistoryContexts({ filePath: file, sourceVersion: await version(), offsets: [offset], classify: classifyMessage })).contexts.get(offset)?.suppressed).toBe(false) +}) diff --git a/src/server/services/sessionHistoryContext.ts b/src/server/services/sessionHistoryContext.ts index bf6c6443..91e2fa3a 100644 --- a/src/server/services/sessionHistoryContext.ts +++ b/src/server/services/sessionHistoryContext.ts @@ -4,7 +4,8 @@ import { mkdtemp, open, rm, stat } from 'node:fs/promises' import { rmSync } from 'node:fs' import { tmpdir } from 'node:os' import { join } from 'node:path' -import { HISTORY_SEMANTIC_RECORD_BYTES, streamBoundedHistory, withHistoryReadBudget } from './boundedSessionHistory.js' +import { withHistoryReadBudget } from './boundedSessionHistory.js' +import { isSessionMetadataTextTruncated, streamSessionMetadata } from './sessionMetadataReader.js' import { ApiError } from '../middleware/errorHandler.js' type Context = { owner?: string; suppressed: boolean } @@ -88,20 +89,25 @@ export async function readHistoryContexts(options: { try { const fingerprint = await sourceAnchors(options.filePath, targetSize, signal) state.database.exec('BEGIN') - const result = await streamBoundedHistory(options.filePath, (entry, completeLine, offset) => { + const result = await streamSessionMetadata(options.filePath, (entry, completeLine, offset) => { const classification = options.classify(entry) const inherited = typeof entry.parentUuid === 'string' ? (getParent.get(entry.parentUuid) as { chain?: string } | null)?.chain : undefined const explicit = typeof entry.parent_tool_use_id === 'string' && entry.parent_tool_use_id ? entry.parent_tool_use_id : undefined const owner = explicit ?? (entry.isSidechain === true ? inherited : undefined) const chain = classification.agentToolId ?? inherited if (typeof entry.uuid === 'string') saveParent.run(entry.uuid, chain ?? null) - if (classification.notification) suppressed = true + const message = entry.message as { role?: unknown } | undefined + // A bounded text preview cannot prove whether a user record contains a + // task notification beyond its prefix. Keep uncertainty fail-closed; + // images and tool payloads do not affect this textual classification. + if (message?.role === 'user' && !entry.isMeta && isSessionMetadataTextTruncated(entry)) suppressed = null + else if (classification.notification) suppressed = true else if (classification.reset) suppressed = false // Keep root-only ownership filtering separate from notification state: // a dedicated child transcript legitimately lacks its parent's Agent call. saveContext.run(offset, owner ?? null, suppressed !== false ? 1 : 0, entry.isSidechain === true && !owner ? 1 : 0) if (completeLine) completeSuppression = suppressed - }, signal, { startOffset: originalOffset, endOffset: targetSize, maxRecordBytes: HISTORY_SEMANTIC_RECORD_BYTES, onSkipped: () => { suppressed = null; completeSuppression = null } }) + }, signal, { startOffset: originalOffset, endOffset: targetSize, onSkipped: () => { suppressed = null; completeSuppression = null } }) if (fingerprint !== await sourceAnchors(options.filePath, targetSize, signal)) throw new ApiError(409, 'History was rewritten during context scan', 'HISTORY_CHANGED') state.database.exec('COMMIT') state.fingerprint = fingerprint diff --git a/src/server/services/sessionMetadataReader.test.ts b/src/server/services/sessionMetadataReader.test.ts new file mode 100644 index 00000000..0338d425 --- /dev/null +++ b/src/server/services/sessionMetadataReader.test.ts @@ -0,0 +1,122 @@ +import { afterEach, beforeEach, expect, test } from 'bun:test' +import { mkdtemp, readFile, rm, writeFile } from 'node:fs/promises' +import { join } from 'node:path' +import { tmpdir } from 'node:os' +import { HISTORY_SEMANTIC_RECORD_BYTES } from './boundedSessionHistory.js' +import { isSessionMetadataTextTruncated, streamSessionMetadata } from './sessionMetadataReader.js' + +let directory: string +let file: string +const large = 'x'.repeat(HISTORY_SEMANTIC_RECORD_BYTES + 1) +beforeEach(async () => { + directory = await mkdtemp(join(tmpdir(), 'metadata-reader-')) + file = join(directory, 'history.jsonl') +}) +afterEach(async () => { + await rm(directory, { recursive: true, force: true }) +}) + +async function read(content: string) { + await writeFile(file, content) + const entries: { + entry: Record + complete: boolean + offset: number + }[] = [] + let skipped = 0 + const scan = await streamSessionMetadata(file, (entry, complete, offset) => entries.push({ entry, complete, offset }), undefined, { onSkipped: () => skipped++ }) + return { entries, scan, skipped } +} + +test('large image bodies retain launch and context fields, including metadata after the body', async () => { + const entry = { type: 'user', message: { content: [{ type: 'image', source: { data: large } }, { type: 'text', text: '你好 🌎' }], role: 'user' }, cwd: '/tmp/正确', uuid: 'turn', parentUuid: 'previous', parent_tool_use_id: 'agent', isSidechain: true, isMeta: false } + const content = JSON.stringify(entry) + '\n' + const result = await read(content) + expect(result.scan).toMatchObject({ omittedRecords: 0, projectedRecords: 1 }) + expect(result.entries[0]).toMatchObject({ complete: true, offset: 0, entry: { ...entry, message: { content: [{ type: 'image' }, { type: 'text', text: '你好 🌎' }], role: 'user' } } }) + expect(isSessionMetadataTextTruncated(result.entries[0]!.entry)).toBe(false) + expect(await readFile(file, 'utf8')).toBe(content) +}) + +test('ordinary records remain byte-semantically complete, including long notification text', async () => { + const entry = { message: { role: 'user', content: 'x'.repeat(70_000) }, custom: { untouched: true } } + const { entries, scan } = await read(JSON.stringify(entry)) + expect(entries[0]).toEqual({ entry, complete: false, offset: 0 }) + expect(scan.projectedRecords).toBe(0) + expect(isSessionMetadataTextTruncated(entries[0]!.entry)).toBe(false) +}) + +test('long text is bounded and explicitly marked as incomplete for semantic classifiers', async () => { + const { entries } = await read(JSON.stringify({ type: 'user', message: { role: 'user', content: [{ type: 'text', text: large }] } })) + const entry = entries[0]!.entry + expect(isSessionMetadataTextTruncated(entry)).toBe(true) + expect(JSON.stringify(entry).length).toBeLessThan(5000) +}) + +test('tool payloads are discarded while all Agent identity blocks remain available', async () => { + const { entries } = await read(JSON.stringify({ type: 'assistant', message: { role: 'assistant', content: [{ type: 'tool_use', input: { prompt: large }, id: 'task-1', name: 'Agent' }, { type: 'tool_result', tool_use_id: 'tool-2', content: large }] } })) + expect(entries[0]!.entry).toEqual({ type: 'assistant', message: { role: 'assistant', content: [{ type: 'tool_use', id: 'task-1', name: 'Agent' }, { type: 'tool_result', tool_use_id: 'tool-2' }] } }) + expect(isSessionMetadataTextTruncated(entries[0]!.entry)).toBe(false) +}) + +test('duplicate keys use JSON last-value semantics and prototype keys cannot mutate objects', async () => { + const { entries } = await read('{"unused":"' + large + '","type":"user","type":"system","repository":{"__proto__":{"polluted":true}},"message":{"role":"user"},"message":null}') + expect(entries[0]!.entry.type).toBe('system') + expect(entries[0]!.entry.message).toBe(null) + expect(Object.prototype.hasOwnProperty.call(entries[0]!.entry.repository, '__proto__')).toBe(true) + expect(({} as { polluted?: boolean }).polluted).toBeUndefined() +}) + +test('malformed large records are skipped without accepting their partial projected metadata', async () => { + const { entries, scan, skipped } = await read('{"type":"session-meta","workDir":"/wrong","body":"' + large + '",}\n' + JSON.stringify({ type: 'session-meta', workDir: '/right' }) + '\n') + expect(scan.omittedRecords).toBe(1) + expect(skipped).toBe(1) + expect(entries.map(item => item.entry)).toEqual([{ type: 'session-meta', workDir: '/right' }]) +}) + +test('unterminated large JSON is skipped but a valid non-newline tail is retained', async () => { + const bad = await read('{"message":{"content":"' + large) + expect(bad.entries).toEqual([]) + expect(bad.scan.omittedRecords).toBe(1) + const good = await read(JSON.stringify({ unused: large, cwd: '/tail' })) + expect(good.entries).toEqual([{ entry: { cwd: '/tail' }, complete: false, offset: 0 }]) +}) + +test('real oversized launch metadata and excessive structure fail explicitly', async () => { + await writeFile(file, JSON.stringify({ unused: large, workDir: 'x'.repeat(128 * 1024 + 1) })) + await expect(streamSessionMetadata(file, () => {})).rejects.toMatchObject({ code: 'SESSION_METADATA_TOO_LARGE' }) + await writeFile(file, '{"unused":"' + large + '","repository":' + '['.repeat(130) + '0' + ']'.repeat(130) + '}') + await expect(streamSessionMetadata(file, () => {})).rejects.toMatchObject({ code: 'SESSION_METADATA_TOO_LARGE' }) +}) + +test('honors byte ranges and propagates consumer errors and cancellation', async () => { + const prefix = JSON.stringify({ type: 'before' }) + '\n' + const body = JSON.stringify({ type: 'user', unused: large }) + '\n' + await writeFile(file, prefix + body + '{"type":"after"}\n') + const offsets: number[] = [] + const scan = await streamSessionMetadata(file, (_entry, _complete, offset) => offsets.push(offset), undefined, { startOffset: Buffer.byteLength(prefix), endOffset: Buffer.byteLength(prefix + body) }) + expect(offsets).toEqual([Buffer.byteLength(prefix)]) + expect(scan.nextOffset).toBe(Buffer.byteLength(prefix + body)) + await expect(streamSessionMetadata(file, () => { + throw new Error('consumer failure') + })).rejects.toThrow('consumer failure') + await expect(streamSessionMetadata(file, () => {}, AbortSignal.abort())).rejects.toThrow() +}) + +test('overlong retained metadata keys are rejected instead of being silently renamed', async () => { + await writeFile(file, JSON.stringify({ unused: large, repository: { ['x'.repeat(257)]: 'value' } })) + await expect(streamSessionMetadata(file, () => {})).rejects.toMatchObject({ code: 'SESSION_METADATA_TOO_LARGE' }) + const { entries } = await read(JSON.stringify({ ['x'.repeat(HISTORY_SEMANTIC_RECORD_BYTES + 1)]: 'unused', cwd: '/right' })) + expect(entries[0]!.entry).toEqual({ cwd: '/right' }) +}) + +test('duplicate message, content, and text keys discard obsolete truncation flags', async () => { + for (const record of [ + '{"message":{"content":"' + large + '"},"message":{"role":"user","content":"final"}}', + '{"message":{"role":"user","content":"' + large + '","content":"final"}}', + '{"message":{"role":"user","content":[{"type":"text","text":"' + large + '","text":"final"}]}}', + ]) { + const { entries } = await read(record) + expect(isSessionMetadataTextTruncated(entries[0]!.entry)).toBe(false) + } +}) diff --git a/src/server/services/sessionMetadataReader.ts b/src/server/services/sessionMetadataReader.ts new file mode 100644 index 00000000..e3694717 --- /dev/null +++ b/src/server/services/sessionMetadataReader.ts @@ -0,0 +1,293 @@ +import { open } from 'node:fs/promises' +import { StringDecoder } from 'node:string_decoder' +import parser, { type Token } from 'stream-json/parser.js' +import { ApiError } from '../middleware/errorHandler.js' +import { HISTORY_SEMANTIC_RECORD_BYTES } from './boundedSessionHistory.js' + +const METADATA_BYTES = 128 * 1024 +const TEXT_PREVIEW_CHARS = 4096 +const MAX_DEPTH = 128 +const truncatedTextEntries = new WeakSet() + +/** A projected long text is unsuitable for authoritative notification classification. */ +export function isSessionMetadataTextTruncated(entry: object): boolean { + return truncatedTextEntries.has(entry) +} + +const rootFields = new Set([ + 'type', 'subtype', 'uuid', 'parentUuid', 'parent_tool_use_id', 'isSidechain', 'isMeta', + 'cwd', 'timestamp', 'workDir', 'repository', 'worktreeSession', 'permissionMode', + 'runtimeProviderId', 'runtimeModelId', 'effortLevel', 'customTitle', 'aiTitle', + 'entrypoint', 'message', 'content', +]) +const blockFields = new Set(['type', 'name', 'id', 'tool_use_id', 'text']) +type Frame = { + path: string[] + value?: Record | unknown[] + key: string + index: number +} + +function limit(): ApiError { + return new ApiError(413, 'Session metadata exceeds its resource budget', 'SESSION_METADATA_TOO_LARGE') +} + +/** Assemble only the small structural envelope; the parser still validates every byte. */ +function createProjection() { + const frames: Frame[] = [] + let root: unknown + let bytes = 0 + let scalar = '' + let scalarPath: string[] = [] + let scalarMode: 'key' | 'string' | 'number' | undefined + let scalarOverflow = false + const truncatedSlots = new WeakMap>() + const charge = (size: number) => { + bytes += size + if (bytes > METADATA_BYTES) throw limit() + } + const nextPath = (): string[] => { + const parent = frames.at(-1) + return parent ? [...parent.path, Array.isArray(parent.value) ? String(parent.index) : parent.key] : [] + } + const selected = (path: string[]): boolean => { + if (!path.length) return true + if (path.includes('\0unselected')) return false + if (!rootFields.has(path[0]!)) return false + if (path[0] !== 'message') return true + if (path.length === 1) return true + if (path[1] === 'role') return path.length === 2 + if (path[1] !== 'content') return false + return path.length <= 3 || path.length === 4 && blockFields.has(path[3]!) + } + const preview = (path: string[]) => path.length === 1 && path[0] === 'content' + || path[0] === 'message' && path[1] === 'content' && (path.length === 2 || path.length === 4 && path[3] === 'text') + const attach = (value: unknown, path: string[], truncated = false) => { + const parent = frames.at(-1) + if (selected(path)) { + if (!parent) root = value + else if (parent.value) { + charge(16 + (Array.isArray(parent.value) ? 0 : Buffer.byteLength(parent.key))) + const slot = Array.isArray(parent.value) ? String(parent.value.length) : parent.key + const slots = truncatedSlots.get(parent.value) ?? new Set() + if (truncated) slots.add(slot) + else slots.delete(slot) + truncatedSlots.set(parent.value, slots) + if (Array.isArray(parent.value)) parent.value.push(value) + else Object.defineProperty(parent.value, parent.key, { value, writable: true, configurable: true, enumerable: true }) + } + } + if (parent) parent.index++ + } + const token = (token: Token) => { + switch (token.name) { + case 'startObject': + case 'startArray': { + if (frames.length >= MAX_DEPTH) throw limit() + const path = nextPath() + const value = selected(path) ? token.name === 'startArray' ? [] : Object.create(null) : undefined + attach(value, path) + frames.push({ path, value, key: '', index: 0 }) + break + } + case 'endObject': + case 'endArray': + frames.pop() + break + case 'startKey': + scalarMode = 'key' + scalar = '' + scalarOverflow = false + break + case 'startString': + case 'startNumber': + scalarMode = token.name === 'startString' ? 'string' : 'number' + scalarPath = nextPath() + scalar = '' + scalarOverflow = false + break + case 'stringChunk': + case 'numberChunk': { + if (scalarMode !== 'key' && !selected(scalarPath)) break + const bound = scalarMode === 'key' ? 256 : preview(scalarPath) ? TEXT_PREVIEW_CHARS : METADATA_BYTES + const available = bound - scalar.length + if (token.value.length > available) { + scalarOverflow = true + if (scalarMode !== 'key' && !preview(scalarPath)) throw limit() + } + scalar += token.value.slice(0, Math.max(0, available)) + break + } + case 'endKey': { + // A retained metadata map must never silently rename a key. Unknown + // root/body fields can be discarded without affecting launch state. + const parent = frames.at(-1)! + if (scalarOverflow && parent.value && parent.path.length && parent.path[0] !== 'message') throw limit() + frames.at(-1)!.key = scalarOverflow ? '\0unselected' : scalar + scalarMode = undefined + break + } + case 'endString': + case 'endNumber': { + if (selected(scalarPath)) { + charge(Buffer.byteLength(scalar)) + } + attach(token.name === 'endNumber' ? Number(scalar) : scalar, scalarPath, scalarOverflow && scalarPath[0] === 'message') + scalarMode = undefined + break + } + case 'trueValue': + case 'falseValue': + case 'nullValue': + attach(token.value, nextPath()) + break + } + } + return { + token, + result: () => { + if (root && typeof root === 'object' && !Array.isArray(root)) { + const message = (root as Record).message as Record | undefined + if (message && typeof message === 'object' && (truncatedSlots.get(message)?.has('content') + || Array.isArray(message.content) && message.content.some(block => block && typeof block === 'object' && truncatedSlots.get(block)?.has('text')))) { + truncatedTextEntries.add(root) + } + return root as Record + } + return undefined + }, + } +} + +/** + * Preserve ordinary JSONL semantics. Oversized bodies use a validated streaming + * projection instead of making otherwise small launch metadata unavailable. + * This is not a semantic replay reader: large message payloads are not returned. + */ +export async function streamSessionMetadata( + filePath: string, + onEntry: (entry: Record, completeLine: boolean, byteStart: number) => void, + signal?: AbortSignal, + options: { + startOffset?: number + endOffset?: number + onSkipped?: () => void + } = {}, +): Promise<{ + sourceVersion: string + omittedRecords: number + projectedRecords: number + scannedBytes: number + nextOffset: number +}> { + const handle = await open(filePath, 'r') + const check = () => { + if (signal?.aborted) throw signal.reason ?? new DOMException('Aborted', 'AbortError') + } + let active: ReturnType | undefined + try { + const stat = await handle.stat({ bigint: true }) + const size = Math.min(Number(stat.size), options.endOffset ?? Number(stat.size)) + const firstOffset = options.startOffset ?? 0 + let lineStart = firstOffset + let nextOffset = firstOffset + let parts: Buffer[] = [] + let length = 0 + let projection: ReturnType | undefined + let parseError: Error | undefined + let projectionError: unknown + let decoder: StringDecoder | undefined + let omittedRecords = 0 + let projectedRecords = 0 + const feed = async (part: Buffer) => { + if (!active || parseError || projectionError) return + const text = decoder!.write(part) + await new Promise(resolve => active!.write(text, () => resolve())) + } + const flush = async (completeLine: boolean) => { + let entry: unknown + if (active) { + if (!parseError && !projectionError) { + await new Promise(resolve => { + active!.once('end', resolve) + active!.once('error', resolve) + active!.end(decoder!.end()) + }) + } + active.destroy() + active = undefined + if (projectionError) throw projectionError + if (!parseError) { + entry = projection!.result() + projectedRecords++ + } + } else if (length) { + try { + entry = JSON.parse((parts.length === 1 ? parts[0]! : Buffer.concat(parts, length)).toString('utf8')) + } catch (error) { + parseError = error as Error + } + } + if (parseError) { + omittedRecords++ + options.onSkipped?.() + } + else if (entry && typeof entry === 'object' && !Array.isArray(entry)) onEntry(entry as Record, completeLine, lineStart) + parts = [] + length = 0 + projection = undefined + parseError = undefined + projectionError = undefined + } + const chunk = Buffer.allocUnsafe(64 * 1024) + for (let offset = firstOffset; offset < size;) { + check() + const { bytesRead } = await handle.read(chunk, 0, Math.min(chunk.length, size - offset), offset) + if (!bytesRead) throw new ApiError(409, 'Session history changed during recovery', 'HISTORY_CHANGED') + offset += bytesRead + let start = 0 + while (start < bytesRead) { + check() + const found = chunk.indexOf(10, start) + const end = found >= 0 && found < bytesRead ? found : bytesRead + const part = chunk.subarray(start, end) + length += part.length + if (!active && length > HISTORY_SEMANTIC_RECORD_BYTES) { + projection = createProjection() + decoder = new StringDecoder('utf8') + active = parser.asStream({ packValues: false, streamValues: true }) + active.on('error', error => { + parseError = error + }) + active.on('data', (token: Token) => { + if (!projectionError) { + try { + projection!.token(token) + } catch (error) { + projectionError = error + } + } + }) + for (const saved of parts) await feed(saved) + parts = [] + } + if (active) await feed(part) + else parts.push(Buffer.from(part)) + start = end + 1 + if (end < bytesRead) { + await flush(true) + lineStart = offset - bytesRead + end + 1 + nextOffset = lineStart + await new Promise(resolve => setImmediate(resolve)) + } + } + } + if (length) await flush(false) + const after = await handle.stat({ bigint: true }) + if (after.size < stat.size || after.size === stat.size && after.mtimeNs !== stat.mtimeNs) throw new ApiError(409, 'Session changed during recovery; retry', 'HISTORY_CHANGED') + return { sourceVersion: `${stat.dev}:${stat.ino}:${stat.size}:${stat.mtimeNs}`, omittedRecords, projectedRecords, scannedBytes: size - firstOffset, nextOffset } + } finally { + active?.destroy() + await handle.close() + } +} diff --git a/src/server/services/sessionService.metadata.test.ts b/src/server/services/sessionService.metadata.test.ts new file mode 100644 index 00000000..d3cb522b --- /dev/null +++ b/src/server/services/sessionService.metadata.test.ts @@ -0,0 +1,74 @@ +import { afterEach, beforeEach, expect, test } from 'bun:test' +import { mkdtemp, mkdir, readFile, rm, writeFile } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { SessionService } from './sessionService.js' +import { resetSettingsCache } from '../../utils/settings/settingsCache.js' +import { sanitizePath } from '../../utils/sessionStoragePortable.js' +import { HISTORY_SEMANTIC_RECORD_BYTES } from './boundedSessionHistory.js' + +const sessionId = 'aaaaaaaa-bbbb-cccc-dddd-eeeeeeeeeeee' +let directory: string +let file: string +let oldHome: string | undefined +let oldConfig: string | undefined +beforeEach(async () => { + directory = await mkdtemp(join(tmpdir(), 'session-metadata-')) + oldHome = process.env.HOME + oldConfig = process.env.CLAUDE_CONFIG_DIR + process.env.HOME = directory + process.env.CLAUDE_CONFIG_DIR = directory + resetSettingsCache() + const project = join(directory, 'projects', sanitizePath(directory)) + await mkdir(project, { recursive: true }) + file = join(project, `${sessionId}.jsonl`) +}) +afterEach(async () => { + if (oldHome === undefined) delete process.env.HOME + else process.env.HOME = oldHome + if (oldConfig === undefined) delete process.env.CLAUDE_CONFIG_DIR + else process.env.CLAUDE_CONFIG_DIR = oldConfig + resetSettingsCache() + await rm(directory, { recursive: true, force: true }) +}) +const body = 'x'.repeat(HISTORY_SEMANTIC_RECORD_BYTES + 1) +async function seed(entries: unknown[]) { + await writeFile(file, entries.map(entry => JSON.stringify(entry)).join('\n') + '\n') +} + +test('cold session metadata and restart survive a large image and retain the latest runtime selection', async () => { + await seed([ + { type: 'session-meta', workDir: directory, runtimeProviderId: 'old', runtimeModelId: 'old', effortLevel: 'high' }, + { type: 'user', cwd: directory, message: { role: 'user', content: [{ type: 'image', source: { data: body } }, { type: 'text', text: 'Image task' }] } }, + { type: 'session-meta', workDir: directory, runtimeProviderId: 'new', runtimeModelId: 'new' }, + ]) + const before = await readFile(file) + for (let restart = 0; restart < 2; restart++) { + const service = new SessionService() + expect(await service.getSessionWorkDir(sessionId)).toBe(directory) + expect(await service.getSessionLaunchInfo(sessionId)).toMatchObject({ workDir: directory, transcriptMessageCount: 1, runtimeProviderId: 'new', runtimeModelId: 'new' }) + expect((await service.getSessionLaunchInfo(sessionId))?.effortLevel).toBeUndefined() + } + expect(await readFile(file)).toEqual(before) +}) + +test('a sole oversized legacy turn remains resumable and cannot be mistaken for an empty placeholder', async () => { + await seed([{ type: 'assistant', message: { role: 'assistant', content: [{ type: 'tool_use', name: 'Bash', input: { command: body } }] }, cwd: directory, repository: { repoName: 'fixture', projectRoot: directory } }]) + const info = await new SessionService().getSessionLaunchInfo(sessionId) + expect(info).toMatchObject({ transcriptMessageCount: 1, workDir: directory, repository: { repoName: 'fixture', projectRoot: directory } }) +}) + +test('large bodies preserve system turn counts, isMeta exclusion, and explicit worktree clearing', async () => { + await seed([ + { type: 'worktree-state', worktreeSession: { worktreePath: '/fixture/worktree', worktreeName: 'fixture' } }, + { type: 'system', message: { role: 'system', content: body }, cwd: directory }, + { type: 'user', isMeta: true, message: { role: 'user', content: body } }, + { type: 'worktree-state', worktreeSession: null }, + ]) + expect(await new SessionService().getSessionLaunchInfo(sessionId)).toMatchObject({ transcriptMessageCount: 1, workDir: directory, worktreeSession: null }) +}) + +test('actual over-budget launch metadata still blocks unsafe launch', async () => { + await seed([{ type: 'session-meta', workDir: body }]) + await expect(new SessionService().getSessionLaunchInfo(sessionId)).rejects.toMatchObject({ code: 'SESSION_METADATA_TOO_LARGE' }) +}) diff --git a/src/server/services/sessionService.ts b/src/server/services/sessionService.ts index 8631b66e..1f6889cf 100644 --- a/src/server/services/sessionService.ts +++ b/src/server/services/sessionService.ts @@ -11,6 +11,7 @@ import { recoverBoundedSessionHistory, type SessionHistoryRecovery } from './ses * 确保 Desktop App 与 CLI 的数据完全互通。 */ +import { streamSessionMetadata } from './sessionMetadataReader.js' import { HISTORY_SEMANTIC_RECORD_BYTES, HISTORY_PAGE_BYTES, displayPreview, readBoundedHistoryPage, streamBoundedHistory, withHistoryReadBudget, type HistoryPageInfo } from './boundedSessionHistory.js' import { constants, createReadStream, createWriteStream, type Stats } from 'node:fs' import { createHash } from 'node:crypto' @@ -691,8 +692,8 @@ export class SessionService { private readonly subagentLookupCache = new Map() private readonly historyRecoveryCache = new Map() - private readonly metadataProjectionCache = new Map() - private readonly metadataProjectionRequests = new Map>() + private readonly metadataProjectionCache = new Map() + private readonly metadataProjectionRequests = new Map>() private readonly sessionHistoryRequests = new Map { const stat = await fs.stat(filePath, { bigint: true }) const signature = `${stat.dev}:${stat.ino}:${stat.size}:${stat.mtimeNs}` @@ -1306,10 +1306,10 @@ export class SessionService { // a handful of maliciously large scalar values or repository fields. if (Buffer.byteLength(JSON.stringify(state)) > 128 * 1024) throw new ApiError(413, 'Session metadata exceeds its resource budget', 'SESSION_METADATA_TOO_LARGE') } - const scan = await streamBoundedHistory(filePath, (entry, completeLine) => { + const scan = await streamSessionMetadata(filePath, (entry, completeLine) => { apply(launch, entry as RawEntry) if (completeLine) apply(summary, entry as RawEntry) - }, undefined, { maxRecordBytes: HISTORY_SEMANTIC_RECORD_BYTES }) + }) const shared = (state: typeof summary) => ({ ...(state.permissionMode ? { permissionMode: state.permissionMode } : {}), ...(state.runtimeProviderId !== undefined ? { runtimeProviderId: state.runtimeProviderId } : {}), @@ -1336,7 +1336,6 @@ export class SessionService { ...shared(launch), }, customTitle: launch.nonemptyCustomTitle, - complete: scan.oversizedRecords === 0, } this.metadataProjectionCache.delete(key) this.metadataProjectionCache.set(key, { signature: scan.sourceVersion, ...result }) @@ -4608,7 +4607,6 @@ export class SessionService { if (!found) return null const projection = await this.getMetadataProjection(found.filePath, found.projectDir) - if (!projection.complete) throw new ApiError(413, 'Session metadata contains oversized records', 'SESSION_METADATA_INCOMPLETE') return projection.launchInfo.workDir } @@ -4638,7 +4636,6 @@ export class SessionService { if (!found) return memory ? { ...memory, transcriptMessageCount: 0 } : null const projection = await this.getMetadataProjection(found.filePath, found.projectDir) - if (!projection.complete) throw new ApiError(413, 'Session metadata contains oversized records', 'SESSION_METADATA_INCOMPLETE') const projected = projection.launchInfo return { ...projected, ...memory, transcriptMessageCount: projected.transcriptMessageCount } }