diff --git a/src/codex/native-residue.ts b/src/codex/native-residue.ts index 1f3bd6551..cc5425215 100644 --- a/src/codex/native-residue.ts +++ b/src/codex/native-residue.ts @@ -1,5 +1,15 @@ import { createHash } from "node:crypto"; -import { lstatSync, readFileSync, readdirSync, realpathSync, statSync } from "node:fs"; +import { + closeSync, + fstatSync, + lstatSync, + openSync, + readFileSync, + readdirSync, + realpathSync, + readSync, + statSync, +} from "node:fs"; import type { Stats } from "node:fs"; import { basename, dirname, join, resolve } from "node:path"; @@ -69,6 +79,8 @@ const MODELS_CACHE_FILE_NAME = basename(CODEX_MODELS_CACHE_PATH); const JOURNAL_FILE_NAME = "opencodex-journal.json"; const HISTORY_DATABASE_FILE_NAME = "state_5.sqlite"; const ROUTED_CATALOG_DESCRIPTION_PREFIX = "Routed via opencodex → "; +const MAX_ROLLOUT_INSPECTION_BYTES = 64 * 1024 * 1024; +const ROLLOUT_READ_CHUNK_BYTES = 64 * 1024; function errorCode(error: unknown): string | undefined { return (error as NodeJS.ErrnoException | undefined)?.code; @@ -141,6 +153,53 @@ function indeterminate( return { kind: "indeterminate", surface, path, reason }; } +function rolloutSessionMetaPayload( + line: string, +): { kind: "payload"; payload: Record | null } | { kind: "malformed"; reason: string } { + let parsed: unknown; + try { + parsed = JSON.parse(line); + } catch (error) { + return { kind: "malformed", reason: `malformed rollout JSONL: ${errorReason(error)}` }; + } + if (!parsed || typeof parsed !== "object" || Array.isArray(parsed)) { + return { kind: "malformed", reason: "rollout JSONL record is not an object" }; + } + const record = parsed as Record; + if (record.type !== "session_meta") return { kind: "payload", payload: null }; + if (!record.payload || typeof record.payload !== "object" || Array.isArray(record.payload)) { + return { kind: "malformed", reason: "session_meta payload has an unknown shape" }; + } + return { kind: "payload", payload: record.payload as Record }; +} + +function consumeRolloutLines( + surface: "history" | "history-backup", + path: string, + partial: string, + first: Record | undefined, + latest: Record | undefined, +): NativeRoutedResidueResult | { kind: "continue"; partial: string; first: Record | undefined; latest: Record | undefined } { + let rest = partial; + let newline = rest.indexOf("\n"); + while (newline !== -1) { + const line = rest.slice(0, newline); + rest = rest.slice(newline + 1); + if (line.trim()) { + const payload = rolloutSessionMetaPayload(line); + if (payload.kind === "malformed") { + return indeterminate(surface, path, payload.reason); + } + if (payload.payload !== null) { + first ??= payload.payload; + latest = payload.payload; + } + } + newline = rest.indexOf("\n"); + } + return { kind: "continue", partial: rest, first, latest }; +} + function classifyToml( surface: "config" | "profile", path: string, @@ -358,50 +417,104 @@ function classifyReferencedRollout( surface: "history" | "history-backup", reference: RolloutReference, ): NativeRoutedResidueResult { - const read = readRegularFile(reference.path); - if (read.kind === "absent") { + const resolved = resolveRegularFile(reference.path); + if (resolved.kind === "absent") { return indeterminate(surface, reference.path, "referenced rollout is absent"); } - if (read.kind === "indeterminate") return indeterminate(surface, reference.path, read.reason); + if (resolved.kind === "indeterminate") return indeterminate(surface, reference.path, resolved.reason); + + let handle: number; + try { + handle = openSync(resolved.path, "r"); + } catch (error) { + return indeterminate(surface, resolved.path, `unreadable rollout: ${errorReason(error)}`); + } let first: Record | undefined; let latest: Record | undefined; - for (const line of read.content.split("\n")) { - if (!line.trim()) continue; - let parsed: unknown; - try { - parsed = JSON.parse(line); - } catch (error) { - return indeterminate(surface, read.path, `malformed rollout JSONL: ${errorReason(error)}`); + let partial = ""; + let totalRead = 0; + try { + const opened = fstatSync(handle); + if (opened.size > MAX_ROLLOUT_INSPECTION_BYTES) { + return indeterminate( + surface, + resolved.path, + `referenced rollout exceeds the ${MAX_ROLLOUT_INSPECTION_BYTES} byte inspection limit`, + ); } - if (!parsed || typeof parsed !== "object" || Array.isArray(parsed)) { - return indeterminate(surface, read.path, "rollout JSONL record is not an object"); + const decoder = new TextDecoder("utf-8", { ignoreBOM: true }); + const buffer = Buffer.allocUnsafe(ROLLOUT_READ_CHUNK_BYTES); + while (totalRead < opened.size) { + const remaining = Math.min(buffer.length, opened.size - totalRead); + const count = readSync(handle, buffer, 0, remaining, totalRead); + if (count === 0) { + return indeterminate(surface, resolved.path, "rollout read ended before the observed size"); + } + if (count < 0) { + return indeterminate(surface, resolved.path, "rollout read ended before the observed size"); + } + totalRead += count; + partial += decoder.decode(buffer.subarray(0, count), { stream: true }); + const consumed = consumeRolloutLines(surface, resolved.path, partial, first, latest); + if (consumed.kind !== "continue") return consumed; + partial = consumed.partial; + first = consumed.first; + latest = consumed.latest; } - const record = parsed as Record; - if (record.type !== "session_meta") continue; - if (!record.payload || typeof record.payload !== "object" || Array.isArray(record.payload)) { - return indeterminate(surface, read.path, "session_meta payload has an unknown shape"); + partial += decoder.decode(); + const consumed = consumeRolloutLines(surface, resolved.path, partial, first, latest); + if (consumed.kind !== "continue") return consumed; + partial = consumed.partial; + first = consumed.first; + latest = consumed.latest; + if (partial.trim()) { + const payload = rolloutSessionMetaPayload(partial); + if (payload.kind === "malformed") { + return indeterminate(surface, resolved.path, payload.reason); + } + if (payload.payload !== null) { + first ??= payload.payload; + latest = payload.payload; + } + } + const after = fstatSync(handle); + if (!sameStat(resolved.stat, after)) { + return indeterminate(surface, resolved.path, "rollout changed while it was being observed"); + } + const pathAfter = statSync(resolved.path); + if (!sameStat(resolved.stat, pathAfter)) { + return indeterminate(surface, resolved.path, "rollout pathname was replaced while it was being observed"); + } + } catch (error) { + if (errorCode(error) === "ENOENT") { + return indeterminate(surface, resolved.path, "referenced rollout is absent"); + } + return indeterminate(surface, resolved.path, `unreadable rollout: ${errorReason(error)}`); + } finally { + try { + closeSync(handle); + } catch { + // Closing an already-closed descriptor cannot affect the classification. } - const payload = record.payload as Record; - first ??= payload; - latest = payload; } if (!first || !latest) { - return indeterminate(surface, read.path, "referenced rollout has no session_meta metadata"); + return indeterminate(surface, resolved.path, "referenced rollout has no session_meta metadata"); } + let hasOpenCodexProvider = false; for (const [position, payload] of [["first", first], ["latest", latest]] as const) { if (payload.id !== reference.id) { - return indeterminate(surface, read.path, `${position} session_meta does not identify the referenced thread`); + return indeterminate(surface, resolved.path, `${position} session_meta does not identify the referenced thread`); } if (typeof payload.model_provider !== "string" || !payload.model_provider) { - return indeterminate(surface, read.path, `${position} session_meta has no provider metadata`); - } - if (payload.model_provider === "opencodex") { - return { kind: "residue", surface, path: read.path }; + return indeterminate(surface, resolved.path, `${position} session_meta has no provider metadata`); } + hasOpenCodexProvider = hasOpenCodexProvider || payload.model_provider === "opencodex"; } - return { kind: "clean" }; + return hasOpenCodexProvider + ? { kind: "residue", surface, path: resolved.path } + : { kind: "clean" }; } function classifyReferencedRollouts( diff --git a/tests/codex-native-residue.test.ts b/tests/codex-native-residue.test.ts index b086e3d5f..b8d6e9ffa 100644 --- a/tests/codex-native-residue.test.ts +++ b/tests/codex-native-residue.test.ts @@ -9,6 +9,7 @@ import { realpathSync, rmSync, symlinkSync, + truncateSync, writeFileSync, } from "node:fs"; import { tmpdir } from "node:os"; @@ -637,12 +638,80 @@ test("routed first rollout metadata is residue even when the latest metadata is }); }); +test("an opencodex first rollout with invalid latest metadata is indeterminate", () => { + createHistoryDatabase("openai"); + writeFileSync( + pathInCodexHome("rollout.jsonl"), + sessionMeta("thread-1", "opencodex") + "\n" + sessionMeta("thread-2", "openai") + "\n", + ); + + expect(classifyNativeRoutedResidue()).toMatchObject({ + kind: "indeterminate", + surface: "history", + path: pathInCodexHome("rollout.jsonl"), + reason: expect.stringContaining("latest session_meta"), + }); +}); + test("a referenced rollout with native first and latest metadata is clean", () => { createHistoryDatabase("openai", ["openai", "openai"]); expect(classifyNativeRoutedResidue()).toEqual({ kind: "clean" }); }); +test("a routed rollout without a trailing newline is residue", () => { + createHistoryDatabase("openai"); + writeFileSync(pathInCodexHome("rollout.jsonl"), sessionMeta("thread-1", "opencodex")); + + expect(classifyNativeRoutedResidue()).toMatchObject({ + kind: "residue", + surface: "history", + path: pathInCodexHome("rollout.jsonl"), + }); +}); + +test("a routed rollout with a non-ASCII id split across the read chunk is residue", () => { + createHistoryDatabase("openai"); + const boundary = 64 * 1024; + const prefix = `{"timestamp":"2026-08-04T00:00:00.000Z","type":"session_meta","payload":{"description":"`; + const suffix = `","id":"thread-1","model_provider":"opencodex","source":"cli"}}\n`; + const paddingLength = boundary - Buffer.byteLength(prefix) - 1; // 🚀 starts at byte 65535, straddling 64 KiB + const content = `${prefix}${"x".repeat(paddingLength)}🚀${suffix}`; + const emojiByteOffset = Buffer.from(content, "utf8").indexOf(Buffer.from("🚀", "utf8")); + expect(emojiByteOffset).toBe(boundary - 1); + writeFileSync(pathInCodexHome("rollout.jsonl"), content); + + expect(classifyNativeRoutedResidue()).toMatchObject({ + kind: "residue", + surface: "history", + path: pathInCodexHome("rollout.jsonl"), + }); +}); + +test("an oversized referenced rollout is indeterminate without being loaded", () => { + createHistoryDatabase("openai"); + truncateSync(pathInCodexHome("rollout.jsonl"), 64 * 1024 * 1024 + 1); + + expect(classifyNativeRoutedResidue()).toMatchObject({ + kind: "indeterminate", + surface: "history", + path: pathInCodexHome("rollout.jsonl"), + reason: expect.stringContaining("inspection limit"), + }); +}); + +test("a BOM-prefixed rollout record is indeterminate", () => { + createHistoryDatabase("openai"); + writeFileSync(pathInCodexHome("rollout.jsonl"), `\uFEFF${sessionMeta("thread-1", "opencodex")}\n`); + + expect(classifyNativeRoutedResidue()).toMatchObject({ + kind: "indeterminate", + surface: "history", + path: pathInCodexHome("rollout.jsonl"), + reason: expect.stringContaining("malformed rollout JSONL"), + }); +}); + for (const fixture of [ { name: "missing",