refactor(ai-vault): split the session scanner into a transcript reader and consumers

The scanner's only output was the Session History summary; a second reader
of the same transcripts (a search index) had nowhere to plug in without
hooking the parse itself. Extract a reader that owns each file read, keeps
the resumable cursor and publishes every decoded message to registered
consumers. The session list stays a fold inside the parser and is the only
consumer here. Parsers take an optional message sink instead of a scope.

Also: read Cursor chats/<md5>/<uuid>/meta.json for cwd, title and
timestamps (Cursor transcripts carry only role and message); share the
lazily spawned worker-thread host between the OpenCode SQLite reader and
the port-scan probe; probe OpenCode's schema before querying; keep the
newest-N discovery set with a bounded insert instead of sort+slice.

Session list output is byte-identical to main across all 18 providers
cold and append-resumed; the one Cursor session gains cwd/timestamps
from meta.json.
This commit is contained in:
Jinwoo-H
2026-09-08 22:33:31 -04:00
parent e80fae0c4d
commit dae2b4981e
40 changed files with 2178 additions and 665 deletions
@@ -0,0 +1,31 @@
import { expect, it } from 'vitest'
import { SessionNewestFiles } from './session-newest-files'
import type { FileWithMtime } from './session-scanner-types'
function file(i: number): FileWithMtime {
const mtimeMs = (i * 7919) % 997
return { path: String(i), mtimeMs, modifiedAt: new Date(mtimeMs).toISOString() }
}
it('retains at most 12 of 100,000 candidates with stable newest-first ties', () => {
const all = Array.from({ length: 100_000 }, (_, i) => file(i))
const retained = new SessionNewestFiles(12)
let peak = 0
for (const candidate of all) {
retained.add(candidate)
peak = Math.max(peak, retained.size)
}
expect(peak).toBe(12)
expect(retained.newest()).toEqual(all.sort((a, b) => b.mtimeMs - a.mtimeMs).slice(0, 12))
})
it('supports full backfill and empty requests', () => {
const all = new SessionNewestFiles(Infinity)
const none = new SessionNewestFiles(0)
for (let i = 0; i < 100; i++) {
all.add(file(i))
none.add(file(i))
}
expect(all.newest()).toHaveLength(100)
expect(none.newest()).toEqual([])
})
+50
View File
@@ -0,0 +1,50 @@
import type { FileWithMtime } from './session-scanner-types'
/** Retain only the requested newest files, preserving traversal order on ties. */
export class SessionNewestFiles {
private readonly files: FileWithMtime[] = []
private readonly limit: number
constructor(limit: number) {
this.limit = Math.max(0, Math.trunc(limit) || 0)
}
add(file: FileWithMtime): void {
// The backfill enumerates with no limit; skip the insert search entirely.
if (!Number.isFinite(this.limit)) {
this.files.push(file)
return
}
if (this.limit <= 0) {
return
}
const last = this.files.at(-1)
if (this.files.length >= this.limit && last && file.mtimeMs <= last.mtimeMs) {
return
}
let low = 0
let high = this.files.length
while (low < high) {
const middle = (low + high) >>> 1
if (this.files[middle].mtimeMs >= file.mtimeMs) {
low = middle + 1
} else {
high = middle
}
}
this.files.splice(low, 0, file)
if (this.files.length > this.limit) {
this.files.pop()
}
}
get size(): number {
return this.files.length
}
/** The unbounded path appends in traversal order, so the sort is not redundant. */
newest(): FileWithMtime[] {
return [...this.files].sort((a, b) => b.mtimeMs - a.mtimeMs)
}
}
@@ -0,0 +1,98 @@
import type { AiVaultSession } from '../../shared/ai-vault-types'
import type { ResumableSessionParseState } from './session-scanner-types'
import type { TranscriptMessageChannel } from './session-transcript-channel'
// Sized past the default recency cap (1000) plus the in-scope cap (2000) so a
// full steady-state result set stays resident between forced rescans.
const MAX_CACHE_ENTRIES = 4096
export type SessionParseResumePoint = {
state: ResumableSessionParseState
// Byte offset just past the last complete ('\n'-terminated) line consumed;
// a trailing unterminated line is deliberately left before this point.
byteOffset: number
// Bound to the cached state, which keeps the reference its parsers were built
// with; a resumed read re-points this channel instead of replacing it.
channel: TranscriptMessageChannel
}
export type SessionParseCacheEntry = {
mtimeMs: number
sizeBytes: number | null
platform: NodeJS.Platform
session: AiVaultSession | null
resume: SessionParseResumePoint | null
}
const cache = new Map<string, SessionParseCacheEntry>()
export function resetSessionParseCacheForTests(): void {
cache.clear()
}
// Drops one entry after its file is deleted. Cleanliness, not correctness:
// discovery walks disk first, so a trashed file is never rediscovered anyway.
export function invalidateSessionParseCacheEntry(path: string): void {
cache.delete(path)
}
// Persisted subset of a cache entry: the non-serializable `resume` parser
// state is dropped (see session-parse-cache-persistence.ts).
export type PersistedSessionParseCacheEntry = Omit<SessionParseCacheEntry, 'resume'>
export function snapshotSessionParseCacheForPersistence(): [
string,
PersistedSessionParseCacheEntry
][] {
return [...cache].map(([path, entry]): [string, PersistedSessionParseCacheEntry] => [
path,
{
mtimeMs: entry.mtimeMs,
sizeBytes: entry.sizeBytes,
platform: entry.platform,
session: entry.session
}
])
}
// Seeded entries carry `resume: null`: after a restart an unchanged file is a
// cache hit; a file that changed while the app was closed pays one full
// (not incremental) re-parse.
export function seedSessionParseCache(
entries: Iterable<[string, PersistedSessionParseCacheEntry]>
): void {
const list = [...entries]
// Snapshot order is oldest→newest (LRU); an over-cap list keeps the newest
// tail rather than seeding the oldest entries and dropping the tail.
for (const [path, entry] of list.slice(Math.max(0, list.length - MAX_CACHE_ENTRIES))) {
if (cache.size >= MAX_CACHE_ENTRIES) {
return
}
// In-process entries are always fresher than persisted ones; never clobber.
if (cache.has(path)) {
continue
}
cache.set(path, {
mtimeMs: entry.mtimeMs,
sizeBytes: entry.sizeBytes,
platform: entry.platform,
session: entry.session,
resume: null
})
}
}
export function getSessionParseCacheEntry(path: string): SessionParseCacheEntry | undefined {
return cache.get(path)
}
export function storeSessionParseCacheEntry(path: string, entry: SessionParseCacheEntry): void {
cache.delete(path)
cache.set(path, entry)
if (cache.size > MAX_CACHE_ENTRIES) {
const oldest = cache.keys().next()
if (!oldest.done) {
cache.delete(oldest.value)
}
}
}
@@ -23,6 +23,12 @@ import {
normalizePreviewText,
timestampMs
} from './session-scanner-values'
import { NO_TRANSCRIPT_MESSAGES, type TranscriptMessageSink } from './session-transcript-consumers'
import {
boundedText,
transcriptMessageRole,
transcriptMessagesFromContent
} from './session-transcript-message-content'
const SESSION_PREVIEW_MESSAGE_LIMIT = 5
@@ -30,9 +36,12 @@ export function createAccumulator(args: {
agent: AiVaultAgent
file: FileWithMtime
sessionId: string
// Where every decoded message goes; absent for one-shot parses with no reader.
messages?: TranscriptMessageSink
}): SessionAccumulator {
return {
agent: args.agent,
messages: args.messages ?? NO_TRANSCRIPT_MESSAGES,
sessionId: args.sessionId,
title: null,
fallbackTitle: null,
@@ -64,19 +73,29 @@ export function cloneSessionAccumulator(accumulator: SessionAccumulator): Sessio
// closure state (claude, codex) build their own ResumableSessionParseState.
export function accumulatorFoldResumeState(
accumulator: SessionAccumulator,
consumeRecordLine: (accumulator: SessionAccumulator, line: string) => void
consumeRecordLine: (accumulator: SessionAccumulator, line: string) => void,
// Runs per finalize, for agents whose metadata lives in a sibling file the
// fold never sees; it may only fill fields the transcript left empty.
enrichBeforeFinalize?: (accumulator: SessionAccumulator) => Promise<void>
): ResumableSessionParseState {
return {
consumeLine: (line) => consumeRecordLine(accumulator, line),
clone: () =>
accumulatorFoldResumeState(cloneSessionAccumulator(accumulator), consumeRecordLine),
accumulatorFoldResumeState(
cloneSessionAccumulator(accumulator),
consumeRecordLine,
enrichBeforeFinalize
),
touchFile: (file) => {
accumulator.modifiedAt = file.modifiedAt
},
// Finalize a snapshot: the live accumulator (and its preview array) keeps
// accumulating appended lines after this session object is handed out.
finalize: (platform, options) =>
finalizeSession(cloneSessionAccumulator(accumulator), platform, options)
finalize: async (platform, options) => {
const snapshot = cloneSessionAccumulator(accumulator)
await enrichBeforeFinalize?.(snapshot)
return finalizeSession(snapshot, platform, options)
}
}
}
@@ -161,8 +180,13 @@ export function addPreviewMessage(
// Why: Claude meta/injected turns still preview, but must not seed the
// copyable first-prompt row.
seedFirstUserPrompt?: boolean
// Set false by callers that already published this record's messages.
publishMessage?: boolean
}
): void {
if (args.publishMessage !== false && accumulator.messages.active) {
publishTranscriptMessage(accumulator, args.role, args.text, args.timestamp)
}
// Seeded before the preview-empty return so the copy body never depends on
// preview-only normalization rules.
seedFullFirstUserPrompt(
@@ -199,15 +223,41 @@ export function addPreviewContent(
() => extractFullFirstUserPromptText(content),
options?.seedFirstUserPrompt
)
// Published from the content value, not the preview string: a consumer needs
// the whole turn, including the tool blocks the 220-char preview drops.
if (accumulator.messages.active) {
for (const message of transcriptMessagesFromContent(role, content, timestampIso(timestamp))) {
accumulator.messages.push(message)
}
}
addPreviewMessage(accumulator, {
role,
text: extractPreviewContentText(content),
timestamp,
// Content path already seeded above when capture is enabled.
seedFirstUserPrompt: false
seedFirstUserPrompt: false,
publishMessage: false
})
}
/** One already-flattened turn; the content path publishes per block instead. */
function publishTranscriptMessage(
accumulator: SessionAccumulator,
role: AiVaultSessionPreviewMessage['role'],
text: string | null,
timestamp: unknown
): void {
const messageRole = transcriptMessageRole(role)
const messageText = text === null ? null : boundedText(text)
if (messageRole && messageText) {
accumulator.messages.push({
role: messageRole,
text: messageText,
timestamp: timestampIso(timestamp)
})
}
}
/**
* Seed the copyable first prompt from the first real user turn. `fullText` is a
* thunk so list scans (capture mode `none`) never pay the extraction cost.
@@ -16,6 +16,7 @@ import { parseCursorSessionFile } from './session-scanner-cursor-parser'
import { parseHermesSessionFile } from './session-scanner-hermes-parser'
import { parseOpenCodeSessionFile } from './session-scanner-opencode-parser'
import type { SessionFileCandidate } from './session-scanner-types'
import type { TranscriptMessageSink } from './session-transcript-consumers'
/**
* Parse a single agent session file into an `AiVaultSession`. Routes to the
@@ -24,25 +25,33 @@ import type { SessionFileCandidate } from './session-scanner-types'
* `parseOpenCodeSqliteSession` instead of the legacy JSON parser.
* @param candidate - The session file candidate to parse.
* @param platform - The platform to use for resume command generation.
* @param messages - Where the parser publishes every decoded message.
* @returns The parsed `AiVaultSession`, or `null` if parsing fails.
*/
export async function parseAgentSessionFile(
candidate: SessionFileCandidate,
platform: NodeJS.Platform
platform: NodeJS.Platform,
messages?: TranscriptMessageSink
): Promise<AiVaultSession | null> {
switch (candidate.agent) {
case 'claude':
return parseClaudeSessionFile(candidate.file, platform)
return parseClaudeSessionFile(candidate.file, platform, messages)
case 'codex':
return parseCodexSessionFile(candidate.file, platform, candidate.codexHome)
return parseCodexSessionFile(
candidate.file,
platform,
candidate.codexHome,
undefined,
messages
)
case 'gemini':
return parseGeminiSessionFile(candidate.file, platform)
return parseGeminiSessionFile(candidate.file, platform, messages)
case 'antigravity':
return parseAntigravitySessionFile(candidate.file, platform)
return parseAntigravitySessionFile(candidate.file, platform, messages)
case 'copilot':
return parseCopilotSessionFile(candidate.file, platform)
return parseCopilotSessionFile(candidate.file, platform, messages)
case 'cursor':
return parseCursorSessionFile(candidate.file, platform)
return parseCursorSessionFile(candidate.file, platform, messages)
case 'opencode': {
// Why: OpenCode 1.17.x sessions are read from SQLite via a synthetic
// <dbPath>#<sessionId> candidate path. Legacy file-based sessions use
@@ -55,29 +64,29 @@ export async function parseAgentSessionFile(
platform
})
}
return parseOpenCodeSessionFile(candidate.file, platform)
return parseOpenCodeSessionFile(candidate.file, platform, messages)
}
case 'grok':
return parseGrokSessionFile(candidate.file, platform)
return parseGrokSessionFile(candidate.file, platform, messages)
case 'hermes':
return parseHermesSessionFile(candidate.file, platform)
return parseHermesSessionFile(candidate.file, platform, messages)
case 'rovo':
return parseRovoSessionFile(candidate.file, platform)
return parseRovoSessionFile(candidate.file, platform, messages)
case 'openclaw':
return parseMessageGraphSessionFile('openclaw', candidate.file, platform)
return parseMessageGraphSessionFile('openclaw', candidate.file, platform, messages)
case 'pi':
return parseMessageGraphSessionFile('pi', candidate.file, platform)
return parseMessageGraphSessionFile('pi', candidate.file, platform, messages)
case 'omp':
return parseMessageGraphSessionFile('omp', candidate.file, platform)
return parseMessageGraphSessionFile('omp', candidate.file, platform, messages)
case 'prime-agent':
return parseMessageGraphSessionFile('prime-agent', candidate.file, platform)
return parseMessageGraphSessionFile('prime-agent', candidate.file, platform, messages)
case 'droid':
return parseDroidSessionFile(candidate.file, platform)
return parseDroidSessionFile(candidate.file, platform, messages)
case 'cline':
return parseClineSessionFile(candidate.file, platform)
return parseClineSessionFile(candidate.file, platform, messages)
case 'devin':
return parseDevinSessionFile(candidate.file, platform)
return parseDevinSessionFile(candidate.file, platform, messages)
case 'kimi':
return parseKimiSessionFile(candidate.file, platform)
return parseKimiSessionFile(candidate.file, platform, messages)
}
}
@@ -8,6 +8,7 @@ import {
clineMessagesPathForMetadata,
isClineSessionMetadataPath
} from './session-scanner-cline-parser'
import { cursorChatMetaPath } from './session-scanner-cursor-chat-meta'
import { resolveKimiSessionsDir } from './session-scanner-kimi-paths'
import { OMP_SESSION_ARTIFACT_DIR_PATTERN } from './session-scanner-omp-subagent-transcripts'
import { claudeProjectsRootDirs, OMP_SESSIONS_DIR, sessionRootDirs } from './session-scanner-roots'
@@ -61,8 +62,9 @@ export type AiVaultAgentSource = {
rootDirs: (options: AiVaultScanOptions, wslHomeDirs: readonly string[]) => string[]
extensions: readonly string[]
filePredicate?: (filePath: string) => boolean
// A sibling whose stat participates in candidate freshness and recency.
contentDependencyPath?: (filePath: string) => string
// A sibling whose stat participates in candidate freshness and recency; async
// for agents that have to look the sibling up rather than derive its path.
contentDependencyPath?: (filePath: string) => string | undefined | Promise<string | undefined>
// Return false to skip a directory; depth 0 is a child of the root.
directoryPredicate?: (name: string, depth: number) => boolean
// Roots that are alternates for one install rather than distinct locations,
@@ -126,7 +128,8 @@ export const AI_VAULT_AGENT_SOURCES: AiVaultAgentSourceTable = {
'projects'
]),
extensions: ['.jsonl'],
filePredicate: (filePath) => pathSegments(filePath).includes('agent-transcripts')
filePredicate: (filePath) => pathSegments(filePath).includes('agent-transcripts'),
contentDependencyPath: cursorChatMetaPath
},
grok: {
rootDirs: (options, wslHomeDirs) =>
@@ -15,6 +15,7 @@ import type {
ResumableSessionParseState,
SessionAccumulator
} from './session-scanner-types'
import type { TranscriptMessageSink } from './session-transcript-consumers'
import { extractString, normalizeTitleText, parseJsonObject } from './session-scanner-values'
type ParserSessionOptions = {
@@ -24,12 +25,13 @@ type ParserSessionOptions = {
export async function parseAntigravitySessionFile(
file: FileWithMtime,
platform: NodeJS.Platform = process.platform
platform: NodeJS.Platform = process.platform,
messages?: TranscriptMessageSink
): Promise<AiVaultSession | null> {
const input = openTranscriptReadStream(file.path, { encoding: 'utf-8' }, 'scan')
const lines = createInterface({ input, crlfDelay: Infinity })
try {
return await parseAntigravitySessionLines({ file, lines, platform })
return await parseAntigravitySessionLines({ file, lines, platform, messages })
} finally {
// readline.close() leaves the underlying stream open; destroy it so a
// mid-parse throw cannot leak the gated transcript handle.
@@ -54,13 +56,14 @@ export async function parseAntigravitySessionContent(
}
export function createAntigravitySessionResumeState(
file: FileWithMtime
file: FileWithMtime,
messages?: TranscriptMessageSink
): ResumableSessionParseState {
const sessionId = antigravityConversationIdFromTranscriptPath(file.path) ?? ''
// Why: the transcript has no cwd/model fields. Workspace enrichment is a
// separate, conservative history join; protobuf/SQLite blobs are unstable.
return accumulatorFoldResumeState(
createAccumulator({ agent: 'antigravity', file, sessionId }),
createAccumulator({ agent: 'antigravity', file, sessionId, messages }),
consumeAntigravityRecordLine
)
}
@@ -70,8 +73,9 @@ async function parseAntigravitySessionLines(args: {
lines: AsyncIterable<string> | Iterable<string>
platform: NodeJS.Platform
options?: ParserSessionOptions
messages?: TranscriptMessageSink
}): Promise<AiVaultSession | null> {
const state = createAntigravitySessionResumeState(args.file)
const state = createAntigravitySessionResumeState(args.file, args.messages)
for await (const line of args.lines) {
state.consumeLine(line)
}
@@ -0,0 +1,46 @@
import { readCodexRolloutSessionMetaId } from '../codex/codex-rollout-session-meta'
import { codexRolloutHardlinkIdentity, dedupeCodexRolloutAliases } from './codex-session-root-dedup'
import { antigravityHistoryPathForBrainDir } from './session-scanner-antigravity-paths'
import { codexHomeForSessionsDir } from './session-scanner-codex-paths'
import { DEFAULT_CODEX_HOME_DIR } from './session-scanner-source-discovery'
import type {
AiVaultScanOptions,
SessionFileCandidate,
SessionFileDiscovery
} from './session-scanner-types'
/** Newest-first parse candidates for a discovery set, with Codex hardlink aliases collapsed. */
export async function sessionCandidatesFromDiscoveries(
discoveries: SessionFileDiscovery[],
options: AiVaultScanOptions
): Promise<SessionFileCandidate[]> {
return dedupeCodexRolloutAliases(
discoveries
.flatMap((discovery) =>
discovery.files.map((file): SessionFileCandidate => ({
agent: discovery.agent,
file,
codexHome:
discovery.agent === 'codex'
? codexHomeForSessionsDir(
discovery.rootDir,
options.defaultCodexHomeDir ?? DEFAULT_CODEX_HOME_DIR
)
: null,
antigravityHistoryPath:
discovery.agent === 'antigravity'
? antigravityHistoryPathForBrainDir(discovery.rootDir)
: undefined
}))
)
.sort((left, right) => right.file.mtimeMs - left.file.mtimeMs),
{
isCodex: (candidate) => candidate.agent === 'codex',
getFilePath: (candidate) => candidate.file.path,
getCodexHome: (candidate) => candidate.codexHome,
getHardlinkIdentity: (candidate) => codexRolloutHardlinkIdentity(candidate.file)
},
(filePath) => readCodexRolloutSessionMetaId(filePath, options.signal, 'scan'),
options.signal
)
}
@@ -9,6 +9,7 @@ import {
updateTimeline
} from './session-scanner-accumulator'
import type { FileWithMtime } from './session-scanner-types'
import type { TranscriptMessageSink } from './session-transcript-consumers'
import {
arrayValue,
asRecord,
@@ -35,7 +36,8 @@ export function clineMessagesPathForMetadata(filePath: string): string {
export async function parseClineSessionFile(
file: FileWithMtime,
platform: NodeJS.Platform = process.platform
platform: NodeJS.Platform = process.platform,
messageSink?: TranscriptMessageSink
): Promise<AiVaultSession | null> {
const metadataContent = await wslGatedReadFile(file.path, 'utf-8', 'scan')
let messagesContent: string | null = null
@@ -52,7 +54,7 @@ export async function parseClineSessionFile(
throw error
}
}
return parseClineSessionContent(file, metadataContent, messagesContent, platform)
return parseClineSessionContent(file, metadataContent, messagesContent, platform, {}, messageSink)
}
function isMissingSessionPathError(error: unknown): boolean {
@@ -68,7 +70,8 @@ export function parseClineSessionContent(
metadataContent: string,
messagesContent: string | null,
platform: NodeJS.Platform = process.platform,
options: ParserSessionOptions = {}
options: ParserSessionOptions = {},
messageSink?: TranscriptMessageSink
): AiVaultSession | null {
const metadata = parseJsonRecord(metadataContent)
if (!metadata) {
@@ -76,7 +79,12 @@ export function parseClineSessionContent(
}
const pathSegments = file.path.replace(/\\/g, '/').split('/').filter(Boolean)
const sessionId = extractString(metadata.session_id) ?? pathSegments.at(-2) ?? ''
const accumulator = createAccumulator({ agent: 'cline', file, sessionId })
const accumulator = createAccumulator({
agent: 'cline',
file,
sessionId,
messages: messageSink
})
accumulator.cwd = extractString(metadata.cwd) ?? extractString(metadata.workspace_root)
accumulator.model = extractString(metadata.model)
updateTimeline(accumulator, metadata.started_at)
@@ -22,6 +22,7 @@ import type {
ResumableSessionParseState,
SessionAccumulator
} from './session-scanner-types'
import type { TranscriptMessageSink } from './session-transcript-consumers'
import {
addCodexUsage,
asRecord,
@@ -29,18 +30,22 @@ import {
extractModel,
extractString,
normalizeCodexUsage,
normalizeTitleText,
parseJsonObject,
subtractCodexUsage
} from './session-scanner-values'
import { remoteSessionContentLines } from './remote-session-content-lines'
import { readCodexTimelineOnlyRecord } from './session-scanner-codex-record-fast-path'
import {
extractCodexSessionMetadataTitle,
isCodexWorkerSession
} from './session-scanner-codex-session-meta'
export async function parseCodexSessionFile(
file: FileWithMtime,
platform: NodeJS.Platform = process.platform,
codexHome: string | null = null,
executionHostId?: ExecutionHostId
executionHostId?: ExecutionHostId,
messages?: TranscriptMessageSink
): Promise<AiVaultSession | null> {
const lines = createInterface({
input: openTranscriptReadStream(file.path, { encoding: 'utf-8' }, 'scan'),
@@ -53,6 +58,7 @@ export async function parseCodexSessionFile(
platform,
codexHome,
executionHostId,
messages,
titleReader: (sessionId) => readCodexSessionIndexTitle(file.path, codexHome, sessionId)
})
}
@@ -89,12 +95,16 @@ type CodexSessionParseState = {
titleSource: 'meta' | 'user' | null
}
function createCodexParseState(file: FileWithMtime): CodexSessionParseState {
function createCodexParseState(
file: FileWithMtime,
messages?: TranscriptMessageSink
): CodexSessionParseState {
return {
accumulator: createAccumulator({
agent: 'codex',
file,
sessionId: sessionIdFromFileName(file.path)
sessionId: sessionIdFromFileName(file.path),
messages
}),
previousTotals: null,
rejectedWorkerSession: false,
@@ -257,10 +267,13 @@ async function finalizeCodexParseState(
export function createCodexSessionResumeState(
file: FileWithMtime,
codexHome: string | null
codexHome: string | null,
messages?: TranscriptMessageSink
): ResumableSessionParseState {
return codexResumeStateFromParseState(createCodexParseState(file), codexHome, (sessionId) =>
readCodexSessionIndexTitle(file.path, codexHome, sessionId)
return codexResumeStateFromParseState(
createCodexParseState(file, messages),
codexHome,
(sessionId) => readCodexSessionIndexTitle(file.path, codexHome, sessionId)
)
}
@@ -298,8 +311,9 @@ async function parseCodexSessionLines(args: {
executionHostId?: ExecutionHostId
executionHostPlatform?: NodeJS.Platform | null
titleReader?: (sessionId: string) => Promise<string | null>
messages?: TranscriptMessageSink
}): Promise<AiVaultSession | null> {
const state = createCodexParseState(args.file)
const state = createCodexParseState(args.file, args.messages)
for await (const line of args.lines) {
consumeCodexRecordLine(state, line)
if (state.rejectedWorkerSession) {
@@ -314,21 +328,3 @@ async function parseCodexSessionLines(args: {
executionHostPlatform: args.executionHostPlatform
})
}
function isCodexWorkerSession(payload: Record<string, unknown>): boolean {
const threadSource = extractString(payload.thread_source) ?? extractString(payload.threadSource)
if (threadSource) {
return threadSource.toLowerCase() !== 'user'
}
const source = asRecord(payload.source)
return Boolean(asRecord(source?.subagent))
}
function extractCodexSessionMetadataTitle(payload: Record<string, unknown>): string | null {
return (
normalizeTitleText(extractString(payload.title) ?? '') ??
normalizeTitleText(extractString(payload.thread_name) ?? '') ??
normalizeTitleText(extractString(payload.threadName) ?? '')
)
}
@@ -0,0 +1,22 @@
import { asRecord, extractString, normalizeTitleText } from './session-scanner-values'
// Field readers for Codex's `session_meta` record, whose key spelling has drifted
// across Codex releases (snake_case rollouts, camelCase app-server rollouts).
export function isCodexWorkerSession(payload: Record<string, unknown>): boolean {
const threadSource = extractString(payload.thread_source) ?? extractString(payload.threadSource)
if (threadSource) {
return threadSource.toLowerCase() !== 'user'
}
const source = asRecord(payload.source)
return Boolean(asRecord(source?.subagent))
}
export function extractCodexSessionMetadataTitle(payload: Record<string, unknown>): string | null {
return (
normalizeTitleText(extractString(payload.title) ?? '') ??
normalizeTitleText(extractString(payload.thread_name) ?? '') ??
normalizeTitleText(extractString(payload.threadName) ?? '')
)
}
@@ -8,6 +8,7 @@ import type {
ResumableSessionParseState,
SessionAccumulator
} from './session-scanner-types'
import type { TranscriptMessageSink } from './session-transcript-consumers'
import {
accumulatorFoldResumeState,
addPreviewMessage,
@@ -32,13 +33,14 @@ type ParserSessionOptions = {
export async function parseCopilotSessionFile(
file: FileWithMtime,
platform: NodeJS.Platform = process.platform
platform: NodeJS.Platform = process.platform,
messages?: TranscriptMessageSink
): Promise<AiVaultSession | null> {
const lines = createInterface({
input: openTranscriptReadStream(file.path, { encoding: 'utf-8' }, 'scan'),
crlfDelay: Infinity
})
return parseCopilotSessionLines({ file, lines, platform })
return parseCopilotSessionLines({ file, lines, platform, messages })
}
export async function parseCopilotSessionContent(
@@ -107,9 +109,17 @@ function consumeCopilotRecordLine(accumulator: SessionAccumulator, line: string)
}
}
export function createCopilotSessionResumeState(file: FileWithMtime): ResumableSessionParseState {
export function createCopilotSessionResumeState(
file: FileWithMtime,
messages?: TranscriptMessageSink
): ResumableSessionParseState {
return accumulatorFoldResumeState(
createAccumulator({ agent: 'copilot', file, sessionId: sessionIdFromFileName(file.path) }),
createAccumulator({
agent: 'copilot',
file,
sessionId: sessionIdFromFileName(file.path),
messages
}),
consumeCopilotRecordLine
)
}
@@ -119,8 +129,9 @@ async function parseCopilotSessionLines(args: {
lines: AsyncIterable<string> | Iterable<string>
platform: NodeJS.Platform
options?: ParserSessionOptions
messages?: TranscriptMessageSink
}): Promise<AiVaultSession | null> {
const state = createCopilotSessionResumeState(args.file)
const state = createCopilotSessionResumeState(args.file, args.messages)
for await (const line of args.lines) {
state.consumeLine(line)
}
@@ -0,0 +1,300 @@
import { mkdir, mkdtemp, rm, utimes, writeFile } from 'node:fs/promises'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { afterEach, describe, expect, it, vi } from 'vitest'
import { WslTranscriptFsError } from '../native-chat/wsl-transcript-fs-error'
// Why: a refused WSL read is the one build failure that must not be cached.
let failNextChatsReaddir = false
let chatsRootReads = 0
vi.mock('../native-chat/wsl-transcript-fs-access', async (importOriginal) => {
const actual = await importOriginal<typeof WslTranscriptFsAccess>()
return {
...actual,
wslGatedReaddir: (
...args: Parameters<typeof actual.wslGatedReaddir>
): ReturnType<typeof actual.wslGatedReaddir> => {
if (args[0].endsWith('chats')) {
chatsRootReads += 1
}
if (failNextChatsReaddir && args[0].includes('workspace-hash')) {
failNextChatsReaddir = false
return Promise.reject(new WslTranscriptFsError('timeout', 'wsl fs timed out'))
}
return actual.wslGatedReaddir(...args)
}
}
})
import type * as WslTranscriptFsAccess from '../native-chat/wsl-transcript-fs-access'
import {
cursorChatMetaPath,
readCursorChatMeta,
resetCursorChatMetaIndexCacheForTests,
withCursorChatMetaScan
} from './session-scanner-cursor-chat-meta'
import {
createCursorSessionResumeState,
parseCursorSessionContent,
parseCursorSessionFile
} from './session-scanner-cursor-parser'
import type { AiVaultScanIssue } from '../../shared/ai-vault-types'
import { AI_VAULT_AGENT_SOURCES } from './session-scanner-agent-sources'
import { discoverFiles } from './session-scanner-discovery'
import type { FileWithMtime, SessionFileDiscovery } from './session-scanner-types'
// Cursor's real meta.json keys (~/.cursor/chats/<md5 of cwd>/<uuid>/meta.json, 2026-09).
type CursorMetaFixture = {
schemaVersion: number
createdAtMs: number
updatedAtMs: number
cwd: string
hasConversation: boolean
title?: string
}
const CREATED_AT_MS = 1_787_039_612_017
const UPDATED_AT_MS = 1_787_039_640_532
let tempRoots: string[] = []
afterEach(async () => {
resetCursorChatMetaIndexCacheForTests()
await Promise.all(tempRoots.map((root) => rm(root, { recursive: true, force: true })))
tempRoots = []
})
async function createCursorHome(): Promise<string> {
const root = await mkdtemp(join(tmpdir(), 'orca-cursor-chat-meta-'))
tempRoots.push(root)
const cursorHome = join(root, '.cursor')
await mkdir(cursorHome, { recursive: true })
return cursorHome
}
async function writeTranscript(
cursorHome: string,
projectSlug: string,
chatId: string,
lines: string[]
): Promise<string> {
const chatDir = join(cursorHome, 'projects', projectSlug, 'agent-transcripts', chatId)
await mkdir(chatDir, { recursive: true })
const transcriptPath = join(chatDir, `${chatId}.jsonl`)
await writeFile(transcriptPath, lines.map((line) => `${line}\n`).join(''))
return transcriptPath
}
async function writeChatMeta(
cursorHome: string,
workspaceHash: string,
chatId: string,
meta: Partial<CursorMetaFixture> = {}
): Promise<string> {
const chatDir = join(cursorHome, 'chats', workspaceHash, chatId)
await mkdir(chatDir, { recursive: true })
const metaPath = join(chatDir, 'meta.json')
await writeFile(
metaPath,
JSON.stringify({
schemaVersion: 1,
createdAtMs: CREATED_AT_MS,
updatedAtMs: UPDATED_AT_MS,
cwd: '/private/tmp/workspace',
hasConversation: true,
...meta
} satisfies CursorMetaFixture)
)
return metaPath
}
function fileWithMtime(path: string): FileWithMtime {
return { path, mtimeMs: 1, modifiedAt: new Date(1).toISOString() }
}
describe('cursor chat meta', () => {
it('resolves the meta.json under the workspace hash that holds the chat id', async () => {
const cursorHome = await createCursorHome()
await writeChatMeta(cursorHome, 'aa37220647fb7ce5eb044aa4bda60807', 'other-chat')
const metaPath = await writeChatMeta(cursorHome, '96fa26ac0f433670ebec73ecef20b47b', 'chat-1', {
title: 'Shell Command Hostname'
})
const transcriptPath = await writeTranscript(cursorHome, 'private-tmp-workspace', 'chat-1', [])
expect(await cursorChatMetaPath(transcriptPath)).toBe(metaPath)
expect(await readCursorChatMeta(transcriptPath)).toEqual({
title: 'Shell Command Hostname',
cwd: '/private/tmp/workspace',
createdAt: new Date(CREATED_AT_MS).toISOString(),
updatedAt: new Date(UPDATED_AT_MS).toISOString()
})
})
it('re-indexes after a chat appears under an already indexed workspace', async () => {
const cursorHome = await createCursorHome()
await writeChatMeta(cursorHome, 'workspace-hash', 'chat-first')
const firstTranscript = await writeTranscript(cursorHome, 'slug', 'chat-first', [])
expect(await cursorChatMetaPath(firstTranscript)).toBeDefined()
const laterMetaPath = await writeChatMeta(cursorHome, 'workspace-hash', 'chat-later')
const laterTranscript = await writeTranscript(cursorHome, 'slug', 'chat-later', [])
expect(await cursorChatMetaPath(laterTranscript)).toBe(laterMetaPath)
})
it('does not cache a metadata index whose build was refused by the WSL gate', async () => {
const cursorHome = await createCursorHome()
const metaPath = await writeChatMeta(cursorHome, 'workspace-hash', 'chat-refused')
const transcriptPath = await writeTranscript(cursorHome, 'slug', 'chat-refused', [])
failNextChatsReaddir = true
await expect(cursorChatMetaPath(transcriptPath)).rejects.toBeInstanceOf(WslTranscriptFsError)
// The next scan rebuilds instead of replaying the rejected promise.
await expect(cursorChatMetaPath(transcriptPath)).resolves.toBe(metaPath)
})
it('validates the index once per scan, not once per transcript', async () => {
const cursorHome = await createCursorHome()
const transcripts: string[] = []
for (const chatId of ['chat-a', 'chat-b', 'chat-c']) {
await writeChatMeta(cursorHome, 'workspace-hash', chatId)
transcripts.push(await writeTranscript(cursorHome, 'slug', chatId, []))
}
chatsRootReads = 0
const inScan = await withCursorChatMetaScan(() =>
Promise.all(transcripts.map((path) => cursorChatMetaPath(path)))
)
expect(inScan.every(Boolean)).toBe(true)
expect(chatsRootReads).toBe(1)
// Outside a scan every lookup re-validates, which is what the parse path needs.
chatsRootReads = 0
await Promise.all(transcripts.map((path) => cursorChatMetaPath(path)))
expect(chatsRootReads).toBe(3)
})
it('yields nothing and does not throw when there is no chats tree', async () => {
const cursorHome = await createCursorHome()
const transcriptPath = await writeTranscript(cursorHome, 'slug', 'chat-orphan', [])
await expect(cursorChatMetaPath(transcriptPath)).resolves.toBeUndefined()
await expect(readCursorChatMeta(transcriptPath)).resolves.toBeNull()
await expect(readCursorChatMeta('/nowhere/near/cursor/chat.jsonl')).resolves.toBeNull()
})
it('yields nothing and does not throw when meta.json is malformed', async () => {
const cursorHome = await createCursorHome()
const chatDir = join(cursorHome, 'chats', 'workspace-hash', 'chat-bad')
await mkdir(chatDir, { recursive: true })
await writeFile(join(chatDir, 'meta.json'), '{ not json')
const transcriptPath = await writeTranscript(cursorHome, 'slug', 'chat-bad', [])
await expect(readCursorChatMeta(transcriptPath)).resolves.toBeNull()
})
})
describe('cursor discovery meta dependency', () => {
it('folds meta.json into candidate freshness so a rewrite invalidates the parse cache', async () => {
const cursorHome = await createCursorHome()
const metaPath = await writeChatMeta(cursorHome, 'workspace-hash', 'chat-7')
await writeTranscript(cursorHome, 'slug', 'chat-7', [])
const issues: AiVaultScanIssue[] = []
const discover = (): Promise<SessionFileDiscovery> =>
discoverFiles({
rootDir: join(cursorHome, 'projects'),
limit: 10,
agent: 'cursor',
issues,
extensions: [...AI_VAULT_AGENT_SOURCES.cursor.extensions],
filePredicate: AI_VAULT_AGENT_SOURCES.cursor.filePredicate,
contentDependencyPath: AI_VAULT_AGENT_SOURCES.cursor.contentDependencyPath
})
const before = (await discover()).files[0]
const future = new Date(Date.now() + 10_000)
await utimes(metaPath, future, future)
const after = (await discover()).files[0]
expect(after?.mtimeMs).toBeGreaterThan(before?.mtimeMs ?? 0)
expect(issues).toEqual([])
})
})
describe('cursor parser chat meta fallback', () => {
it('fills cwd, timestamps and title from meta.json', async () => {
const cursorHome = await createCursorHome()
await writeChatMeta(cursorHome, 'workspace-hash', 'chat-2', { title: 'Named From Meta' })
const transcriptPath = await writeTranscript(cursorHome, 'slug', 'chat-2', [
JSON.stringify({ role: 'assistant', message: { content: 'hello' } })
])
const session = await parseCursorSessionFile(fileWithMtime(transcriptPath), 'darwin')
expect(session?.cwd).toBe('/private/tmp/workspace')
expect(session?.title).toBe('Named From Meta')
expect(session?.createdAt).toBe(new Date(CREATED_AT_MS).toISOString())
expect(session?.updatedAt).toBe(new Date(UPDATED_AT_MS).toISOString())
})
it('keeps a transcript title and timestamps over meta.json', async () => {
const cursorHome = await createCursorHome()
await writeChatMeta(cursorHome, 'workspace-hash', 'chat-3', { title: 'Meta Title' })
const transcriptPath = await writeTranscript(cursorHome, 'slug', 'chat-3', [
JSON.stringify({
role: 'user',
message: { content: 'transcript first prompt' },
timestamp: '2026-01-01T00:00:00.000Z'
})
])
const session = await parseCursorSessionFile(fileWithMtime(transcriptPath), 'darwin')
expect(session?.title).toBe('transcript first prompt')
expect(session?.createdAt).toBe('2026-01-01T00:00:00.000Z')
// cwd is never in the transcript, so it still comes from meta.json.
expect(session?.cwd).toBe('/private/tmp/workspace')
})
it('builds the resume command from the meta.json cwd', async () => {
const cursorHome = await createCursorHome()
await writeChatMeta(cursorHome, 'workspace-hash', 'chat-4', { cwd: '/repo/from-meta' })
const transcriptPath = await writeTranscript(cursorHome, 'slug', 'chat-4', [
JSON.stringify({ role: 'user', message: { content: 'hi' } })
])
const session = await parseCursorSessionFile(fileWithMtime(transcriptPath), 'darwin')
expect(session?.resumeCommand).toContain('/repo/from-meta')
})
it('leaves remote content parses to the transcript alone', async () => {
const cursorHome = await createCursorHome()
await writeChatMeta(cursorHome, 'workspace-hash', 'chat-5', { title: 'Meta Title' })
const transcriptPath = await writeTranscript(cursorHome, 'slug', 'chat-5', [])
const session = await parseCursorSessionContent(
fileWithMtime(transcriptPath),
`${JSON.stringify({ role: 'assistant', message: { content: 'remote' } })}\n`,
'linux'
)
expect(session?.cwd).toBeNull()
expect(session?.title).not.toBe('Meta Title')
})
it('applies the fallback on every finalize of a resumed parse', async () => {
const cursorHome = await createCursorHome()
await writeChatMeta(cursorHome, 'workspace-hash', 'chat-6', { title: 'Resumed Meta' })
const transcriptPath = await writeTranscript(cursorHome, 'slug', 'chat-6', [])
const state = createCursorSessionResumeState(fileWithMtime(transcriptPath))
state.consumeLine(JSON.stringify({ role: 'assistant', message: { content: 'first' } }))
const first = await state.finalize('darwin')
state.consumeLine(JSON.stringify({ role: 'assistant', message: { content: 'second' } }))
const second = await state.finalize('darwin')
expect(first?.cwd).toBe('/private/tmp/workspace')
expect(second?.title).toBe('Resumed Meta')
expect(second?.messageCount).toBe(2)
})
})
@@ -0,0 +1,221 @@
import { AsyncLocalStorage } from 'node:async_hooks'
import { basename, dirname, join } from 'node:path'
import { wslGatedReaddir, wslGatedStat } from '../native-chat/wsl-transcript-fs-access'
import { WslTranscriptFsError } from '../native-chat/wsl-transcript-fs-gate'
import { timestampIso } from './session-scanner-accumulator'
import { extractString, normalizeTitleText, readJsonObjectIfExists } from './session-scanner-values'
// Cursor keeps a chat's transcript and its metadata in two unrelated trees:
// <cursor>/projects/<slug>/agent-transcripts/<uuid>/<uuid>.jsonl holds the
// messages, while <cursor>/chats/<md5 of cwd>/<uuid>/meta.json holds the cwd,
// title and timestamps. The md5 hashes the very cwd we are looking for, so the
// only way across is an index of the chat directories.
const CURSOR_CHATS_DIR = 'chats'
const CURSOR_CHAT_META_FILE = 'meta.json'
const CURSOR_TRANSCRIPTS_DIR = 'agent-transcripts'
const CURSOR_PROJECTS_DIR = 'projects'
// Why: custom and WSL Cursor homes can vary over a long-lived main process.
const CURSOR_CHAT_META_INDEX_CACHE_MAX = 8
export type CursorChatMeta = {
title: string | null
cwd: string | null
createdAt: string | null
updatedAt: string | null
}
type CursorChatMetaIndexEntry = {
signature: string
metaPathByChatId: Map<string, string>
}
const cursorChatMetaIndexCache = new Map<string, Promise<CursorChatMetaIndexEntry>>()
// Why: validating the cache costs a readdir plus a stat per workspace, and
// discovery asks once per transcript; one scan sees the tree once instead.
const scanScopedIndex = new AsyncLocalStorage<Map<string, Promise<Map<string, string>>>>()
export function resetCursorChatMetaIndexCacheForTests(): void {
cursorChatMetaIndexCache.clear()
}
/** Runs one discovery scan; every Cursor transcript inside it shares a single validated index read. */
export function withCursorChatMetaScan<T>(fn: () => Promise<T>): Promise<T> {
return scanScopedIndex.run(new Map(), fn)
}
/** Path a discovery stat can watch so a rewritten meta.json invalidates the parse cache. */
export async function cursorChatMetaPath(transcriptPath: string): Promise<string | undefined> {
const chatsRoot = cursorChatsRootFromTranscriptPath(transcriptPath)
const chatId = cursorChatIdFromTranscriptPath(transcriptPath)
if (!chatsRoot || !chatId) {
return undefined
}
const index = await readCursorChatMetaIndexOncePerScan(chatsRoot)
return index.get(chatId)
}
function readCursorChatMetaIndexOncePerScan(chatsRoot: string): Promise<Map<string, string>> {
const scan = scanScopedIndex.getStore()
if (!scan) {
return readCursorChatMetaIndex(chatsRoot)
}
let pending = scan.get(chatsRoot)
if (!pending) {
pending = readCursorChatMetaIndex(chatsRoot)
scan.set(chatsRoot, pending)
}
return pending
}
export async function readCursorChatMeta(transcriptPath: string): Promise<CursorChatMeta | null> {
const metaPath = await cursorChatMetaPath(transcriptPath)
if (!metaPath) {
return null
}
const record = await readJsonObjectIfExists(metaPath)
if (!record) {
return null
}
return {
title: normalizeTitleText(extractString(record.title) ?? ''),
cwd: extractString(record.cwd),
createdAt: timestampIso(record.createdAtMs),
updatedAt: timestampIso(record.updatedAtMs)
}
}
function cursorChatIdFromTranscriptPath(transcriptPath: string): string | null {
const chatDir = dirname(transcriptPath)
return basename(dirname(chatDir)) === CURSOR_TRANSCRIPTS_DIR ? basename(chatDir) : null
}
function cursorChatsRootFromTranscriptPath(transcriptPath: string): string | null {
let currentDir = dirname(transcriptPath)
while (currentDir && dirname(currentDir) !== currentDir) {
// The chats tree is a sibling of the projects tree, custom Cursor homes included.
if (basename(currentDir) === CURSOR_PROJECTS_DIR) {
return join(dirname(currentDir), CURSOR_CHATS_DIR)
}
currentDir = dirname(currentDir)
}
return null
}
async function readCursorChatMetaIndex(chatsRoot: string): Promise<Map<string, string>> {
let workspaceDirs: string[]
try {
workspaceDirs = (await wslGatedReaddir(chatsRoot, 'scan'))
.filter((entry) => entry.isDirectory())
.map((entry) => entry.name)
.sort()
} catch (error) {
// Why: a refused WSL read is not "no chats"; letting it through keeps the
// session out of the parse cache instead of caching it without metadata.
if (error instanceof WslTranscriptFsError) {
throw error
}
return new Map()
}
const signature = await readCursorChatsSignature(chatsRoot, workspaceDirs)
const cached = await readCachedCursorChatMetaIndex(chatsRoot, signature)
if (cached) {
return cached
}
const pending = buildCursorChatMetaIndex(chatsRoot, workspaceDirs).then((metaPathByChatId) => ({
signature,
metaPathByChatId
}))
storeCursorChatMetaIndexEntry(chatsRoot, pending)
// Why: a rejected build (a refused WSL read) must not be served from the
// cache forever; the next scan rebuilds while this one still sees the error.
pending.catch(() => {
if (cursorChatMetaIndexCache.get(chatsRoot) === pending) {
cursorChatMetaIndexCache.delete(chatsRoot)
}
})
return (await pending).metaPathByChatId
}
// Why: a new chat only bumps its own workspace directory, so the chats root's
// own mtime would keep serving an index that is missing the newest sessions.
async function readCursorChatsSignature(
chatsRoot: string,
workspaceDirs: string[]
): Promise<string> {
const parts = await Promise.all(
workspaceDirs.map(async (name) => {
try {
const dirStat = await wslGatedStat(join(chatsRoot, name), 'scan')
return `${name}:${dirStat.mtimeMs}`
} catch {
return `${name}:?`
}
})
)
return parts.join('|')
}
async function buildCursorChatMetaIndex(
chatsRoot: string,
workspaceDirs: string[]
): Promise<Map<string, string>> {
const metaPathByChatId = new Map<string, string>()
for (const workspaceDir of workspaceDirs) {
let chatDirs
try {
chatDirs = await wslGatedReaddir(join(chatsRoot, workspaceDir), 'scan')
} catch (error) {
if (error instanceof WslTranscriptFsError) {
throw error
}
continue
}
for (const chatDir of chatDirs) {
// Why: the same chat id never appears under two workspace hashes, so the
// first hit wins and a duplicate would only cost a wasted read.
if (chatDir.isDirectory() && !metaPathByChatId.has(chatDir.name)) {
metaPathByChatId.set(
chatDir.name,
join(chatsRoot, workspaceDir, chatDir.name, CURSOR_CHAT_META_FILE)
)
}
}
}
return metaPathByChatId
}
async function readCachedCursorChatMetaIndex(
chatsRoot: string,
signature: string
): Promise<Map<string, string> | undefined> {
const cached = cursorChatMetaIndexCache.get(chatsRoot)
if (!cached) {
return undefined
}
const entry = await cached
if (entry.signature !== signature) {
return undefined
}
// Why: a concurrent scan can replace this Promise while it resolves; only the
// still-current entry may refresh recency without bypassing the cap.
if (cursorChatMetaIndexCache.get(chatsRoot) === cached) {
cursorChatMetaIndexCache.delete(chatsRoot)
cursorChatMetaIndexCache.set(chatsRoot, cached)
}
return entry.metaPathByChatId
}
function storeCursorChatMetaIndexEntry(
chatsRoot: string,
pending: Promise<CursorChatMetaIndexEntry>
): void {
cursorChatMetaIndexCache.delete(chatsRoot)
cursorChatMetaIndexCache.set(chatsRoot, pending)
if (cursorChatMetaIndexCache.size > CURSOR_CHAT_META_INDEX_CACHE_MAX) {
const oldest = cursorChatMetaIndexCache.keys().next()
if (!oldest.done) {
cursorChatMetaIndexCache.delete(oldest.value)
}
}
}
@@ -8,6 +8,7 @@ import type {
ResumableSessionParseState,
SessionAccumulator
} from './session-scanner-types'
import type { TranscriptMessageSink } from './session-transcript-consumers'
import {
accumulatorFoldResumeState,
addPreviewContent,
@@ -22,6 +23,7 @@ import {
extractString,
parseJsonObject
} from './session-scanner-values'
import { readCursorChatMeta } from './session-scanner-cursor-chat-meta'
type ParserSessionOptions = {
executionHostId?: ExecutionHostId
@@ -30,13 +32,14 @@ type ParserSessionOptions = {
export async function parseCursorSessionFile(
file: FileWithMtime,
platform: NodeJS.Platform = process.platform
platform: NodeJS.Platform = process.platform,
messages?: TranscriptMessageSink
): Promise<AiVaultSession | null> {
const lines = createInterface({
input: openTranscriptReadStream(file.path, { encoding: 'utf-8' }, 'scan'),
crlfDelay: Infinity
})
return parseCursorSessionLines({ file, lines, platform })
return parseCursorSessionLines({ file, lines, platform, messages })
}
export async function parseCursorSessionContent(
@@ -50,7 +53,8 @@ export async function parseCursorSessionContent(
file,
lines: remoteSessionContentLines(content, signal),
platform,
options
options,
enrichFromChatMeta: false
})
}
@@ -75,20 +79,55 @@ function consumeCursorRecordLine(accumulator: SessionAccumulator, line: string):
}
}
export function createCursorSessionResumeState(file: FileWithMtime): ResumableSessionParseState {
export function createCursorSessionResumeState(
file: FileWithMtime,
// Remote hosts stream transcript content only, with no sibling meta.json to read.
enrichFromChatMeta = true,
messages?: TranscriptMessageSink
): ResumableSessionParseState {
return accumulatorFoldResumeState(
createAccumulator({ agent: 'cursor', file, sessionId: sessionIdFromFileName(file.path) }),
consumeCursorRecordLine
createAccumulator({
agent: 'cursor',
file,
sessionId: sessionIdFromFileName(file.path),
messages
}),
consumeCursorRecordLine,
enrichFromChatMeta ? (accumulator) => applyCursorChatMeta(accumulator, file.path) : undefined
)
}
/** Fills only what the transcript never recorded; its own records always win. */
async function applyCursorChatMeta(
accumulator: SessionAccumulator,
transcriptPath: string
): Promise<void> {
if (accumulator.cwd && accumulator.createdAt && accumulator.updatedAt && accumulator.title) {
return
}
const meta = await readCursorChatMeta(transcriptPath)
if (!meta) {
return
}
accumulator.title ??= meta.title
accumulator.cwd ??= meta.cwd
accumulator.createdAt ??= meta.createdAt
accumulator.updatedAt ??= meta.updatedAt
}
async function parseCursorSessionLines(args: {
file: FileWithMtime
lines: AsyncIterable<string> | Iterable<string>
platform: NodeJS.Platform
options?: ParserSessionOptions
enrichFromChatMeta?: boolean
messages?: TranscriptMessageSink
}): Promise<AiVaultSession | null> {
const state = createCursorSessionResumeState(args.file)
const state = createCursorSessionResumeState(
args.file,
args.enrichFromChatMeta ?? true,
args.messages
)
for await (const line of args.lines) {
state.consumeLine(line)
}
@@ -2,6 +2,7 @@ import { wslGatedReadFile } from '../native-chat/wsl-transcript-fs-access'
import type { AiVaultSession } from '../../shared/ai-vault-types'
import type { ExecutionHostId } from '../../shared/execution-host'
import type { FileWithMtime } from './session-scanner-types'
import type { TranscriptMessageSink } from './session-transcript-consumers'
import {
addPreviewContent,
createAccumulator,
@@ -25,20 +26,28 @@ type ParserSessionOptions = {
export async function parseDevinSessionFile(
file: FileWithMtime,
platform: NodeJS.Platform = process.platform
platform: NodeJS.Platform = process.platform,
messages?: TranscriptMessageSink
): Promise<AiVaultSession | null> {
return parseDevinSessionContent(
file,
await wslGatedReadFile(file.path, 'utf-8', 'scan'),
platform
platform,
{},
undefined,
messages
)
}
/** Remote content parses have no reader attached; the sink stays undefined. */
export function parseDevinSessionContent(
file: FileWithMtime,
content: string,
platform: NodeJS.Platform = process.platform,
options: ParserSessionOptions = {}
options: ParserSessionOptions = {},
_signal?: AbortSignal,
messages?: TranscriptMessageSink
): AiVaultSession | null {
const record = asRecord(JSON.parse(content) as unknown)
if (!record) {
@@ -48,7 +57,7 @@ export function parseDevinSessionContent(
extractString(record.session_id) ??
extractString(record.sessionId) ??
sessionIdFromFileName(file.path)
const accumulator = createAccumulator({ agent: 'devin', file, sessionId })
const accumulator = createAccumulator({ agent: 'devin', file, sessionId, messages })
const agentRecord = asRecord(record.agent)
accumulator.model =
extractString(agentRecord?.model_name) ??
+70 -51
View File
@@ -1,10 +1,11 @@
import type { Dirent } from 'node:fs'
import { extname, join } from 'node:path'
import { SessionNewestFiles } from './session-newest-files'
import type { AiVaultAgent, AiVaultScanIssue } from '../../shared/ai-vault-types'
import { wslGatedReaddir, wslGatedStat } from '../native-chat/wsl-transcript-fs-access'
import { WslTranscriptFsError } from '../native-chat/wsl-transcript-fs-gate'
import { recordSessionScanIssue } from './session-scan-issues'
import type { FileWithMtime, SessionFileDiscovery } from './session-scanner-types'
import type { SessionFileDiscovery } from './session-scanner-types'
import { errorMessage } from './session-scanner-values'
export async function discoverFiles(args: {
@@ -14,16 +15,45 @@ export async function discoverFiles(args: {
issues: AiVaultScanIssue[]
extensions: string[]
filePredicate?: (path: string) => boolean
contentDependencyPath?: (path: string) => string
contentDependencyPath?: (path: string) => string | undefined | Promise<string | undefined>
directoryPredicate?: (name: string, depth: number) => boolean
}): Promise<SessionFileDiscovery> {
let paths: string[]
const files = new SessionNewestFiles(args.limit)
try {
paths = await walkSessionFiles(args.rootDir, args.agent, args.issues, {
extensions: new Set(args.extensions),
filePredicate: args.filePredicate,
directoryPredicate: args.directoryPredicate
})
await forEachSessionFile(
args.rootDir,
args.agent,
args.issues,
{
extensions: new Set(args.extensions),
filePredicate: args.filePredicate,
directoryPredicate: args.directoryPredicate
},
async (path) => {
try {
const fileStat = await wslGatedStat(path, 'scan')
const dependencyStat = await optionalContentDependencyStat(
await args.contentDependencyPath?.(path)
)
const mtimeMs = Math.max(fileStat.mtimeMs, dependencyStat?.mtimeMs ?? 0)
files.add({
path,
mtimeMs,
modifiedAt: new Date(mtimeMs).toISOString(),
sizeBytes: fileStat.size + (dependencyStat?.size ?? 0),
dev: fileStat.dev,
ino: fileStat.ino,
nlink: fileStat.nlink
})
} catch (err) {
recordSessionScanIssue(args.issues, {
agent: args.agent,
path,
message: errorMessage(err)
})
}
}
)
} catch (err) {
// Why: discoverAiVaultSessionSources fans out with Promise.all, so one
// stalled distro would otherwise reject the whole vault scan — including
@@ -38,34 +68,7 @@ export async function discoverFiles(args: {
})
return { agent: args.agent, rootDir: args.rootDir, files: [] }
}
const files: FileWithMtime[] = []
for (const path of paths) {
try {
const fileStat = await wslGatedStat(path, 'scan')
const dependencyStat = await optionalContentDependencyStat(args.contentDependencyPath?.(path))
const mtimeMs = Math.max(fileStat.mtimeMs, dependencyStat?.mtimeMs ?? 0)
files.push({
path,
mtimeMs,
modifiedAt: new Date(mtimeMs).toISOString(),
sizeBytes: fileStat.size + (dependencyStat?.size ?? 0),
dev: fileStat.dev,
ino: fileStat.ino,
nlink: fileStat.nlink
})
} catch (err) {
recordSessionScanIssue(args.issues, {
agent: args.agent,
path,
message: errorMessage(err)
})
}
}
return {
agent: args.agent,
rootDir: args.rootDir,
files: files.sort((left, right) => right.mtimeMs - left.mtimeMs).slice(0, args.limit)
}
return { agent: args.agent, rootDir: args.rootDir, files: files.newest() }
}
async function optionalContentDependencyStat(
@@ -85,21 +88,39 @@ async function optionalContentDependencyStat(
}
}
export type SessionFileWalkOptions = {
extensions: Set<string>
filePredicate?: (path: string) => boolean
// Return false to skip descending into a directory; depth 0 is a child of
// rootDir, so pruned subtrees are never stat'd or parsed.
directoryPredicate?: (name: string, depth: number) => boolean
readDirectory?: (dirPath: string) => Promise<Dirent[]>
signal?: AbortSignal
}
/** Collecting form for callers that want every match; bounded scans stream. */
export async function walkSessionFiles(
dirPath: string,
agent: AiVaultAgent,
issues: AiVaultScanIssue[],
options: {
extensions: Set<string>
filePredicate?: (path: string) => boolean
// Return false to skip descending into a directory; depth 0 is a child of
// rootDir, so pruned subtrees are never stat'd or parsed.
directoryPredicate?: (name: string, depth: number) => boolean
readDirectory?: (dirPath: string) => Promise<Dirent[]>
signal?: AbortSignal
},
depth = 0
options: SessionFileWalkOptions
): Promise<string[]> {
const files: string[] = []
await forEachSessionFile(dirPath, agent, issues, options, async (path) => {
files.push(path)
})
return files
}
/** Streams matches to `onFile` so a bounded consumer never retains the whole tree. */
export async function forEachSessionFile(
dirPath: string,
agent: AiVaultAgent,
issues: AiVaultScanIssue[],
options: SessionFileWalkOptions,
onFile: (path: string) => Promise<void>,
depth = 0
): Promise<void> {
options.signal?.throwIfAborted()
let entries
try {
@@ -113,10 +134,9 @@ export async function walkSessionFiles(
if (error instanceof WslTranscriptFsError) {
throw error
}
return []
return
}
const files: string[] = []
for (const entry of entries) {
options.signal?.throwIfAborted()
const fullPath = join(dirPath, entry.name)
@@ -124,7 +144,7 @@ export async function walkSessionFiles(
// Skip whole subtrees an agent never wants (e.g. subagent transcripts),
// avoiding the readdir cost of descending into them.
if (options.directoryPredicate?.(entry.name, depth) ?? true) {
files.push(...(await walkSessionFiles(fullPath, agent, issues, options, depth + 1)))
await forEachSessionFile(fullPath, agent, issues, options, onFile, depth + 1)
}
continue
}
@@ -133,8 +153,7 @@ export async function walkSessionFiles(
options.extensions.has(extname(entry.name).toLowerCase()) &&
(options.filePredicate?.(fullPath) ?? true)
) {
files.push(fullPath)
await onFile(fullPath)
}
}
return files
}
@@ -8,6 +8,7 @@ import type {
ResumableSessionParseState,
SessionAccumulator
} from './session-scanner-types'
import type { TranscriptMessageSink } from './session-transcript-consumers'
import {
accumulatorFoldResumeState,
addPreviewMessage,
@@ -32,12 +33,13 @@ type ParserSessionOptions = {
export async function parseDroidSessionFile(
file: FileWithMtime,
platform: NodeJS.Platform = process.platform
platform: NodeJS.Platform = process.platform,
messages?: TranscriptMessageSink
): Promise<AiVaultSession | null> {
const input = openTranscriptReadStream(file.path, { encoding: 'utf-8' }, 'scan')
const lines = createInterface({ input, crlfDelay: Infinity })
try {
return await parseDroidSessionLines({ file, lines, platform })
return await parseDroidSessionLines({ file, lines, platform, messages })
} finally {
// readline.close() leaves the underlying stream open; destroy it so a
// mid-parse throw cannot leak the gated transcript handle.
@@ -94,9 +96,17 @@ function consumeDroidRecordLine(accumulator: SessionAccumulator, line: string):
}
}
export function createDroidSessionResumeState(file: FileWithMtime): ResumableSessionParseState {
export function createDroidSessionResumeState(
file: FileWithMtime,
messages?: TranscriptMessageSink
): ResumableSessionParseState {
return accumulatorFoldResumeState(
createAccumulator({ agent: 'droid', file, sessionId: sessionIdFromFileName(file.path) }),
createAccumulator({
agent: 'droid',
file,
sessionId: sessionIdFromFileName(file.path),
messages
}),
consumeDroidRecordLine
)
}
@@ -106,8 +116,9 @@ async function parseDroidSessionLines(args: {
lines: AsyncIterable<string> | Iterable<string>
platform: NodeJS.Platform
options?: ParserSessionOptions
messages?: TranscriptMessageSink
}): Promise<AiVaultSession | null> {
const state = createDroidSessionResumeState(args.file)
const state = createDroidSessionResumeState(args.file, args.messages)
for await (const line of args.lines) {
state.consumeLine(line)
}
@@ -8,6 +8,7 @@ import type {
ResumableSessionParseState,
SessionAccumulator
} from './session-scanner-types'
import type { TranscriptMessageSink } from './session-transcript-consumers'
import {
accumulatorFoldResumeState,
addPreviewContent,
@@ -27,16 +28,19 @@ import {
export async function parseGeminiSessionFile(
file: FileWithMtime,
platform: NodeJS.Platform = process.platform
platform: NodeJS.Platform = process.platform,
messages?: TranscriptMessageSink
): Promise<AiVaultSession | null> {
if (file.path.endsWith('.jsonl')) {
return parseGeminiJsonlSessionFile(file, platform)
return parseGeminiJsonlSessionFile(file, platform, messages)
}
return parseGeminiJsonSessionContent(
file,
await wslGatedReadFile(file.path, 'utf-8', 'scan'),
platform
platform,
{},
messages
)
}
@@ -62,7 +66,8 @@ function parseGeminiJsonSessionContent(
file: FileWithMtime,
content: string,
platform: NodeJS.Platform,
options: ResumableParseFinalizeOptions = {}
options: ResumableParseFinalizeOptions = {},
messages?: TranscriptMessageSink
): AiVaultSession | null {
const record = asRecord(JSON.parse(content) as unknown)
if (!record) {
@@ -71,7 +76,8 @@ function parseGeminiJsonSessionContent(
const accumulator = createAccumulator({
agent: 'gemini',
file,
sessionId: extractString(record.sessionId) ?? sessionIdFromFileName(file.path)
sessionId: extractString(record.sessionId) ?? sessionIdFromFileName(file.path),
messages
})
updateTimeline(accumulator, extractString(record.startTime))
updateTimeline(accumulator, extractString(record.lastUpdated))
@@ -83,13 +89,14 @@ function parseGeminiJsonSessionContent(
export async function parseGeminiJsonlSessionFile(
file: FileWithMtime,
platform: NodeJS.Platform
platform: NodeJS.Platform,
messages?: TranscriptMessageSink
): Promise<AiVaultSession | null> {
const lines = createInterface({
input: openTranscriptReadStream(file.path, { encoding: 'utf-8' }, 'scan'),
crlfDelay: Infinity
})
return parseGeminiJsonlSessionLines({ file, lines, platform })
return parseGeminiJsonlSessionLines({ file, lines, platform, messages })
}
function consumeGeminiJsonlRecordLine(accumulator: SessionAccumulator, line: string): void {
@@ -114,10 +121,16 @@ function consumeGeminiJsonlRecordLine(accumulator: SessionAccumulator, line: str
// Resumable only for the JSONL log format; Gemini's legacy single-JSON
// session documents are rewritten in place and must be re-read whole.
export function createGeminiJsonlSessionResumeState(
file: FileWithMtime
file: FileWithMtime,
messages?: TranscriptMessageSink
): ResumableSessionParseState {
return accumulatorFoldResumeState(
createAccumulator({ agent: 'gemini', file, sessionId: sessionIdFromFileName(file.path) }),
createAccumulator({
agent: 'gemini',
file,
sessionId: sessionIdFromFileName(file.path),
messages
}),
consumeGeminiJsonlRecordLine
)
}
@@ -127,8 +140,9 @@ async function parseGeminiJsonlSessionLines(args: {
lines: AsyncIterable<string> | Iterable<string>
platform: NodeJS.Platform
options?: ResumableParseFinalizeOptions
messages?: TranscriptMessageSink
}): Promise<AiVaultSession | null> {
const state = createGeminiJsonlSessionResumeState(args.file)
const state = createGeminiJsonlSessionResumeState(args.file, args.messages)
for await (const line of args.lines) {
state.consumeLine(line)
}
@@ -10,6 +10,7 @@ import type {
ResumableSessionParseState,
SessionAccumulator
} from './session-scanner-types'
import type { TranscriptMessageSink } from './session-transcript-consumers'
import {
accumulatorFoldResumeState,
addPreviewContent,
@@ -38,7 +39,8 @@ type ParserSessionOptions = {
export async function parseRovoSessionFile(
file: FileWithMtime,
platform: NodeJS.Platform = process.platform
platform: NodeJS.Platform = process.platform,
messages?: TranscriptMessageSink
): Promise<AiVaultSession | null> {
const metadata = asRecord(
JSON.parse(await wslGatedReadFile(file.path, 'utf-8', 'scan')) as unknown
@@ -49,7 +51,8 @@ export async function parseRovoSessionFile(
const accumulator = createAccumulator({
agent: 'rovo',
file,
sessionId: basename(dirname(file.path))
sessionId: basename(dirname(file.path)),
messages
})
accumulator.title = firstString(metadata, ['title', 'name', 'summary'])
accumulator.cwd = firstString(metadata, [
@@ -174,12 +177,13 @@ export type MessageGraphAgent = 'openclaw' | 'pi' | 'omp' | 'prime-agent'
export async function parseMessageGraphSessionFile(
agent: MessageGraphAgent,
file: FileWithMtime,
platform: NodeJS.Platform = process.platform
platform: NodeJS.Platform = process.platform,
messages?: TranscriptMessageSink
): Promise<AiVaultSession | null> {
const input = openTranscriptReadStream(file.path, { encoding: 'utf-8' }, 'scan')
const lines = createInterface({ input, crlfDelay: Infinity })
try {
return await parseMessageGraphSessionLines({ agent, file, lines, platform })
return await parseMessageGraphSessionLines({ agent, file, lines, platform, messages })
} finally {
// readline.close() leaves the underlying stream open; destroy it so a
// mid-parse throw cannot leak the gated transcript handle.
@@ -245,10 +249,11 @@ function consumeMessageGraphRecordLine(accumulator: SessionAccumulator, line: st
export function createMessageGraphSessionResumeState(
agent: MessageGraphAgent,
file: FileWithMtime
file: FileWithMtime,
messages?: TranscriptMessageSink
): ResumableSessionParseState {
const state = accumulatorFoldResumeState(
createAccumulator({ agent, file, sessionId: sessionIdFromFileName(file.path) }),
createAccumulator({ agent, file, sessionId: sessionIdFromFileName(file.path), messages }),
consumeMessageGraphRecordLine
)
// Why: only OMP materializes task-subagent transcripts beside its sessions
@@ -263,8 +268,9 @@ async function parseMessageGraphSessionLines(args: {
lines: AsyncIterable<string> | Iterable<string>
platform: NodeJS.Platform
options?: ParserSessionOptions
messages?: TranscriptMessageSink
}): Promise<AiVaultSession | null> {
const state = createMessageGraphSessionResumeState(args.agent, args.file)
const state = createMessageGraphSessionResumeState(args.agent, args.file, args.messages)
for await (const line of args.lines) {
state.consumeLine(line)
}
@@ -4,6 +4,7 @@ import { dirname, join } from 'node:path'
import { createInterface } from 'node:readline'
import type { AiVaultSession } from '../../shared/ai-vault-types'
import type { FileWithMtime, SessionAccumulator } from './session-scanner-types'
import type { TranscriptMessageSink } from './session-transcript-consumers'
import {
addPreviewMessage,
createAccumulator,
@@ -34,7 +35,8 @@ const GROK_USER_QUERY_PREVIEW_SCAN_LIMIT = 4096
export async function parseGrokSessionFile(
file: FileWithMtime,
platform: NodeJS.Platform = process.platform
platform: NodeJS.Platform = process.platform,
messages?: TranscriptMessageSink
): Promise<AiVaultSession | null> {
const record = asRecord(JSON.parse(await wslGatedReadFile(file.path, 'utf-8', 'scan')) as unknown)
if (!record) {
@@ -42,7 +44,7 @@ export async function parseGrokSessionFile(
}
const info = asRecord(record.info)
const sessionId = extractString(info?.id) ?? sessionIdFromFileName(dirname(file.path))
const accumulator = createAccumulator({ agent: 'grok', file, sessionId })
const accumulator = createAccumulator({ agent: 'grok', file, sessionId, messages })
accumulator.cwd = extractString(info?.cwd)
accumulator.title =
normalizeTitleText(extractString(record.generated_title) ?? '') ??
@@ -2,6 +2,7 @@ import { wslGatedReadFile } from '../native-chat/wsl-transcript-fs-access'
import type { AiVaultSession } from '../../shared/ai-vault-types'
import type { ExecutionHostId } from '../../shared/execution-host'
import type { FileWithMtime } from './session-scanner-types'
import type { TranscriptMessageSink } from './session-transcript-consumers'
import {
addPreviewContent,
createAccumulator,
@@ -24,20 +25,28 @@ type ParserSessionOptions = {
export async function parseHermesSessionFile(
file: FileWithMtime,
platform: NodeJS.Platform = process.platform
platform: NodeJS.Platform = process.platform,
messages?: TranscriptMessageSink
): Promise<AiVaultSession | null> {
return parseHermesSessionContent(
file,
await wslGatedReadFile(file.path, 'utf-8', 'scan'),
platform
platform,
{},
undefined,
messages
)
}
/** Remote content parses have no reader attached; the sink stays undefined. */
export async function parseHermesSessionContent(
file: FileWithMtime,
content: string,
platform: NodeJS.Platform = process.platform,
options: ParserSessionOptions = {}
options: ParserSessionOptions = {},
_signal?: AbortSignal,
messages?: TranscriptMessageSink
): Promise<AiVaultSession | null> {
const record = asRecord(JSON.parse(content) as unknown)
if (!record) {
@@ -46,7 +55,8 @@ export async function parseHermesSessionContent(
const accumulator = createAccumulator({
agent: 'hermes',
file,
sessionId: extractString(record.session_id) ?? sessionIdFromFileName(file.path)
sessionId: extractString(record.session_id) ?? sessionIdFromFileName(file.path),
messages
})
accumulator.model = extractString(record.model)
accumulator.cwd = extractString(record.cwd)
@@ -16,6 +16,7 @@ import {
readKimiWorkDirBySessionId
} from './session-scanner-kimi-paths'
import type { FileWithMtime, SessionAccumulator } from './session-scanner-types'
import type { TranscriptMessageSink } from './session-transcript-consumers'
import {
asRecord,
extractContentText,
@@ -32,7 +33,8 @@ import {
// session_index.jsonl; model/messages/tokens come from the wire transcript.
export async function parseKimiSessionFile(
file: FileWithMtime,
platform: NodeJS.Platform = process.platform
platform: NodeJS.Platform = process.platform,
messages?: TranscriptMessageSink
): Promise<AiVaultSession | null> {
let stateRecord: Record<string, unknown> | null
try {
@@ -53,7 +55,7 @@ export async function parseKimiSessionFile(
}
const sessionId = kimiSessionIdFromStatePath(file.path)
const accumulator = createAccumulator({ agent: 'kimi', file, sessionId })
const accumulator = createAccumulator({ agent: 'kimi', file, sessionId, messages })
// Why: Kimi sessions are work-dir-scoped — the resume command must `cd` into
// the original directory or the CLI rejects it. That path lives only in the
@@ -3,6 +3,7 @@ import { WslTranscriptFsError } from '../native-chat/wsl-transcript-fs-gate'
import { join } from 'node:path'
import type { AiVaultSession } from '../../shared/ai-vault-types'
import type { FileWithMtime, SessionAccumulator } from './session-scanner-types'
import type { TranscriptMessageSink } from './session-transcript-consumers'
import {
addPreviewMessage,
createAccumulator,
@@ -25,14 +26,15 @@ import {
export async function parseOpenCodeSessionFile(
file: FileWithMtime,
platform: NodeJS.Platform = process.platform
platform: NodeJS.Platform = process.platform,
messages?: TranscriptMessageSink
): Promise<AiVaultSession | null> {
const record = asRecord(JSON.parse(await wslGatedReadFile(file.path, 'utf-8', 'scan')) as unknown)
if (!record) {
return null
}
const sessionId = extractString(record.id) ?? sessionIdFromFileName(file.path)
const accumulator = createAccumulator({ agent: 'opencode', file, sessionId })
const accumulator = createAccumulator({ agent: 'opencode', file, sessionId, messages })
accumulator.title = normalizeTitleText(extractString(record.title) ?? '')
accumulator.cwd = extractString(record.directory)
updateTimeline(accumulator, timeObjectValue(record.time, 'created'))
@@ -0,0 +1,29 @@
import type SyncDatabase from '../sqlite/sync-database'
import { columnExists, tableExists } from '../opencode-usage/schema-helpers'
// Why: OpenCode's schema has moved more than once, so every read probes for the
// columns it names. These are the two shapes the session parser depends on;
// keeping them here stops each reader from inventing its own partial gate.
/** Enough of `message` to count a session's turns. */
export function canCountOpenCodeMessages(db: SyncDatabase): boolean {
return (
tableExists(db, 'message') &&
columnExists(db, 'message', 'session_id') &&
columnExists(db, 'message', 'data')
)
}
/** Enough of `message`×`part` to read a session's parts in turn order. */
export function canReadOpenCodeMessageParts(db: SyncDatabase): boolean {
return (
canCountOpenCodeMessages(db) &&
columnExists(db, 'message', 'id') &&
// Every parts read orders by it; unprobed, a schema without it throws mid-read.
columnExists(db, 'message', 'time_created') &&
tableExists(db, 'part') &&
columnExists(db, 'part', 'message_id') &&
columnExists(db, 'part', 'time_created') &&
columnExists(db, 'part', 'data')
)
}
@@ -1,4 +1,4 @@
import type { Worker } from 'node:worker_threads'
import { LazyWorkerThreadHost, type WorkerThreadFactory } from '../lazy-worker-thread-host'
import type { AiVaultScanIssue, AiVaultSession } from '../../shared/ai-vault-types'
import type {
OpenCodeSqliteListRequest,
@@ -11,9 +11,10 @@ import type { SessionFileCandidate } from './session-scanner-types'
import { errorMessage } from './session-scanner-values'
// Why (#8864): a lazily-spawned, unref'd worker runs OpenCode SQLite reads off
// the main-process event loop. Lifecycle (idle teardown, FIFO one-at-a-time
// dispatch, per-call timeouts, respawn-on-fault) mirrors src/main/speech/
// stt-service.ts. The default spawn + shared singleton live in
// the main-process event loop. This module owns the request half (FIFO
// one-at-a-time dispatch, per-call timeouts, respawn-on-fault); the thread's
// own lifetime belongs to LazyWorkerThreadHost, shared with the port-scan probe
// client. The default spawn + shared singleton live in
// session-scanner-opencode-sqlite-worker-spawn.ts.
export const LIST_TIMEOUT_MS = 30_000
@@ -25,8 +26,6 @@ export const IDLE_TEARDOWN_MS = 30_000
// fresh scan burst starts from idle (so the cap is per-scan, not process-wide).
export const MAX_CONSECUTIVE_DEATHS = 3
export type WorkerFactory = () => Worker
// Omit<union, 'id'> collapses to the shared keys, so omit each member and let
// the client stamp the correlation id.
type OpenCodeSqliteRequestBody =
@@ -53,20 +52,27 @@ class OpenCodeSqliteWorkerUnavailableError extends Error {}
* no worker can be spawned rather than moving SQLite work onto the main thread.
*/
export class OpenCodeSqliteWorkerClient {
private worker: Worker | null = null
private active: PendingCall | null = null
private queue: PendingCall[] = []
private idleTimer: NodeJS.Timeout | null = null
private consecutiveDeaths = 0
private nextId = 1
private loggedWorkerUnavailable = false
private cleanupWorkerListeners: (() => void) | null = null
private readonly workerFactory: WorkerFactory
private readonly log: (message: string) => void
private readonly host: LazyWorkerThreadHost<OpenCodeSqliteWorkerResponse>
constructor(options: { workerFactory: WorkerFactory; log?: (message: string) => void }) {
this.workerFactory = options.workerFactory
this.log = options.log ?? ((message) => console.warn(message))
constructor(options: { workerFactory: WorkerThreadFactory; log?: (message: string) => void }) {
const log = options.log ?? ((message: string) => console.warn(message))
this.host = new LazyWorkerThreadHost<OpenCodeSqliteWorkerResponse>({
factory: options.workerFactory,
idleTeardownMs: IDLE_TEARDOWN_MS,
onMessage: (response) => this.onMessage(response),
onError: (error) => this.onWorkerFault(error),
onExit: (code) => this.onWorkerExit(code),
isIdle: () => !this.active && this.queue.length === 0,
// Why (#8864): never fall back to synchronous SQLite reads here; a missing
// bundle or resource-exhausted spawn must omit OpenCode history rather than
// reintroduce the main-process hang this worker boundary prevents.
onUnavailable: (err) =>
log(`OpenCode SQLite worker unavailable; skipping its history. ${errorMessage(err)}`)
})
}
/**
@@ -167,7 +173,7 @@ export class OpenCodeSqliteWorkerClient {
if (this.active || this.queue.length === 0) {
return
}
const worker = this.ensureWorker()
const worker = this.host.ensure()
if (!worker) {
this.failQueuedAsUnavailable()
return
@@ -177,7 +183,7 @@ export class OpenCodeSqliteWorkerClient {
return
}
this.active = call
this.clearIdleTimer()
this.host.clearIdleTimer()
// Timeout clock starts at dispatch (not enqueue): a batch may enqueue up to
// 8 parses at once, and a queue-inclusive timeout would fire falsely.
call.timer = setTimeout(() => this.onTimeout(call), call.timeoutMs)
@@ -185,39 +191,6 @@ export class OpenCodeSqliteWorkerClient {
worker.postMessage(call.request)
}
private ensureWorker(): Worker | null {
if (this.worker) {
return this.worker
}
try {
const worker = this.workerFactory()
const onMessage = (response: OpenCodeSqliteWorkerResponse): void => this.onMessage(response)
const onError = (error: Error): void => this.onWorkerFault(error)
const onExit = (code: number): void => this.onWorkerExit(code)
worker.on('message', onMessage)
worker.on('error', onError)
worker.on('exit', onExit)
this.cleanupWorkerListeners = () => {
worker.off('message', onMessage)
worker.off('error', onError)
worker.off('exit', onExit)
}
// Never keep the app alive for a scan worker.
worker.unref?.()
this.worker = worker
return worker
} catch (err) {
// Why (#8864): never fall back to synchronous SQLite reads here; a missing
// bundle or resource-exhausted spawn must omit OpenCode history rather than
// reintroduce the main-process hang this worker boundary prevents.
if (!this.loggedWorkerUnavailable) {
this.loggedWorkerUnavailable = true
this.log(`OpenCode SQLite worker unavailable; skipping its history. ${errorMessage(err)}`)
}
return null
}
}
private onMessage(response: OpenCodeSqliteWorkerResponse): void {
const call = this.active
if (!call || call.request.id !== response.id) {
@@ -243,7 +216,7 @@ export class OpenCodeSqliteWorkerClient {
// A clean self-exit is not a death, but the stale handle must be dropped
// or the next dispatch would post into the dead worker and stall to timeout.
if (code === 0 && !this.active && this.queue.length === 0) {
this.destroyWorker()
this.host.destroy()
return
}
this.onWorkerFault(new Error(`OpenCode SQLite worker exited with code ${code}`))
@@ -251,7 +224,7 @@ export class OpenCodeSqliteWorkerClient {
private onWorkerFault(error: Error): void {
const failed = this.active
this.destroyWorker()
this.host.destroy()
this.consecutiveDeaths++
if (failed) {
this.settle(failed, () => failed.reject(error))
@@ -302,46 +275,7 @@ export class OpenCodeSqliteWorkerClient {
if (this.queue.length > 0) {
this.pump()
} else {
this.scheduleIdleTeardown()
this.host.scheduleIdleTeardown()
}
}
private scheduleIdleTeardown(): void {
this.clearIdleTimer()
if (!this.worker) {
return
}
this.idleTimer = setTimeout(() => this.teardownIfIdle(), IDLE_TEARDOWN_MS)
this.idleTimer.unref?.()
}
private teardownIfIdle(): void {
this.idleTimer = null
// Only tear down with nothing active AND nothing queued: a request arriving
// as the timer fires must never be lost to a self-exiting worker.
if (this.active || this.queue.length > 0) {
return
}
this.destroyWorker()
}
private clearIdleTimer(): void {
if (this.idleTimer) {
clearTimeout(this.idleTimer)
this.idleTimer = null
}
}
private destroyWorker(): void {
this.clearIdleTimer()
const worker = this.worker
this.worker = null
if (!worker) {
return
}
this.cleanupWorkerListeners?.()
this.cleanupWorkerListeners = null
worker.removeAllListeners()
void worker.terminate().catch(() => undefined)
}
}
@@ -10,6 +10,10 @@ import {
shouldCaptureFullFirstUserPrompt
} from './session-scanner-first-user-prompt'
import { readOpenCodeDatabase } from './session-scanner-opencode-sqlite-open'
import {
canCountOpenCodeMessages,
canReadOpenCodeMessageParts
} from './session-scanner-opencode-sqlite-schema'
import { normalizeTitleText } from './session-scanner-values'
import type SyncDatabase from '../sqlite/sync-database'
import { columnExists, tableExists } from '../opencode-usage/schema-helpers'
@@ -70,14 +74,6 @@ function sessionNumberColumnSelect(db: SyncDatabase, columnName: string): string
return columnExists(db, 'session', columnName) ? `s.${columnName}` : '0'
}
function canCountOpenCodeMessages(db: SyncDatabase): boolean {
return (
tableExists(db, 'message') &&
columnExists(db, 'message', 'session_id') &&
columnExists(db, 'message', 'data')
)
}
function buildSessionQuery(db: SyncDatabase): string {
const messageCountSubquery = canCountOpenCodeMessages(db)
? `(SELECT COUNT(*) FROM message m
@@ -154,14 +150,7 @@ function extractPartText(partData: string): string | null {
}
function readFirstUserPromptFromOpenCodeDb(db: SyncDatabase, sessionId: string): string | null {
if (
!canCountOpenCodeMessages(db) ||
!tableExists(db, 'part') ||
!columnExists(db, 'message', 'id') ||
!columnExists(db, 'part', 'message_id') ||
!columnExists(db, 'part', 'time_created') ||
!columnExists(db, 'part', 'data')
) {
if (!canReadOpenCodeMessageParts(db)) {
return null
}
@@ -206,14 +195,7 @@ function readFirstUserPromptFromOpenCodeDb(db: SyncDatabase, sessionId: string):
}
function buildPreviewQuery(db: SyncDatabase): string | null {
if (
!canCountOpenCodeMessages(db) ||
!tableExists(db, 'part') ||
!columnExists(db, 'message', 'id') ||
!columnExists(db, 'part', 'message_id') ||
!columnExists(db, 'part', 'time_created') ||
!columnExists(db, 'part', 'data')
) {
if (!canReadOpenCodeMessageParts(db)) {
return null
}
return `SELECT json_extract(m.data, '$.role') AS role,
+83 -228
View File
@@ -1,7 +1,5 @@
import { readTranscriptSlice } from '../native-chat/wsl-transcript-fs-access'
import type { AiVaultSession } from '../../shared/ai-vault-types'
import { createAntigravitySessionResumeState } from './session-scanner-antigravity-parser'
import { parseAgentSessionFile } from './session-scanner-agent-parser'
import { createCodexSessionResumeState } from './session-scanner-codex-parser'
import { createDroidSessionResumeState } from './session-scanner-droid-parser'
import { createMessageGraphSessionResumeState } from './session-scanner-graph-parsers'
@@ -13,27 +11,25 @@ import { countSubagentTranscripts } from './session-scanner-subagent-transcripts
import { countOmpSubagentTranscripts } from './session-scanner-omp-subagent-transcripts'
import type { ResumableSessionParseState, SessionFileCandidate } from './session-scanner-types'
import { refreshCachedCodexTitle } from './session-scanner-codex-cached-title'
import { consumeCompleteJsonlLines } from './session-scanner-jsonl-reader'
import {
getSessionParseCacheEntry,
storeSessionParseCacheEntry,
type SessionParseCacheEntry
} from './session-parse-cache-store'
import type { TranscriptMessageSink } from './session-transcript-consumers'
import {
readResumableTranscript,
readWholeTranscript,
type TranscriptReadStats
} from './session-transcript-reader'
// Sized past the default recency cap (1000) plus the in-scope cap (2000) so a
// full steady-state result set stays resident between forced rescans.
const MAX_CACHE_ENTRIES = 4096
const NEWLINE_BYTE = 0x0a
type ResumePoint = {
state: ResumableSessionParseState
// Byte offset just past the last complete ('\n'-terminated) line consumed;
// a trailing unterminated line is deliberately left before this point.
byteOffset: number
}
type SessionParseCacheEntry = {
mtimeMs: number
sizeBytes: number | null
platform: NodeJS.Platform
session: AiVaultSession | null
resume: ResumePoint | null
}
export {
invalidateSessionParseCacheEntry,
resetSessionParseCacheForTests,
seedSessionParseCache,
snapshotSessionParseCacheForPersistence,
type PersistedSessionParseCacheEntry
} from './session-parse-cache-store'
// Incremental append-parsing applies only to transcripts that are append-only
// JSONL line-folds. Whole-JSON documents (grok/rovo/devin/hermes/gemini-json)
@@ -44,31 +40,32 @@ type SessionParseCacheEntry = {
// cached state instead, never pay for a throwaway accumulator.
function resumableStateFactoryFor(
candidate: SessionFileCandidate
): (() => ResumableSessionParseState) | null {
): ((messages: TranscriptMessageSink) => ResumableSessionParseState) | null {
switch (candidate.agent) {
case 'claude':
return () => createClaudeSessionResumeState(candidate.file)
return (messages) => createClaudeSessionResumeState(candidate.file, messages)
case 'codex':
return () => createCodexSessionResumeState(candidate.file, candidate.codexHome)
return (messages) =>
createCodexSessionResumeState(candidate.file, candidate.codexHome, messages)
case 'cursor':
return () => createCursorSessionResumeState(candidate.file)
return (messages) => createCursorSessionResumeState(candidate.file, true, messages)
case 'copilot':
return () => createCopilotSessionResumeState(candidate.file)
return (messages) => createCopilotSessionResumeState(candidate.file, messages)
case 'droid':
return () => createDroidSessionResumeState(candidate.file)
return (messages) => createDroidSessionResumeState(candidate.file, messages)
case 'openclaw':
case 'pi':
case 'omp':
case 'prime-agent': {
const agent = candidate.agent
return () => createMessageGraphSessionResumeState(agent, candidate.file)
return (messages) => createMessageGraphSessionResumeState(agent, candidate.file, messages)
}
case 'gemini':
return candidate.file.path.endsWith('.jsonl')
? () => createGeminiJsonlSessionResumeState(candidate.file)
? (messages) => createGeminiJsonlSessionResumeState(candidate.file, messages)
: null
case 'antigravity':
return () => createAntigravitySessionResumeState(candidate.file)
return (messages) => createAntigravitySessionResumeState(candidate.file, messages)
case 'devin':
case 'grok':
case 'hermes':
@@ -80,98 +77,23 @@ function resumableStateFactoryFor(
}
}
export type SessionParseStats = {
export type SessionParseStats = TranscriptReadStats & {
reused: number
incremental: number
fullParses: number
// Transcripts the parser already excluded (Codex workers), re-listed after a
// write and dismissed without reading. Counted apart from `incremental` so a
// scan span still shows how much work the early stop actually removed.
earlyStopped: number
bytesRead: number
}
export function createSessionParseStats(): SessionParseStats {
return { reused: 0, incremental: 0, fullParses: 0, earlyStopped: 0, bytesRead: 0 }
}
const cache = new Map<string, SessionParseCacheEntry>()
export function resetSessionParseCacheForTests(): void {
cache.clear()
}
// Drops one entry after its file is deleted. Cleanliness, not correctness:
// discovery walks disk first, so a trashed file is never rediscovered anyway.
export function invalidateSessionParseCacheEntry(path: string): void {
cache.delete(path)
}
// Persisted subset of a cache entry: the non-serializable `resume` parser
// state is dropped (see session-parse-cache-persistence.ts).
export type PersistedSessionParseCacheEntry = Omit<SessionParseCacheEntry, 'resume'>
export function snapshotSessionParseCacheForPersistence(): [
string,
PersistedSessionParseCacheEntry
][] {
return [...cache].map(([path, entry]): [string, PersistedSessionParseCacheEntry] => [
path,
{
mtimeMs: entry.mtimeMs,
sizeBytes: entry.sizeBytes,
platform: entry.platform,
session: entry.session
}
])
}
// Seeded entries carry `resume: null`: after a restart an unchanged file is a
// cache hit; a file that changed while the app was closed pays one full
// (not incremental) re-parse.
export function seedSessionParseCache(
entries: Iterable<[string, PersistedSessionParseCacheEntry]>
): void {
const list = [...entries]
// Snapshot order is oldest→newest (LRU); an over-cap list keeps the newest
// tail rather than seeding the oldest entries and dropping the tail.
for (const [path, entry] of list.slice(Math.max(0, list.length - MAX_CACHE_ENTRIES))) {
if (cache.size >= MAX_CACHE_ENTRIES) {
return
}
// In-process entries are always fresher than persisted ones; never clobber.
if (cache.has(path)) {
continue
}
cache.set(path, {
mtimeMs: entry.mtimeMs,
sizeBytes: entry.sizeBytes,
platform: entry.platform,
session: entry.session,
resume: null
})
}
}
function storeEntry(path: string, entry: SessionParseCacheEntry): void {
cache.delete(path)
cache.set(path, entry)
if (cache.size > MAX_CACHE_ENTRIES) {
const oldest = cache.keys().next()
if (!oldest.done) {
cache.delete(oldest.value)
}
}
}
/**
* Parse a session file, reusing prior work where the file is provably
* unchanged (mtime+size) and, for append-only JSONL transcripts (Claude,
* Codex, Cursor, Copilot, Droid, OpenClaw/Pi/OMP, Gemini-JSONL), resuming the
* parse from the last consumed byte when the file only grew. This is what
* keeps the renderer's ~5s forced rescans from re-reading gigabytes of
* transcripts (STA-1278/STA-1417: main process pegging one core during
* multi-agent workloads).
* The session list's cursor over the transcript reader: it remembers what each
* file looked like when it was last listed, reuses that work where the file is
* provably unchanged (mtime+size), and otherwise asks the reader to resume from
* the last consumed byte or re-read the file whole. This is what keeps the
* renderer's ~5s forced rescans from re-reading gigabytes of transcripts
* (STA-1278/STA-1417: main process pegging one core during multi-agent
* workloads). Other consumers of the reader keep their own equivalent cursor
* and never consult this one.
*/
export async function parseAgentSessionFileCached(
candidate: SessionFileCandidate,
@@ -179,7 +101,7 @@ export async function parseAgentSessionFileCached(
stats?: SessionParseStats
): Promise<AiVaultSession | null> {
const { file } = candidate
const entry = cache.get(file.path)
const entry = getSessionParseCacheEntry(file.path)
const unchanged =
entry !== undefined &&
@@ -187,56 +109,30 @@ export async function parseAgentSessionFileCached(
entry.mtimeMs === file.mtimeMs &&
(entry.sizeBytes === null || file.sizeBytes === undefined || entry.sizeBytes === file.sizeBytes)
if (unchanged) {
if (stats) {
stats.reused++
}
// A zero-turn transcript usually never changes again, but its sibling
// subagent dir (Claude `<session>/subagents/`, OMP's same-named artifact
// dir) can gain files after the parent's last write (a still-running
// subagent finishing). The mtime+size key can't see that, so refresh the
// cheap directory count on reuse.
if (entry.session && entry.session.messageCount === 0) {
const subagentTranscriptCount =
candidate.agent === 'claude'
? await countSubagentTranscripts(file.path)
: candidate.agent === 'omp'
? await countOmpSubagentTranscripts(file.path)
: null
if (
subagentTranscriptCount !== null &&
subagentTranscriptCount !== entry.session.subagentTranscriptCount
) {
entry.session = { ...entry.session, subagentTranscriptCount }
}
}
// Codex titles come from session_index.jsonl, which mtime+size can't see.
// Remote counterpart: remote-session-scanner.ts's reusedCodexTitleRefresh.
if (entry.session && candidate.agent === 'codex') {
entry.session = await refreshCachedCodexTitle(candidate, entry.session)
}
storeEntry(file.path, entry)
return entry.session
return reuseCachedSession(candidate, entry, stats)
}
const stateFactory = resumableStateFactoryFor(candidate)
if (stateFactory) {
const parsed = await parseResumableCandidate({
const read = await readResumableTranscript({
candidate,
platform,
entry,
stats,
stateFactory
resume: entry?.platform === platform ? entry.resume : null,
stateFactory,
stats
})
storeEntry(file.path, parsed)
return parsed.session
storeSessionParseCacheEntry(file.path, {
mtimeMs: file.mtimeMs,
sizeBytes: file.sizeBytes ?? null,
platform,
session: read.session,
resume: read.resume
})
return read.session
}
if (stats) {
stats.fullParses++
stats.bytesRead += file.sizeBytes ?? 0
}
const session = await parseAgentSessionFile(candidate, platform)
storeEntry(file.path, {
const session = await readWholeTranscript({ candidate, platform, stats })
storeSessionParseCacheEntry(file.path, {
mtimeMs: file.mtimeMs,
sizeBytes: file.sizeBytes ?? null,
platform,
@@ -246,79 +142,38 @@ export async function parseAgentSessionFileCached(
return session
}
async function parseResumableCandidate(args: {
candidate: SessionFileCandidate
platform: NodeJS.Platform
entry: SessionParseCacheEntry | undefined
async function reuseCachedSession(
candidate: SessionFileCandidate,
entry: SessionParseCacheEntry,
stats?: SessionParseStats
stateFactory: () => ResumableSessionParseState
}): Promise<SessionParseCacheEntry> {
const { file } = args.candidate
const resume = args.entry?.platform === args.platform ? args.entry.resume : null
const canResume =
resume !== null &&
resume !== undefined &&
typeof file.sizeBytes === 'number' &&
file.sizeBytes >= resume.byteOffset &&
(resume.byteOffset === 0 || (await endsWithNewlineAt(file.path, resume.byteOffset)))
// Clone before consuming: a failed read must not corrupt the cached state,
// or the next resume would double-count the lines applied before the error.
const state = canResume ? resume.state.clone() : args.stateFactory()
const startOffset = canResume ? resume.byteOffset : 0
// Mirrors the reader's entry guard so a dismissed transcript is not reported
// as an incremental parse that read nothing.
const stoppedBeforeRead = state.shouldStop?.() === true
if (args.stats) {
if (stoppedBeforeRead) {
args.stats.earlyStopped++
} else if (canResume) {
args.stats.incremental++
} else {
args.stats.fullParses++
): Promise<AiVaultSession | null> {
if (stats) {
stats.reused++
}
// A zero-turn transcript usually never changes again, but its sibling
// subagent dir (Claude `<session>/subagents/`, OMP's same-named artifact
// dir) can gain files after the parent's last write (a still-running
// subagent finishing). The mtime+size key can't see that, so refresh the
// cheap directory count on reuse.
if (entry.session && entry.session.messageCount === 0) {
const subagentTranscriptCount =
candidate.agent === 'claude'
? await countSubagentTranscripts(candidate.file.path)
: candidate.agent === 'omp'
? await countOmpSubagentTranscripts(candidate.file.path)
: null
if (
subagentTranscriptCount !== null &&
subagentTranscriptCount !== entry.session.subagentTranscriptCount
) {
entry.session = { ...entry.session, subagentTranscriptCount }
}
}
const readResult = await consumeCompleteJsonlLines({
path: file.path,
start: startOffset,
onLine: (line) => state.consumeLine(line),
// Bound: the optional hooks are declared as methods, so a parser written
// with method syntax must not lose `this` on the way into the reader.
onLineBytes: state.consumeLineBytes?.bind(state),
shouldStop: state.shouldStop?.bind(state)
})
if (args.stats) {
args.stats.bytesRead += readResult.bytesRead
}
// The stat this scan displays is current even when nothing new was consumed.
state.touchFile(file)
// Keep parity with the one-shot parser: a final unterminated line is shown,
// but stays out of the resumable state so the (possibly still-growing) line
// is re-read once complete instead of being half-counted.
let displayState = state
if (readResult.trailingPartialLine !== null) {
displayState = state.clone()
displayState.consumeLine(readResult.trailingPartialLine)
}
return {
mtimeMs: file.mtimeMs,
sizeBytes: file.sizeBytes ?? null,
platform: args.platform,
session: await displayState.finalize(args.platform),
resume: { state, byteOffset: readResult.consumedThrough }
// Codex titles come from session_index.jsonl, which mtime+size can't see.
// Remote counterpart: remote-session-scanner.ts's reusedCodexTitleRefresh.
if (entry.session && candidate.agent === 'codex') {
entry.session = await refreshCachedCodexTitle(candidate, entry.session)
}
}
// A resume point is only valid if it still sits just past a line break;
// anything else means the file was rewritten, not appended. Heuristic: a
// grown rewrite keeping '\n' at exactly this byte would slip through, but
// agent transcripts are append-only so that trade is accepted (worst case is
// a stale vault row until the file is next truncated or the app restarts).
async function endsWithNewlineAt(path: string, offset: number): Promise<boolean> {
const slice = await readTranscriptSlice(path, offset - 1, 1, 'scan')
return slice.length === 1 && slice[0] === NEWLINE_BYTE
storeSessionParseCacheEntry(candidate.file.path, entry)
return entry.session
}
@@ -10,6 +10,7 @@ import type {
ResumableSessionParseState,
SessionAccumulator
} from './session-scanner-types'
import type { TranscriptMessageSink } from './session-transcript-consumers'
import {
addPreviewContent,
createAccumulator,
@@ -42,12 +43,16 @@ export type ClaudeSessionParseState = {
firstUserTitle: string | null
}
export function createClaudeSessionParseState(file: FileWithMtime): ClaudeSessionParseState {
export function createClaudeSessionParseState(
file: FileWithMtime,
messages?: TranscriptMessageSink
): ClaudeSessionParseState {
return {
accumulator: createAccumulator({
agent: 'claude',
file,
sessionId: sessionIdFromFileName(file.path)
sessionId: sessionIdFromFileName(file.path),
messages
}),
metaTitle: null,
generatedTitle: null,
@@ -188,8 +193,11 @@ export async function finalizeClaudeSessionParseState(
return finalizeSession(snapshot.accumulator, platform, options)
}
export function createClaudeSessionResumeState(file: FileWithMtime): ResumableSessionParseState {
return claudeResumeStateFromParseState(createClaudeSessionParseState(file))
export function createClaudeSessionResumeState(
file: FileWithMtime,
messages?: TranscriptMessageSink
): ResumableSessionParseState {
return claudeResumeStateFromParseState(createClaudeSessionParseState(file, messages))
}
function claudeResumeStateFromParseState(
@@ -207,13 +215,14 @@ function claudeResumeStateFromParseState(
export async function parseClaudeSessionFile(
file: FileWithMtime,
platform: NodeJS.Platform = process.platform
platform: NodeJS.Platform = process.platform,
messages?: TranscriptMessageSink
): Promise<AiVaultSession | null> {
const lines = createInterface({
input: openTranscriptReadStream(file.path, { encoding: 'utf-8' }, 'scan'),
crlfDelay: Infinity
})
return parseClaudeSessionLines({ file, lines, platform })
return parseClaudeSessionLines({ file, lines, platform, messages })
}
export async function parseClaudeSessionContent(
@@ -236,8 +245,9 @@ async function parseClaudeSessionLines(args: {
lines: AsyncIterable<string> | Iterable<string>
platform: NodeJS.Platform
options?: ParserSessionOptions
messages?: TranscriptMessageSink
}): Promise<AiVaultSession | null> {
const state = createClaudeSessionParseState(args.file)
const state = createClaudeSessionParseState(args.file, args.messages)
for await (const line of args.lines) {
consumeClaudeSessionLine(state, line)
}
@@ -4,6 +4,7 @@ import { discoverFiles } from './session-scanner-discovery'
import { opencodeDiscoveries } from './session-scanner-opencode-sources'
import { antigravityDiscoveries } from './session-scanner-antigravity-sources'
import { AI_VAULT_AGENT_SOURCES, type AiVaultAgentSource } from './session-scanner-agent-sources'
import { withCursorChatMetaScan } from './session-scanner-cursor-chat-meta'
import { normalizedWslHomeDirs } from './session-scanner-roots'
import type { AiVaultScanOptions, SessionFileDiscovery } from './session-scanner-types'
@@ -17,25 +18,27 @@ export async function discoverAiVaultSessionSources(args: {
const { options, limitPerAgent, issues } = args
const wslHomeDirs = normalizedWslHomeDirs(options.wslHomeDirs)
return Promise.all([
// Why: OpenCode 1.17.x migrated sessions from per-session JSON files to a
// SQLite DB. discoverOpenCodeSessions runs both the file scanner (legacy)
// and the SQLite scanner (1.17.x); dedup by sessionId happens inside.
...opencodeDiscoveries(options, wslHomeDirs, limitPerAgent, issues),
...antigravityDiscoveries(options, wslHomeDirs, limitPerAgent, issues),
...Object.entries(AI_VAULT_AGENT_SOURCES).flatMap(([agent, source]) =>
source
? agentDiscoveries(
agent as AiVaultAgent,
source,
options,
wslHomeDirs,
limitPerAgent,
issues
)
: []
)
])
return withCursorChatMetaScan(() =>
Promise.all([
// Why: OpenCode 1.17.x migrated sessions from per-session JSON files to a
// SQLite DB. discoverOpenCodeSessions runs both the file scanner (legacy)
// and the SQLite scanner (1.17.x); dedup by sessionId happens inside.
...opencodeDiscoveries(options, wslHomeDirs, limitPerAgent, issues),
...antigravityDiscoveries(options, wslHomeDirs, limitPerAgent, issues),
...Object.entries(AI_VAULT_AGENT_SOURCES).flatMap(([agent, source]) =>
source
? agentDiscoveries(
agent as AiVaultAgent,
source,
options,
wslHomeDirs,
limitPerAgent,
issues
)
: []
)
])
)
}
function agentDiscoveries(
@@ -5,6 +5,7 @@ import type {
AiVaultSessionPreviewMessage
} from '../../shared/ai-vault-types'
import type { ExecutionHostId } from '../../shared/execution-host'
import type { TranscriptMessageSink } from './session-transcript-consumers'
export type AiVaultScanOptions = {
claudeProjectsDir?: string
@@ -108,6 +109,9 @@ export type ResumableSessionParseState = {
export type SessionAccumulator = {
agent: AiVaultAgent
// Every decoded message this fold sees also goes here, for the reader's
// consumers. Shared by clones on purpose: one read, one message stream.
messages: TranscriptMessageSink
sessionId: string
title: string | null
fallbackTitle: string | null
+4 -41
View File
@@ -6,19 +6,13 @@ import type {
import { LOCAL_EXECUTION_HOST_ID, type ExecutionHostId } from '../../shared/execution-host'
import { withSpan } from '../observability/tracer'
import { sessionSortTime } from './session-scanner-accumulator'
import {
codexRolloutHardlinkIdentity,
dedupeCodexRolloutAliases,
dedupeCodexSessionsBySessionId
} from './codex-session-root-dedup'
import { readCodexRolloutSessionMetaId } from '../codex/codex-rollout-session-meta'
import { dedupeCodexSessionsBySessionId } from './codex-session-root-dedup'
import {
createAntigravityWorkspaceResolver,
readLocalAntigravityHistory,
type AntigravityWorkspaceResolver
} from './session-scanner-antigravity-history'
import { antigravityHistoryPathForBrainDir } from './session-scanner-antigravity-paths'
import { codexHomeForSessionsDir } from './session-scanner-codex-paths'
import { sessionCandidatesFromDiscoveries } from './session-scanner-candidates'
import {
ensureSessionParseCacheLoaded,
scheduleSessionParseCachePersist
@@ -30,10 +24,7 @@ import {
} from './session-scanner-parse-cache'
import { recordSessionScanIssue } from './session-scan-issues'
import { discoverInScopeClaudeFiles } from './session-scanner-scope-discovery'
import {
DEFAULT_CODEX_HOME_DIR,
discoverAiVaultSessionSources
} from './session-scanner-source-discovery'
import { discoverAiVaultSessionSources } from './session-scanner-source-discovery'
import type {
AiVaultScanOptions,
SessionFileCandidate,
@@ -83,35 +74,7 @@ export async function scanAiVaultSessions(
const discoveries = await discoverAiVaultSessionSources({ options, limitPerAgent, issues })
throwIfAiVaultScanCancelled(options.signal)
const candidates = await dedupeCodexRolloutAliases(
discoveries
.flatMap((discovery) =>
discovery.files.map((file): SessionFileCandidate => ({
agent: discovery.agent,
file,
codexHome:
discovery.agent === 'codex'
? codexHomeForSessionsDir(
discovery.rootDir,
options.defaultCodexHomeDir ?? DEFAULT_CODEX_HOME_DIR
)
: null,
antigravityHistoryPath:
discovery.agent === 'antigravity'
? antigravityHistoryPathForBrainDir(discovery.rootDir)
: undefined
}))
)
.sort((left, right) => right.file.mtimeMs - left.file.mtimeMs),
{
isCodex: (candidate) => candidate.agent === 'codex',
getFilePath: (candidate) => candidate.file.path,
getCodexHome: (candidate) => candidate.codexHome,
getHardlinkIdentity: (candidate) => codexRolloutHardlinkIdentity(candidate.file)
},
(filePath) => readCodexRolloutSessionMetaId(filePath, options.signal, 'scan'),
options.signal
)
const candidates = await sessionCandidatesFromDiscoveries(discoveries, options)
const parsedSessions = await parseSessionCandidates({
candidates: candidates.slice(0, limit * SESSION_PARSE_CANDIDATE_MULTIPLIER),
@@ -0,0 +1,92 @@
import {
hasTranscriptConsumers,
transcriptConsumers,
type TranscriptMessage,
type TranscriptMessageSink,
type TranscriptReadConsumer,
type TranscriptReadOutcome,
type TranscriptReadStart
} from './session-transcript-consumers'
/**
* The sink a parser pushes into, and the fan-out to every registered consumer.
*
* One channel belongs to one file for as long as its resumable parse state
* lives, because the cached state (and every clone of it) holds this reference.
* A read re-points the channel at that read's consumers instead of replacing it.
*/
export class TranscriptMessageChannel implements TranscriptMessageSink {
private readers: TranscriptReadConsumer[] = []
private muted = false
/** True while a read is open with at least one consumer attached. */
get active(): boolean {
return this.readers.length > 0
}
beginRead(start: TranscriptReadStart): void {
this.muted = false
this.readers = []
// Keeps a scan with no consumers allocation-free on its hottest path.
if (!hasTranscriptConsumers()) {
return
}
for (const consumer of transcriptConsumers()) {
try {
const reader = consumer.beginRead(start)
if (reader) {
this.readers.push(reader)
}
} catch {
// A consumer that cannot open this read simply does not see it.
}
}
}
push(message: TranscriptMessage): void {
if (this.muted || this.readers.length === 0) {
return
}
// A throwing consumer is dropped for the rest of the read rather than
// failing the parse; it then gets no `finish`, so it never records a cursor
// for a stream it did not see in full.
let index = 0
while (index < this.readers.length) {
try {
this.readers[index].message(message)
index++
} catch {
this.readers.splice(index, 1)
}
}
}
/**
* Suppresses emission for a display-only re-read: the trailing unterminated
* line is shown in the list but is re-read once complete, so emitting it here
* would hand every consumer the same line twice. `fn` must be synchronous.
*/
mute<T>(fn: () => T): T {
const previous = this.muted
this.muted = true
try {
return fn()
} finally {
this.muted = previous
}
}
finishRead(outcome: TranscriptReadOutcome): void {
const readers = this.readers
this.readers = []
this.muted = false
for (const reader of readers) {
try {
reader.finish(outcome)
} catch {
// A consumer failure must never fail the session list.
}
}
}
}
@@ -0,0 +1,212 @@
import { appendFile, mkdir, mkdtemp, rm, writeFile } from 'node:fs/promises'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { afterEach, expect, it } from 'vitest'
import { scanAiVaultSessions } from './session-scanner'
import { resetSessionParseCacheForTests } from './session-scanner-parse-cache'
import { isolatedScanRoots, jsonLines } from './session-scanner-test-fixtures'
import {
registerTranscriptConsumer,
resetTranscriptConsumersForTests,
type TranscriptMessage,
type TranscriptReadOutcome,
type TranscriptReadStart
} from './session-transcript-consumers'
type RecordedRead = {
start: TranscriptReadStart
messages: TranscriptMessage[]
outcome: TranscriptReadOutcome | null
}
function recordingConsumer(): { reads: RecordedRead[]; unregister: () => void } {
const reads: RecordedRead[] = []
const unregister = registerTranscriptConsumer({
beginRead: (start) => {
const read: RecordedRead = { start, messages: [], outcome: null }
reads.push(read)
return {
message: (message) => read.messages.push(message),
finish: (outcome) => {
read.outcome = outcome
}
}
}
})
return { reads, unregister }
}
function textsFor(reads: RecordedRead[], agent: string): string[] {
return reads
.filter((read) => read.start.candidate.agent === agent)
.flatMap((read) => read.messages.map((message) => `${message.role}:${message.text}`))
}
let tempRoots: string[] = []
afterEach(async () => {
resetTranscriptConsumersForTests()
resetSessionParseCacheForTests()
await Promise.all(tempRoots.map((root) => rm(root, { recursive: true, force: true })))
tempRoots = []
})
function claudeTurns(from: number, to: number): unknown[] {
const records: unknown[] = []
for (let index = from; index <= to; index++) {
records.push({
type: 'user',
sessionId: 'claude-session',
timestamp: `2026-05-01T10:0${index}:00.000Z`,
cwd: '/tmp/claude',
message: { role: 'user', content: `ask ${index}` }
})
records.push({
type: 'assistant',
sessionId: 'claude-session',
timestamp: `2026-05-01T10:0${index}:01.000Z`,
message: {
role: 'assistant',
content: [
{ type: 'text', text: `reply ${index}` },
{ type: 'tool_use', name: 'Bash', input: { command: `ls ${index}` } }
]
}
})
}
return records
}
async function writeClaudeFixture(): Promise<{
root: string
roots: ReturnType<typeof isolatedScanRoots>
transcript: string
}> {
const root = await mkdtemp(join(tmpdir(), 'orca-transcript-consumers-'))
tempRoots.push(root)
const roots = isolatedScanRoots(root)
const transcript = join(roots.claudeProjectsDir, 'project', 'claude-session.jsonl')
await mkdir(join(roots.claudeProjectsDir, 'project'), { recursive: true })
await writeFile(transcript, `${jsonLines(claudeTurns(1, 4))}\n`)
return { root, roots, transcript }
}
it('delivers one message stream to every registered consumer', async () => {
const { roots } = await writeClaudeFixture()
const first = recordingConsumer()
const second = recordingConsumer()
const result = await scanAiVaultSessions({ ...roots, platform: 'darwin', limit: 20 })
expect(result.issues).toEqual([])
const stream = textsFor(first.reads, 'claude')
expect(stream).toEqual(textsFor(second.reads, 'claude'))
expect(stream).toEqual([
'user:ask 1',
'assistant:reply 1',
'tool:Bash: ls 1',
'user:ask 2',
'assistant:reply 2',
'tool:Bash: ls 2',
'user:ask 3',
'assistant:reply 3',
'tool:Bash: ls 3',
'user:ask 4',
'assistant:reply 4',
'tool:Bash: ls 4'
])
// The list's own fold keeps only the newest five preview turns, so the stream
// is demonstrably the reader's, not a projection of the session row.
const session = result.sessions.find((entry) => entry.agent === 'claude')
expect(session?.previewMessages).toHaveLength(5)
expect(session?.messageCount).toBe(8)
})
it('leaves the session list identical whether or not a consumer is registered', async () => {
const withoutConsumer = await writeClaudeFixture()
const bare = await scanAiVaultSessions({
...withoutConsumer.roots,
platform: 'darwin',
limit: 20
})
resetSessionParseCacheForTests()
recordingConsumer()
const observed = await scanAiVaultSessions({
...withoutConsumer.roots,
platform: 'darwin',
limit: 20
})
expect(observed.sessions).toEqual(bare.sessions)
})
it('replays only the appended lines on a resumed read', async () => {
const { roots, transcript } = await writeClaudeFixture()
const consumer = recordingConsumer()
await scanAiVaultSessions({ ...roots, platform: 'darwin', limit: 20 })
const firstRead = consumer.reads.at(-1)
expect(firstRead?.start.mode).toBe('replace')
expect(firstRead?.start.previousByteOffset).toBe(0)
expect(firstRead?.outcome?.incomplete).toBe(false)
await appendFile(transcript, `${jsonLines(claudeTurns(5, 5))}\n`)
consumer.reads.length = 0
await scanAiVaultSessions({ ...roots, platform: 'darwin', limit: 20 })
const resumed = consumer.reads.find((read) => read.start.candidate.agent === 'claude')
expect(resumed?.start.mode).toBe('append')
expect(resumed?.start.previousByteOffset).toBe(firstRead?.outcome?.byteOffset)
expect(textsFor(consumer.reads, 'claude')).toEqual([
'user:ask 5',
'assistant:reply 5',
'tool:Bash: ls 5'
])
})
it('publishes a trailing unterminated line once, when it is complete', async () => {
const { roots, transcript } = await writeClaudeFixture()
const consumer = recordingConsumer()
await scanAiVaultSessions({ ...roots, platform: 'darwin', limit: 20 })
// A half-written record: the list shows it, the stream must not carry it yet.
const [partial] = claudeTurns(5, 5)
await appendFile(transcript, JSON.stringify(partial))
consumer.reads.length = 0
await scanAiVaultSessions({ ...roots, platform: 'darwin', limit: 20 })
expect(textsFor(consumer.reads, 'claude')).toEqual([])
await appendFile(transcript, '\n')
consumer.reads.length = 0
await scanAiVaultSessions({ ...roots, platform: 'darwin', limit: 20 })
expect(textsFor(consumer.reads, 'claude')).toEqual(['user:ask 5'])
})
it('keeps the session list working when a consumer throws', async () => {
const { roots } = await writeClaudeFixture()
registerTranscriptConsumer({
beginRead: () => ({
message: () => {
throw new Error('consumer exploded')
},
finish: () => undefined
})
})
const healthy = recordingConsumer()
const result = await scanAiVaultSessions({ ...roots, platform: 'darwin', limit: 20 })
expect(result.issues).toEqual([])
expect(result.sessions.find((entry) => entry.agent === 'claude')?.messageCount).toBe(8)
expect(textsFor(healthy.reads, 'claude')).toHaveLength(12)
})
it('skips a read a consumer declines without disturbing the others', async () => {
const { roots } = await writeClaudeFixture()
registerTranscriptConsumer({ beginRead: () => null })
const healthy = recordingConsumer()
await scanAiVaultSessions({ ...roots, platform: 'darwin', limit: 20 })
expect(textsFor(healthy.reads, 'claude')).toHaveLength(12)
})
@@ -0,0 +1,81 @@
import type { AiVaultSession } from '../../shared/ai-vault-types'
import type { SessionFileCandidate } from './session-scanner-types'
// Why: the transcript reader owns discovery, per-file cursors and decoding; a
// consumer only folds the message stream. Registering a second consumer (a
// search index, a digest) must not require touching the reader or the parse
// cache, so the reader publishes reads rather than knowing who reads them.
export type TranscriptMessageRole = 'user' | 'assistant' | 'tool'
export type TranscriptMessage = {
role: TranscriptMessageRole
/** Untruncated decoded text; caps and redaction are consumer policy. */
text: string
timestamp: string | null
}
/** Where a parser hands its decoded messages; the reader supplies the instance. */
export type TranscriptMessageSink = {
/** False when nobody is listening: parsers skip the extraction entirely. */
readonly active: boolean
push(message: TranscriptMessage): void
}
export const NO_TRANSCRIPT_MESSAGES: TranscriptMessageSink = {
active: false,
push: () => undefined
}
export type TranscriptReadStart = {
candidate: SessionFileCandidate
/** `replace`: the whole file is being re-read; `append`: a resumed read. */
mode: 'replace' | 'append'
/** Byte offset the messages of this read continue from. */
previousByteOffset: number
}
export type TranscriptReadOutcome = {
/** Null when the parser rejected the file (an excluded Codex worker transcript). */
session: AiVaultSession | null
/** Byte offset just past the last complete line this read consumed. */
byteOffset: number
/** The read did not cover the whole span; its messages are not the whole file. */
incomplete: boolean
}
/** One consumer's view of one file read. */
export type TranscriptReadConsumer = {
message(message: TranscriptMessage): void
finish(outcome: TranscriptReadOutcome): void
}
export type TranscriptConsumer = {
/**
* Open this read, or return null to ignore it. A consumer whose own cursor is
* behind `previousByteOffset` declines here and re-reads on its own schedule;
* it must never ask another consumer where it is.
*/
beginRead(start: TranscriptReadStart): TranscriptReadConsumer | null
}
const consumers = new Set<TranscriptConsumer>()
export function registerTranscriptConsumer(consumer: TranscriptConsumer): () => void {
consumers.add(consumer)
return () => {
consumers.delete(consumer)
}
}
export function transcriptConsumers(): readonly TranscriptConsumer[] {
return [...consumers]
}
export function hasTranscriptConsumers(): boolean {
return consumers.size > 0
}
export function resetTranscriptConsumersForTests(): void {
consumers.clear()
}
@@ -0,0 +1,61 @@
import { expect, it } from 'vitest'
import { transcriptMessagesFromContent } from './session-transcript-message-content'
const AT = '2026-05-01T10:00:00.000Z'
it('keeps a plain string turn under the record role', () => {
expect(transcriptMessagesFromContent('user', 'just words', AT)).toEqual([
{ role: 'user', text: 'just words', timestamp: AT }
])
})
it('drops turns whose role a consumer cannot use', () => {
expect(transcriptMessagesFromContent('system', 'boot', AT)).toEqual([])
expect(transcriptMessagesFromContent('unknown', 'noise', AT)).toEqual([])
})
it('joins text blocks and appends tool blocks as their own messages', () => {
expect(
transcriptMessagesFromContent(
'assistant',
[
{ type: 'text', text: 'first' },
{ type: 'tool_use', name: 'Bash', input: { command: 'ls -la', description: 'ignored' } },
{ type: 'thinking', text: 'second' },
{ type: 'image', source: {} }
],
AT
)
).toEqual([
{ role: 'assistant', text: 'first\nsecond', timestamp: AT },
{ role: 'tool', text: 'Bash: ls -la', timestamp: AT }
])
})
it('reads a tool result carried on a user record as a tool message', () => {
expect(
transcriptMessagesFromContent(
'user',
[{ type: 'tool_result', content: [{ type: 'text', text: 'exit 0' }] }],
AT
)
).toEqual([{ role: 'tool', text: 'exit 0', timestamp: AT }])
})
it('names a tool call even with no recognisable argument', () => {
expect(
transcriptMessagesFromContent('assistant', [{ type: 'tool_use', name: 'Read', input: {} }], AT)
).toEqual([{ role: 'tool', text: 'Read', timestamp: AT }])
})
it('emits nothing for blank or absent content', () => {
expect(transcriptMessagesFromContent('user', ' ', AT)).toEqual([])
expect(transcriptMessagesFromContent('user', null, AT)).toEqual([])
expect(transcriptMessagesFromContent('assistant', [{ type: 'tool_use' }], AT)).toEqual([])
})
it('does not apply the list preview cap', () => {
const long = 'x'.repeat(5000)
const [message] = transcriptMessagesFromContent('user', [{ type: 'text', text: long }], AT)
expect(message.text).toHaveLength(5000)
})
@@ -0,0 +1,138 @@
import { asRecord } from './session-scanner-record-value'
import { sliceAtCodeUnitLimit } from './session-scanner-text-normalization'
import type { AiVaultSessionPreviewMessage } from '../../shared/ai-vault-types'
import type { TranscriptMessage, TranscriptMessageRole } from './session-transcript-consumers'
// Safety bound only: a consumer applies its own caps. Matches the first-prompt
// copy path's ceiling so one pathological paste cannot dominate a scan.
const TRANSCRIPT_MESSAGE_TEXT_LIMIT = 256 * 1024
const TOOL_ARGUMENT_SCAN_LIMIT = 2000
const TEXT_BLOCK_TYPES = new Set(['text', 'input_text', 'output_text', 'thinking', 'reasoning'])
// The argument that identifies what a tool call actually did.
const TOOL_INPUT_KEYS = ['command', 'cmd', 'file_path', 'path', 'pattern', 'query', 'description']
type PreviewRole = AiVaultSessionPreviewMessage['role']
/** Only conversational roles reach consumers; system/unknown turns are noise. */
export function transcriptMessageRole(role: PreviewRole): TranscriptMessageRole | null {
return role === 'user' || role === 'assistant' || role === 'tool' ? role : null
}
export function toolCallText(name: unknown, input: unknown): string | null {
const toolName = typeof name === 'string' && name.trim() ? name.trim() : null
const inputRecord = asRecord(input)
let argument: string | null = null
if (inputRecord) {
for (const key of TOOL_INPUT_KEYS) {
const value = inputRecord[key]
if (typeof value === 'string' && value.trim()) {
argument = value
break
}
}
} else if (typeof input === 'string' && input.trim()) {
argument = input
}
if (!toolName && !argument) {
return null
}
const bounded = argument ? sliceAtCodeUnitLimit(argument, TOOL_ARGUMENT_SCAN_LIMIT) : null
return toolName && bounded ? `${toolName}: ${bounded}` : (toolName ?? bounded)
}
/** Flattens a tool_result body (a string, or an array of text blocks). */
function toolResultText(content: unknown): string | null {
if (typeof content === 'string') {
return content.trim() ? content : null
}
if (!Array.isArray(content)) {
return null
}
const parts: string[] = []
let length = 0
for (const item of content) {
const text = typeof item === 'string' ? item : asRecord(item)?.text
if (typeof text === 'string' && text) {
parts.push(text)
length += text.length
if (length >= TRANSCRIPT_MESSAGE_TEXT_LIMIT) {
break
}
}
}
const joined = parts.join('\n')
return joined.trim() ? joined : null
}
/**
* Splits one provider content value into the messages it decodes to. Text
* blocks keep the record's role; tool_use and tool_result blocks become `tool`
* messages whichever record carried them (Claude stores tool results on user
* records), so a consumer never has to know a provider's record shapes.
*/
export function transcriptMessagesFromContent(
role: PreviewRole,
content: unknown,
timestamp: string | null
): TranscriptMessage[] {
const messages: TranscriptMessage[] = []
const textRole = transcriptMessageRole(role)
if (typeof content === 'string') {
const text = boundedText(content)
return text && textRole ? [{ role: textRole, text, timestamp }] : []
}
const blocks = Array.isArray(content) ? content : content != null ? [content] : []
const textParts: string[] = []
for (const block of blocks) {
if (typeof block === 'string') {
textParts.push(block)
continue
}
const item = asRecord(block)
if (!item) {
continue
}
const type = typeof item.type === 'string' ? item.type : null
if (type === 'tool_use') {
pushMessage(messages, 'tool', toolCallText(item.name, item.input), timestamp)
continue
}
if (type === 'tool_result') {
pushMessage(messages, 'tool', toolResultText(item.content), timestamp)
continue
}
if (type !== null && !TEXT_BLOCK_TYPES.has(type)) {
continue
}
const text = typeof item.text === 'string' ? item.text : item.content
if (typeof text === 'string' && text) {
textParts.push(text)
}
}
if (textRole && textParts.length > 0) {
// The record's own words lead; its tool blocks follow in transcript order.
const text = boundedText(textParts.join('\n'))
if (text) {
messages.unshift({ role: textRole, text, timestamp })
}
}
return messages
}
function pushMessage(
messages: TranscriptMessage[],
role: TranscriptMessageRole,
text: string | null,
timestamp: string | null
): void {
const bounded = text === null ? null : boundedText(text)
if (bounded) {
messages.push({ role, text: bounded, timestamp })
}
}
export function boundedText(value: string): string | null {
const bounded = sliceAtCodeUnitLimit(value, TRANSCRIPT_MESSAGE_TEXT_LIMIT)
return bounded.trim() ? bounded : null
}
@@ -0,0 +1,149 @@
import { readTranscriptSlice } from '../native-chat/wsl-transcript-fs-access'
import type { AiVaultSession } from '../../shared/ai-vault-types'
import { parseAgentSessionFile } from './session-scanner-agent-parser'
import { consumeCompleteJsonlLines } from './session-scanner-jsonl-reader'
import type { ResumableSessionParseState, SessionFileCandidate } from './session-scanner-types'
import type { SessionParseResumePoint } from './session-parse-cache-store'
import { TranscriptMessageChannel } from './session-transcript-channel'
const NEWLINE_BYTE = 0x0a
// Why: this layer owns reading a transcript and nothing else. It decides where
// a read starts, drives the parser, publishes the decoded messages to every
// registered consumer, and reports where the read ended. Which of those results
// are cached, listed or indexed belongs to the callers.
export type TranscriptReadStats = {
incremental: number
fullParses: number
// Transcripts the parser already excluded (Codex workers), re-listed after a
// write and dismissed without reading. Counted apart from `incremental` so a
// scan span still shows how much work the early stop actually removed.
earlyStopped: number
bytesRead: number
}
export type ResumableTranscriptRead = {
session: AiVaultSession | null
/** The fold to resume from next time, and the channel bound to it. */
resume: SessionParseResumePoint
}
/**
* Read an append-only transcript, resuming from `resume` when the file only
* grew and the recorded offset still sits on a line boundary. Anything else
* (a rewrite, a truncation, a platform change) re-reads the whole file.
*/
export async function readResumableTranscript(args: {
candidate: SessionFileCandidate
platform: NodeJS.Platform
resume: SessionParseResumePoint | null
stateFactory: (messages: TranscriptMessageChannel) => ResumableSessionParseState
stats?: TranscriptReadStats
}): Promise<ResumableTranscriptRead> {
const { file } = args.candidate
const resume = args.resume
const canResume =
resume !== null &&
typeof file.sizeBytes === 'number' &&
file.sizeBytes >= resume.byteOffset &&
(resume.byteOffset === 0 || (await endsWithNewlineAt(file.path, resume.byteOffset)))
// Clone before consuming: a failed read must not corrupt the cached state,
// or the next resume would double-count the lines applied before the error.
const channel = canResume ? resume.channel : new TranscriptMessageChannel()
const state = canResume ? resume.state.clone() : args.stateFactory(channel)
const startOffset = canResume ? resume.byteOffset : 0
// Mirrors the reader's entry guard so a dismissed transcript is not reported
// as an incremental parse that read nothing.
const stoppedBeforeRead = state.shouldStop?.() === true
if (args.stats) {
if (stoppedBeforeRead) {
args.stats.earlyStopped++
} else if (canResume) {
args.stats.incremental++
} else {
args.stats.fullParses++
}
}
channel.beginRead({
candidate: args.candidate,
mode: canResume ? 'append' : 'replace',
previousByteOffset: startOffset
})
try {
const readResult = await consumeCompleteJsonlLines({
path: file.path,
start: startOffset,
onLine: (line) => state.consumeLine(line),
// Bound: the optional hooks are declared as methods, so a parser written
// with method syntax must not lose `this` on the way into the reader.
onLineBytes: state.consumeLineBytes?.bind(state),
shouldStop: state.shouldStop?.bind(state)
})
if (args.stats) {
args.stats.bytesRead += readResult.bytesRead
}
// The stat this scan displays is current even when nothing new was consumed.
state.touchFile(file)
// Keep parity with the one-shot parser: a final unterminated line is shown,
// but stays out of the resumable state so the (possibly still-growing) line
// is re-read once complete instead of being half-counted.
let displayState = state
if (readResult.trailingPartialLine !== null) {
const partialLine = readResult.trailingPartialLine
displayState = state.clone()
channel.mute(() => displayState.consumeLine(partialLine))
}
const session = await displayState.finalize(args.platform)
channel.finishRead({ session, byteOffset: readResult.consumedThrough, incomplete: false })
return {
session,
resume: { state, byteOffset: readResult.consumedThrough, channel }
}
} catch (error) {
channel.finishRead({ session: null, byteOffset: startOffset, incomplete: true })
throw error
}
}
/**
* Read a transcript whose format is rewritten in place rather than appended
* (whole-JSON documents, Kimi's state doc, OpenCode). There is no cursor to
* keep, so every read is a whole-file `replace`.
*/
export async function readWholeTranscript(args: {
candidate: SessionFileCandidate
platform: NodeJS.Platform
stats?: TranscriptReadStats
}): Promise<AiVaultSession | null> {
const { file } = args.candidate
if (args.stats) {
args.stats.fullParses++
args.stats.bytesRead += file.sizeBytes ?? 0
}
const channel = new TranscriptMessageChannel()
channel.beginRead({ candidate: args.candidate, mode: 'replace', previousByteOffset: 0 })
try {
const session = await parseAgentSessionFile(args.candidate, args.platform, channel)
channel.finishRead({ session, byteOffset: file.sizeBytes ?? 0, incomplete: false })
return session
} catch (error) {
channel.finishRead({ session: null, byteOffset: 0, incomplete: true })
throw error
}
}
// A resume point is only valid if it still sits just past a line break;
// anything else means the file was rewritten, not appended. Heuristic: a
// grown rewrite keeping '\n' at exactly this byte would slip through, but
// agent transcripts are append-only so that trade is accepted (worst case is
// a stale vault row until the file is next truncated or the app restarts).
async function endsWithNewlineAt(path: string, offset: number): Promise<boolean> {
const slice = await readTranscriptSlice(path, offset - 1, 1, 'scan')
return slice.length === 1 && slice[0] === NEWLINE_BYTE
}
+102
View File
@@ -0,0 +1,102 @@
import type { Worker } from 'node:worker_threads'
export type WorkerThreadFactory = () => Worker
/**
* Owns the lifetime of one lazily-spawned worker thread: spawn on demand,
* listener wiring, teardown, and idle expiry. It holds no request state, so
* every decision about which call a message belongs to stays with its client,
* and a failed spawn is reported rather than thrown so the client can fail its
* queued calls closed instead of moving the work back onto the main thread.
*/
export class LazyWorkerThreadHost<TResponse> {
private worker: Worker | null = null
private idleTimer: NodeJS.Timeout | null = null
private cleanupListeners: (() => void) | null = null
private reportedUnavailable = false
constructor(
private readonly options: {
factory: WorkerThreadFactory
idleTeardownMs: number
onMessage: (response: TResponse) => void
onError: (error: Error) => void
onExit: (code: number) => void
/** Nothing active and nothing queued, checked again when the idle timer fires. */
isIdle: () => boolean
/** First spawn failure only: a repeating one must not repeat the log. */
onUnavailable: (error: unknown) => void
}
) {}
get current(): Worker | null {
return this.worker
}
/** The live worker, spawning one if needed; null when no worker can be had. */
ensure(): Worker | null {
if (this.worker) {
return this.worker
}
try {
const worker = this.options.factory()
const onMessage = (response: TResponse): void => this.options.onMessage(response)
const onError = (error: Error): void => this.options.onError(error)
const onExit = (code: number): void => this.options.onExit(code)
worker.on('message', onMessage)
worker.on('error', onError)
worker.on('exit', onExit)
this.cleanupListeners = () => {
worker.off('message', onMessage)
worker.off('error', onError)
worker.off('exit', onExit)
}
// Never keep the app alive for background work.
worker.unref?.()
this.worker = worker
return worker
} catch (err) {
if (!this.reportedUnavailable) {
this.reportedUnavailable = true
this.options.onUnavailable(err)
}
return null
}
}
destroy(): void {
this.clearIdleTimer()
const worker = this.worker
this.worker = null
if (!worker) {
return
}
this.cleanupListeners?.()
this.cleanupListeners = null
worker.removeAllListeners()
void worker.terminate().catch(() => undefined)
}
scheduleIdleTeardown(): void {
this.clearIdleTimer()
if (!this.worker) {
return
}
this.idleTimer = setTimeout(() => {
this.idleTimer = null
// Re-checked here: a request arriving as the timer fires must never be
// lost to a self-exiting worker.
if (this.options.isIdle()) {
this.destroy()
}
}, this.options.idleTeardownMs)
this.idleTimer.unref?.()
}
clearIdleTimer(): void {
if (this.idleTimer) {
clearTimeout(this.idleTimer)
this.idleTimer = null
}
}
}
+28 -93
View File
@@ -2,6 +2,7 @@ import { existsSync } from 'node:fs'
import { getAppEnvironment, hasAppEnvironment } from '../../shared/app-environment'
import { join } from 'node:path'
import { Worker } from 'node:worker_threads'
import { LazyWorkerThreadHost, type WorkerThreadFactory } from '../lazy-worker-thread-host'
import {
PORT_SCAN_COMMAND_TIMEOUT_MS,
PortScanCommandTimeoutError,
@@ -11,11 +12,10 @@ import {
// Why (#11161): a lazily-spawned, unref'd worker runs the port scan's probe
// spawns off the Electron main-process event loop, because libuv performs
// process creation inline on the calling thread. Lifecycle (FIFO one-at-a-time
// dispatch, per-call deadlines, respawn-on-fault, idle teardown, fail-closed)
// mirrors src/main/ai-vault/session-scanner-opencode-sqlite-worker-client.ts;
// the duplicated ~150 lines are cheaper than a premature shared abstraction, so
// a third adopter should extract one.
// process creation inline on the calling thread. This module owns the request
// half (FIFO one-at-a-time dispatch, per-call deadlines, respawn-on-fault); the
// thread's own lifetime belongs to LazyWorkerThreadHost, shared with
// src/main/ai-vault/session-scanner-opencode-sqlite-worker-client.ts.
//
// This module used to contain the literal text require('electron'), which fails the
// plain-Node entry guard even inside a try/catch. It reads the AppEnvironment port
@@ -36,7 +36,7 @@ export const MAX_CONSECUTIVE_DEATHS = 3
export const MAX_QUEUED_CALLS = 8
export type PortScanCommandResult = { stdout: string; spawnMs: number }
export type PortScanWorkerFactory = () => Worker
export type PortScanWorkerFactory = WorkerThreadFactory
// Distinguishes "no worker at all" from a timeout or crash so the scanner can
// log it once and callers never mistake it for a command timeout.
@@ -62,20 +62,27 @@ type PendingCall = {
* be spawned rather than moving process creation back onto the main thread.
*/
export class PortScanCommandClient {
private worker: Worker | null = null
private active: PendingCall | null = null
private queue: PendingCall[] = []
private idleTimer: NodeJS.Timeout | null = null
private consecutiveDeaths = 0
private nextId = 1
private loggedWorkerUnavailable = false
private cleanupWorkerListeners: (() => void) | null = null
private readonly workerFactory: PortScanWorkerFactory
private readonly log: (message: string) => void
private readonly host: LazyWorkerThreadHost<PortScanCommandResponse>
constructor(options: { workerFactory: PortScanWorkerFactory; log?: (message: string) => void }) {
this.workerFactory = options.workerFactory
this.log = options.log ?? ((message) => console.warn(message))
const log = options.log ?? ((message: string) => console.warn(message))
this.host = new LazyWorkerThreadHost<PortScanCommandResponse>({
factory: options.workerFactory,
idleTeardownMs: IDLE_TEARDOWN_MS,
onMessage: (response) => this.onMessage(response),
onError: (error) => this.onWorkerFault(error),
onExit: (code) => this.onWorkerExit(code),
isIdle: () => !this.active && this.queue.length === 0,
// Why (#11161): never fall back to in-process execFile here; a missing
// bundle must report port scanning as unavailable rather than reintroduce
// the main-thread freeze this worker boundary exists to prevent.
onUnavailable: (err) =>
log(`[workspace-ports] probe worker unavailable. ${errorMessage(err)}`)
})
}
/**
@@ -109,7 +116,7 @@ export class PortScanCommandClient {
if (this.active || this.queue.length === 0) {
return
}
const worker = this.ensureWorker()
const worker = this.host.ensure()
if (!worker) {
this.failQueuedAsUnavailable()
return
@@ -119,7 +126,7 @@ export class PortScanCommandClient {
return
}
this.active = call
this.clearIdleTimer()
this.host.clearIdleTimer()
// Why (#11161): one at a time. uv_spawn blocks the worker's own loop, so a
// second concurrent request would have its deadline armed while the first
// spawn is still stalling the thread, producing a false timeout.
@@ -128,39 +135,6 @@ export class PortScanCommandClient {
worker.postMessage(call.request)
}
private ensureWorker(): Worker | null {
if (this.worker) {
return this.worker
}
try {
const worker = this.workerFactory()
const onMessage = (response: PortScanCommandResponse): void => this.onMessage(response)
const onError = (error: Error): void => this.onWorkerFault(error)
const onExit = (code: number): void => this.onWorkerExit(code)
worker.on('message', onMessage)
worker.on('error', onError)
worker.on('exit', onExit)
this.cleanupWorkerListeners = () => {
worker.off('message', onMessage)
worker.off('error', onError)
worker.off('exit', onExit)
}
// Never keep the app alive for a port scan.
worker.unref?.()
this.worker = worker
return worker
} catch (err) {
// Why (#11161): never fall back to in-process execFile here; a missing
// bundle must report port scanning as unavailable rather than reintroduce
// the main-thread freeze this worker boundary exists to prevent.
if (!this.loggedWorkerUnavailable) {
this.loggedWorkerUnavailable = true
this.log(`[workspace-ports] probe worker unavailable. ${errorMessage(err)}`)
}
return null
}
}
private onMessage(response: PortScanCommandResponse): void {
const call = this.active
if (!call || call.request.id !== response.id) {
@@ -191,7 +165,7 @@ export class PortScanCommandClient {
// A clean self-exit is not a death, but the stale handle must be dropped or
// the next dispatch would post into a dead worker and stall to its deadline.
if (code === 0 && !this.active && this.queue.length === 0) {
this.destroyWorker()
this.host.destroy()
return
}
this.onWorkerFault(new Error(`Port scan probe worker exited with code ${code}`))
@@ -199,7 +173,7 @@ export class PortScanCommandClient {
private onWorkerFault(error: Error): void {
const failed = this.active
this.destroyWorker()
this.host.destroy()
this.consecutiveDeaths++
if (failed) {
this.settle(failed, () => failed.reject(error))
@@ -248,50 +222,11 @@ export class PortScanCommandClient {
if (this.queue.length > 0) {
this.pump()
} else {
this.scheduleIdleTeardown()
// Terminating can orphan a probe child mid-spawn; the worker reaps what it
// can on exit, and every probe here is short-lived.
this.host.scheduleIdleTeardown()
}
}
private scheduleIdleTeardown(): void {
this.clearIdleTimer()
if (!this.worker) {
return
}
this.idleTimer = setTimeout(() => this.teardownIfIdle(), IDLE_TEARDOWN_MS)
this.idleTimer.unref?.()
}
private teardownIfIdle(): void {
this.idleTimer = null
// Only tear down with nothing active AND nothing queued: a request arriving
// as the timer fires must never be lost to a self-exiting worker.
if (this.active || this.queue.length > 0) {
return
}
this.destroyWorker()
}
private clearIdleTimer(): void {
if (this.idleTimer) {
clearTimeout(this.idleTimer)
this.idleTimer = null
}
}
private destroyWorker(): void {
this.clearIdleTimer()
const worker = this.worker
this.worker = null
if (!worker) {
return
}
this.cleanupWorkerListeners?.()
this.cleanupWorkerListeners = null
worker.removeAllListeners()
// Terminating can orphan a probe child mid-spawn; the worker reaps what it
// can on exit, and every probe here is short-lived.
void worker.terminate().catch(() => undefined)
}
}
function errorMessage(error: unknown): string {