From 0e360f0a48c813440fc7c706a3d3880fccaafae9 Mon Sep 17 00:00:00 2001 From: weishu Date: Mon, 5 Jan 2026 15:53:30 +0800 Subject: [PATCH] refactor: extract session scanner base class for code reuse Create BaseSessionScanner abstract class to consolidate common session scanning logic shared between Claude and Codex session scanners. Both scanners now inherit from this base, reducing duplication and providing consistent patterns for file watching, event processing, and state management. --- cli/src/claude/utils/sessionScanner.ts | 282 ++++++------ cli/src/codex/utils/codexSessionScanner.ts | 411 +++++++++++------- .../common/session/BaseSessionScanner.ts | 189 ++++++++ 3 files changed, 581 insertions(+), 301 deletions(-) create mode 100644 cli/src/modules/common/session/BaseSessionScanner.ts diff --git a/cli/src/claude/utils/sessionScanner.ts b/cli/src/claude/utils/sessionScanner.ts index cae2634c..d97dd223 100644 --- a/cli/src/claude/utils/sessionScanner.ts +++ b/cli/src/claude/utils/sessionScanner.ts @@ -1,10 +1,9 @@ -import { InvalidateSync } from "@/utils/sync"; import { RawJSONLines, RawJSONLinesSchema } from "../types"; -import { join } from "node:path"; +import { basename, join } from "node:path"; import { readFile } from "node:fs/promises"; import { logger } from "@/ui/logger"; -import { startFileWatcher } from "@/modules/watcher/startFileWatcher"; import { getProjectPath } from "./path"; +import { BaseSessionScanner, SessionFileScanEntry, SessionFileScanResult, SessionFileScanStats } from "@/modules/common/session/BaseSessionScanner"; /** * Known internal Claude Code event types that should be silently skipped. @@ -18,133 +17,139 @@ const INTERNAL_CLAUDE_EVENT_TYPES = new Set([ ]); export async function createSessionScanner(opts: { - sessionId: string | null, - workingDirectory: string - onMessage: (message: RawJSONLines) => void + sessionId: string | null; + workingDirectory: string; + onMessage: (message: RawJSONLines) => void; }) { - - // Resolve project directory - const projectDir = getProjectPath(opts.workingDirectory); - - // Finished, pending finishing and current session - let finishedSessions = new Set(); - let pendingSessions = new Set(); - let currentSessionId: string | null = null; - let watchers = new Map void)>(); - let processedMessageKeys = new Set(); - - // Mark existing messages as processed and start watching the initial session - if (opts.sessionId) { - let messages = await readSessionLog(projectDir, opts.sessionId); - logger.debug(`[SESSION_SCANNER] Marking ${messages.length} existing messages as processed from session ${opts.sessionId}`); - for (let m of messages) { - processedMessageKeys.add(messageKey(m)); - } - // IMPORTANT: Also start watching the initial session file because Claude Code - // may continue writing to it even after creating a new session with --resume - // (agent tasks and other updates can still write to the original session file) - currentSessionId = opts.sessionId; - } - - // Main sync function - const sync = new InvalidateSync(async () => { - // logger.debug(`[SESSION_SCANNER] Syncing...`); - - // Collect session ids - include ALL sessions that have watchers - // This ensures we continue processing sessions that Claude Code may still write to - let sessions: string[] = []; - for (let p of pendingSessions) { - sessions.push(p); - } - if (currentSessionId && !pendingSessions.has(currentSessionId)) { - sessions.push(currentSessionId); - } - // Also process sessions that have active watchers (they may still receive updates) - for (let [sessionId] of watchers) { - if (!sessions.includes(sessionId)) { - sessions.push(sessionId); - } - } - - // Process sessions - for (let session of sessions) { - const sessionMessages = await readSessionLog(projectDir, session); - let skipped = 0; - let sent = 0; - for (let file of sessionMessages) { - let key = messageKey(file); - if (processedMessageKeys.has(key)) { - skipped++; - continue; - } - processedMessageKeys.add(key); - logger.debug(`[SESSION_SCANNER] Sending new message: type=${file.type}, uuid=${file.type === 'summary' ? file.leafUuid : file.uuid}`); - opts.onMessage(file); - sent++; - } - if (sessionMessages.length > 0) { - logger.debug(`[SESSION_SCANNER] Session ${session}: found=${sessionMessages.length}, skipped=${skipped}, sent=${sent}`); - } - } - - // Move pending sessions to finished sessions (but keep processing them via watchers) - for (let p of sessions) { - if (pendingSessions.has(p)) { - pendingSessions.delete(p); - finishedSessions.add(p); - } - } - - // Update watchers for all sessions - for (let p of sessions) { - if (!watchers.has(p)) { - logger.debug(`[SESSION_SCANNER] Starting watcher for session: ${p}`); - watchers.set(p, startFileWatcher(join(projectDir, `${p}.jsonl`), () => { sync.invalidate(); })); - } - } + const scanner = new ClaudeSessionScanner({ + sessionId: opts.sessionId, + workingDirectory: opts.workingDirectory, + onMessage: opts.onMessage }); - await sync.invalidateAndAwait(); - // Periodic sync - const intervalId = setInterval(() => { sync.invalidate(); }, 3000); + await scanner.start(); - // Public interface return { cleanup: async () => { - clearInterval(intervalId); - for (let w of watchers.values()) { - w(); - } - watchers.clear(); - await sync.invalidateAndAwait(); - sync.stop(); + await scanner.cleanup(); }, onNewSession: (sessionId: string) => { - if (currentSessionId === sessionId) { - logger.debug(`[SESSION_SCANNER] New session: ${sessionId} is the same as the current session, skipping`); - return; - } - if (finishedSessions.has(sessionId)) { - logger.debug(`[SESSION_SCANNER] New session: ${sessionId} is already finished, skipping`); - return; - } - if (pendingSessions.has(sessionId)) { - logger.debug(`[SESSION_SCANNER] New session: ${sessionId} is already pending, skipping`); - return; - } - if (currentSessionId) { - pendingSessions.add(currentSessionId); - } - logger.debug(`[SESSION_SCANNER] New session: ${sessionId}`) - currentSessionId = sessionId; - sync.invalidate(); - }, - } + scanner.onNewSession(sessionId); + } + }; } export type SessionScanner = ReturnType; +class ClaudeSessionScanner extends BaseSessionScanner { + private readonly projectDir: string; + private readonly onMessage: (message: RawJSONLines) => void; + private readonly finishedSessions = new Set(); + private readonly pendingSessions = new Set(); + private currentSessionId: string | null; + private readonly scannedSessions = new Set(); + + constructor(opts: { sessionId: string | null; workingDirectory: string; onMessage: (message: RawJSONLines) => void }) { + super({ intervalMs: 3000 }); + this.projectDir = getProjectPath(opts.workingDirectory); + this.onMessage = opts.onMessage; + this.currentSessionId = opts.sessionId; + } + + public onNewSession(sessionId: string): void { + if (this.currentSessionId === sessionId) { + logger.debug(`[SESSION_SCANNER] New session: ${sessionId} is the same as the current session, skipping`); + return; + } + if (this.finishedSessions.has(sessionId)) { + logger.debug(`[SESSION_SCANNER] New session: ${sessionId} is already finished, skipping`); + return; + } + if (this.pendingSessions.has(sessionId)) { + logger.debug(`[SESSION_SCANNER] New session: ${sessionId} is already pending, skipping`); + return; + } + if (this.currentSessionId) { + this.pendingSessions.add(this.currentSessionId); + } + logger.debug(`[SESSION_SCANNER] New session: ${sessionId}`); + this.currentSessionId = sessionId; + this.invalidate(); + } + + protected async initialize(): Promise { + if (!this.currentSessionId) { + return; + } + const sessionFile = this.sessionFilePath(this.currentSessionId); + const { events, totalLines } = await readSessionLog(sessionFile, 0); + logger.debug(`[SESSION_SCANNER] Marking ${events.length} existing messages as processed from session ${this.currentSessionId}`); + const keys = events.map((entry) => messageKey(entry.event)); + this.seedProcessedKeys(keys); + this.setCursor(sessionFile, totalLines); + } + + protected async beforeScan(): Promise { + this.scannedSessions.clear(); + } + + protected async findSessionFiles(): Promise { + const files = new Set(); + for (const sessionId of this.pendingSessions) { + files.add(this.sessionFilePath(sessionId)); + } + if (this.currentSessionId && !this.pendingSessions.has(this.currentSessionId)) { + files.add(this.sessionFilePath(this.currentSessionId)); + } + for (const watched of this.getWatchedFiles()) { + files.add(watched); + } + return [...files]; + } + + protected async parseSessionFile(filePath: string, cursor: number): Promise> { + const sessionId = sessionIdFromPath(filePath); + if (sessionId) { + this.scannedSessions.add(sessionId); + } + const { events, totalLines } = await readSessionLog(filePath, cursor); + return { + events, + nextCursor: totalLines + }; + } + + protected generateEventKey(event: RawJSONLines): string { + return messageKey(event); + } + + protected async handleFileScan(stats: SessionFileScanStats): Promise { + for (const message of stats.events) { + const id = message.type === 'summary' ? message.leafUuid : message.uuid; + logger.debug(`[SESSION_SCANNER] Sending new message: type=${message.type}, uuid=${id}`); + this.onMessage(message); + } + if (stats.parsedCount > 0) { + const sessionId = sessionIdFromPath(stats.filePath) ?? 'unknown'; + logger.debug(`[SESSION_SCANNER] Session ${sessionId}: found=${stats.parsedCount}, skipped=${stats.skippedCount}, sent=${stats.newCount}`); + } + } + + protected async afterScan(): Promise { + for (const sessionId of this.scannedSessions) { + if (this.pendingSessions.has(sessionId)) { + this.pendingSessions.delete(sessionId); + this.finishedSessions.add(sessionId); + } + } + } + + private sessionFilePath(sessionId: string): string { + return join(this.projectDir, `${sessionId}.jsonl`); + } +} + // // Helpers // @@ -164,22 +169,28 @@ function messageKey(message: RawJSONLines): string { } /** - * Read and parse session log file - * Returns only valid conversation messages, silently skipping internal events + * Read and parse session log file. + * Returns only valid conversation messages, silently skipping internal events. */ -async function readSessionLog(projectDir: string, sessionId: string): Promise { - const expectedSessionFile = join(projectDir, `${sessionId}.jsonl`); - logger.debug(`[SESSION_SCANNER] Reading session file: ${expectedSessionFile}`); +async function readSessionLog(filePath: string, startLine: number): Promise<{ events: SessionFileScanEntry[]; totalLines: number }> { + logger.debug(`[SESSION_SCANNER] Reading session file: ${filePath}`); let file: string; try { - file = await readFile(expectedSessionFile, 'utf-8'); + file = await readFile(filePath, 'utf-8'); } catch (error) { - logger.debug(`[SESSION_SCANNER] Session file not found: ${expectedSessionFile}`); - return []; + logger.debug(`[SESSION_SCANNER] Session file not found: ${filePath}`); + return { events: [], totalLines: startLine }; } - let lines = file.split('\n'); - let messages: RawJSONLines[] = []; - for (let l of lines) { + const lines = file.split('\n'); + const hasTrailingEmpty = lines.length > 0 && lines[lines.length - 1] === ''; + const totalLines = hasTrailingEmpty ? lines.length - 1 : lines.length; + let effectiveStartLine = startLine; + if (effectiveStartLine > totalLines) { + effectiveStartLine = 0; + } + const messages: SessionFileScanEntry[] = []; + for (let index = effectiveStartLine; index < lines.length; index += 1) { + const l = lines[index]; try { if (l.trim() === '') { continue; @@ -194,15 +205,22 @@ async function readSessionLog(projectDir: string, sessionId: string): Promise { - const codexHomeDir = process.env.CODEX_HOME || join(homedir(), '.codex'); - const sessionsRoot = join(codexHomeDir, 'sessions'); - - const processedLineCounts = new Map(); - const watchers = new Map void>(); - const sessionIdByFile = new Map(); - const sessionCwdByFile = new Map(); - const sessionTimestampByFile = new Map(); - const pendingEventsByFile = new Map(); - const sessionMetaParsed = new Set(); - - let activeSessionId: string | null = opts.sessionId; - let reportedSessionId: string | null = opts.sessionId; - let isClosing = false; - let matchFailed = false; - const targetCwd = opts.cwd && opts.cwd.trim().length > 0 ? normalizePath(opts.cwd) : null; - const referenceTimestampMs = opts.startupTimestampMs ?? Date.now(); - const sessionStartWindowMs = opts.sessionStartWindowMs ?? DEFAULT_SESSION_START_WINDOW_MS; - const matchDeadlineMs = referenceTimestampMs + sessionStartWindowMs; - const sessionDatePrefixes = targetCwd - ? getSessionDatePrefixes(referenceTimestampMs, sessionStartWindowMs) - : null; - logger.debug(`[CODEX_SESSION_SCANNER] Init: targetCwd=${targetCwd ?? 'none'} startupTs=${new Date(referenceTimestampMs).toISOString()} windowMs=${sessionStartWindowMs}`); if (!targetCwd && !opts.sessionId) { - matchFailed = true; const message = 'No cwd provided for Codex session matching; refusing to fallback.'; logger.warn(`[CODEX_SESSION_SCANNER] ${message}`); opts.onSessionMatchFailed?.(message); @@ -70,35 +45,212 @@ export async function createCodexSessionScanner(opts: CodexSessionScannerOptions }; } - function reportSessionId(sessionId: string): void { - if (reportedSessionId === sessionId) { + const scanner = new CodexSessionScannerImpl(opts, targetCwd); + await scanner.start(); + + return { + cleanup: async () => { + await scanner.cleanup(); + }, + onNewSession: (sessionId: string) => { + scanner.onNewSession(sessionId); + } + }; +} + +class CodexSessionScannerImpl extends BaseSessionScanner { + private readonly sessionsRoot: string; + private readonly onEvent: (event: CodexSessionEvent) => void; + private readonly onSessionFound?: (sessionId: string) => void; + private readonly onSessionMatchFailed?: (message: string) => void; + private readonly sessionIdByFile = new Map(); + private readonly sessionCwdByFile = new Map(); + private readonly sessionTimestampByFile = new Map(); + private readonly pendingEventsByFile = new Map(); + private readonly sessionMetaParsed = new Set(); + private readonly fileEpochByPath = new Map(); + private readonly targetCwd: string | null; + private readonly referenceTimestampMs: number; + private readonly sessionStartWindowMs: number; + private readonly matchDeadlineMs: number; + private readonly sessionDatePrefixes: Set | null; + + private activeSessionId: string | null; + private reportedSessionId: string | null; + private matchFailed = false; + private bestWithinWindow: Candidate | null = null; + + constructor(opts: CodexSessionScannerOptions, targetCwd: string | null) { + super({ intervalMs: 2000 }); + const codexHomeDir = process.env.CODEX_HOME || join(homedir(), '.codex'); + this.sessionsRoot = join(codexHomeDir, 'sessions'); + this.onEvent = opts.onEvent; + this.onSessionFound = opts.onSessionFound; + this.onSessionMatchFailed = opts.onSessionMatchFailed; + this.activeSessionId = opts.sessionId; + this.reportedSessionId = opts.sessionId; + this.targetCwd = targetCwd; + this.referenceTimestampMs = opts.startupTimestampMs ?? Date.now(); + this.sessionStartWindowMs = opts.sessionStartWindowMs ?? DEFAULT_SESSION_START_WINDOW_MS; + this.matchDeadlineMs = this.referenceTimestampMs + this.sessionStartWindowMs; + this.sessionDatePrefixes = this.targetCwd + ? getSessionDatePrefixes(this.referenceTimestampMs, this.sessionStartWindowMs) + : null; + + logger.debug(`[CODEX_SESSION_SCANNER] Init: targetCwd=${this.targetCwd ?? 'none'} startupTs=${new Date(this.referenceTimestampMs).toISOString()} windowMs=${this.sessionStartWindowMs}`); + } + + public onNewSession(sessionId: string): void { + if (this.activeSessionId === sessionId) { return; } - reportedSessionId = sessionId; - opts.onSessionFound?.(sessionId); + logger.debug(`[CODEX_SESSION_SCANNER] Switching to new session: ${sessionId}`); + this.setActiveSessionId(sessionId); + this.invalidate(); } - function setActiveSessionId(sessionId: string): void { - activeSessionId = sessionId; - reportSessionId(sessionId); - if (targetCwd) { - flushPendingEventsForSession(sessionId); - } else { - pendingEventsByFile.clear(); + protected shouldScan(): boolean { + return !this.matchFailed; + } + + protected shouldWatchFile(filePath: string): boolean { + if (!this.activeSessionId) { + if (!this.targetCwd) { + return false; + } + return this.getCandidateForFile(filePath) !== null; + } + const fileSessionId = this.sessionIdByFile.get(filePath); + if (fileSessionId) { + return fileSessionId === this.activeSessionId; + } + return filePath.endsWith(`-${this.activeSessionId}.jsonl`); + } + + protected async initialize(): Promise { + const files = await this.listSessionFiles(this.sessionsRoot); + for (const filePath of files) { + const { nextCursor } = await this.readSessionFile(filePath, 0); + this.setCursor(filePath, nextCursor); + if (this.shouldWatchFile(filePath)) { + this.ensureWatcher(filePath); + } } } - async function listSessionFiles(dir: string): Promise { + protected async beforeScan(): Promise { + this.bestWithinWindow = null; + } + + protected async findSessionFiles(): Promise { + const files = await this.listSessionFiles(this.sessionsRoot); + return sortFilesByMtime(files); + } + + protected async parseSessionFile(filePath: string, cursor: number): Promise> { + if (this.shouldSkipFile(filePath)) { + return { events: [], nextCursor: cursor }; + } + return this.readSessionFile(filePath, cursor); + } + + protected generateEventKey(event: CodexSessionEvent, context: { filePath: string; lineIndex?: number }): string { + const epoch = this.fileEpochByPath.get(context.filePath) ?? 0; + const lineIndex = context.lineIndex ?? -1; + return `${context.filePath}:${epoch}:${lineIndex}`; + } + + protected async handleFileScan(stats: SessionFileScanStats): Promise { + const filePath = stats.filePath; + const fileSessionId = this.sessionIdByFile.get(filePath) ?? null; + + if (!this.activeSessionId && this.targetCwd) { + this.appendPendingEvents(filePath, stats.events, fileSessionId); + const candidate = this.getCandidateForFile(filePath); + if (candidate) { + if (!this.bestWithinWindow || candidate.score < this.bestWithinWindow.score) { + this.bestWithinWindow = candidate; + } + } + if (stats.newCount > 0) { + logger.debug(`[CODEX_SESSION_SCANNER] Buffered ${stats.newCount} pending events from ${filePath}`); + } + return; + } + + const emittedForFile = this.emitEvents(stats.events, fileSessionId); + if (emittedForFile > 0) { + logger.debug(`[CODEX_SESSION_SCANNER] Emitted ${emittedForFile} new events from ${filePath}`); + } + } + + protected async afterScan(): Promise { + if (!this.activeSessionId && this.targetCwd) { + if (this.bestWithinWindow) { + logger.debug(`[CODEX_SESSION_SCANNER] Selected session ${this.bestWithinWindow.sessionId} within start window`); + this.setActiveSessionId(this.bestWithinWindow.sessionId); + } else if (Date.now() > this.matchDeadlineMs) { + this.matchFailed = true; + this.pendingEventsByFile.clear(); + const message = `No Codex session found within ${this.sessionStartWindowMs}ms for cwd ${this.targetCwd}; refusing fallback.`; + logger.warn(`[CODEX_SESSION_SCANNER] ${message}`); + this.onSessionMatchFailed?.(message); + } else if (this.pendingEventsByFile.size > 0) { + logger.debug('[CODEX_SESSION_SCANNER] No session candidate matched yet; pending events buffered'); + } + } + } + + private shouldSkipFile(filePath: string): boolean { + if (!this.activeSessionId) { + return false; + } + const fileSessionId = this.sessionIdByFile.get(filePath); + if (fileSessionId && fileSessionId !== this.activeSessionId) { + return true; + } + if (!fileSessionId && !filePath.endsWith(`-${this.activeSessionId}.jsonl`)) { + return true; + } + return false; + } + + private reportSessionId(sessionId: string): void { + if (this.reportedSessionId === sessionId) { + return; + } + this.reportedSessionId = sessionId; + this.onSessionFound?.(sessionId); + } + + private setActiveSessionId(sessionId: string): void { + this.activeSessionId = sessionId; + this.reportSessionId(sessionId); + const candidateFiles = this.getFilesForSession(sessionId); + for (const filePath of candidateFiles) { + if (this.shouldWatchFile(filePath)) { + this.ensureWatcher(filePath); + } + } + this.pruneWatchers(this.getWatchedFiles().filter((filePath) => this.shouldWatchFile(filePath))); + if (this.targetCwd) { + this.flushPendingEventsForSession(sessionId); + } else { + this.pendingEventsByFile.clear(); + } + } + + private async listSessionFiles(dir: string): Promise { try { const entries = await readdir(dir, { withFileTypes: true }); const results: string[] = []; for (const entry of entries) { const full = join(dir, entry.name); - if (!shouldIncludeSessionPath(full, sessionsRoot, sessionDatePrefixes)) { + if (!shouldIncludeSessionPath(full, this.sessionsRoot, this.sessionDatePrefixes)) { continue; } if (entry.isDirectory()) { - results.push(...await listSessionFiles(full)); + results.push(...await this.listSessionFiles(full)); } else if (entry.isFile() && entry.name.endsWith('.jsonl')) { results.push(full); } @@ -109,24 +261,26 @@ export async function createCodexSessionScanner(opts: CodexSessionScannerOptions } } - async function readSessionFile(filePath: string, startLine: number): Promise<{ events: CodexSessionEvent[]; totalLines: number }> { + private async readSessionFile(filePath: string, startLine: number): Promise> { let content: string; try { content = await readFile(filePath, 'utf-8'); } catch (error) { - return { events: [], totalLines: startLine }; + return { events: [], nextCursor: startLine }; } - const events: CodexSessionEvent[] = []; + const events: SessionFileScanEntry[] = []; const lines = content.split('\n'); const hasTrailingEmpty = lines.length > 0 && lines[lines.length - 1] === ''; const totalLines = hasTrailingEmpty ? lines.length - 1 : lines.length; let effectiveStartLine = startLine; if (effectiveStartLine > totalLines) { effectiveStartLine = 0; + const nextEpoch = (this.fileEpochByPath.get(filePath) ?? 0) + 1; + this.fileEpochByPath.set(filePath, nextEpoch); } - const hasSessionMeta = sessionMetaParsed.has(filePath); + const hasSessionMeta = this.sessionMetaParsed.has(filePath); const parseFrom = hasSessionMeta ? effectiveStartLine : 0; for (let index = parseFrom; index < lines.length; index += 1) { @@ -135,66 +289,55 @@ export async function createCodexSessionScanner(opts: CodexSessionScannerOptions continue; } try { - const parsed = JSON.parse(trimmed); + const parsed = JSON.parse(trimmed) as CodexSessionEvent; if (parsed?.type === 'session_meta') { const payload = asRecord(parsed.payload); const sessionId = payload ? asString(payload.id) : null; if (sessionId) { - sessionIdByFile.set(filePath, sessionId); + this.sessionIdByFile.set(filePath, sessionId); } const sessionCwd = payload ? asString(payload.cwd) : null; const normalizedCwd = sessionCwd ? normalizePath(sessionCwd) : null; if (normalizedCwd) { - sessionCwdByFile.set(filePath, normalizedCwd); + this.sessionCwdByFile.set(filePath, normalizedCwd); } const rawTimestamp = payload ? payload.timestamp : null; const sessionTimestamp = payload ? parseTimestamp(payload.timestamp) : null; if (sessionTimestamp !== null) { - sessionTimestampByFile.set(filePath, sessionTimestamp); + this.sessionTimestampByFile.set(filePath, sessionTimestamp); } logger.debug(`[CODEX_SESSION_SCANNER] Session meta: file=${filePath} cwd=${sessionCwd ?? 'none'} normalizedCwd=${normalizedCwd ?? 'none'} timestamp=${rawTimestamp ?? 'none'} parsedTs=${sessionTimestamp ?? 'none'}`); - sessionMetaParsed.add(filePath); + this.sessionMetaParsed.add(filePath); } if (index >= effectiveStartLine) { - events.push(parsed); + events.push({ event: parsed, lineIndex: index }); } } catch (error) { logger.debug(`[CODEX_SESSION_SCANNER] Failed to parse line: ${error}`); } } - return { events, totalLines }; + return { events, nextCursor: totalLines }; } - async function initializeProcessedMessages(): Promise { - const files = await listSessionFiles(sessionsRoot); - for (const filePath of files) { - const { totalLines } = await readSessionFile(filePath, 0); - processedLineCounts.set(filePath, totalLines); - if (!isClosing && !watchers.has(filePath)) { - watchers.set(filePath, startFileWatcher(filePath, () => sync.invalidate())); - } - } - } - - function getCandidateForFile(filePath: string): Candidate | null { - const sessionId = sessionIdByFile.get(filePath); + private getCandidateForFile(filePath: string): Candidate | null { + const sessionId = this.sessionIdByFile.get(filePath); if (!sessionId) { return null; } - const fileCwd = sessionCwdByFile.get(filePath); - if (targetCwd && fileCwd !== targetCwd) { + const fileCwd = this.sessionCwdByFile.get(filePath); + if (this.targetCwd && fileCwd !== this.targetCwd) { return null; } - const sessionTimestamp = sessionTimestampByFile.get(filePath); + const sessionTimestamp = this.sessionTimestampByFile.get(filePath); if (sessionTimestamp === undefined) { return null; } - const diff = Math.abs(sessionTimestamp - referenceTimestampMs); - if (diff > sessionStartWindowMs) { + const diff = Math.abs(sessionTimestamp - this.referenceTimestampMs); + if (diff > this.sessionStartWindowMs) { return null; } @@ -204,11 +347,25 @@ export async function createCodexSessionScanner(opts: CodexSessionScannerOptions }; } - function appendPendingEvents(filePath: string, events: CodexSessionEvent[], fileSessionId: string | null): void { + private getFilesForSession(sessionId: string): string[] { + const matches: string[] = []; + for (const [filePath, storedSessionId] of this.sessionIdByFile.entries()) { + if (storedSessionId === sessionId) { + matches.push(filePath); + } + } + if (matches.length > 0) { + return matches; + } + const suffix = `-${sessionId}.jsonl`; + return this.getWatchedFiles().filter((filePath) => filePath.endsWith(suffix)); + } + + private appendPendingEvents(filePath: string, events: CodexSessionEvent[], fileSessionId: string | null): void { if (events.length === 0) { return; } - const existing = pendingEventsByFile.get(filePath); + const existing = this.pendingEventsByFile.get(filePath); if (existing) { existing.events.push(...events); if (!existing.fileSessionId && fileSessionId) { @@ -216,131 +373,47 @@ export async function createCodexSessionScanner(opts: CodexSessionScannerOptions } return; } - pendingEventsByFile.set(filePath, { + this.pendingEventsByFile.set(filePath, { events: [...events], fileSessionId }); } - function emitEvents(events: CodexSessionEvent[], fileSessionId: string | null): number { + private emitEvents(events: CodexSessionEvent[], fileSessionId: string | null): number { let emittedForFile = 0; for (const event of events) { const payload = asRecord(event.payload); const payloadSessionId = payload ? asString(payload.id) : null; const eventSessionId = payloadSessionId ?? fileSessionId ?? null; - if (activeSessionId && eventSessionId && eventSessionId !== activeSessionId) { + if (this.activeSessionId && eventSessionId && eventSessionId !== this.activeSessionId) { continue; } - opts.onEvent(event); + this.onEvent(event); emittedForFile += 1; } return emittedForFile; } - function flushPendingEventsForSession(sessionId: string): void { - if (pendingEventsByFile.size === 0) { + private flushPendingEventsForSession(sessionId: string): void { + if (this.pendingEventsByFile.size === 0) { return; } let emitted = 0; - for (const [filePath, pending] of pendingEventsByFile.entries()) { + for (const [filePath, pending] of this.pendingEventsByFile.entries()) { const matches = (pending.fileSessionId && pending.fileSessionId === sessionId) || filePath.endsWith(`-${sessionId}.jsonl`); if (!matches) { continue; } - emitted += emitEvents(pending.events, pending.fileSessionId); + emitted += this.emitEvents(pending.events, pending.fileSessionId); } - pendingEventsByFile.clear(); + this.pendingEventsByFile.clear(); if (emitted > 0) { logger.debug(`[CODEX_SESSION_SCANNER] Emitted ${emitted} pending events for session ${sessionId}`); } } - - const sync = new InvalidateSync(async () => { - if (isClosing || matchFailed) { - return; - } - const files = await listSessionFiles(sessionsRoot); - const sortedFiles = await sortFilesByMtime(files); - let bestWithinWindow: Candidate | null = null; - - for (const filePath of sortedFiles) { - if (isClosing) { - return; - } - if (!watchers.has(filePath)) { - watchers.set(filePath, startFileWatcher(filePath, () => sync.invalidate())); - } - - const fileSessionId = sessionIdByFile.get(filePath); - if (activeSessionId && fileSessionId && fileSessionId !== activeSessionId) { - continue; - } - if (activeSessionId && !fileSessionId && !filePath.endsWith(`-${activeSessionId}.jsonl`)) { - continue; - } - - const lastProcessedLine = processedLineCounts.get(filePath) ?? 0; - const { events, totalLines } = await readSessionFile(filePath, lastProcessedLine); - processedLineCounts.set(filePath, totalLines); - const candidate = !activeSessionId && targetCwd ? getCandidateForFile(filePath) : null; - if (!activeSessionId && targetCwd) { - appendPendingEvents(filePath, events, fileSessionId ?? null); - if (candidate) { - if (!bestWithinWindow || candidate.score < bestWithinWindow.score) { - bestWithinWindow = candidate; - } - } - continue; - } - - const emittedForFile = emitEvents(events, fileSessionId ?? null); - if (emittedForFile > 0) { - logger.debug(`[CODEX_SESSION_SCANNER] Emitted ${emittedForFile} new events from ${filePath}`); - } - } - - if (!activeSessionId && targetCwd) { - if (bestWithinWindow) { - logger.debug(`[CODEX_SESSION_SCANNER] Selected session ${bestWithinWindow.sessionId} within start window`); - setActiveSessionId(bestWithinWindow.sessionId); - } else if (Date.now() > matchDeadlineMs) { - matchFailed = true; - pendingEventsByFile.clear(); - const message = `No Codex session found within ${sessionStartWindowMs}ms for cwd ${targetCwd}; refusing fallback.`; - logger.warn(`[CODEX_SESSION_SCANNER] ${message}`); - opts.onSessionMatchFailed?.(message); - } else if (pendingEventsByFile.size > 0) { - logger.debug('[CODEX_SESSION_SCANNER] No session candidate matched yet; pending events buffered'); - } - } - }); - - await initializeProcessedMessages(); - await sync.invalidateAndAwait(); - const intervalId = setInterval(() => sync.invalidate(), 2000); - - return { - cleanup: async () => { - isClosing = true; - clearInterval(intervalId); - sync.stop(); - for (const stop of watchers.values()) { - stop(); - } - watchers.clear(); - }, - onNewSession: (sessionId: string) => { - if (activeSessionId === sessionId) { - return; - } - logger.debug(`[CODEX_SESSION_SCANNER] Switching to new session: ${sessionId}`); - setActiveSessionId(sessionId); - sync.invalidate(); - } - }; } async function sortFilesByMtime(files: string[]): Promise { diff --git a/cli/src/modules/common/session/BaseSessionScanner.ts b/cli/src/modules/common/session/BaseSessionScanner.ts new file mode 100644 index 00000000..e19d0e75 --- /dev/null +++ b/cli/src/modules/common/session/BaseSessionScanner.ts @@ -0,0 +1,189 @@ +import { InvalidateSync } from "@/utils/sync"; +import { startFileWatcher } from "@/modules/watcher/startFileWatcher"; + +export type SessionFileScanEntry = { + event: TEvent; + lineIndex?: number; +}; + +export type SessionFileScanResult = { + events: SessionFileScanEntry[]; + nextCursor: number; +}; + +export type SessionFileScanStats = { + filePath: string; + events: TEvent[]; + parsedCount: number; + newCount: number; + skippedCount: number; + cursor: number; + nextCursor: number; +}; + +type BaseSessionScannerOptions = { + intervalMs: number; +}; + +export abstract class BaseSessionScanner { + private readonly sync: InvalidateSync; + private readonly watchers = new Map void>(); + private readonly processedEventKeys = new Set(); + private readonly fileCursors = new Map(); + private intervalId: ReturnType | null = null; + private stopped = false; + private scanPromise: Promise | null = null; + + protected constructor(private readonly options: BaseSessionScannerOptions) { + this.sync = new InvalidateSync(() => this.scan()); + } + + protected abstract findSessionFiles(): Promise; + protected abstract parseSessionFile(filePath: string, cursor: number): Promise>; + protected abstract generateEventKey(event: TEvent, context: { filePath: string; lineIndex?: number }): string; + + protected async handleFileScan(_stats: SessionFileScanStats): Promise { + } + + protected async initialize(): Promise { + } + + protected async beforeScan(): Promise { + } + + protected async afterScan(): Promise { + } + + protected shouldScan(): boolean { + return true; + } + + protected shouldWatchFile(_filePath: string): boolean { + return true; + } + + protected ensureWatcher(filePath: string): void { + if (this.watchers.has(filePath)) { + return; + } + this.watchers.set(filePath, startFileWatcher(filePath, () => this.sync.invalidate())); + } + + protected invalidate(): void { + this.sync.invalidate(); + } + + protected getCursor(filePath: string): number { + return this.fileCursors.get(filePath) ?? 0; + } + + protected setCursor(filePath: string, cursor: number): void { + this.fileCursors.set(filePath, cursor); + } + + protected seedProcessedKeys(keys: Iterable): void { + for (const key of keys) { + this.recordProcessedKey(key); + } + } + + protected getWatchedFiles(): string[] { + return [...this.watchers.keys()]; + } + + protected pruneWatchers(keepFiles: Iterable): void { + const keep = new Set(keepFiles); + for (const [filePath, stop] of this.watchers.entries()) { + if (keep.has(filePath)) { + continue; + } + stop(); + this.watchers.delete(filePath); + } + } + + public async start(): Promise { + await this.initialize(); + await this.sync.invalidateAndAwait(); + this.intervalId = setInterval(() => this.sync.invalidate(), this.options.intervalMs); + } + + public async cleanup(): Promise { + this.stopped = true; + if (this.intervalId) { + clearInterval(this.intervalId); + this.intervalId = null; + } + this.sync.stop(); + const pendingScan = this.scanPromise; + for (const stop of this.watchers.values()) { + stop(); + } + this.watchers.clear(); + if (pendingScan) { + await pendingScan.catch(() => {}); + } + } + + private async scan(): Promise { + if (this.stopped || !this.shouldScan()) { + return; + } + if (this.scanPromise) { + return this.scanPromise; + } + this.scanPromise = this.runScan(); + try { + await this.scanPromise; + } finally { + this.scanPromise = null; + } + } + + private async runScan(): Promise { + if (this.stopped || !this.shouldScan()) { + return; + } + await this.beforeScan(); + const files = await this.findSessionFiles(); + for (const filePath of files) { + if (this.stopped || !this.shouldScan()) { + return; + } + if (this.shouldWatchFile(filePath)) { + this.ensureWatcher(filePath); + } + const cursor = this.getCursor(filePath); + const { events, nextCursor } = await this.parseSessionFile(filePath, cursor); + const newEvents: TEvent[] = []; + const newKeys: string[] = []; + for (const entry of events) { + const key = this.generateEventKey(entry.event, { filePath, lineIndex: entry.lineIndex }); + if (this.processedEventKeys.has(key)) { + this.recordProcessedKey(key); + continue; + } + newKeys.push(key); + newEvents.push(entry.event); + } + await this.handleFileScan({ + filePath, + events: newEvents, + parsedCount: events.length, + newCount: newEvents.length, + skippedCount: events.length - newEvents.length, + cursor, + nextCursor + }); + this.setCursor(filePath, nextCursor); + for (const key of newKeys) { + this.recordProcessedKey(key); + } + } + await this.afterScan(); + } + + private recordProcessedKey(key: string): void { + this.processedEventKeys.add(key); + } +}