diff --git a/cli/src/claude/utils/sessionScanner.test.ts b/cli/src/claude/utils/sessionScanner.test.ts index 40841a46..e9f02181 100644 --- a/cli/src/claude/utils/sessionScanner.test.ts +++ b/cli/src/claude/utils/sessionScanner.test.ts @@ -1,5 +1,5 @@ import { describe, it, expect, beforeEach, afterEach } from 'vitest' -import { createSessionScanner } from './sessionScanner' +import { createSessionScanner, readSessionLog } from './sessionScanner' import { RawJSONLines } from '../types' import { mkdir, writeFile, appendFile, rm, readFile } from 'node:fs/promises' import { join } from 'node:path' @@ -144,4 +144,83 @@ describe('sessionScanner', () => { expect(content).toContain('readme.md') } }) + + it('reads only new bytes on a subsequent scan (incremental, no reparse of prior region)', async () => { + const filePath = join(testDir, 'incremental.jsonl') + const line1 = JSON.stringify({ type: 'user', uuid: 'u1', message: { content: 'first' } }) + const line2 = JSON.stringify({ type: 'user', uuid: 'u2', message: { content: 'second' } }) + await writeFile(filePath, line1 + '\n') + + const first = await readSessionLog(filePath, 0) + expect(first.events).toHaveLength(1) + expect(first.nextCursor).toBe(Buffer.byteLength(line1 + '\n')) + + await appendFile(filePath, line2 + '\n') + const second = await readSessionLog(filePath, first.nextCursor) + + // Starting from the cursor returned by the first scan, only the newly + // appended line comes back — the already-seen line is not re-parsed. + expect(second.events).toHaveLength(1) + expect(second.events[0].event.type).toBe('user') + if (second.events[0].event.type === 'user') { + expect(second.events[0].event.uuid).toBe('u2') + } + expect(second.nextCursor).toBe(Buffer.byteLength(line1 + '\n' + line2 + '\n')) + }) + + it('holds back an unterminated trailing line until its newline arrives', async () => { + const filePath = join(testDir, 'partial.jsonl') + const complete = JSON.stringify({ type: 'user', uuid: 'u1', message: { content: 'complete' } }) + const partial = JSON.stringify({ type: 'user', uuid: 'u2', message: { content: 'partial' } }) + + // Simulate a write in progress: a full line followed by a truncated one + // with no trailing newline yet. + await writeFile(filePath, complete + '\n' + partial.slice(0, 10)) + + const result = await readSessionLog(filePath, 0) + expect(result.events).toHaveLength(1) + if (result.events[0].event.type === 'user') { + expect(result.events[0].event.uuid).toBe('u1') + } + // Cursor sits right after the complete line's newline, not at EOF — the + // partial trailing bytes are left unconsumed. + expect(result.nextCursor).toBe(Buffer.byteLength(complete + '\n')) + + // Completing the line makes it available on the next scan. + await appendFile(filePath, partial.slice(10) + '\n') + const followUp = await readSessionLog(filePath, result.nextCursor) + expect(followUp.events).toHaveLength(1) + if (followUp.events[0].event.type === 'user') { + expect(followUp.events[0].event.uuid).toBe('u2') + } + }) + + it('forwards a complete final record written without a trailing newline', async () => { + const filePath = join(testDir, 'no-trailing-newline.jsonl') + const line1 = JSON.stringify({ type: 'user', uuid: 'u1', message: { content: 'first' } }) + const line2 = JSON.stringify({ type: 'user', uuid: 'u2', message: { content: 'last' } }) + + // A prior line terminated by a newline, then a complete final record that + // was flushed without its terminating newline (shutdown/import). + await writeFile(filePath, line1 + '\n' + line2) + + const result = await readSessionLog(filePath, 0) + expect(result.events).toHaveLength(2) + expect(result.events.map((e) => e.event.type === 'user' ? e.event.uuid : null)).toEqual(['u1', 'u2']) + // The whole file is consumed even though it does not end in a newline. + expect(result.nextCursor).toBe(Buffer.byteLength(line1 + '\n' + line2)) + }) + + it('forwards a single complete record with no newline at all', async () => { + const filePath = join(testDir, 'single-no-newline.jsonl') + const only = JSON.stringify({ type: 'user', uuid: 'u1', message: { content: 'only' } }) + await writeFile(filePath, only) + + const result = await readSessionLog(filePath, 0) + expect(result.events).toHaveLength(1) + if (result.events[0].event.type === 'user') { + expect(result.events[0].event.uuid).toBe('u1') + } + expect(result.nextCursor).toBe(Buffer.byteLength(only)) + }) }) diff --git a/cli/src/claude/utils/sessionScanner.ts b/cli/src/claude/utils/sessionScanner.ts index d97dd223..7417cc6c 100644 --- a/cli/src/claude/utils/sessionScanner.ts +++ b/cli/src/claude/utils/sessionScanner.ts @@ -1,6 +1,6 @@ import { RawJSONLines, RawJSONLinesSchema } from "../types"; import { basename, join } from "node:path"; -import { readFile } from "node:fs/promises"; +import { open, stat } from "node:fs/promises"; import { logger } from "@/ui/logger"; import { getProjectPath } from "./path"; import { BaseSessionScanner, SessionFileScanEntry, SessionFileScanResult, SessionFileScanStats } from "@/modules/common/session/BaseSessionScanner"; @@ -83,11 +83,11 @@ class ClaudeSessionScanner extends BaseSessionScanner { return; } const sessionFile = this.sessionFilePath(this.currentSessionId); - const { events, totalLines } = await readSessionLog(sessionFile, 0); + const { events, nextCursor } = 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); + this.setCursor(sessionFile, nextCursor); } protected async beforeScan(): Promise { @@ -113,10 +113,10 @@ class ClaudeSessionScanner extends BaseSessionScanner { if (sessionId) { this.scannedSessions.add(sessionId); } - const { events, totalLines } = await readSessionLog(filePath, cursor); + const { events, nextCursor } = await readSessionLog(filePath, cursor); return { events, - nextCursor: totalLines + nextCursor }; } @@ -169,52 +169,119 @@ function messageKey(message: RawJSONLines): string { } /** - * Read and parse session log file. - * Returns only valid conversation messages, silently skipping internal events. + * Whether a trailing segment (after the last newline) is already a complete + * JSON value. A record still being written parses as incomplete, so this + * distinguishes a flushed final record with no terminating newline from a + * genuinely partial line. */ -async function readSessionLog(filePath: string, startLine: number): Promise<{ events: SessionFileScanEntry[]; totalLines: number }> { - logger.debug(`[SESSION_SCANNER] Reading session file: ${filePath}`); - let file: string; +function isCompleteJsonLine(segment: Buffer): boolean { try { - file = await readFile(filePath, 'utf-8'); + JSON.parse(segment.toString('utf-8')); + return true; + } catch { + return false; + } +} + +/** + * Incrementally read and parse a session log file. + * + * The cursor is a BYTE OFFSET into the (append-only) JSONL. Each scan stats the + * file and reads only the bytes after the cursor — so the cost is O(new content) + * regardless of how large the conversation has grown, instead of re-reading the + * whole file on every scan, poll- or watch-driven. A trailing partial line (a + * write in progress) is left unconsumed until its newline arrives — unless it + * already forms a complete record flushed without a terminating newline, which + * is consumed rather than stranded. If the file shrank, the cursor resets to 0 + * and the whole file is re-read (dedup by uuid in the base scanner absorbs any + * re-sent events). + */ +export async function readSessionLog(filePath: string, startByte: number): Promise<{ events: SessionFileScanEntry[]; nextCursor: number }> { + let size: number; + try { + size = (await stat(filePath)).size; } catch (error) { logger.debug(`[SESSION_SCANNER] Session file not found: ${filePath}`); - return { events: [], totalLines: startLine }; + return { events: [], nextCursor: startByte }; } - 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; + + let from = startByte; + if (from > size) { + from = 0; // file was truncated/rewritten — re-read from the top } - const messages: SessionFileScanEntry[] = []; - for (let index = effectiveStartLine; index < lines.length; index += 1) { - const l = lines[index]; + if (from >= size) { + return { events: [], nextCursor: size }; // no new bytes + } + + let chunk: Buffer; + try { + const length = size - from; + const buffer = Buffer.allocUnsafe(length); + let bytesRead = 0; + const fd = await open(filePath, 'r'); try { - if (l.trim() === '') { - continue; + // A single read may return fewer bytes than requested, so loop until + // the range is filled or EOF is hit. + while (bytesRead < length) { + const result = await fd.read(buffer, bytesRead, length - bytesRead, from + bytesRead); + if (result.bytesRead === 0) { + break; + } + bytesRead += result.bytesRead; } - let message = JSON.parse(l); - - // Silently skip known internal Claude Code events - // These are state/tracking events, not conversation messages + } finally { + await fd.close(); + } + // The tail of an allocUnsafe buffer is uninitialized heap, so only the + // first `bytesRead` bytes are valid. Operating past them would let a stray + // 0x0a in garbage advance the cursor past never-read data → dropped lines. + chunk = buffer.subarray(0, bytesRead); + } catch (error) { + logger.debug(`[SESSION_SCANNER] Failed to read session file ${filePath}: ${error}`); + return { events: [], nextCursor: startByte }; + } + + // Everything up to and including the last newline is complete lines. A + // segment after it is normally a partial write, held back until its newline + // arrives on a later scan. But a final record can be flushed without a + // trailing newline (e.g. at shutdown or on import); the previous whole-file + // reader parsed such a record, so if the trailing segment already parses as + // a complete JSON value, consume it now instead of stranding it until the + // next append. + let readableEnd = chunk.lastIndexOf(0x0a) + 1; // 0 when no newline yet + const trailing = chunk.subarray(readableEnd); + if (trailing.length > 0 && isCompleteJsonLine(trailing)) { + readableEnd = chunk.length; + } + if (readableEnd === 0) { + return { events: [], nextCursor: from }; + } + const nextCursor = from + readableEnd; + const text = chunk.subarray(0, readableEnd).toString('utf-8'); + + const messages: SessionFileScanEntry[] = []; + for (const l of text.split('\n')) { + if (l.trim() === '') { + continue; + } + try { + const message = JSON.parse(l); + // Silently skip known internal Claude Code state/tracking events. if (message.type && INTERNAL_CLAUDE_EVENT_TYPES.has(message.type)) { continue; } - - let parsed = RawJSONLinesSchema.safeParse(message); + const parsed = RawJSONLinesSchema.safeParse(message); if (!parsed.success) { // Unknown message types are silently skipped. continue; } - messages.push({ event: parsed.data, lineIndex: index }); + messages.push({ event: parsed.data }); } catch (e) { logger.debug(`[SESSION_SCANNER] Error processing message: ${e}`); continue; } } - return { events: messages, totalLines }; + return { events: messages, nextCursor }; } function sessionIdFromPath(filePath: string): string | null {