Skip to content

Commit 6070834

Browse files
committed
Stop null-padding turns.jsonl and recover poisoned resume loads
When keepBytes exceeds on-disk size, rebuild the segment instead of truncate-past-EOF (which pads null bytes). On load, if the base store fails, recover usable turns with null-strip parse and soft-default metadata; unrecoverable errors name the file.
1 parent 2c4643b commit 6070834

4 files changed

Lines changed: 206 additions & 15 deletions

File tree

‎src/session/incremental-jsonl.test.ts‎

Lines changed: 28 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -237,6 +237,34 @@ describe("createSegmentedJSONLWriter", () => {
237237
});
238238
});
239239

240+
describe("createSegmentedJSONLWriter stale keepBytes", () => {
241+
test("does not pad null bytes when on-disk file shrank below keepBytes", async () => {
242+
const dir = tempDir();
243+
const write = createSegmentedJSONLWriter(dir, BASE);
244+
245+
const a = { id: 1, text: "first-record" };
246+
const b = { id: 2, text: "second-record" };
247+
const c = { id: 3, text: "third-record" };
248+
await write([a, b, c]);
249+
250+
const full = path.join(dir, BASE);
251+
// Simulate external shrink/compaction that left the in-memory offsets stale:
252+
// file is shorter than the writer's remembered keepBytes for a shared prefix.
253+
const keptOnDisk = fullSnapshot([a]);
254+
fs.writeFileSync(full, keptOnDisk);
255+
expect(fs.statSync(full).size).toBeLessThan(Buffer.byteLength(fullSnapshot([a, b, c])));
256+
257+
// Shared prefix [a, b] would compute keepBytes past the shrunken file size.
258+
// Writer must rebuild rather than truncate-extend with null padding.
259+
const d = { id: 4, text: "after-shrink" };
260+
await write([a, b, d]);
261+
262+
const onDisk = fs.readFileSync(full);
263+
expect(onDisk.includes(0)).toBe(false);
264+
expect(await combined(dir)).toBe(fullSnapshot([a, b, d]));
265+
});
266+
});
267+
240268
describe("segment readers", () => {
241269
test("readExtraSegmentTexts returns tail segments in order", async () => {
242270
const dir = tempDir();

‎src/session/incremental-jsonl.ts‎

Lines changed: 24 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -211,12 +211,31 @@ export function createSegmentedJSONLWriter(
211211
const full = path.join(dir, name);
212212
const truncateInPlace = isFirst && state !== null && entry.keepBytes > 0;
213213
if (truncateInPlace) {
214-
const handle = await fs.promises.open(full, "r+");
214+
// Stale keepBytes (e.g. after external shrink/compaction) can exceed the
215+
// on-disk size. POSIX truncate-past-EOF pads with null bytes, which
216+
// poisons the JSONL and breaks resume with `\u0000` parse errors.
217+
// Never extend via truncate — rewrite the full segment instead.
218+
let existingSize = 0;
215219
try {
216-
await handle.truncate(entry.keepBytes);
217-
if (entry.text.length > 0) await handle.write(entry.text, entry.keepBytes);
218-
} finally {
219-
await handle.close();
220+
existingSize = (await fs.promises.stat(full)).size;
221+
} catch {
222+
existingSize = 0;
223+
}
224+
if (entry.keepBytes > existingSize) {
225+
// Offsets are wrong relative to disk. Rebuild the kept prefix from
226+
// the in-memory records that belong in this segment, then append
227+
// the planned text (the post-prefix lines for this segment).
228+
const keptRecords = records.slice(firstSegStartRecord, prefix);
229+
const fullText = keptRecords.map((r) => lineFor(r)).join("") + entry.text;
230+
await fs.promises.writeFile(full, fullText);
231+
} else {
232+
const handle = await fs.promises.open(full, "r+");
233+
try {
234+
await handle.truncate(entry.keepBytes);
235+
if (entry.text.length > 0) await handle.write(entry.text, entry.keepBytes);
236+
} finally {
237+
await handle.close();
238+
}
220239
}
221240
} else {
222241
await fs.promises.writeFile(full, entry.text);

‎src/session/optimized-context-store.test.ts‎

Lines changed: 52 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -61,6 +61,58 @@ describe("createOptimizedContextStore load", () => {
6161
expect(loaded.turns.map((t) => (t.content[0] as { text: string }).text)).toEqual(["a", "b"]);
6262
});
6363

64+
test("recovers usable turns when turns.jsonl has a mid-file null-byte hole", async () => {
65+
const dir = tempDir();
66+
const store = await createOptimizedContextStore(dir);
67+
68+
const head = jsonl([turn("a"), turn("b")]);
69+
const tail = jsonl([turn("c")]);
70+
// Simulate truncate-past-EOF null padding between valid JSONL records.
71+
const poisoned = Buffer.concat([
72+
Buffer.from(head, "utf8"),
73+
Buffer.alloc(64, 0),
74+
Buffer.from(tail, "utf8"),
75+
]);
76+
fs.writeFileSync(path.join(dir, TURNS_FILE), poisoned);
77+
78+
const loaded = await store.load();
79+
expect(loaded.turns.map((t) => (t.content[0] as { text: string }).text)).toEqual([
80+
"a",
81+
"b",
82+
"c",
83+
]);
84+
});
85+
86+
test("soft-defaults metadata when metadata.json is corrupt but turns load", async () => {
87+
const dir = tempDir();
88+
const store = await createOptimizedContextStore(dir);
89+
90+
fs.writeFileSync(path.join(dir, TURNS_FILE), jsonl([turn("kept")]));
91+
// Corrupt metadata alone must not abort resume when turns are fine.
92+
// Base load parses turns first then metadata — if metadata throws, recovery
93+
// path soft-defaults and still returns turns.
94+
fs.writeFileSync(path.join(dir, "metadata.json"), "{not-json\x00");
95+
96+
const loaded = await store.load();
97+
expect(loaded.turns).toHaveLength(1);
98+
expect((loaded.turns[0]!.content[0] as { text: string }).text).toBe("kept");
99+
expect(loaded.pendingOperations).toEqual([]);
100+
expect(loaded.connectorState).toBeNull();
101+
});
102+
103+
test("unrecoverable turns.jsonl names the file in the error", async () => {
104+
const dir = tempDir();
105+
const store = await createOptimizedContextStore(dir);
106+
107+
// Mid-file garbage that is not null padding and not a torn tail — unrecoverable.
108+
fs.writeFileSync(
109+
path.join(dir, TURNS_FILE),
110+
jsonl([turn("a")]) + "THIS IS NOT JSON\n" + jsonl([turn("b")]),
111+
);
112+
113+
await expect(store.load()).rejects.toThrow(/turns\.jsonl/);
114+
});
115+
64116
// Compacted head rewrites segment 0 while a prior multi-segment history's
65117
// tails stay on disk. Concatenating them reintroduces tool_call ids that the
66118
// compact head already kept — drop the orphan tails so the session can resume.

‎src/session/optimized-context-store.ts‎

Lines changed: 102 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -68,31 +68,86 @@ function sanitizeCallId(callId: string): string {
6868
* Parse conversation turns out of one JSONL segment. A crash can tear the final
6969
* line of the active (last) segment mid-write; when `tolerateTornTail` is set a
7070
* final line that fails to parse is dropped rather than aborting the resume.
71+
*
72+
* Null bytes (truncate-past-EOF padding from a stale keepBytes write) are stripped
73+
* so a poisoned segment can still yield its usable turns on resume. Errors name
74+
* `fileName` when provided so diagnostics point at the on-disk file, not a bare
75+
* Bun JSON token.
7176
*/
72-
function parseSegmentTurns(text: string, tolerateTornTail: boolean): ConversationTurn[] {
77+
function parseSegmentTurns(
78+
text: string,
79+
tolerateTornTail: boolean,
80+
fileName = "turns segment",
81+
): ConversationTurn[] {
7382
if (text.length === 0) return [];
74-
const lines = text.split("\n");
83+
// POSIX truncate past EOF pads with `\0`. Strip them so the rest of the JSONL
84+
// remains parseable instead of dying on Unrecognized token '\u0000'.
85+
const cleaned = text.includes("\0") ? text.replaceAll("\0", "") : text;
86+
if (cleaned.length === 0) return [];
87+
const lines = cleaned.split("\n");
7588
if (lines[lines.length - 1] === "") lines.pop();
7689

7790
const turns: ConversationTurn[] = [];
7891
for (let i = 0; i < lines.length; i++) {
92+
const line = lines[i]!;
93+
if (line.length === 0) continue;
7994
const isLast = i === lines.length - 1;
8095
let raw: unknown;
8196
try {
82-
raw = JSON.parse(lines[i]!);
97+
raw = JSON.parse(line);
8398
} catch (cause) {
8499
if (tolerateTornTail && isLast) break;
85-
throw new Error("turns segment has malformed JSON", { cause });
100+
throw new Error(`${fileName} has malformed JSON at line ${i + 1}`, { cause });
86101
}
87102
const result = ConversationTurnSchema(raw);
88103
if (result instanceof type.errors) {
89-
throw new Error(`turns segment has unexpected structure: ${result.summary}`);
104+
throw new Error(`${fileName} has unexpected structure at line ${i + 1}: ${result.summary}`);
90105
}
91106
turns.push(result);
92107
}
93108
return turns;
94109
}
95110

111+
const EMPTY_TOKEN_USAGE = {
112+
input: 0,
113+
output: 0,
114+
cacheRead: 0,
115+
cacheWrite: 0,
116+
thinking: 0,
117+
} as const;
118+
119+
function emptyMetadata(): {
120+
pendingOperations: never[];
121+
tokenUsage: typeof EMPTY_TOKEN_USAGE;
122+
connectorState: null;
123+
} {
124+
return {
125+
pendingOperations: [],
126+
tokenUsage: { ...EMPTY_TOKEN_USAGE },
127+
connectorState: null,
128+
};
129+
}
130+
131+
/**
132+
* Soft-default metadata when the recovery path cannot use the base store.
133+
* Corrupt or missing metadata.json must not abort resume of usable turns.
134+
*/
135+
async function loadMetadataSoft(dir: string): Promise<ReturnType<typeof emptyMetadata>> {
136+
const metadataPath = path.join(dir, METADATA_FILE);
137+
try {
138+
if (!(await pathExists(metadataPath))) return emptyMetadata();
139+
const text = await fs.promises.readFile(metadataPath, "utf-8");
140+
JSON.parse(text);
141+
// Schema lives in the base store; recovery only needs a safe shell.
142+
return emptyMetadata();
143+
} catch (cause) {
144+
log.warn("metadata.json unreadable during resilient load; using empty defaults", {
145+
cause: cause instanceof Error ? cause.message : String(cause),
146+
});
147+
return emptyMetadata();
148+
}
149+
}
150+
96151
// Mirrors assertWellFormedToolSequence without throwing. Used to choose the
97152
// longest segment prefix the reactor will accept after a load. Unpaired
98153
// trailing tool_calls are allowed; dups and orphan results fail.
@@ -333,12 +388,49 @@ export async function createOptimizedContextStore(dir: string): Promise<ContextS
333388
// the complete turn history is the actual live conversation state, not an
334389
// optional convenience — callers that only need a recent tail (e.g. TUI
335390
// resume hydration) should use `loadRecentTurns` instead.
391+
//
392+
// When the base isogit store hard-fails (e.g. null-padded turns.jsonl from
393+
// a stale truncate), recover usable turns via resilient segment parse and
394+
// soft-default metadata so resume does not die on a bare Bun JSON token.
336395
async load(signal) {
337-
const baseResult = await base.load(signal);
338-
const extraTexts = await readExtraSegmentTexts(dir, TURNS_FILE);
339-
if (extraTexts.length === 0) return baseResult;
340-
const turns = await loadTurnsWithoutMalformedToolSequence(baseResult.turns, extraTexts);
341-
return { ...baseResult, turns };
396+
try {
397+
const baseResult = await base.load(signal);
398+
const extraTexts = await readExtraSegmentTexts(dir, TURNS_FILE);
399+
if (extraTexts.length === 0) return baseResult;
400+
const turns = await loadTurnsWithoutMalformedToolSequence(baseResult.turns, extraTexts);
401+
return { ...baseResult, turns };
402+
} catch (cause) {
403+
log.warn(
404+
"base context store load failed; recovering turns from disk segments",
405+
{ cause: cause instanceof Error ? cause.message : String(cause) },
406+
);
407+
let baseTurns: ConversationTurn[];
408+
try {
409+
// Prefer resilient parse of segment 0 alone so orphan-tail heal still runs.
410+
const basePath = path.join(dir, TURNS_FILE);
411+
if (await pathExists(basePath)) {
412+
const text = await fs.promises.readFile(basePath, "utf-8");
413+
baseTurns = parseSegmentTurns(text, false, TURNS_FILE);
414+
} else {
415+
baseTurns = [];
416+
}
417+
} catch (parseCause) {
418+
// Unrecoverable: rethrow with the file name in the message.
419+
throw new Error(
420+
`failed to load ${TURNS_FILE}: ${
421+
parseCause instanceof Error ? parseCause.message : String(parseCause)
422+
}`,
423+
{ cause: parseCause },
424+
);
425+
}
426+
const extraTexts = await readExtraSegmentTexts(dir, TURNS_FILE);
427+
const turns =
428+
extraTexts.length === 0
429+
? baseTurns
430+
: await loadTurnsWithoutMalformedToolSequence(baseTurns, extraTexts);
431+
const metadata = await loadMetadataSoft(dir);
432+
return { turns, ...metadata };
433+
}
342434
},
343435
setConnectorState: (state) => base.setConnectorState(state),
344436
branch: (name, signal) => base.branch(name, signal),

0 commit comments

Comments
 (0)