From a60b6fa3eea93d106adffa3aa7b2e3119c9c0699 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=86=BC=E5=81=A5=E8=81=AA?= Date: Mon, 24 Aug 2026 14:53:23 +0800 Subject: [PATCH] fix(deepseek): fence retired host turns on local failure --- .../deepseek-harness/deepseek-harness.test.ts | 87 +++++++++++++++++++ server/drivers/deepseek-harness/index.ts | 6 +- 2 files changed, 91 insertions(+), 2 deletions(-) diff --git a/server/drivers/deepseek-harness/deepseek-harness.test.ts b/server/drivers/deepseek-harness/deepseek-harness.test.ts index a13a156ad..317ecfd03 100644 --- a/server/drivers/deepseek-harness/deepseek-harness.test.ts +++ b/server/drivers/deepseek-harness/deepseek-harness.test.ts @@ -748,6 +748,93 @@ describe("DeepSeek Harness turns", () => { events.stop(); }); + it("does not replay a retired terminal after queue submission overflow", async () => { + const fake = await host(); + let holdPrompt = true; + fake.onRawRequest = (request) => holdPrompt && dshClientRequestSchema.safeParse(request.body).data?.method === "session.prompt"; + fake.onRequest = ({ body }) => officialResponse(dshClientRequestSchema.parse(body), "queue-overflow-session"); + const instance = await DeepSeekHarnessDriver.create({ instanceId: "deepseekHarness", displayName: undefined, environment: {}, enabled: true, config: { baseUrl: fake.baseUrl, transport: "direct" } }); + const events = recordEvents(instance.adapter); + const model = encodeDshModelId("deepseek", "chat"); + const sessionEvent = (rpcId: string, type: string, seq: number, data: DshJsonValue) => ({ + type: "server-request", + rpcId, + method: "session/event", + payload: { type: "session/event", sessionId: "queue-overflow-session", event: { type, seq, time: seq, data } }, + }); + + const first = instance.adapter.sendTurn({ threadId: "queue-overflow", text: "first", model }); + await fake.waitForRawResponse(); + for (let seq = 1; seq <= 257; seq++) { + fake.send("mux", sessionEvent(`overflow-frame-${seq}`, "assistant/chunk.text-delta", seq, { turn: 1, delta: "x" })); + } + await fake.waitForStreamRoundTrip("mux"); + holdPrompt = false; + fake.releaseRawResponses(); + await expect(first).rejects.toThrow("event buffer overflowed"); + + holdPrompt = true; + const replacementStart = instance.adapter.sendTurn({ threadId: "queue-overflow", text: "replacement", model }); + await fake.waitForRawResponse(); + fake.send("mux", sessionEvent("old-overflow-end", "turn/end", 2, { turn: 1, reason: { kind: "completed" } })); + await fake.waitForStreamRoundTrip("mux"); + holdPrompt = false; + fake.releaseRawResponses(); + + const replacement = await replacementStart; + expect(events.events.filter((event) => event.type === "turn.completed" && event.turnId === replacement.turnId)).toHaveLength(0); + fake.send("mux", sessionEvent("replacement-start", "turn/start", 3, { turn: 2 })); + fake.send("mux", sessionEvent("replacement-output", "assistant/chunk", 4, { turn: 2, step: 1, chunk: { type: "text-delta", index: 0, text: "replacement output" } })); + fake.send("mux", sessionEvent("replacement-end", "turn/end", 5, { turn: 2, reason: { kind: "completed" } })); + + await events.until((event) => event.type === "turn.completed" && event.turnId === replacement.turnId, 500); + expect(events.events).toContainEqual(expect.objectContaining({ type: "content.delta", turnId: replacement.turnId, delta: "replacement output" })); + await instance.dispose(); + events.stop(); + }); + + it("does not replay a retired terminal after stream-loss recovery reuses a session", async () => { + const fake = await host(); + let holdPrompt = false; + fake.onRawRequest = (request) => holdPrompt && dshClientRequestSchema.safeParse(request.body).data?.method === "session.prompt"; + fake.onRequest = ({ body }) => officialResponse(dshClientRequestSchema.parse(body), "stream-reuse-session"); + const instance = await DeepSeekHarnessDriver.create({ instanceId: "deepseekHarness", displayName: undefined, environment: {}, enabled: true, config: { baseUrl: fake.baseUrl, transport: "direct" } }); + const events = recordEvents(instance.adapter); + const model = encodeDshModelId("deepseek", "chat"); + const sessionEvent = (rpcId: string, type: string, seq: number, data: DshJsonValue) => ({ + type: "server-request", + rpcId, + method: "session/event", + payload: { type: "session/event", sessionId: "stream-reuse-session", event: { type, seq, time: seq, data } }, + }); + + await instance.adapter.sendTurn({ threadId: "stream-reuse", text: "first", model }); + await fake.waitForStream("mux"); + fake.send("mux", sessionEvent("old-stream-start", "turn/start", 1, { turn: 1 })); + await fake.waitForStreamRoundTrip("mux"); + fake.send("host", { type: "server-request", rpcId: "stream-loss", method: "stream/error", payload: { type: "stream/error", error: { code: "internal", message: "lost", details: {} } } }); + await events.until((event) => event.type === "turn.completed" && event.threadId === "stream-reuse" && event.stopReason === "stream_lost", 2_000); + + holdPrompt = true; + const replacementStart = instance.adapter.sendTurn({ threadId: "stream-reuse", text: "replacement", model }); + await fake.waitForRawResponse(); + fake.send("mux", sessionEvent("old-stream-end", "turn/end", 2, { turn: 1, reason: { kind: "completed" } })); + await fake.waitForStreamRoundTrip("mux"); + holdPrompt = false; + fake.releaseRawResponses(); + + const replacement = await replacementStart; + expect(events.events.filter((event) => event.type === "turn.completed" && event.turnId === replacement.turnId)).toHaveLength(0); + fake.send("mux", sessionEvent("replacement-stream-start", "turn/start", 3, { turn: 2 })); + fake.send("mux", sessionEvent("replacement-stream-output", "assistant/chunk", 4, { turn: 2, step: 1, chunk: { type: "text-delta", index: 0, text: "replacement output" } })); + fake.send("mux", sessionEvent("replacement-stream-end", "turn/end", 5, { turn: 2, reason: { kind: "completed" } })); + + await events.until((event) => event.type === "turn.completed" && event.turnId === replacement.turnId, 500); + expect(events.events).toContainEqual(expect.objectContaining({ type: "content.delta", turnId: replacement.turnId, delta: "replacement output" })); + await instance.dispose(); + events.stop(); + }); + it.each([ ["stopAll", "session.create"], ["stopAll", "session.selectModel"], diff --git a/server/drivers/deepseek-harness/index.ts b/server/drivers/deepseek-harness/index.ts index 04ceb0546..eaf3a39af 100644 --- a/server/drivers/deepseek-harness/index.ts +++ b/server/drivers/deepseek-harness/index.ts @@ -238,14 +238,14 @@ export const DeepSeekHarnessDriver: ProviderDriver = { || method === "tool/call" || method === "tool/result" ); - const retireHostTurn = (running: ActiveTurn) => { + const retireHostTurn = (running: ActiveTurn, includeUnknown = true) => { if (!running.sessionId) return; const previous = retiredHostTurns.get(running.sessionId); if (running.hostTurn !== undefined) { retiredHostTurns.set(running.sessionId, previous === null || previous === undefined ? running.hostTurn : Math.max(previous, running.hostTurn)); - } else if (!retiredHostTurns.has(running.sessionId)) { + } else if (includeUnknown && !retiredHostTurns.has(running.sessionId)) { retiredHostTurns.set(running.sessionId, null); } }; @@ -601,6 +601,7 @@ export const DeepSeekHarnessDriver: ProviderDriver = { running.queueFrames = []; running.queueFrameBytes = 0; running.settleStartup(); + retireHostTurn(running); await cancelRunning(running); throw new Error("DeepSeek Harness queue submission event buffer overflowed"); } @@ -695,6 +696,7 @@ export const DeepSeekHarnessDriver: ProviderDriver = { running.completed = true; running.cancelled = true; running.failReady(); + retireHostTurn(running, false); if (active.get(threadId) === running) active.delete(threadId); } await Promise.allSettled(lost.map(([, running]) => running.sessionId