diff --git a/src/server/services/localIndex/sessionIndex.ts b/src/server/services/localIndex/sessionIndex.ts index 0fd9077c..82868adc 100644 --- a/src/server/services/localIndex/sessionIndex.ts +++ b/src/server/services/localIndex/sessionIndex.ts @@ -183,50 +183,44 @@ function boundedInteger(value: number | undefined, fallback: number): number { return Math.max(0, Math.trunc(value!)) } +/** A picker query must not scan an unbounded session table synchronously. + * Past this many rows the page is a lower bound: exact matches already seen + * stay ranked first, and the caller reports the total as incomplete. */ +const UNICODE_METADATA_SCAN_ROWS = 4096 + function searchUnicodeSessionMetadata(operation: LocalIndexReadOperation, needle: string, limit: number, offset: number): SessionIndexPage { const rank = (row: SessionRow): number => { const names = [row.title.toLowerCase(), row.session_id.toLowerCase()] if (![...names, row.work_dir?.toLowerCase() ?? '', row.project_path.toLowerCase()].some(value => value.includes(needle))) return -1 return names.includes(needle) ? 3 : names.some(value => value.startsWith(needle)) ? 2 : names.some(value => value.includes(needle)) ? 1 : 0 } - const scan = (visit: (row: SessionRow) => void) => { + const scan = (visit: (row: SessionRow) => boolean) => { let last: SessionRow | undefined - while (true) { + let scanned = 0 + while (scanned < UNICODE_METADATA_SCAN_ROWS) { const rows = operation.all(` SELECT * FROM sessions ${last ? 'WHERE modified_at_ms < ? OR (modified_at_ms = ? AND session_id > ?) OR (modified_at_ms = ? AND session_id = ? AND transcript_path > ?)' : ''} ORDER BY modified_at_ms DESC, session_id ASC, transcript_path ASC LIMIT 256 `, ...(last ? [last.modified_at_ms, last.modified_at_ms, last.session_id, last.modified_at_ms, last.session_id, last.transcript_path] : [])) - for (const row of rows) visit(row) - if (rows.length < 256) break + for (const row of rows) { + if (!visit(row)) return + if (++scanned >= UNICODE_METADATA_SCAN_ROWS) return + } + if (rows.length < 256) return last = rows[rows.length - 1] } } - // Count rank buckets before selecting the requested page. Two bounded scans - // avoid retaining offset + limit rows for arbitrarily deep pagination. - const counts = [0, 0, 0, 0] - scan(row => { const value = rank(row); if (value >= 0) counts[value]!++ }) - const total = counts.reduce((sum, value) => sum + value, 0) - if (offset >= total) return { sessions: [], total } - const skips = [0, 0, 0, 0] - const takes = [0, 0, 0, 0] - let remainingSkip = offset - let remainingTake = limit - for (let value = 3; value >= 0; value--) { - skips[value] = Math.min(remainingSkip, counts[value]!) - remainingSkip -= skips[value]! - takes[value] = Math.min(remainingTake, counts[value]! - skips[value]!) - remainingTake -= takes[value]! - } + // One pass. The scan cap bounds how many matches can be retained, so the + // page can be sliced after ranking instead of scanning the table twice. const buckets: SessionRow[][] = [[], [], [], []] scan(row => { const value = rank(row) - if (value < 0 || takes[value] === 0) return - if (skips[value]! > 0) { skips[value]!--; return } - buckets[value]!.push(row) - takes[value]!-- + if (value >= 0) buckets[value]!.push(row) + return true }) - return { sessions: buckets.reverse().flat().map(sessionFromRow), total } + const ordered = [3, 2, 1, 0].flatMap(value => buckets[value]!) + return { sessions: ordered.slice(offset, offset + limit).map(sessionFromRow), total: ordered.length } } function sessionFromRow(row: SessionRow): IndexedSessionRow { diff --git a/src/server/services/searchService.ts b/src/server/services/searchService.ts index 632f60e7..1653a58a 100644 --- a/src/server/services/searchService.ts +++ b/src/server/services/searchService.ts @@ -307,9 +307,15 @@ export class SearchService { // --------------------------------------------------------------------------- /** Menu suggestions deliberately trust the disposable index, never opening canonical transcripts. */ - async searchSessionSuggestions(query: string, options: { limit?: number; signal?: AbortSignal } = {}) { + async searchSessionSuggestions(query: string, options: { limit?: number; signal?: AbortSignal; deadlineMs?: number } = {}) { throwIfAborted(options.signal) if (!query.trim()) return { sessions: [], truncated: false, indexUnavailable: false } + // The FTS lookup is one synchronous statement, so the deadline is checked + // before it starts. A picker that already spent its budget on metadata + // returns that page instead of stalling every other request. + if (options.deadlineMs !== undefined && Date.now() > options.deadlineMs) { + return { sessions: [], truncated: true, indexUnavailable: false } + } const projected = this.suggestIndexedSessions(query, options) throwIfAborted(options.signal) return { diff --git a/src/server/services/sessionCollaborationHost.test.ts b/src/server/services/sessionCollaborationHost.test.ts index e8c3b87b..530e9f5f 100644 --- a/src/server/services/sessionCollaborationHost.test.ts +++ b/src/server/services/sessionCollaborationHost.test.ts @@ -112,7 +112,7 @@ test('a full metadata page skips transcript search and returns immediately from const service = await getSessionCollaborationService() const result = await service.candidates('修复') expect(result.sessions.map(item => item.sessionId)).toEqual(sessions.map(item => item.id)) - expect(metadata).toHaveBeenCalledWith('修复', { limit: 30, offset: 0, signal: undefined }) + expect(metadata).toHaveBeenCalledWith('修复', expect.objectContaining({ limit: 30, offset: 0, signal: undefined })) expect(await service.list({ query: '修复' })).toMatchObject({ total: 20_000, totalIsLowerBound: true, truncated: true }) expect(list).not.toHaveBeenCalled() expect(fullText).not.toHaveBeenCalled() diff --git a/src/server/services/sessionCollaborationHost.ts b/src/server/services/sessionCollaborationHost.ts index 7a3ee72a..a5f170a9 100644 --- a/src/server/services/sessionCollaborationHost.ts +++ b/src/server/services/sessionCollaborationHost.ts @@ -42,13 +42,19 @@ export async function getSessionCollaborationService(): Promise deadlineMs) { + return { ...metadata, truncated: true, totalIsLowerBound: true } + } + const content = await searchService.searchSessionSuggestions(needle, { limit: 100, signal, deadlineMs }) signal?.throwIfAborted() const details = new Map(sessionService.getSessionSuggestionMetadata(content.sessions.map(item => item.sessionId)).map(item => [item.id, item])) const normalized = needle.toLowerCase() diff --git a/src/server/services/sessionService.ts b/src/server/services/sessionService.ts index a1d0e6dd..21066f1a 100644 --- a/src/server/services/sessionService.ts +++ b/src/server/services/sessionService.ts @@ -10,7 +10,7 @@ import { recoverBoundedSessionHistory, type SessionHistoryRecovery } from './ses */ import { HISTORY_SEMANTIC_RECORD_BYTES, HISTORY_PAGE_BYTES, displayPreview, readBoundedHistoryPage, streamBoundedHistory, withHistoryReadBudget, type HistoryPageInfo } from './boundedSessionHistory.js' -import { constants, createReadStream, type Stats } from 'node:fs' +import { constants, createReadStream, createWriteStream, type Stats } from 'node:fs' import { createHash } from 'node:crypto' import * as fs from 'node:fs/promises' import * as path from 'node:path' @@ -3375,7 +3375,7 @@ export class SessionService { } /** Metadata-only reference lookup: no per-result transcript or workspace hydration. */ - async searchSessionMetadata(query: string, options: { limit?: number; offset?: number; signal?: AbortSignal } = {}): Promise<{ sessions: Array<{ id: string; title: string; workDir: string | null; projectPath: string; modifiedAt: string }>; total: number }> { + async searchSessionMetadata(query: string, options: { limit?: number; offset?: number; signal?: AbortSignal; deadlineMs?: number } = {}): Promise<{ sessions: Array<{ id: string; title: string; workDir: string | null; projectPath: string; modifiedAt: string }>; total: number; truncated?: boolean }> { options.signal?.throwIfAborted() this.syncSharedMutationEpoch() const limit = Math.min(100, Math.max(1, options.limit ?? 30)) @@ -3394,8 +3394,12 @@ export class SessionService { const scope = this.getConfigDir() this.prepareSessionListCaches(scope) const rows: Array<{ id: string; title: string; workDir: string | null; projectPath: string; modifiedAt: string }> = [] + let truncated = false for (const file of await this.discoverSessionFiles(undefined, scope)) { options.signal?.throwIfAborted() + // The picker shares this process with every other request. Stop walking + // transcripts once its budget is spent and return the rows already read. + if (options.deadlineMs !== undefined && Date.now() > options.deadlineMs) { truncated = true; break } try { const summary = await this.getCachedSessionListSummary(file.filePath, file.projectDir, await fs.stat(file.filePath), scope) rows.push({ id: file.sessionId, title: summary.title, workDir: summary.workDir, projectPath: file.projectDir, modifiedAt: summary.modifiedAt }) @@ -3409,7 +3413,7 @@ export class SessionService { } const matches = rows.filter(row => [row.title, row.id, row.workDir ?? '', row.projectPath].some(value => value.toLowerCase().includes(needle))) matches.sort((a, b) => rank(b) - rank(a) || Date.parse(b.modifiedAt) - Date.parse(a.modifiedAt) || a.id.localeCompare(b.id) || a.projectPath.localeCompare(b.projectPath)) - return { sessions: matches.slice(offset, offset + limit), total: matches.length } + return { sessions: matches.slice(offset, offset + limit), total: matches.length, ...(truncated ? { truncated } : {}) } } /** List all sessions, optionally filtered by physical project path. */ @@ -4608,15 +4612,33 @@ export class SessionService { throw ApiError.notFound(`Session not found: ${sessionId}`) } - const entries = await this.readJsonlFile(found.filePath) - const workDir = this.resolveWorkDirFromEntries(entries, found.projectDir) || fallbackWorkDir || process.cwd() - const repository = this.resolveRepositoryFromEntries(entries) + // Only the newest metadata survives a clear. Walk the transcript one + // record at a time so a large session is not parsed into one array. + const preserved = { workDir: undefined as string | undefined, cwd: undefined as string | undefined, + repository: undefined as PreparedSessionWorkspace['repository'] | undefined, + permissionMode: undefined as string | undefined } + await streamBoundedHistory(found.filePath, entry => { + const record = entry as RawEntry + if (record.type === 'session-meta') { + if (typeof (record as Record).workDir === 'string') preserved.workDir = (record as Record).workDir as string + if (typeof record.permissionMode === 'string' && VALID_SESSION_PERMISSION_MODES.has(record.permissionMode)) preserved.permissionMode = record.permissionMode + } + if (typeof record.cwd === 'string' && record.cwd.trim()) preserved.cwd = record.cwd + const repository = (record as Record).repository + if (repository && typeof repository === 'object') preserved.repository = repository as PreparedSessionWorkspace['repository'] + }, undefined, { maxRecordBytes: HISTORY_SEMANTIC_RECORD_BYTES }).catch(error => { + if ((error as NodeJS.ErrnoException).code !== 'ENOENT') throw error + }) + const workDir = (preserved.workDir && normalizeDriveRootPathForPlatform(preserved.workDir)) + || (preserved.cwd && normalizeDriveRootPathForPlatform(preserved.cwd)) + || this.desanitizePath(found.projectDir) || fallbackWorkDir || process.cwd() + const repository = preserved.repository const permissionMode = ( preservedPermissionMode && VALID_SESSION_PERMISSION_MODES.has(preservedPermissionMode) ) ? preservedPermissionMode - : this.resolvePermissionModeFromEntries(entries) + : preserved.permissionMode const now = new Date().toISOString() const initialEntry = { @@ -4829,33 +4851,51 @@ export class SessionService { } const removedIds = new Set(removedMessageIds) - const filteredEntries = entries.filter( - (entry) => { - if (typeof entry.uuid !== 'string') return true - if (removedIds.has(entry.uuid)) return false - if ( - entry.message?.role && - (entry.type === 'user' || entry.type === 'assistant' || entry.type === 'system') - ) { - return remainingMessageIds.has(entry.uuid) - } - return true - }, - ) - - const content = - filteredEntries.length > 0 - ? filteredEntries.map((entry) => JSON.stringify(entry)).join('\n') + '\n' - : '' + const kept = (entry: RawEntry): boolean => { + if (typeof entry.uuid !== 'string') return true + if (removedIds.has(entry.uuid)) return false + if ( + entry.message?.role && + (entry.type === 'user' || entry.type === 'assistant' || entry.type === 'system') + ) { + return remainingMessageIds.has(entry.uuid) + } + return true + } + // Copy the original lines that survive. Re-serializing every retained entry + // would hold the whole transcript as one string on the request thread. const transcriptStats = await fs.stat(found.filePath) const tempFilePath = `${found.filePath}.rewind-${crypto.randomUUID()}.tmp` + const output = createWriteStream(tempFilePath, { mode: transcriptStats.mode }) + let failed = false + const fail = (error: Error) => { if (!failed) { failed = true; output.destroy(error) } } try { - await fs.writeFile(tempFilePath, content, { - encoding: 'utf-8', - mode: transcriptStats.mode, + await new Promise((resolve, reject) => { + output.on('error', reject) + output.on('finish', resolve) + void (async () => { + const input = createReadStream(found.filePath, { encoding: 'utf8' }) + try { + for await (const line of createInterface({ input, crlfDelay: Infinity })) { + if (line.trim()) { + try { + if (!kept(JSON.parse(line) as RawEntry)) continue + } catch { /* Keep a line the transcript reader would also keep. */ } + } + if (!output.write(`${line}\n`)) await new Promise(resume => output.once('drain', resume)) + } + output.end() + } catch (error) { + fail(error instanceof Error ? error : new Error(String(error))) + } finally { + input.destroy() + } + })() }) + if (!this.shouldPersistSession()) return { removedCount: 0, removedMessageIds: [] } await fs.rename(tempFilePath, found.filePath) } finally { + output.destroy() await fs.rm(tempFilePath, { force: true }) } this.invalidateSessionListCache()