mirror of
https://github.com/NanmiCoder/claude-code-haha.git
synced 2026-10-10 11:53:10 +08:00
fix(sessions): keep mention search and transcript edits off the request thread
A Chinese @ query scanned the whole session table twice and then ran a body search, so every other request queued behind it. Cap that scan, skip the body lookup once the budget is spent, and stream session clears and trims instead of holding the transcript in memory.
This commit is contained in:
@@ -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<SessionRow>(`
|
||||
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 {
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -42,13 +42,19 @@ export async function getSessionCollaborationService(): Promise<SessionCollabora
|
||||
sessions: {
|
||||
async list({ query, limit, offset, signal }) {
|
||||
const needle = query?.trim() ?? ''
|
||||
const metadata = await sessionService.searchSessionMetadata(needle, { limit, offset, signal })
|
||||
// Metadata ranking and the body lookup are both synchronous SQLite.
|
||||
// Whichever is still pending once this elapses is skipped so the
|
||||
// picker cannot hold the process while other requests wait.
|
||||
const deadlineMs = Date.now() + 300
|
||||
const metadata = await sessionService.searchSessionMetadata(needle, { limit, offset, signal, deadlineMs })
|
||||
signal?.throwIfAborted()
|
||||
// A full metadata page already answers the picker. Do not wait for
|
||||
// full-text matches (or canonical transcript validation) to display it.
|
||||
if (!needle) return metadata
|
||||
if (metadata.sessions.length === limit) return { ...metadata, truncated: true, totalIsLowerBound: true }
|
||||
const content = await searchService.searchSessionSuggestions(needle, { limit: 100, signal })
|
||||
if (metadata.truncated || metadata.sessions.length === limit || Date.now() > 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()
|
||||
|
||||
@@ -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<string, unknown>).workDir === 'string') preserved.workDir = (record as Record<string, unknown>).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<string, unknown>).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<void>((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<void>(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()
|
||||
|
||||
Reference in New Issue
Block a user