From 842a929691b1e89afff7a9d5bab35cf96ad571dd Mon Sep 17 00:00:00 2001 From: Leo Date: Sat, 22 Aug 2026 21:31:38 +0800 Subject: [PATCH 1/3] =?UTF-8?q?=E2=9A=A1=20fix(sidecar):=20settleToolCall?= =?UTF-8?q?=20=E6=94=B9=20append-only=20=E4=BE=A7=E8=BD=A6=E8=A1=8C?= =?UTF-8?q?=EF=BC=8C=E6=B6=88=E6=AF=8F=20tool=5Fresult=20=E5=85=A8?= =?UTF-8?q?=E9=87=8F=E9=87=8D=E5=86=99=E7=9A=84=20O(N=C2=B2)=20IO=20(#411?= =?UTF-8?q?=E2=91=A0)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 每条 tool_result 触发的整文件 JSON.parse+原子重写在长 run 下是 O(N²) IO。 改为追加 tool_settled 侧车行(O(1)),get() 投影时物化回 tool_call 行, finalize compactModelStreamItems 一次重写完成裁剪+物化收敛;countItems 排除侧车行。运行中语义不变(#256 终态契约测试保持通过)。 --- .../runner/run-state-store.test.ts | 39 +++++++++++ .../agent-runtime/runner/run-state-store.ts | 64 ++++++++++++++----- 2 files changed, 88 insertions(+), 15 deletions(-) diff --git a/apps/sidecar/src/services/agent-runtime/runner/run-state-store.test.ts b/apps/sidecar/src/services/agent-runtime/runner/run-state-store.test.ts index 3ceaa881e..3b5aecb96 100644 --- a/apps/sidecar/src/services/agent-runtime/runner/run-state-store.test.ts +++ b/apps/sidecar/src/services/agent-runtime/runner/run-state-store.test.ts @@ -171,6 +171,45 @@ describe("run-state-store", () => { expect(after).toEqual(before); }); + test("settleToolCall 侧车行:get 投影物化终态,countItems 排除,compact 收敛", async () => { + const dir = mkdtempSync(join(tmpdir(), "lume-run-state-store-settle-")); + const store = createFileBackedLumeRunStateStore(dir); + await store.create(makeState("run-settle", "running")); + await store.appendItem("run-settle", { + type: "tool_call", id: "tu-1", toolName: "bash", input: {}, parentAgentId: "root", status: "pending", createdAt: "2026-04-29T00:00:01.000Z" + }); + await store.appendItem("run-settle", { + type: "tool_call", id: "tu-2", toolName: "edit", input: {}, parentAgentId: "root", status: "pending", createdAt: "2026-04-29T00:00:02.000Z" + }); + + // 侧车 append(O(1)):get 投影物化终态回 tool_call 行 + await store.settleToolCall("run-settle", "tu-1", "completed", "2026-04-29T00:00:03.000Z"); + await store.settleToolCall("run-settle", "tu-2", "failed", "2026-04-29T00:00:04.000Z"); + + const settled = await store.get("run-settle"); + const byId = new Map(settled?.generatedItems.map((item) => [item.type === "tool_call" ? item.id : item.id, item])); + expect((byId.get("tu-1") as { status?: string; endedAt?: string }).status).toBe("completed"); + expect((byId.get("tu-1") as { endedAt?: string }).endedAt).toBe("2026-04-29T00:00:03.000Z"); + expect((byId.get("tu-2") as { status?: string }).status).toBe("failed"); + // 侧车行不作为独立 item 暴露 + expect(settled?.generatedItems.every((item) => item.type !== ("tool_settled" as never))).toBe(true); + // 计数排除侧车行 + expect(await store.countItems("run-settle")).toBe(2); + + // 孤儿侧车行(目标 tool_call 不存在)投影时静默丢弃 + await store.settleToolCall("run-settle", "tu-orphan", "completed", "2026-04-29T00:00:05.000Z"); + const orphaned = await store.get("run-settle"); + expect(orphaned?.generatedItems).toHaveLength(2); + + // compact 一并物化收敛:重写后文件无侧车行且 tool_call 自带终态 + await store.appendItem("run-settle", { type: "model_stream", id: "d1", event: { type: "stream_event" } as never, createdAt: "2026-04-29T00:00:06.000Z" }); + await store.appendItem("run-settle", { type: "assistant_message", id: "a1", content: [{ type: "text", text: "final" }] as never, createdAt: "2026-04-29T00:00:07.000Z" }); + await store.compactModelStreamItems("run-settle"); + const compacted = await store.get("run-settle"); + expect(compacted?.generatedItems.map((item) => item.type)).toEqual(["tool_call", "tool_call", "assistant_message"]); + expect(((compacted?.generatedItems[0] ?? {}) as { status?: string }).status).toBe("completed"); + }); + test("todo 快照 save/read 往返,readLatestTodoState 优先快照", async () => { const dir = mkdtempSync(join(tmpdir(), "lume-run-state-store-")); const store = createFileBackedLumeRunStateStore(dir); diff --git a/apps/sidecar/src/services/agent-runtime/runner/run-state-store.ts b/apps/sidecar/src/services/agent-runtime/runner/run-state-store.ts index 6c0940483..6ea5d400c 100644 --- a/apps/sidecar/src/services/agent-runtime/runner/run-state-store.ts +++ b/apps/sidecar/src/services/agent-runtime/runner/run-state-store.ts @@ -72,6 +72,38 @@ function readJsonlFile(path: string): T[] { .filter((item): item is T => item !== null); } +/** + * tool_call 终态侧车行:settleToolCall 只 append 此记录(O(1)), + * 读取投影时物化回对应 tool_call 行,finalize compact 时随重写一次性收敛掉。 + */ +interface ToolSettleRecord { + type: "tool_settled"; + toolCallId: string; + status: "completed" | "failed"; + endedAt: string; +} + +/** 物化侧车终态到 tool_call 行并剔除侧车行;无侧车行时原数组原样返回(零分配快路)。 */ +function projectSettledItems(items: LumeRunItem[]): LumeRunItem[] { + let settles: Map | null = null; + const retained: LumeRunItem[] = []; + for (const item of items) { + const record = item as unknown as Partial; + if (record.type === "tool_settled" && typeof record.toolCallId === "string") { + // 同 id 多条取最新(后到终态胜出) + (settles ??= new Map()).set(record.toolCallId, record as ToolSettleRecord); + } else { + retained.push(item); + } + } + if (!settles) return items; + return retained.map((item) => { + if (item.type !== "tool_call") return item; + const settle = settles.get(item.id); + return settle ? { ...item, status: settle.status, endedAt: settle.endedAt } : item; + }); +} + class FileBackedLumeRunStateStore implements LumeRunStateStore { private readonly sessionDir: string; private readonly runsDir: string; @@ -101,7 +133,7 @@ class FileBackedLumeRunStateStore implements LumeRunStateStore { if (!state) return null; return { ...state, - generatedItems: readJsonlFile(this.itemsPath(runId)) + generatedItems: projectSettledItems(readJsonlFile(this.itemsPath(runId))) }; } @@ -136,17 +168,12 @@ class FileBackedLumeRunStateStore implements LumeRunStateStore { } async settleToolCall(runId: string, toolCallId: string, status: "completed" | "failed", endedAt: string): Promise { - const path = this.itemsPath(runId); - if (!existsSync(path)) return; - const items = readJsonlFile(path); - let settled = false; - const next = items.map((item) => { - if (item.type !== "tool_call" || item.id !== toolCallId || settled) return item; - settled = true; - return { ...item, status, endedAt } as LumeRunItem; - }); - if (!settled) return; - writeTextAtomic(path, next.map((item) => JSON.stringify(item)).join("\n") + "\n"); + // append-only 侧车行(O(1)):不再全量读改写 items.jsonl——每条 tool_result 触发的 + // 整文件重写在长 run 下是 O(N²) IO。物化时机:读取投影(get)与 finalize compact。 + // 目标行不存在时记录成孤儿侧车行,投影时自然丢弃——与原"找不到静默"语义一致。 + if (!existsSync(this.itemsPath(runId))) return; + const record: ToolSettleRecord = { type: "tool_settled", toolCallId, status, endedAt }; + appendFileSync(this.itemsPath(runId), JSON.stringify(record) + "\n", "utf-8"); } async listByThread(threadId: string): Promise { @@ -180,7 +207,11 @@ class FileBackedLumeRunStateStore implements LumeRunStateStore { async countItems(runId: string): Promise { const path = this.itemsPath(runId); if (!existsSync(path)) return 0; - return readFileSync(path, "utf-8").split("\n").filter((line) => line.trim().length > 0).length; + // 侧车行不是用户可见 item,排除之(JSON.stringify 无空格序列化,子串匹配可靠) + return readFileSync(path, "utf-8") + .split("\n") + .filter((line) => line.trim().length > 0 && !line.includes('"type":"tool_settled"')) + .length; } async compactModelStreamItems(runId: string): Promise { @@ -189,8 +220,11 @@ class FileBackedLumeRunStateStore implements LumeRunStateStore { const items = readJsonlFile(path); // 与 hydrate 投影同一判定:无 assistant_message 的 run 依赖 model_stream 重建文本,不裁 if (items.length === 0 || !runHasAssistantMessage(items)) return; - const retained = items.filter((item) => item.type !== "model_stream"); - if (retained.length === items.length) return; + // finalize 收敛一并物化侧车终态:重写后文件回归纯净形态(tool_call 行自带终态,无侧车行) + const filtered = items.filter((item) => item.type !== "model_stream"); + const retained = projectSettledItems(filtered); + // 无 delta 可裁且无侧车行待物化(投影快路原引用返回)才免重写 + if (filtered.length === items.length && retained === filtered) return; writeTextAtomic(path, retained.map((item) => JSON.stringify(item)).join("\n") + "\n"); } From 253ec7c441d9192dbffb80eb65bd1d87b4ffdf4b Mon Sep 17 00:00:00 2001 From: Leo Date: Sat, 22 Aug 2026 21:31:50 +0800 Subject: [PATCH 2/3] =?UTF-8?q?=E2=9A=A1=20fix(sidecar):=20=E4=BA=8B?= =?UTF-8?q?=E4=BB=B6=E6=80=BB=E7=BA=BF=E8=AF=BB=E8=B7=AF=E5=BE=84=E5=A2=9E?= =?UTF-8?q?=E9=87=8F=E5=8C=96=E2=80=94=E2=80=94hasEvents=20=E5=8F=AA?= =?UTF-8?q?=E8=AF=BB=E9=A6=96=E5=9D=97=EF=BC=8Cread=20=E6=8C=89=E6=96=87?= =?UTF-8?q?=E4=BB=B6=20offset=20=E5=A2=9E=E9=87=8F=20(#411=E2=91=A1)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit hasEvents 读全量只为看首行:改为定位读头部 64KB,头部无完整非空行时 回退全量判定,F4 分叉语义不变。GET_EVENTS 每次 afterSeq 增量轮询仍整文件 读+parse:ThreadState 记录字节水位(readOffset/maxSeqSeen),afterSeq 不早于 水位时只读新字节段;早于水位或 releaseThread 后回退全量,任意 afterSeq 语义保持。毒行截断与残尾处理与 readFile 严格一致。 --- .../events/thread-event-bus.test.ts | 20 +++ .../agent-runtime/events/thread-event-bus.ts | 125 +++++++++++++++++- 2 files changed, 143 insertions(+), 2 deletions(-) diff --git a/apps/sidecar/src/services/agent-runtime/events/thread-event-bus.test.ts b/apps/sidecar/src/services/agent-runtime/events/thread-event-bus.test.ts index 41955d0e3..1fcea56ac 100644 --- a/apps/sidecar/src/services/agent-runtime/events/thread-event-bus.test.ts +++ b/apps/sidecar/src/services/agent-runtime/events/thread-event-bus.test.ts @@ -119,6 +119,26 @@ describe("ThreadEventBus", () => { expect(received.map((e) => e.phase)).toEqual(["update", "end"]) }) + test("增量 read(#411):水位内 afterSeq 走快路,结果与全量一致;早于水位回退全量", async () => { + dir = mkdtempSync(join(tmpdir(), "bus-incr-")) + const bus = getThreadEventBus(dir) + await bus.publish("th1", "r1", skeletonEvent("run", "start")) + await bus.publish("th1", "r1", skeletonEvent("message", "end")) + // 首次 read 建立水位(全量) + expect((await bus.read("th1")).map((e) => e.seq)).toEqual([1, 2]) + // 水位内增量:只读新字节段 + await bus.publish("th1", "r1", skeletonEvent("turn", "end")) + const inc = await bus.read("th1", 2) + expect(inc.map((e) => e.seq)).toEqual([3]) + // 早于水位的 afterSeq:回退全量仍可答 + const old = await bus.read("th1", 0) + expect(old.map((e) => e.seq)).toEqual([1, 2, 3]) + // releaseThread 后水位随 state 失效,重建后全量续读不丢 + releaseThreadEventBus(dir, "th1") + await bus.publish("th1", "r2", skeletonEvent("run", "end")) + expect((await bus.read("th1")).map((e) => e.seq)).toEqual([1, 2, 3, 4]) + }) + test("hasEvents 与 readFile 截断语义严格一致(F4 分叉判空捷径)", async () => { dir = mkdtempSync(join(tmpdir(), "bus-")) const bus = new ThreadEventBus(dir) diff --git a/apps/sidecar/src/services/agent-runtime/events/thread-event-bus.ts b/apps/sidecar/src/services/agent-runtime/events/thread-event-bus.ts index 17e96ac88..bf1a4c63c 100644 --- a/apps/sidecar/src/services/agent-runtime/events/thread-event-bus.ts +++ b/apps/sidecar/src/services/agent-runtime/events/thread-event-bus.ts @@ -5,7 +5,7 @@ * 后续非 update 相位、PERSIST_COALESCE_MS 窗口或 releaseThread 落盘。 * 此前每个流式 delta 的累计全文都同步写盘,长会话 events 体积近二次增长(~74% 字节浪费)。 */ -import { appendFileSync, existsSync, mkdirSync, readFileSync, writeFileSync } from "node:fs" +import { appendFileSync, closeSync, existsSync, fstatSync, mkdirSync, openSync, readFileSync, readSync, writeFileSync } from "node:fs" import { join } from "node:path" import type { SdkEventEnvelope, SdkLifecycleEvent } from "@lume/shared" @@ -15,6 +15,9 @@ const UPDATE_COALESCE_MS = 16 /** update 相位持久折叠窗口(ms):崩溃时最多丢该窗口内的流式过渡态(终值由非 update 相位兜底落盘) */ const PERSIST_COALESCE_MS = 500 +/** hasEvents 判空只读文件头部字节数;首条事件通常是小事件(run.started),超出即回退全量判定 */ +const HAS_EVENTS_HEAD_BYTES = 64 * 1024 + interface ThreadState { /** 下一个待分配的 seq(初始化 = 文件最后一条合法行 seq + 1) */ nextSeq: number @@ -25,6 +28,10 @@ interface ThreadState { /** 持久折叠缓冲:与 coalesceBuffer 同 key 语义,但负责落盘——折叠后磁盘每个 key 每窗口只有最新累计态 */ persistBuffer: Map persistTimer: ReturnType | null + /** 增量读水位:已消费到的文件字节偏移(null=未建立,read 走全量);releaseThread 后随 state 失效 */ + readOffset: number | null + /** readOffset 之前的最大 seq——afterSeq ≥ 该值才可走增量快路(更早的 afterSeq 需要重读全量) */ + maxSeqSeen: number } export class ThreadEventBus { @@ -90,17 +97,96 @@ export class ThreadEventBus { /** 快照/续传:seq > afterSeq 的全部事件(文件 + 未落盘的持久折叠缓冲,按 seq 归并)。 */ async read(threadId: string, afterSeq?: number): Promise { - const all = [...this.readFile(threadId), ...this.pendingEnvelopes(threadId)].sort((a, b) => a.seq - b.seq) + const st = this.threads.get(threadId) + let fileEvents: SdkEventEnvelope[] + if ( + st?.readOffset != null + && (afterSeq ?? 0) >= st.maxSeqSeen + && existsSync(this.file(threadId)) + ) { + // 增量快路:活跃线程(publish 过)且 afterSeq 不早于已消费水位——只读新字节段。 + const segment = this.readSegmentFrom(this.file(threadId), st.readOffset) + if (segment) { + fileEvents = segment.envelopes + st.readOffset = segment.nextOffset + const last = fileEvents[fileEvents.length - 1] + if (last) st.maxSeqSeen = Math.max(st.maxSeqSeen, last.seq) + } else { + // 文件被替换/截断:回退全量并重立水位 + fileEvents = this.readFile(threadId) + this.resetReadWatermark(st, threadId) + } + } else { + fileEvents = this.readFile(threadId) + if (st) this.resetReadWatermark(st, threadId) + } + const all = [...fileEvents, ...this.pendingEnvelopes(threadId)].sort((a, b) => a.seq - b.seq) return afterSeq === undefined ? all : all.filter((e) => e.seq > afterSeq) } + /** 全量读后建立增量水位:offset=文件末尾,maxSeqSeen=最后一条合法行 seq。 */ + private resetReadWatermark(st: ThreadState, threadId: string): void { + const envelopes = this.readFile(threadId) + st.maxSeqSeen = envelopes[envelopes.length - 1]?.seq ?? 0 + let fd: number | undefined + try { + fd = openSync(this.file(threadId), "r") + st.readOffset = fstatSync(fd).size + } catch { + st.readOffset = null + } finally { + if (fd !== undefined) closeSync(fd) + } + } + + /** + * 从 start 字节偏移读到 EOF 并解析完整行;返回事件与新偏移(停在最后一个 \n 之后, + * 残尾留待下次——append-only 下只会被补全)。文件短于 start(被替换/截断)返回 null。 + */ + private readSegmentFrom(file: string, start: number): { envelopes: SdkEventEnvelope[]; nextOffset: number } | null { + let fd: number | undefined + try { + fd = openSync(file, "r") + const size = fstatSync(fd).size + if (size < start) return null + const buffer = Buffer.alloc(size - start) + let total = 0 + while (total < buffer.length) { + const n = readSync(fd, buffer, total, buffer.length - total, start + total) + if (n <= 0) break + total += n + } + const text = buffer.toString("utf8", 0, total) + const lastNewline = text.lastIndexOf("\n") + if (lastNewline === -1) return { envelopes: [], nextOffset: start } + const complete = text.slice(0, lastNewline + 1) + const out: SdkEventEnvelope[] = [] + for (const line of complete.split("\n")) { + if (!line) continue + try { + out.push(JSON.parse(line) as SdkEventEnvelope) + } catch { + break + } + } + return { envelopes: out, nextOffset: start + Buffer.byteLength(complete, "utf8") } + } catch { + return null + } finally { + if (fd !== undefined) closeSync(fd) + } + } + /** * 判空捷径(F4 分叉用):与 readFile 的截断语义严格一致——逐行找到第一条非空行, * JSON.parse 成功即有事件、失败(毒行)即无;全空行/文件缺失为无。不做全量对象分配。 + * 只读头部 HAS_EVENTS_HEAD_BYTES:首条事件通常是小事件,判空 O(1);头部无定论时回退全量。 */ hasEvents(threadId: string): boolean { const file = this.file(threadId) if (existsSync(file)) { + const verdict = this.headHasEventVerdict(file) + if (verdict !== "unknown") return verdict === "yes" for (const line of readFileSync(file, "utf8").split("\n")) { if (!line) continue try { @@ -115,6 +201,39 @@ export class ThreadEventBus { return this.pendingEnvelopes(threadId).length > 0 } + /** 头部判定:"yes"=首条非空行是合法事件;"no"=毒行;"unknown"=头部无完整非空行。 */ + private headHasEventVerdict(file: string): "yes" | "no" | "unknown" { + let fd: number | undefined + try { + fd = openSync(file, "r") + const size = fstatSync(fd).size + const buffer = Buffer.alloc(Math.min(size, HAS_EVENTS_HEAD_BYTES)) + let total = 0 + while (total < buffer.length) { + const n = readSync(fd, buffer, total, buffer.length - total, total) + if (n <= 0) break + total += n + } + const text = buffer.toString("utf8", 0, total) + // 末段可能被截断(头部边界切在行中),仅判定以 \n 收尾的完整行 + const lines = text.endsWith("\n") ? text.split("\n") : text.slice(0, text.lastIndexOf("\n") + 1).split("\n") + for (const line of lines) { + if (!line) continue + try { + JSON.parse(line) + return "yes" + } catch { + return "no" + } + } + return "unknown" + } catch { + return "unknown" + } finally { + if (fd !== undefined) closeSync(fd) + } + } + private file(threadId: string): string { return join(this.sessionDir, `${threadId}.events.jsonl`) } @@ -132,6 +251,8 @@ export class ThreadEventBus { coalesceTimer: null, persistBuffer: new Map(), persistTimer: null, + readOffset: null, + maxSeqSeen: 0, } this.threads.set(threadId, st) } From 5309700f0426eaf9b24c57cafbb6964685c2ba87 Mon Sep 17 00:00:00 2001 From: Leo Date: Sat, 22 Aug 2026 21:32:22 +0800 Subject: [PATCH 3/3] =?UTF-8?q?=F0=9F=90=9B=20fix(sidecar):=20waiting=5Fba?= =?UTF-8?q?ckground=20checkpoint=20=E5=B4=A9=E6=BA=83=E6=81=A2=E5=A4=8D?= =?UTF-8?q?=E6=AD=BB=E8=B7=AF=E2=80=94=E2=80=94=E6=8C=89=E6=8C=81=E4=B9=85?= =?UTF-8?q?=E5=8C=96=E7=BB=88=E6=80=81=E9=80=9A=E7=9F=A5=E7=9C=9F=E6=AD=A3?= =?UTF-8?q?=E7=BB=AD=E8=B7=91=20(#411=E2=91=A2)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 崩溃后 live 终态监听(run.handleAsyncEvent)已不在,checkpoint 永停 waiting_background;恢复横幅出现但 RESUME_RUN 命中 no-op 分支,文案谎称 "已重新附着"实际什么都没做。现经注入的 resolveBackgroundNotification 查 background-process-recovery 落盘的终态通知:有则与 live 同形转换 (ready_to_resume + syntheticToolResult)真正续跑,无则如实返回等待态。 --- apps/sidecar/src/rpc/agent-handlers.ts | 12 +++ .../interruption/resume-service.test.ts | 89 +++++++++++++++++++ .../interruption/resume-service.ts | 48 ++++++++-- 3 files changed, 144 insertions(+), 5 deletions(-) diff --git a/apps/sidecar/src/rpc/agent-handlers.ts b/apps/sidecar/src/rpc/agent-handlers.ts index eb6720809..2f5af626f 100644 --- a/apps/sidecar/src/rpc/agent-handlers.ts +++ b/apps/sidecar/src/rpc/agent-handlers.ts @@ -14,6 +14,7 @@ import { deleteAgentThread, getAgentThreadMeta, getAgentThreadMessages, + getAgentThreadSDKMessages, getRecentAgentThreadMessages, listAllAgentThreads, listAgentThreads, @@ -419,6 +420,17 @@ export function createAgentHandlers( { runStateStore, continuationStore, + // 崩溃恢复(#411③):后台任务终态通知由 background-process-recovery 落盘 + // transcript,按 processJobId 取回供 waiting_background checkpoint 转换 + resolveBackgroundNotification: async (processJobId) => { + const notification = getAgentThreadSDKMessages(input.threadId).find( + (message) => + message.type === "system" && + message.subtype === "task_notification" && + message.task_id === processJobId, + ); + return notification; + }, }, async (checkpoint, state) => { // interrupted(软中止 checkpoint):engine 已为被中断工具补 error 占位 diff --git a/apps/sidecar/src/services/agent-runtime/interruption/resume-service.test.ts b/apps/sidecar/src/services/agent-runtime/interruption/resume-service.test.ts index 32c46b252..da4202a4e 100644 --- a/apps/sidecar/src/services/agent-runtime/interruption/resume-service.test.ts +++ b/apps/sidecar/src/services/agent-runtime/interruption/resume-service.test.ts @@ -310,6 +310,95 @@ describe("LumeResumeService", () => { expect((await continuationStore.get("run-1"))?.status).toBe("resumed"); }); + test("waiting_background(#411③):有持久化终态通知时转换为可续跑 checkpoint", async () => { + const dir = mkdtempSync(join(tmpdir(), "lume-resume-bg-terminal-")); + const runStateStore = createFileBackedLumeRunStateStore(dir); + const continuationStore = createFileBackedRunContinuationStore(dir); + await runStateStore.create(makeRunState()); + await continuationStore.upsert({ + version: 2, + runId: "run-1", + threadId: "thread-1", + status: "waiting_background", + checkpoint: { + step: "waiting_for_tool_result", + toolCallId: "tool-1", + toolName: "Bash", + toolKind: "execute", + processJobId: "job-1" + }, + reason: "后台命令已持久化,恢复时重新附着而不重复执行。", + createdAt: "2026-04-29T00:00:00.000Z", + updatedAt: "2026-04-29T00:00:00.000Z" + }); + + let received: RunContinuationState | undefined; + const result = await new LumeResumeService( + { + runStateStore, + continuationStore, + resolveBackgroundNotification: async (processJobId) => ({ + type: "system", + subtype: "task_notification", + task_id: processJobId, + tool_use_id: "tool-1", + status: "completed", + message: "done", + session_id: "thread-1" + } as any), + }, + async (checkpoint) => { + received = checkpoint; + return { finalOutput: "resumed" }; + } + ).resumeRun({ runId: "run-1" }); + + expect(result.status).toBe("resumed"); + // 与 live handleAsyncEvent 同形:syntheticToolResult 从终态通知构造,step 推进 + expect(received?.status).toBe("ready_to_resume"); + expect(received?.checkpoint.step).toBe("after_tool_result"); + const synthetic = received?.checkpoint.syntheticToolResult as Record; + expect(synthetic.tool_use_id).toBe("tool-1"); + expect(synthetic.content).toBe("done"); + expect(synthetic.is_error).toBeUndefined(); + expect((await continuationStore.get("run-1"))?.status).toBe("resumed"); + }); + + test("waiting_background(#411③):无终态通知时如实返回等待态且不改动 checkpoint", async () => { + const dir = mkdtempSync(join(tmpdir(), "lume-resume-bg-waiting-")); + const runStateStore = createFileBackedLumeRunStateStore(dir); + const continuationStore = createFileBackedRunContinuationStore(dir); + await runStateStore.create(makeRunState()); + await continuationStore.upsert({ + version: 2, + runId: "run-1", + threadId: "thread-1", + status: "waiting_background", + checkpoint: { + step: "waiting_for_tool_result", + toolCallId: "tool-1", + toolName: "Bash", + toolKind: "execute", + processJobId: "job-1" + }, + reason: "后台命令已持久化,恢复时重新附着而不重复执行。", + createdAt: "2026-04-29T00:00:00.000Z", + updatedAt: "2026-04-29T00:00:00.000Z" + }); + + const result = await new LumeResumeService( + { + runStateStore, + continuationStore, + resolveBackgroundNotification: async () => undefined, + }, + async () => ({ finalOutput: "should-not-run" }) + ).resumeRun({ runId: "run-1" }); + + expect(result.status).toBe("waiting_background"); + expect((await continuationStore.get("run-1"))?.status).toBe("waiting_background"); + }); + test("does not replay a V2 side-effect tool with an unknown result", async () => { const dir = mkdtempSync(join(tmpdir(), "lume-resume-v2-unknown-")); const runStateStore = createFileBackedLumeRunStateStore(dir); diff --git a/apps/sidecar/src/services/agent-runtime/interruption/resume-service.ts b/apps/sidecar/src/services/agent-runtime/interruption/resume-service.ts index e7b2de41b..6e8fdd6b4 100644 --- a/apps/sidecar/src/services/agent-runtime/interruption/resume-service.ts +++ b/apps/sidecar/src/services/agent-runtime/interruption/resume-service.ts @@ -1,3 +1,4 @@ +import type { SDKMessage } from "@lume/shared"; import type { RunContinuationState } from "../runner/run-continuation"; import type { LumeRunState } from "../runner/run-state"; import type { RunContinuationStore } from "../runner/run-continuation-store"; @@ -23,6 +24,12 @@ export class LumeResumeService { private readonly stores?: { runStateStore: LumeRunStateStore; continuationStore: RunContinuationStore; + /** + * 崩溃恢复(#411③):按 processJobId 取后台任务已持久化的终态通知 + * (background-process-recovery 服务在重启后落盘 transcript); + * undefined = 尚无持久化终态(任务可能仍活着)。 + */ + resolveBackgroundNotification?: (processJobId: string) => Promise; }, private readonly continueRunFromCheckpoint?: ContinueRunFromCheckpoint ) {} @@ -88,13 +95,44 @@ export class LumeResumeService { if (continuation.version === 2 && continuation.status === "waiting_background") { if (continuation.checkpoint.syntheticToolResult === undefined) { - return { - status: "waiting_background", - error: continuation.reason ?? "后台任务仍在运行,已重新附着且不会重复执行命令。" + // 崩溃后 live 终态监听(run.handleAsyncEvent)已不在:查 recovery 服务落盘的 + // 持久化终态通知,有则与 live 同形转换(ready_to_resume + syntheticToolResult), + // 让恢复真正发生;无则任务可能仍活着(durable 跨重启),如实返回等待态。 + const notification = continuation.checkpoint.processJobId + ? await this.stores.resolveBackgroundNotification?.(continuation.checkpoint.processJobId) + : undefined; + if (!notification || typeof notification !== "object") { + return { + status: "waiting_background", + error: "后台任务尚未产生持久化终态;若任务仍在运行,终态落盘后再次恢复即可继续。" + }; + } + const record = notification as unknown as Record; + const synthetic = { + type: "tool_result", + tool_use_id: + (typeof record.tool_use_id === "string" && record.tool_use_id) + || continuation.checkpoint.toolCallId + || "", + content: record.message ?? record.summary ?? "", + ...(record.status === "failed" || record.status === "stopped" || record.status === "interrupted" + ? { is_error: true } + : {}), + ...(record.execution && typeof record.execution === "object" + ? { _meta: { execution: record.execution } } + : {}), }; + await this.stores.continuationStore.update(input.runId, { + status: "ready_to_resume", + checkpoint: { ...continuation.checkpoint, step: "after_tool_result", syntheticToolResult: synthetic }, + reason: "后台命令已进入终态(崩溃恢复)。" + }); + continuation.status = "ready_to_resume"; + continuation.checkpoint = { ...continuation.checkpoint, step: "after_tool_result", syntheticToolResult: synthetic }; + } else { + await this.stores.continuationStore.update(input.runId, { status: "ready_to_resume" }); + continuation.status = "ready_to_resume"; } - await this.stores.continuationStore.update(input.runId, { status: "ready_to_resume" }); - continuation.status = "ready_to_resume"; } if (