mirror of
https://github.com/wu736139669/hapi.git
synced 2026-08-05 06:24:37 +00:00
fix(codex): scan transcripts incrementally (#1031)
Co-authored-by: zj1123581321 <zj1123581321@users.noreply.github.com>
This commit is contained in:
@@ -1,9 +1,9 @@
|
|||||||
import { afterEach, beforeEach, describe, expect, it } from 'vitest';
|
import { afterEach, beforeEach, describe, expect, it } from 'vitest';
|
||||||
import { appendFile, mkdir, rm, writeFile } from 'node:fs/promises';
|
import { appendFile, mkdir, rename, rm, writeFile } from 'node:fs/promises';
|
||||||
import { existsSync } from 'node:fs';
|
import { existsSync } from 'node:fs';
|
||||||
import { join } from 'node:path';
|
import { join } from 'node:path';
|
||||||
import { tmpdir } from 'node:os';
|
import { tmpdir } from 'node:os';
|
||||||
import { createCodexSessionScanner } from './codexSessionScanner';
|
import { createCodexSessionScanner, readTranscriptRange } from './codexSessionScanner';
|
||||||
import type { CodexSessionEvent } from './codexEventConverter';
|
import type { CodexSessionEvent } from './codexEventConverter';
|
||||||
|
|
||||||
const wait = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms));
|
const wait = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms));
|
||||||
@@ -59,6 +59,21 @@ describe('codexSessionScanner', () => {
|
|||||||
expect(events[0]?.type).toBe('event_msg');
|
expect(events[0]?.type).toBe('event_msg');
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it('reads exactly the requested transcript byte range', async () => {
|
||||||
|
const initial = 'existing transcript\n';
|
||||||
|
const appended = 'new event\n';
|
||||||
|
await writeFile(transcriptPath, initial);
|
||||||
|
await appendFile(transcriptPath, appended);
|
||||||
|
|
||||||
|
const content = await readTranscriptRange(
|
||||||
|
transcriptPath,
|
||||||
|
Buffer.byteLength(initial),
|
||||||
|
Buffer.byteLength(appended)
|
||||||
|
);
|
||||||
|
|
||||||
|
expect(content.toString('utf-8')).toBe(appended);
|
||||||
|
});
|
||||||
|
|
||||||
it('can replay existing transcript history on first attach', async () => {
|
it('can replay existing transcript history on first attach', async () => {
|
||||||
await writeFile(
|
await writeFile(
|
||||||
transcriptPath,
|
transcriptPath,
|
||||||
@@ -164,6 +179,30 @@ describe('codexSessionScanner', () => {
|
|||||||
expect(events[0]?.payload).toEqual({ type: 'agent_message', message: 'after-truncate' });
|
expect(events[0]?.payload).toEqual({ type: 'agent_message', message: 'after-truncate' });
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it('resets the byte cursor when the transcript file is replaced at the same size', async () => {
|
||||||
|
const before = JSON.stringify({
|
||||||
|
type: 'event_msg',
|
||||||
|
payload: { type: 'agent_message', message: 'before-replace' }
|
||||||
|
}) + '\n';
|
||||||
|
const after = before.replace('before-replace', 'after-replace!');
|
||||||
|
expect(Buffer.byteLength(after)).toBe(Buffer.byteLength(before));
|
||||||
|
await writeFile(transcriptPath, before);
|
||||||
|
|
||||||
|
scanner = await createCodexSessionScanner({
|
||||||
|
transcriptPath,
|
||||||
|
onEvent: (event) => events.push(event)
|
||||||
|
});
|
||||||
|
const replacementPath = join(testDir, 'replacement.jsonl');
|
||||||
|
await writeFile(replacementPath, after);
|
||||||
|
await rename(replacementPath, transcriptPath);
|
||||||
|
// A watcher bound to the old inode may not observe an atomic replacement;
|
||||||
|
// the periodic scan must still detect it by file identity.
|
||||||
|
await wait(2200);
|
||||||
|
|
||||||
|
expect(events).toHaveLength(1);
|
||||||
|
expect(events[0]?.payload).toEqual({ type: 'agent_message', message: 'after-replace!' });
|
||||||
|
});
|
||||||
|
|
||||||
it('retries an unterminated final record after it is completed', async () => {
|
it('retries an unterminated final record after it is completed', async () => {
|
||||||
await writeFile(
|
await writeFile(
|
||||||
transcriptPath,
|
transcriptPath,
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
import { readFile } from 'node:fs/promises';
|
import { open, stat } from 'node:fs/promises';
|
||||||
import { BaseSessionScanner, SessionFileScanEntry, SessionFileScanResult, SessionFileScanStats } from '@/modules/common/session/BaseSessionScanner';
|
import { BaseSessionScanner, SessionFileScanEntry, SessionFileScanResult, SessionFileScanStats } from '@/modules/common/session/BaseSessionScanner';
|
||||||
import { logger } from '@/ui/logger';
|
import { logger } from '@/ui/logger';
|
||||||
import type { CodexSessionEvent } from './codexEventConverter';
|
import type { CodexSessionEvent } from './codexEventConverter';
|
||||||
@@ -34,7 +34,12 @@ class CodexSessionScannerImpl extends BaseSessionScanner<CodexSessionEvent> {
|
|||||||
private readonly onEvent: (event: CodexSessionEvent) => void;
|
private readonly onEvent: (event: CodexSessionEvent) => void;
|
||||||
private readonly onSessionId?: (sessionId: string) => void;
|
private readonly onSessionId?: (sessionId: string) => void;
|
||||||
private readonly fileEpochByPath = new Map<string, number>();
|
private readonly fileEpochByPath = new Map<string, number>();
|
||||||
private readonly fileSizeByPath = new Map<string, number>();
|
private readonly fileStateByPath = new Map<string, {
|
||||||
|
device: number;
|
||||||
|
inode: number;
|
||||||
|
partialLine: Buffer;
|
||||||
|
nextLineIndex: number;
|
||||||
|
}>();
|
||||||
private replayExistingHistoryOnNextAttach: boolean;
|
private replayExistingHistoryOnNextAttach: boolean;
|
||||||
private observedSessionId: string | null = null;
|
private observedSessionId: string | null = null;
|
||||||
|
|
||||||
@@ -110,70 +115,94 @@ class CodexSessionScannerImpl extends BaseSessionScanner<CodexSessionEvent> {
|
|||||||
this.setCursor(filePath, nextCursor);
|
this.setCursor(filePath, nextCursor);
|
||||||
}
|
}
|
||||||
|
|
||||||
private async readSessionFile(filePath: string, startLine: number): Promise<SessionFileScanResult<CodexSessionEvent>> {
|
private async readSessionFile(filePath: string, startOffset: number): Promise<SessionFileScanResult<CodexSessionEvent>> {
|
||||||
let content: string;
|
let fileStats;
|
||||||
try {
|
try {
|
||||||
content = await readFile(filePath, 'utf-8');
|
fileStats = await stat(filePath);
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
logger.debug(`[codex-session-scanner] Failed to read transcript ${filePath}: ${error}`);
|
logger.debug(`[codex-session-scanner] Failed to stat transcript ${filePath}: ${error}`);
|
||||||
return { events: [], nextCursor: startLine };
|
return { events: [], nextCursor: startOffset };
|
||||||
}
|
}
|
||||||
|
|
||||||
const lines = content.split('\n');
|
const previousState = this.fileStateByPath.get(filePath);
|
||||||
const hasTrailingEmpty = lines.length > 0 && lines[lines.length - 1] === '';
|
const identityChanged = Boolean(
|
||||||
const totalLines = hasTrailingEmpty ? lines.length - 1 : lines.length;
|
previousState
|
||||||
let nextCursor = totalLines;
|
&& (previousState.device !== fileStats.dev || previousState.inode !== fileStats.ino)
|
||||||
const currentSize = Buffer.byteLength(content);
|
);
|
||||||
const previousSize = this.fileSizeByPath.get(filePath);
|
let effectiveStartOffset = startOffset;
|
||||||
let effectiveStartLine = startLine;
|
let partialLine = previousState?.partialLine ?? Buffer.alloc(0);
|
||||||
|
let nextLineIndex = previousState?.nextLineIndex ?? 0;
|
||||||
|
|
||||||
if ((previousSize !== undefined && currentSize < previousSize) || effectiveStartLine > totalLines) {
|
if (identityChanged || fileStats.size < effectiveStartOffset) {
|
||||||
effectiveStartLine = 0;
|
effectiveStartOffset = 0;
|
||||||
|
partialLine = Buffer.alloc(0);
|
||||||
|
nextLineIndex = 0;
|
||||||
const nextEpoch = (this.fileEpochByPath.get(filePath) ?? 0) + 1;
|
const nextEpoch = (this.fileEpochByPath.get(filePath) ?? 0) + 1;
|
||||||
this.fileEpochByPath.set(filePath, nextEpoch);
|
this.fileEpochByPath.set(filePath, nextEpoch);
|
||||||
}
|
}
|
||||||
this.fileSizeByPath.set(filePath, currentSize);
|
|
||||||
|
|
||||||
const events: SessionFileScanEntry<CodexSessionEvent>[] = [];
|
const bytesToRead = fileStats.size - effectiveStartOffset;
|
||||||
for (let lineIndex = 0; lineIndex < totalLines; lineIndex += 1) {
|
let appended: Buffer = Buffer.alloc(0);
|
||||||
const line = lines[lineIndex];
|
if (bytesToRead > 0) {
|
||||||
if (!line || line.trim().length === 0) {
|
|
||||||
continue;
|
|
||||||
}
|
|
||||||
|
|
||||||
let parsed: unknown;
|
|
||||||
try {
|
try {
|
||||||
parsed = JSON.parse(line);
|
appended = await readTranscriptRange(filePath, effectiveStartOffset, bytesToRead);
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
logger.debug(`[codex-session-scanner] Failed to parse transcript line ${filePath}:${lineIndex + 1}: ${error}`);
|
logger.debug(`[codex-session-scanner] Failed to read transcript ${filePath}: ${error}`);
|
||||||
if (!hasTrailingEmpty && lineIndex === totalLines - 1) {
|
return { events: [], nextCursor: startOffset };
|
||||||
nextCursor = lineIndex;
|
|
||||||
}
|
|
||||||
continue;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
const event = parseCodexSessionEvent(parsed);
|
|
||||||
if (!event) {
|
|
||||||
continue;
|
|
||||||
}
|
|
||||||
|
|
||||||
if (event.type === 'session_meta') {
|
|
||||||
const sessionId = extractSessionId(event);
|
|
||||||
if (sessionId) {
|
|
||||||
this.updateSessionId(sessionId);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
if (lineIndex < effectiveStartLine) {
|
|
||||||
continue;
|
|
||||||
}
|
|
||||||
|
|
||||||
events.push({ event, lineIndex });
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const content = partialLine.length > 0
|
||||||
|
? Buffer.concat([partialLine, appended])
|
||||||
|
: appended;
|
||||||
|
const events: SessionFileScanEntry<CodexSessionEvent>[] = [];
|
||||||
|
|
||||||
|
const parseLine = (lineBuffer: Buffer, lineIndex: number, allowIncomplete: boolean): boolean => {
|
||||||
|
const line = lineBuffer.toString('utf-8');
|
||||||
|
if (!line || line.trim().length === 0) return true;
|
||||||
|
try {
|
||||||
|
const event = parseCodexSessionEvent(JSON.parse(line));
|
||||||
|
if (!event) return true;
|
||||||
|
if (event.type === 'session_meta') {
|
||||||
|
const sessionId = extractSessionId(event);
|
||||||
|
if (sessionId) this.updateSessionId(sessionId);
|
||||||
|
}
|
||||||
|
events.push({ event, lineIndex });
|
||||||
|
return true;
|
||||||
|
} catch (error) {
|
||||||
|
if (!allowIncomplete) {
|
||||||
|
logger.debug(`[codex-session-scanner] Failed to parse transcript line ${filePath}:${lineIndex + 1}: ${error}`);
|
||||||
|
}
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
let lineStart = 0;
|
||||||
|
for (let index = 0; index < content.length; index += 1) {
|
||||||
|
if (content[index] !== 0x0a) continue;
|
||||||
|
parseLine(content.subarray(lineStart, index), nextLineIndex, false);
|
||||||
|
nextLineIndex += 1;
|
||||||
|
lineStart = index + 1;
|
||||||
|
}
|
||||||
|
|
||||||
|
const trailing = content.subarray(lineStart);
|
||||||
|
if (trailing.length > 0 && parseLine(trailing, nextLineIndex, true)) {
|
||||||
|
partialLine = Buffer.alloc(0);
|
||||||
|
nextLineIndex += 1;
|
||||||
|
} else {
|
||||||
|
partialLine = Buffer.from(trailing);
|
||||||
|
}
|
||||||
|
|
||||||
|
this.fileStateByPath.set(filePath, {
|
||||||
|
device: fileStats.dev,
|
||||||
|
inode: fileStats.ino,
|
||||||
|
partialLine,
|
||||||
|
nextLineIndex
|
||||||
|
});
|
||||||
|
|
||||||
return {
|
return {
|
||||||
events,
|
events,
|
||||||
nextCursor
|
nextCursor: effectiveStartOffset + appended.length
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -186,6 +215,22 @@ class CodexSessionScannerImpl extends BaseSessionScanner<CodexSessionEvent> {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export async function readTranscriptRange(filePath: string, startOffset: number, length: number): Promise<Buffer> {
|
||||||
|
const content = Buffer.allocUnsafe(length);
|
||||||
|
let bytesRead = 0;
|
||||||
|
const handle = await open(filePath, 'r');
|
||||||
|
try {
|
||||||
|
while (bytesRead < length) {
|
||||||
|
const result = await handle.read(content, bytesRead, length - bytesRead, startOffset + bytesRead);
|
||||||
|
if (result.bytesRead === 0) break;
|
||||||
|
bytesRead += result.bytesRead;
|
||||||
|
}
|
||||||
|
} finally {
|
||||||
|
await handle.close();
|
||||||
|
}
|
||||||
|
return bytesRead === content.length ? content : content.subarray(0, bytesRead);
|
||||||
|
}
|
||||||
|
|
||||||
function parseCodexSessionEvent(value: unknown): CodexSessionEvent | null {
|
function parseCodexSessionEvent(value: unknown): CodexSessionEvent | null {
|
||||||
if (!value || typeof value !== 'object') {
|
if (!value || typeof value !== 'object') {
|
||||||
return null;
|
return null;
|
||||||
|
|||||||
Reference in New Issue
Block a user