Skip to content

Commit 57b4553

Browse files
committed
Skip mid-file interleaved garbage in turns.jsonl on resume
JSON.parse failures on a turns line warn and skip so a glued truncated manage_tasks write or mid-file junk no longer aborts reactor resume. Schema-invalid turns still fail closed. Warns name turns.jsonl and the line number; glued {...}{...} fragments are salvaged when possible. Fixes CL-7052
1 parent 02a3f85 commit 57b4553

2 files changed

Lines changed: 189 additions & 22 deletions

File tree

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

Lines changed: 89 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -146,19 +146,78 @@ describe("createOptimizedContextStore load", () => {
146146
expect(loaded.connectorState).toBeNull();
147147
});
148148

149-
test("unrecoverable turns.jsonl names the file in the error", async () => {
149+
test("resumes past mid-file interleaved garbage in turns.jsonl", async () => {
150150
const dir = tempDir();
151151
const store = await createOptimizedContextStore(dir);
152152

153-
// Mid-file garbage that is not null padding and not a torn tail — unrecoverable.
153+
// Mid-file garbage that is not null padding — skip the bad line, keep neighbors.
154154
fs.writeFileSync(
155155
path.join(dir, TURNS_FILE),
156156
jsonl([turn("a")]) + "THIS IS NOT JSON\n" + jsonl([turn("b")]),
157157
);
158158

159+
const loaded = await store.load();
160+
expect(loaded.turns.map((t) => (t.content[0] as { text: string }).text)).toEqual([
161+
"a",
162+
"b",
163+
]);
164+
});
165+
166+
test("schema-invalid turns.jsonl still fails closed and names the file", async () => {
167+
const dir = tempDir();
168+
const store = await createOptimizedContextStore(dir);
169+
170+
const badTurn = JSON.stringify({ role: "user", content: "not-an-array", timestamp: 1 });
171+
fs.writeFileSync(
172+
path.join(dir, TURNS_FILE),
173+
jsonl([turn("a")]) + badTurn + "\n" + jsonl([turn("b")]),
174+
);
175+
159176
await expect(store.load()).rejects.toThrow(/turns\.jsonl/);
160177
});
161178

179+
test("salvages a glued truncated manage_tasks record and the next turn", async () => {
180+
const dir = tempDir();
181+
const store = await createOptimizedContextStore(dir);
182+
183+
// Production shape: truncated manage_tasks tool_call JSON glued onto the next
184+
// turn with no newline — JSON.parse of the whole line fails, but salvage keeps
185+
// the complete trailing turn.
186+
const truncatedManageTasks =
187+
'{"role":"assistant","content":[{"type":"tool_call","id":"call-mt-1","name":"manage_tasks","arguments":{"action":"update","updates":[{"id":"t1","status":"do';
188+
const nextTurn = JSON.stringify(turn("after-glue"));
189+
fs.writeFileSync(
190+
path.join(dir, TURNS_FILE),
191+
jsonl([turn("before")]) + truncatedManageTasks + nextTurn + "\n" + jsonl([turn("tail")]),
192+
);
193+
194+
const loaded = await store.load();
195+
expect(loaded.turns.map((t) => (t.content[0] as { text: string }).text)).toEqual([
196+
"before",
197+
"after-glue",
198+
"tail",
199+
]);
200+
});
201+
202+
test("resumes past mid-file garbage plus a torn trailing line", async () => {
203+
const dir = tempDir();
204+
const store = await createOptimizedContextStore(dir);
205+
206+
fs.writeFileSync(
207+
path.join(dir, TURNS_FILE),
208+
jsonl([turn("a")]) +
209+
"GARBAGE\n" +
210+
jsonl([turn("b")]) +
211+
'{"role":"user","content":[{"type":"te',
212+
);
213+
214+
const loaded = await store.load();
215+
expect(loaded.turns.map((t) => (t.content[0] as { text: string }).text)).toEqual([
216+
"a",
217+
"b",
218+
]);
219+
});
220+
162221
// Compacted head rewrites segment 0 while a prior multi-segment history's
163222
// tails stay on disk. Concatenating them reintroduces tool_call ids that the
164223
// compact head already kept — drop the orphan tails so the session can resume.
@@ -388,29 +447,53 @@ describe("loadRecentTurns", () => {
388447
expect(loaded.map((t) => (t.content[0] as { text: string }).text)).toEqual(["a", "b", "c"]);
389448
});
390449

391-
test("the reactor's load() stays strict on the same corrupt fixture and names the segment", async () => {
450+
test("the reactor's load() skips mid-file parse garbage and keeps neighbors", async () => {
392451
const dir = tempDir();
393452
const store = await createOptimizedContextStore(dir);
394453
fs.writeFileSync(
395454
path.join(dir, TURNS_FILE),
396455
jsonl([turn("a")]) + '{"role":"user","content":[{"type":"te\n' + jsonl([turn("b")]),
397456
);
398457

399-
await expect(store.load()).rejects.toThrow(TURNS_FILE);
458+
const loaded = await store.load();
459+
expect(loaded.turns.map((t) => (t.content[0] as { text: string }).text)).toEqual([
460+
"a",
461+
"b",
462+
]);
400463
});
401464

402-
test("reactor's load() stays strict and names an unrecoverable extra segment", async () => {
465+
test("reactor's load() skips mid-file garbage in an extra segment", async () => {
403466
const dir = tempDir();
404467
const store = await createOptimizedContextStore(dir);
405468
const segmentName = segmentFileName(TURNS_FILE, 1);
406469

407470
fs.writeFileSync(path.join(dir, TURNS_FILE), jsonl([turn("a")]));
408-
// Mid-file garbage that is neither null padding nor a torn tail — unrecoverable.
471+
// Mid-file garbage that is neither null padding nor a torn tail — skip it.
409472
fs.writeFileSync(
410473
path.join(dir, segmentName),
411474
jsonl([turn("b")]) + "THIS IS NOT JSON\n" + jsonl([turn("c")]),
412475
);
413476

477+
const loaded = await store.load();
478+
expect(loaded.turns.map((t) => (t.content[0] as { text: string }).text)).toEqual([
479+
"a",
480+
"b",
481+
"c",
482+
]);
483+
});
484+
485+
test("reactor's load() fails closed on schema-invalid lines in an extra segment", async () => {
486+
const dir = tempDir();
487+
const store = await createOptimizedContextStore(dir);
488+
const segmentName = segmentFileName(TURNS_FILE, 1);
489+
490+
fs.writeFileSync(path.join(dir, TURNS_FILE), jsonl([turn("a")]));
491+
const badTurn = JSON.stringify({ role: "user", content: "not-an-array", timestamp: 1 });
492+
fs.writeFileSync(
493+
path.join(dir, segmentName),
494+
jsonl([turn("b")]) + badTurn + "\n" + jsonl([turn("c")]),
495+
);
496+
414497
await expect(store.load()).rejects.toThrow(segmentName);
415498
});
416499
});

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

Lines changed: 100 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -80,11 +80,17 @@ function sanitizeCallId(callId: string): string {
8080
* `fileName` when provided so diagnostics point at the on-disk file, not a bare
8181
* Bun JSON token.
8282
*
83+
* Mid-file lines that fail `JSON.parse` (garbage, glued truncated fragments, or
84+
* interleaved junk) are warned and skipped so resume can continue past them
85+
* (CL-7052). Arktype schema failures stay strict on the reactor path — only
86+
* `skipMalformed` (display-only `loadRecentTurns`) soft-skips those.
87+
*
8388
* `skipMalformed` is for display-only reads (see loadRecentTurns): a bad line
8489
* anywhere in any segment drops that line and keeps the surrounding history,
8590
* because a blank transcript is a worse answer than a transcript with a hole in
8691
* it. The reactor's own load() must never use it — there, history *is* the live
87-
* conversation state and silently dropping a turn would corrupt it (CL-5935).
92+
* conversation state and silently dropping a schema-invalid turn would corrupt
93+
* it (CL-5935).
8894
*/
8995
function parseSegmentTurns(
9096
text: string,
@@ -105,30 +111,108 @@ function parseSegmentTurns(
105111
const line = lines[i]!;
106112
if (line.length === 0) continue;
107113
const isLast = i === lines.length - 1;
108-
let raw: unknown;
114+
const lineNo = i + 1;
115+
116+
let candidates: unknown[];
117+
let fromSalvage = false;
109118
try {
110-
raw = JSON.parse(line);
111-
} catch (cause) {
112-
if (tolerateTornTail && isLast) break;
113-
if (skipMalformed) {
114-
log.warn?.(`skipping malformed JSON at ${fileName} line ${i + 1}`);
119+
candidates = [JSON.parse(line)];
120+
} catch {
121+
candidates = salvageGluedJsonObjects(line);
122+
fromSalvage = true;
123+
if (candidates.length === 0) {
124+
if (tolerateTornTail && isLast) {
125+
log.warn?.(`skipping torn trailing JSON at ${fileName} line ${lineNo}`);
126+
break;
127+
}
128+
log.warn?.(`skipping malformed JSON at ${fileName} line ${lineNo}`);
115129
continue;
116130
}
117-
throw new Error(`${fileName} has malformed JSON at line ${i + 1}`, { cause });
131+
log.warn?.(
132+
`salvaged ${candidates.length} JSON object(s) from glued/malformed line at ${fileName} line ${lineNo}`,
133+
);
118134
}
119-
const result = ConversationTurnSchema(raw);
120-
if (result instanceof type.errors) {
121-
if (skipMalformed) {
122-
log.warn?.(`skipping unexpected structure at ${fileName} line ${i + 1}`);
123-
continue;
135+
136+
for (const raw of candidates) {
137+
const result = ConversationTurnSchema(raw);
138+
if (result instanceof type.errors) {
139+
// Whole-line JSON that fails the turn schema stays strict on the reactor
140+
// path. Salvaged fragments from a glued/garbage line are skipped — they
141+
// are not intentional turn records.
142+
if (skipMalformed || fromSalvage) {
143+
log.warn?.(`skipping unexpected structure at ${fileName} line ${lineNo}`);
144+
continue;
145+
}
146+
throw new Error(
147+
`${fileName} has unexpected structure at line ${lineNo}: ${result.summary}`,
148+
);
124149
}
125-
throw new Error(`${fileName} has unexpected structure at line ${i + 1}: ${result.summary}`);
150+
turns.push(result);
126151
}
127-
turns.push(result);
128152
}
129153
return turns;
130154
}
131155

156+
/**
157+
* Recover zero or more top-level `{...}` values glued on one physical line
158+
* (e.g. a truncated manage_tasks write followed immediately by the next turn
159+
* with no newline). Starts at every `{` so a truncated head that never closes
160+
* does not swallow a later complete object. Incomplete spans and fragments that
161+
* `JSON.parse` rejects are dropped.
162+
*/
163+
function salvageGluedJsonObjects(line: string): unknown[] {
164+
const objects: unknown[] = [];
165+
let searchFrom = 0;
166+
while (searchFrom < line.length) {
167+
const start = line.indexOf("{", searchFrom);
168+
if (start < 0) break;
169+
170+
let depth = 0;
171+
let inString = false;
172+
let escape = false;
173+
let end = -1;
174+
for (let i = start; i < line.length; i++) {
175+
const c = line[i]!;
176+
if (inString) {
177+
if (escape) {
178+
escape = false;
179+
continue;
180+
}
181+
if (c === "\\") {
182+
escape = true;
183+
continue;
184+
}
185+
if (c === '"') inString = false;
186+
continue;
187+
}
188+
if (c === '"') {
189+
inString = true;
190+
continue;
191+
}
192+
if (c === "{") depth++;
193+
else if (c === "}") {
194+
depth--;
195+
if (depth === 0) {
196+
end = i + 1;
197+
break;
198+
}
199+
}
200+
}
201+
if (end < 0) {
202+
// Unclosed from this `{` — try the next candidate start (truncated glue head).
203+
searchFrom = start + 1;
204+
continue;
205+
}
206+
try {
207+
objects.push(JSON.parse(line.slice(start, end)));
208+
searchFrom = end;
209+
} catch {
210+
searchFrom = start + 1;
211+
}
212+
}
213+
return objects;
214+
}
215+
132216
const EMPTY_TOKEN_USAGE = {
133217
input: 0,
134218
output: 0,
@@ -435,7 +519,7 @@ export async function createOptimizedContextStore(dir: string): Promise<ContextS
435519
const basePath = path.join(dir, TURNS_FILE);
436520
if (await pathExists(basePath)) {
437521
const text = await fs.promises.readFile(basePath, "utf-8");
438-
baseTurns = parseSegmentTurns(text, false, TURNS_FILE);
522+
baseTurns = parseSegmentTurns(text, true, TURNS_FILE);
439523
} else {
440524
baseTurns = [];
441525
}

0 commit comments

Comments
 (0)