From 4a899a265fcb23a871444e438a04454c38e9b63e Mon Sep 17 00:00:00 2001 From: Yaowei Zheng Date: Fri, 7 Aug 2026 17:49:40 +0800 Subject: [PATCH] fix(core): serialize Trace appends; tolerate malformed middle lines on read (#234) Co-authored-by: Claude Fable 5 --- packages/core/src/trace/index.ts | 2 +- packages/core/src/trace/resume.ts | 66 +++++++-- packages/core/src/trace/writer.ts | 69 ++++++++-- packages/core/test/replay.test.ts | 128 +++++++++++++++++- packages/core/test/trace.test.ts | 90 +++++++++++- .../docs/content/sessions-and-traces.en.md | 3 +- .../docs/content/sessions-and-traces.zh.md | 3 +- packages/server/src/services/trace-service.ts | 11 +- .../server/test/trace-import-export.test.ts | 5 +- 9 files changed, 344 insertions(+), 33 deletions(-) diff --git a/packages/core/src/trace/index.ts b/packages/core/src/trace/index.ts index a467cb3..0998981 100644 --- a/packages/core/src/trace/index.ts +++ b/packages/core/src/trace/index.ts @@ -7,4 +7,4 @@ export { readTraceTolerant, resumeTrace, } from "./resume.js"; -export type { LocatedTraceFile, ResumeResult } from "./resume.js"; +export type { LocatedTraceFile, ParseTraceLinesOptions, ResumeResult } from "./resume.js"; diff --git a/packages/core/src/trace/resume.ts b/packages/core/src/trace/resume.ts index 46f014c..088c662 100644 --- a/packages/core/src/trace/resume.ts +++ b/packages/core/src/trace/resume.ts @@ -77,13 +77,44 @@ export interface ResumeResult { /** Content of the pairing-backfill placeholder output (the tool hadn't finished and no output was persisted before the process exited). */ const PROCESS_EXIT_PLACEHOLDER = "[interrupted: process exited before the tool finished]"; +/** Options for `parseTraceLines`. */ +export interface ParseTraceLinesOptions { + /** + * Handling of a malformed line **in the middle** of the content (the trailing truncated + * line is always tolerated and ignored, see `parseTraceLines`): + * - "skip" (default): drop that line and keep every parseable record. Reading a Trace back + * is best-effort — the corruption already happened, and throwing away the whole file + * helps no one. Replay is built to stay structurally valid with records missing (crash + * tolerance, pairing backfill), so surviving records still resume. + * - "throw": rethrow the JSON parse error. For validating untrusted input where damage + * should be reported instead of silently repaired (server-side Trace import). + */ + onMalformed?: "skip" | "throw"; + /** Called for each skipped malformed middle line (1-based line number). Not called for the tolerated trailing truncated line. */ + onSkip?: (lineNumber: number, error: unknown) => void; +} + /** * Parse Trace JSONL content. Tolerates a **truncated last line** left behind by an abnormal - * process exit (that line is ignored); corruption in the middle is outside the crash window - * (append-only, single writer), so it throws loudly. + * process exit (that line is silently ignored). A malformed line in the **middle** — e.g. a + * record torn by the pre-fix concurrent-append interleaving (#215) — is skipped by default + * (reported via `onSkip`), keeping every parseable record; pass `onMalformed: "throw"` to + * make middle corruption a hard error instead. */ -export function parseTraceLines(content: string): OmniMessage[] { +export function parseTraceLines(content: string, options?: ParseTraceLinesOptions): OmniMessage[] { + const onMalformed = options?.onMalformed ?? "skip"; const lines = content.split("\n"); + // Index of the last non-empty line, found with one backward scan up front. The malformed + // branch below runs once per skipped line in skip mode, so rescanning the remainder there + // (slice + every) would make widely damaged content O(n^2) — seconds of synchronous CPU on + // the server history path for a large non-JSONL file. + let lastNonEmptyIndex = -1; + for (let i = lines.length - 1; i >= 0; i--) { + if (lines[i]!.trim().length > 0) { + lastNonEmptyIndex = i; + break; + } + } const out: OmniMessage[] = []; for (let i = 0; i < lines.length; i++) { const line = lines[i]!.trim(); @@ -91,17 +122,36 @@ export function parseTraceLines(content: string): OmniMessage[] { try { out.push(JSON.parse(line) as OmniMessage); } catch (err) { - const isLastNonEmpty = lines.slice(i + 1).every((l) => l.trim().length === 0); - if (isLastNonEmpty) break; - throw err; + if (i === lastNonEmptyIndex) break; // truncated tail: expected crash artifact, ignore silently + if (onMalformed === "throw") throw err; + options?.onSkip?.(i + 1, err); } } return out; } -/** Read and parse a Trace file (tolerates a truncated last line). */ +/** Maximum count of skipped line numbers spelled out in the readTraceTolerant diagnostic. */ +const SKIPPED_LINES_SHOWN = 20; + +/** + * Read and parse a Trace file, best-effort: tolerates a truncated last line, and skips + * malformed middle lines (files damaged by the pre-fix concurrent-append bug, #215) while + * keeping every parseable record. Skipped lines are diagnosed once per read on stderr, with + * the spelled-out line-number list capped (a heavily damaged file can skip thousands). + */ export async function readTraceTolerant(path: string): Promise { - return parseTraceLines(await readFile(path, "utf8")); + const skipped: number[] = []; + const messages = parseTraceLines(await readFile(path, "utf8"), { + onSkip: (lineNumber) => skipped.push(lineNumber), + }); + if (skipped.length > 0) { + const shown = skipped.slice(0, SKIPPED_LINES_SHOWN); + const ellipsis = skipped.length > shown.length ? ", …" : ""; + process.stderr.write( + `[trace] skipped ${skipped.length} malformed line(s) [${shown.join(", ")}${ellipsis}] in ${path}\n`, + ); + } + return messages; } /** A located Trace file: its path, containing date-directory name, and index. */ diff --git a/packages/core/src/trace/writer.ts b/packages/core/src/trace/writer.ts index 6df78ba..c14d6e8 100644 --- a/packages/core/src/trace/writer.ts +++ b/packages/core/src/trace/writer.ts @@ -68,9 +68,30 @@ function isRecordable(msg: OmniMessage): boolean { /** * append-only JSONL Trace writer. * - * Single-writer scenario (MVP): concurrency safety isn't required, but every write uses - * `appendFile` (O_APPEND) rather than caching a file handle and seeking to write, avoiding - * overwriting existing content; this also removes the need for an explicit close. + * Every write uses `appendFile` (O_APPEND) rather than caching a file handle and seeking to + * write, avoiding overwriting existing content; this also removes the need for an explicit + * close. + * + * Concurrency: one live Session has exactly one Writer instance, but that instance receives + * appends from **multiple async producers in the same process** — the engine's LLM stream + * driver and each parallel tool execution write independently. `appendFile` splits large + * payloads (multi-MB records such as base64 image Data URLs) into multiple underlying writes, + * so two overlapping appends can interleave mid-record and corrupt the JSONL (#215). All + * mutating operations (`write`, `rotate`) are therefore serialized through one per-instance + * promise chain: each record lands as one uninterrupted append, and a rotation cannot land in + * the middle of an append. Cross-instance/file concurrency does not occur by design **within a + * single server/CLI process**: a Session allows one active run at a time, child sessions write + * their own files, and server-side Trace import only ever creates brand-new files (`wx`). Two + * processes pointed at the same agent data directory are outside this contract (and outside + * supported usage) — the chain cannot cover them. + * + * Error semantics: a failed operation rejects **that caller's** returned promise only; the + * chain itself absorbs the failure so subsequent operations still run (a transient disk error + * must not wedge Trace recording for the rest of the session). A **hung** append is different: + * it blocks this instance's chain until it settles — deliberate head-of-line blocking, because + * ordering cannot be preserved around an append whose outcome is still unknown. Callers + * (context_engine, Session) already treat Trace writes as best-effort and log the surfaced + * error. */ export class Writer { private readonly tracesDir: string; @@ -80,6 +101,8 @@ export class Writer { private index = 1; /** Set true once the date directory has been created for the current file, to avoid a redundant mkdir. */ private ensuredDirForIndex = -1; + /** Serialization chain: mutating operations run strictly in submission order (see class docs). */ + private chain: Promise = Promise.resolve(); constructor(opts: WriterOptions) { this.tracesDir = opts.tracesDir; @@ -94,18 +117,40 @@ export class Writer { return join(this.tracesDir, this.dateDir, fileName); } + /** + * Enqueues one mutating operation on the serialization chain. The returned promise settles + * with that operation's own outcome (so a failure surfaces to its caller), while the chain + * swallows the failure and proceeds with whatever was enqueued next. + */ + private enqueue(op: () => Promise): Promise { + const result = this.chain.then(op); + this.chain = result.then( + () => undefined, + () => undefined, // a failed operation must not wedge the chain + ); + return result; + } + /** * Appends one message. Only written if it's a recordable message; streaming `partial_*` is * skipped. `mkdir -p`s the date directory on the first write to the current file. + * + * The append is serialized on the instance chain, so the record lands as one uninterrupted + * JSONL line even when other writes (or a rotation) are submitted concurrently; the target + * path is resolved when the operation runs, so a write submitted after `rotate()` goes to + * the new file. The returned promise resolves only after this record has been appended + * (callers may rely on write-then-send ordering), and rejects if this append failed. */ async write(msg: OmniMessage): Promise { if (!isRecordable(msg)) return; - const path = this.currentPath(); - if (this.ensuredDirForIndex !== this.index) { - await mkdir(dirname(path), { recursive: true }); - this.ensuredDirForIndex = this.index; - } - await appendFile(path, `${JSON.stringify(msg)}\n`, "utf8"); + return this.enqueue(async () => { + const path = this.currentPath(); + if (this.ensuredDirForIndex !== this.index) { + await mkdir(dirname(path), { recursive: true }); + this.ensuredDirForIndex = this.index; + } + await appendFile(path, `${JSON.stringify(msg)}\n`, "utf8"); + }); } /** Writes multiple messages in sequence. */ @@ -132,10 +177,14 @@ export class Writer { /** * Starts a new Trace file: increments the index, so the next `write` goes to the new file. * Used to split into a separate file when the context is compacted and a new context segment is produced. + * Serialized on the instance chain: appends submitted before the rotation land in the old + * file, appends submitted after it land in the new one — a rotation can never split a record. * Docs: /docs/sessions-and-traces § "Trace design". */ async rotate(): Promise { - this.index += 1; + return this.enqueue(async () => { + this.index += 1; + }); } } diff --git a/packages/core/test/replay.test.ts b/packages/core/test/replay.test.ts index 3294029..cbd5f5f 100644 --- a/packages/core/test/replay.test.ts +++ b/packages/core/test/replay.test.ts @@ -12,15 +12,16 @@ * - Tolerates a truncated trailing line left by an abnormal process exit. * - Round-trip: a Trace written out by the engine, once replayed, matches the history the model actually received. */ -import { mkdtemp, rm } from "node:fs/promises"; +import { mkdir, mkdtemp, rm, writeFile } from "node:fs/promises"; import { tmpdir } from "node:os"; import { join } from "node:path"; -import { afterEach, beforeEach, describe, expect, it } from "vitest"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { assistantText, compactionBegin, compactionEnd, emptyTokenCounts, + imageUrlMessage, requestBegin, requestEnd, sessionMeta, @@ -31,7 +32,7 @@ import { userText, } from "../src/omnimessage/index.js"; import type { OmniMessage, TokenCounts } from "../src/omnimessage/index.js"; -import { parseTraceLines, resumeTrace } from "../src/trace/resume.js"; +import { parseTraceLines, readTraceTolerant, resumeTrace } from "../src/trace/resume.js"; import { ContextEngine } from "../src/engine/context-engine.js"; import { Environment } from "../src/environment/index.js"; import { Writer, readTrace } from "../src/trace/index.js"; @@ -324,15 +325,130 @@ describe("resumeTrace", () => { }); describe("parseTraceLines", () => { + /** + * Content shaped like the #215 corruption: a large record torn in half by a concurrent + * append, with a complete small record and a newline landing in the middle. Line 2 + * (big prefix + small record) and line 3 (big suffix) are both unparseable. + */ + function tornContent(): { content: string; survivors: OmniMessage[] } { + const before = userText("before"); + const after = userText("after"); + const big = JSON.stringify(imageUrlMessage(`data:image/png;base64,${"A".repeat(2048)}`)); + const interleaved = JSON.stringify(tokenUsage(usage(1), usage(1))); + const torn = `${big.slice(0, 100)}${interleaved}\n${big.slice(100)}`; + const content = `${JSON.stringify(before)}\n${torn}\n${JSON.stringify(after)}\n`; + return { content, survivors: [before, after] }; + } + it("tolerates a torn trailing line (process crash mid-write)", () => { const content = `${JSON.stringify(userText("a"))}\n{"timestamp":"2026-07-06T`; - const msgs = parseTraceLines(content); + const skipped: number[] = []; + const msgs = parseTraceLines(content, { onSkip: (line) => skipped.push(line) }); expect(msgs).toHaveLength(1); + // The trailing truncated line is the expected crash artifact — ignored, not diagnosed. + expect(skipped).toEqual([]); }); - it("throws on mid-file corruption", () => { + it("skips a malformed middle line by default, keeping every parseable record (#215)", () => { + const { content, survivors } = tornContent(); + const skipped: number[] = []; + const msgs = parseTraceLines(content, { onSkip: (line) => skipped.push(line) }); + expect(msgs).toEqual(survivors); + expect(skipped).toEqual([2, 3]); + }); + + it('throws on mid-file corruption when onMalformed is "throw" (import validation)', () => { const content = `not-json\n${JSON.stringify(userText("a"))}\n`; - expect(() => parseTraceLines(content)).toThrow(); + expect(() => parseTraceLines(content, { onMalformed: "throw" })).toThrow(); + // The trailing truncated line stays tolerated even in strict mode. + const truncatedTail = `${JSON.stringify(userText("a"))}\n{"timestamp":"2026-07-06T`; + expect(parseTraceLines(truncatedTail, { onMalformed: "throw" })).toHaveLength(1); + }); + + it("handles thousands of malformed lines without quadratic rescans (wide corruption)", () => { + // A large mostly-unparseable file (e.g. a non-JSONL file dropped into the traces dir). + // The tail check must not rescan the remainder per skipped line — with the old + // slice-per-hit check, this content took seconds of synchronous CPU; now it's O(n). + const early = userText("early survivor"); + const late = userText("late survivor"); + const lines: string[] = []; + for (let i = 0; i < 5000; i++) lines.push(`{torn ${i}`); + lines.push(JSON.stringify(early)); + for (let i = 0; i < 5000; i++) lines.push(`{torn ${5000 + i}`); + lines.push(JSON.stringify(late)); + const content = `${lines.join("\n")}\n`; + + const skipped: number[] = []; + const msgs = parseTraceLines(content, { onSkip: (line) => skipped.push(line) }); + expect(msgs).toEqual([early, late]); + // Every malformed line is reported exactly once, in order, none treated as the tail + // (the last non-empty line is the valid "late survivor"). + expect(skipped).toHaveLength(10000); + expect(skipped[0]).toBe(1); + expect(skipped[4999]).toBe(5000); + expect(skipped[5000]).toBe(5002); + expect(skipped[9999]).toBe(10001); + }); +}); + +describe("readTraceTolerant", () => { + let dir: string; + + beforeEach(async () => { + dir = await mkdtemp(join(tmpdir(), "penguin-tolerant-")); + }); + + afterEach(async () => { + await rm(dir, { recursive: true, force: true }); + vi.restoreAllMocks(); + }); + + it("returns the surviving records of a corrupted file without throwing (#215 recovery)", async () => { + // A file damaged by the pre-fix interleaving: a large record torn in half around a + // complete small record. No migration exists — the tolerant read IS the recovery. + const before = userText("before"); + const after = userText("after"); + const big = JSON.stringify(imageUrlMessage(`data:image/png;base64,${"B".repeat(4096)}`)); + const interleaved = JSON.stringify(assistantText("landed mid-record")); + const file = join(dir, "sess_abc_001.jsonl"); + await mkdir(dir, { recursive: true }); + await writeFile( + file, + `${JSON.stringify(before)}\n${big.slice(0, 64)}${interleaved}\n${big.slice(64)}\n${JSON.stringify(after)}\n`, + "utf8", + ); + + const stderr = vi.spyOn(process.stderr, "write").mockImplementation(() => true); + const msgs = await readTraceTolerant(file); + expect(msgs).toEqual([before, after]); + // The skipped lines are diagnosed once, with the file path. + const diagnostics = stderr.mock.calls + .map((c) => String(c[0])) + .filter((s) => s.includes("[trace]")); + expect(diagnostics).toHaveLength(1); + expect(diagnostics[0]).toContain("skipped 2 malformed line(s)"); + expect(diagnostics[0]).toContain(file); + }); + + it("caps the spelled-out line numbers in the diagnostic on heavily damaged files", async () => { + const survivor = userText("survivor"); + const badLines = Array.from({ length: 30 }, (_, i) => `{torn ${i}`); + const file = join(dir, "sess_abc_001.jsonl"); + await mkdir(dir, { recursive: true }); + await writeFile(file, `${badLines.join("\n")}\n${JSON.stringify(survivor)}\n`, "utf8"); + + const stderr = vi.spyOn(process.stderr, "write").mockImplementation(() => true); + const msgs = await readTraceTolerant(file); + expect(msgs).toEqual([survivor]); + const diagnostic = stderr.mock.calls + .map((c) => String(c[0])) + .find((s) => s.includes("[trace]"))!; + // The total is exact, but only the first 20 line numbers are spelled out. + expect(diagnostic).toContain("skipped 30 malformed line(s)"); + const list = /malformed line\(s\) \[([^\]]*)\]/.exec(diagnostic); + expect(list).not.toBeNull(); + const firstTwenty = Array.from({ length: 20 }, (_, i) => i + 1).join(", "); + expect(list![1]).toBe(`${firstTwenty}, …`); }); }); diff --git a/packages/core/test/trace.test.ts b/packages/core/test/trace.test.ts index 377acf5..6c50c20 100644 --- a/packages/core/test/trace.test.ts +++ b/packages/core/test/trace.test.ts @@ -1,4 +1,4 @@ -import { mkdtemp, rm, stat } from "node:fs/promises"; +import { mkdtemp, readFile, rm, stat, writeFile } from "node:fs/promises"; import { tmpdir } from "node:os"; import { join, sep } from "node:path"; @@ -6,6 +6,7 @@ import { afterEach, beforeEach, describe, expect, it } from "vitest"; import { assistantText, + imageUrlMessage, partialText, sessionMeta, subagentEvent, @@ -13,6 +14,7 @@ import { emptyTokenCounts, withOrigin, } from "../src/omnimessage/index.js"; +import type { OmniMessage } from "../src/omnimessage/index.js"; import { Writer, readTrace } from "../src/trace/index.js"; const SESSION_ID = "sess_abc"; @@ -155,4 +157,90 @@ describe("Writer", () => { // append-only verification: the old file's row count does not increase. expect((await readTrace(firstPath)).length).toBe(2); }); + + // #215 regression: appends come from multiple async producers in one process (the LLM stream + // driver plus each parallel tool), and fs.appendFile splits multi-MB payloads into multiple + // underlying writes — without serialization a concurrent append can land mid-record. + describe("concurrency (#215)", () => { + it("keeps a multi-MB record intact under concurrent small appends", async () => { + const writer = new Writer({ + tracesDir, + sessionId: SESSION_ID, + date: new Date(2026, 0, 9), + }); + // Several MB, like a base64 image Data URL — large enough that fs.appendFile splits it + // into multiple underlying writes (Node chunks the data at 512 KiB), small enough to + // keep the test fast. + const big = imageUrlMessage(`data:image/png;base64,${"A".repeat(3 * 1024 * 1024)}`); + const submitted: OmniMessage[] = [ + meta(), + big, + ...[1, 2, 3, 4, 5].map((i) => assistantText(`small ${i}`)), + ]; + // Submit every write before awaiting any: concurrent producers, submission order known. + await Promise.all(submitted.map((msg) => writer.write(msg))); + + // Every non-empty line parses independently, and count and order match submission. + const raw = await readFile(writer.currentPath(), "utf8"); + const lines = raw.split("\n").filter((line) => line.trim().length > 0); + expect(lines).toEqual(submitted.map((msg) => JSON.stringify(msg))); + for (const line of lines) expect(() => JSON.parse(line)).not.toThrow(); + }); + + it("serializes rotate() against in-flight writes: shard boundaries follow submission order", async () => { + const writer = new Writer({ + tracesDir, + sessionId: SESSION_ID, + date: new Date(2026, 0, 9), + }); + const firstPath = writer.currentPath(); + // A large record right before the rotation: the rotation must wait for the whole + // append, and writes submitted after it must land in the next shard. + const big = imageUrlMessage(`data:image/png;base64,${"B".repeat(2 * 1024 * 1024)}`); + const shardOne: OmniMessage[] = [meta(), big]; + const shardTwo: OmniMessage[] = [meta(), assistantText("new context")]; + await Promise.all([ + writer.write(shardOne[0]!), + writer.write(shardOne[1]!), + writer.rotate(), + writer.write(shardTwo[0]!), + writer.write(shardTwo[1]!), + ]); + + const secondPath = writer.currentPath(); + expect(secondPath.endsWith(`${SESSION_ID}_002.jsonl`)).toBe(true); + const firstLines = (await readFile(firstPath, "utf8")) + .split("\n") + .filter((line) => line.trim().length > 0); + const secondLines = (await readFile(secondPath, "utf8")) + .split("\n") + .filter((line) => line.trim().length > 0); + expect(firstLines).toEqual(shardOne.map((msg) => JSON.stringify(msg))); + expect(secondLines).toEqual(shardTwo.map((msg) => JSON.stringify(msg))); + }); + + it("a failed write rejects its own caller without wedging subsequent writes", async () => { + const writer = new Writer({ + tracesDir, + sessionId: SESSION_ID, + date: new Date(2026, 0, 9), + }); + // Occupy the date-dir path with a regular file so mkdir (and the append) must fail. + const dateDirPath = join(tracesDir, "2026-01-09"); + await writeFile(dateDirPath, "blocker", "utf8"); + const results = await Promise.allSettled([ + writer.write(meta()), + writer.write(assistantText("also blocked")), + ]); + // Each failure surfaces to its own caller (callers log it as best-effort)... + expect(results.map((r) => r.status)).toEqual(["rejected", "rejected"]); + + // ...and the chain is not wedged: once the cause clears, later writes land normally. + await rm(dateDirPath); + const recovered = assistantText("after recovery"); + await writer.write(recovered); + const rows = await readTrace(writer.currentPath()); + expect(rows).toEqual([recovered]); + }); + }); }); diff --git a/packages/docs/content/sessions-and-traces.en.md b/packages/docs/content/sessions-and-traces.en.md index cc7bf75..f26dd91 100644 --- a/packages/docs/content/sessions-and-traces.en.md +++ b/packages/docs/content/sessions-and-traces.en.md @@ -54,6 +54,7 @@ A Trace is an append-only JSON Lines file; each line is one OmniMessage envelope - Recorded: `session_meta`, complete `model_msg`, and all `event_msg`. - Not recorded: streaming `partial_*` fragments (the producer appends the complete message once the segment ends), and nested messages tagged with `origin` — a subagent's messages go to the child Session's own Trace, while the parent Trace keeps a single `subagent` pointer event at the spawn site recording the child Session id. - `request_begin` and `request_end(status)` come in pairs delimiting one Request; replay uses `request_end.status === "completed"` as the commit criterion for that turn. +- Appends are serialized inside the writer: records from concurrent producers (the model stream, parallel tool executions) land strictly one after another, each as a single uninterrupted line — a multi-megabyte record such as a base64 image Data URL can never be torn apart by a concurrent append, and file rotation never splits a record. See `packages/core/src/trace/writer.ts` for the implementation. @@ -82,7 +83,7 @@ The Trace is the single source of truth for recovery — there is no separate se 4. Reconstruct the carry-over (undelivered tool outputs, interruption markers) plus turn and Token counters; 5. Continue appending to the same Trace file. -Recovery requires that the Workspace and the model still exist. What recovery guarantees is structural legality: only committed turns are replayed, with `tool_call` / `tool_call_output` pairing intact; incomplete model output (thinking, text) is allowed to be lost. A truncated last line left by an abnormal process exit is tolerated and ignored. See `packages/core/src/trace/resume.ts`. +Recovery requires that the Workspace and the model still exist. What recovery guarantees is structural legality: only committed turns are replayed, with `tool_call` / `tool_call_output` pairing intact; incomplete model output (thinking, text) is allowed to be lost. A truncated last line left by an abnormal process exit is tolerated and ignored; a malformed line in the middle of a file (e.g. in a file damaged before writer appends were serialized) is skipped with a diagnostic on stderr, and every parseable record is kept. See `packages/core/src/trace/resume.ts`. Special case: if the latest Trace file ends with a completed compaction, that context is closed as a whole — resume starts from an empty context; in summarize mode the `[context_summary]` is reconstructed and prepended to the first input after resume (old Traces using the earlier angle-bracket `` form are still understood). diff --git a/packages/docs/content/sessions-and-traces.zh.md b/packages/docs/content/sessions-and-traces.zh.md index ab9e002..ac61e11 100644 --- a/packages/docs/content/sessions-and-traces.zh.md +++ b/packages/docs/content/sessions-and-traces.zh.md @@ -52,6 +52,7 @@ Trace 是 append-only 的 JSON Lines 文件,每行一个 OmniMessage 信封( - 记录的消息:`session_meta`、完整的 `model_msg`、全部 `event_msg`。 - 不记录的消息:流式 `partial_*` 分片(片段结束后由生产方补写完整消息);带 `origin` 标记的嵌套消息——子 Agent 的消息写入子 Session 自己的 Trace,父 Trace 只在派生位置保留一个 `subagent` 指针事件,记录子 Session id。 - `request_begin` 与 `request_end(status)` 成对出现,界定一轮 Request;回放以 `request_end.status === "completed"` 作为该轮已提交的判据。 +- 追加在写入器内部串行执行:并发生产者(模型流、并行工具执行)的记录严格逐条落盘,每条都是一行完整内容——base64 图片 Data URL 之类的多 MB 大记录不会被并发追加撕裂,文件轮转也不会切断任何记录。 实现见 `packages/core/src/trace/writer.ts`。 @@ -80,7 +81,7 @@ Trace 是恢复的唯一事实来源,没有独立的会话数据库需要与 4. 重建 carry-over(未送达的工具输出、中断标记)与轮数、Token 计数器; 5. 继续追加写入同一个 Trace 文件。 -恢复的前提是 Workspace 与模型仍然存在。恢复保证的是结构合法性:只回放已提交的轮次,`tool_call` 与 `tool_call_output` 配对完整;未完成的模型输出(thinking、文本)允许丢失。异常退出留下的截断末行会被容忍并忽略。实现见 `packages/core/src/trace/resume.ts`。 +恢复的前提是 Workspace 与模型仍然存在。恢复保证的是结构合法性:只回放已提交的轮次,`tool_call` 与 `tool_call_output` 配对完整;未完成的模型输出(thinking、文本)允许丢失。异常退出留下的截断末行会被容忍并忽略;文件中间的损坏行(如写入器追加串行化之前受损的存量文件)会被跳过并在 stderr 给出诊断,其余可解析的记录全部保留。实现见 `packages/core/src/trace/resume.ts`。 特殊情形:若最新 Trace 文件以一次完成的压缩收尾,则该上下文已整体关闭——恢复从空上下文开始;summarize 模式下会重建 `[context_summary]` 摘要,前置到恢复后第一轮输入中(旧 Trace 中早期的尖括号 `` 形式仍可识别)。 diff --git a/packages/server/src/services/trace-service.ts b/packages/server/src/services/trace-service.ts index 8125f64..06e2180 100644 --- a/packages/server/src/services/trace-service.ts +++ b/packages/server/src/services/trace-service.ts @@ -2,9 +2,10 @@ * Trace service. * * History messages: all of the Session's index files concatenated in order - * (readTraceTolerant, tolerating a truncated last line), containing only the - * complete messages and events that were actually written to Trace (naturally - * excluding partial_*); in-flight increments are continued by SSE. + * (readTraceTolerant, tolerating a truncated last line and skipping malformed + * middle lines), containing only the complete messages and events that were + * actually written to Trace (naturally excluding partial_*); in-flight + * increments are continued by SSE. * Performance analysis is derived from a single Trace file: nearest-neighbor * pairing of request_begin/end, tool call duration pairing, reconnect / compaction * counts, and Token trend. @@ -1206,7 +1207,9 @@ export class TraceService { const invalid = (message: string) => new HttpError(400, "invalid_trace", message); let records: OmniMessage[]; try { - records = parseTraceLines(content); + // Strict parse: import is the gate for user-supplied files, so a malformed middle + // line is reported as 400 rather than silently dropped (read paths skip it instead). + records = parseTraceLines(content, { onMalformed: "throw" }); } catch { throw invalid("The file is not valid Trace JSONL."); } diff --git a/packages/server/test/trace-import-export.test.ts b/packages/server/test/trace-import-export.test.ts index 58dfa4b..33d6691 100644 --- a/packages/server/test/trace-import-export.test.ts +++ b/packages/server/test/trace-import-export.test.ts @@ -189,7 +189,10 @@ describe("trace-import-export", () => { it("import: malformed middle line → 400 invalid_trace", async () => { const lines = toContent(sampleTrace(SID)).split("\n"); - lines[1] = "{not json"; // corrupt a middle line (non-final: parseTraceLines throws loudly) + // Corrupt a middle (non-final) line. Read paths skip malformed middle lines (#215 + // recovery), but import validates strictly (onMalformed: "throw") — damage in an + // uploaded file is reported, not silently dropped. + lines[1] = "{not json"; const res = await owner.post(`${base()}/import`, { dataBase64: b64(lines.join("\n")) }); expect(res.status).toBe(400); expect(await errorCode(res)).toBe("invalid_trace");