diff --git a/desktop/src/stores/chatStore.test.ts b/desktop/src/stores/chatStore.test.ts index ae246f11..7915b0c1 100644 --- a/desktop/src/stores/chatStore.test.ts +++ b/desktop/src/stores/chatStore.test.ts @@ -1,4 +1,4 @@ -import { beforeEach, describe, expect, it, vi } from 'vitest' +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' import type { MessageEntry } from '../types/session' import { useSessionRuntimeStore } from './sessionRuntimeStore' @@ -1332,6 +1332,50 @@ describe('chatStore history mapping', () => { expect(notifyDesktopMock).not.toHaveBeenCalled() }) + // Same reconnect replay as above, one block type over. A thinking UIMessage + // carries no transcriptMessageId (types/chat.ts:291), so neither the + // appendAssistantTextMessage guard nor dropDuplicateTranscriptTextMessages + // can reach it — a replayed thinking block has nothing identifying it. + it('does not duplicate a hydrated thinking block when live output replays after reconnect', () => { + useChatStore.setState({ + sessions: { + [TEST_SESSION_ID]: makeSession({ + chatState: 'idle', + messages: [ + { + id: 'live-user', + type: 'user_text', + content: 'live prompt', + transcriptMessageId: 'transcript-user-1', + timestamp: 1, + }, + { id: 'live-thinking', type: 'thinking', content: 'weighing the options', timestamp: 2 }, + { + id: 'live-assistant', + type: 'assistant_text', + content: 'live answer', + transcriptMessageId: 'transcript-assistant-1', + timestamp: 3, + }, + ], + }), + }, + }) + + useChatStore.getState().handleServerMessage(TEST_SESSION_ID, { + type: 'thinking', + text: 'weighing the options', + }) + + const session = useChatStore.getState().sessions[TEST_SESSION_ID] + expect(session?.messages).toMatchObject([ + { id: 'live-user', type: 'user_text' }, + { id: 'live-thinking', type: 'thinking', content: 'weighing the options' }, + { id: 'live-assistant', type: 'assistant_text', content: 'live answer' }, + ]) + if (session?.elapsedTimer) clearInterval(session.elapsedTimer) + }) + it('collapses duplicate assistant replies after transcript id hydration', async () => { vi.mocked(sessionsApi.getMessages).mockResolvedValueOnce({ messages: [ @@ -6319,3 +6363,209 @@ describe('chatStore history mapping', () => { vi.useRealTimers() }) }) + +// A desktop window left open for hours re-renders a long-finished session as a +// wall of collapsed "已思考" bubbles with no reply text between them. No new LLM +// call happens — the finished turn's stream events are pushed at the renderer a +// second time on top of already-hydrated history. appendAssistantTextMessage +// guards against exactly that replay, `case 'thinking'` does not. +type FinishedTurnStep = + | { kind: 'thinking'; text: string } + | { kind: 'text'; text: string } + | { kind: 'tool'; id: string; name: string } + +// Verbatim block order of session fd09ae78-297a-427e-a849-5251731bf30c (Kimi +// k3, a Computer Use run): 13 thinking blocks, 5 of them empty. Only the very +// last assistant reply is elided — it is a long markdown table. +const FINISHED_COMPUTER_USE_TURN: FinishedTurnStep[] = [ + { kind: 'thinking', text: 'User wants Computer Use. I need to load the computer-use tools. Let me fetch the tool schemas first.' }, + { kind: 'text', text: '我先加载 Computer Use 工具,然后严格按顺序执行。' }, + { kind: 'tool', id: 'tool_aJGuYYOY4gzHHHFIaTus4M3e', name: 'ToolSearch' }, + { kind: 'thinking', text: 'Start: get_app_state to find search box, then step 1 click it. User wants get_app_state after each step and report the search field content.' }, + { kind: 'tool', id: 'tool_kosJoZ93D4WBqAWnmelI9Rik', name: 'mcp__computer-use__get_app_state' }, + { kind: 'thinking', text: 'Search box at approx (410, 50). Step 1: single click.' }, + { kind: 'text', text: '初始状态:搜索框为空(只有放大镜占位图标),位置约 (410, 50)。\n\n**第 1 步:单击搜索框**' }, + { kind: 'tool', id: 'tool_KlLouIdZHyPcj0Jbmy5nfJOf', name: 'mcp__computer-use__click' }, + { kind: 'thinking', text: '' }, + { kind: 'tool', id: 'tool_lNWmFSm7Er61dyZJcWGkGH6b', name: 'mcp__computer-use__get_app_state' }, + { kind: 'thinking', text: 'Step 1 done: search box still empty (only magnifier icon). Step 2: type "test".' }, + { kind: 'text', text: '第 1 步后:搜索框里**为空**(只有放大镜占位图标,无文字)。\n\n**第 2 步:type_text 输入 "test"**' }, + { kind: 'tool', id: 'tool_hTHyhielyvIaUIWN6YJDz2g5', name: 'mcp__computer-use__type_text' }, + { kind: 'thinking', text: '' }, + { kind: 'tool', id: 'tool_951Y0sQNx8tMZGGMS7XjVeOL', name: 'mcp__computer-use__get_app_state' }, + { kind: 'thinking', text: 'Step 2 done: search box still empty — "test" did not appear. Report honestly. Step 3: type "abc".' }, + { kind: 'text', text: '第 2 步后:搜索框里**为空** —— "test" 没有显示出来。\n\n**第 3 步:type_text 输入 "abc"**' }, + { kind: 'tool', id: 'tool_Y1MzoPySfdwd7tNrFj389CEN', name: 'mcp__computer-use__type_text' }, + { kind: 'thinking', text: 'Continue: get_app_state after step 3.' }, + { kind: 'tool', id: 'tool_8TrKG2qCE6Q4bDYE5hMvngSW', name: 'mcp__computer-use__get_app_state' }, + { kind: 'thinking', text: 'Step 3 done: still empty. Step 4: double-click search box.' }, + { kind: 'text', text: '第 3 步后:搜索框里**依然为空** —— "abc" 也没有显示。\n\n**第 4 步:双击搜索框**' }, + { kind: 'tool', id: 'tool_qguIh8LXvDxOn6wE4aSlys2K', name: 'mcp__computer-use__click' }, + { kind: 'thinking', text: '' }, + { kind: 'tool', id: 'tool_T4513oOu6to3GY1EeEPMLvFr', name: 'mcp__computer-use__get_app_state' }, + { kind: 'thinking', text: '' }, + { kind: 'text', text: '第 4 步后:搜索框里**仍然为空**,界面无任何变化。\n\n**第 5 步:type_text 输入 "xyz"**' }, + { kind: 'tool', id: 'tool_XhJWl5aNuSPaUXctFQ91KSPM', name: 'mcp__computer-use__type_text' }, + { kind: 'thinking', text: '' }, + { kind: 'tool', id: 'tool_b2c5HshyO48rmbYqmlcGj94K', name: 'mcp__computer-use__get_app_state' }, + { kind: 'thinking', text: 'Done. Report all 5 steps faithfully.' }, + { kind: 'text', text: '5 步全部完成,如实汇报每步后搜索框的内容:' }, +] + +const FIRST_THINKING = FINISHED_COMPUTER_USE_TURN[0] as { kind: 'thinking'; text: string } +const LAST_THINKING = FINISHED_COMPUTER_USE_TURN[30] as { kind: 'thinking'; text: string } + +/** The CLI persists one transcript entry per content block; mirror that shape. */ +function buildFinishedTurnHistory(): MessageEntry[] { + const entries: MessageEntry[] = [{ + id: 'cu-prompt', + type: 'user', + timestamp: '2026-07-28T15:00:00.000Z', + content: '用 Computer Use 操作网易云音乐,严格按顺序做,每步做完立刻 get_app_state。', + }] + FINISHED_COMPUTER_USE_TURN.forEach((step, index) => { + const timestamp = new Date(Date.UTC(2026, 6, 28, 15, 1, index)).toISOString() + if (step.kind === 'thinking') { + entries.push({ id: `cu-a-${index}`, type: 'assistant', timestamp, content: [{ type: 'thinking', thinking: step.text }] }) + return + } + if (step.kind === 'text') { + entries.push({ id: `cu-a-${index}`, type: 'assistant', timestamp, content: [{ type: 'text', text: step.text }] }) + return + } + entries.push({ id: `cu-a-${index}`, type: 'assistant', timestamp, content: [{ type: 'tool_use', name: step.name, id: step.id, input: {} }] }) + entries.push({ id: `cu-u-${index}`, type: 'user', timestamp, content: [{ type: 'tool_result', tool_use_id: step.id, content: 'ok', is_error: false }] }) + }) + return entries +} + +/** + * Push the finished turn at the renderer exactly as the server emits it. Both + * emit paths skip falsy thinking (`handler.ts` thinking_delta and the + * no-stream-events assistant fallback), so the five empty transcript blocks + * never reach the wire — `withEmptyThinking` exists only for the defensive case. + */ +function replayFinishedTurn(options: { withText: boolean; withEmptyThinking?: boolean }) { + const store = useChatStore.getState() + for (const step of FINISHED_COMPUTER_USE_TURN) { + if (step.kind === 'thinking') { + if (step.text || options.withEmptyThinking) { + store.handleServerMessage(TEST_SESSION_ID, { type: 'thinking', text: step.text }) + } + continue + } + if (step.kind === 'text') { + if (options.withText) store.handleServerMessage(TEST_SESSION_ID, { type: 'content_delta', text: step.text }) + continue + } + store.handleServerMessage(TEST_SESSION_ID, { + type: 'tool_use_complete', + toolName: step.name, + toolUseId: step.id, + input: {}, + }) + store.handleServerMessage(TEST_SESSION_ID, { + type: 'tool_result', + toolUseId: step.id, + content: 'ok', + isError: false, + }) + } +} + +function thinkingBlocks() { + return (useChatStore.getState().sessions[TEST_SESSION_ID]?.messages ?? []) + .filter((message) => message.type === 'thinking') + .map((message) => message.content) +} + +describe('chatStore wake replay of a finished thinking turn', () => { + beforeEach(async () => { + sendMock.mockReset() + notifyDesktopMock.mockReset() + updateTabStatusMock.mockReset() + getMemberBySessionIdMock.mockReset() + getMemberBySessionIdMock.mockReturnValue(null) + connectionStateHandlers.clear() + vi.mocked(sessionsApi.getMessages).mockReset() + vi.mocked(sessionsApi.getMessages).mockResolvedValue({ messages: buildFinishedTurnHistory() }) + localStorage.clear() + useSettingsStore.setState({ locale: 'en' }) + useChatStore.setState({ + ...initialState, + sessions: { [TEST_SESSION_ID]: makeSession({ chatState: 'idle', messages: [] }) }, + }) + await useChatStore.getState().loadHistory(TEST_SESSION_ID) + }) + + afterEach(() => { + const timer = useChatStore.getState().sessions[TEST_SESSION_ID]?.elapsedTimer + if (timer) clearInterval(timer) + }) + + it('hydrates the finished turn as eight non-empty thinking blocks', () => { + // Sanity check on the fixture: history mapping drops the five empty + // thinking blocks, so anything empty on screen came from the stream path. + expect(thinkingBlocks()).toHaveLength(8) + expect(thinkingBlocks().every((content) => content.trim().length > 0)).toBe(true) + }) + + it('does not re-append hydrated thinking when a finished turn is replayed after wake', () => { + const before = thinkingBlocks() + + replayFinishedTurn({ withText: false }) + + expect(thinkingBlocks()).toEqual(before) + }) + + // Defensive. The transcript really does hold five empty thinking blocks and + // history mapping filters them, but no server emit path puts an empty + // thinking event on the wire today — so a collapsed "已思考" row on screen is + // not by itself proof that an empty block was streamed. + it('does not spawn empty thinking bubbles from replayed empty thinking events', () => { + replayFinishedTurn({ withText: false, withEmptyThinking: true }) + + expect(thinkingBlocks().filter((content) => !content.trim())).toEqual([]) + }) + + it('does not merge the turn tail into the turn head when the replay wraps around', () => { + replayFinishedTurn({ withText: false }) + replayFinishedTurn({ withText: false }) + + // The last thinking block of one replay round and the first of the next + // arrive back to back with nothing between them, so the merge branch glues + // two unrelated reasoning blocks into a single bubble. + expect(thinkingBlocks().filter((content) => + content.includes(LAST_THINKING.text) && content.includes(FIRST_THINKING.text), + )).toEqual([]) + }) + + it('does not grow the thinking wall with every extra replay round', () => { + const hydrated = thinkingBlocks().length + + replayFinishedTurn({ withText: false }) + const afterFirst = thinkingBlocks().length + replayFinishedTurn({ withText: false }) + const afterSecond = thinkingBlocks().length + + expect({ afterFirst, afterSecond }).toEqual({ + afterFirst: hydrated, + afterSecond: hydrated, + }) + }) + + // Secondary gap. The reported screenshot has no reply text between the + // bubbles, so the observed replay carries thinking only — but the existing + // appendAssistantTextMessage guard would not have held either: it only + // matches while the tail still is the hydrated assistant message. + it('does not re-append replayed reply text once the tail moves past the hydrated message', () => { + const assistantTextOf = () => (useChatStore.getState().sessions[TEST_SESSION_ID]?.messages ?? []) + .filter((message) => message.type === 'assistant_text') + .map((message) => message.content) + const before = assistantTextOf() + + replayFinishedTurn({ withText: true }) + + expect(assistantTextOf()).toEqual(before) + }) +}) diff --git a/desktop/src/stores/chatStore.ts b/desktop/src/stores/chatStore.ts index 2c0638e0..b5ad5b13 100644 --- a/desktop/src/stores/chatStore.ts +++ b/desktop/src/stores/chatStore.ts @@ -597,6 +597,21 @@ function appendAssistantTextMessage( ) { return messages } + // 上面那道只在尾部仍是那条 hydrated 消息时才够得着。整轮重放时,正文到达前 + // 尾部早被 thinking / tool_result 顶掉了,于是重复的回复照样追加进来。 + // 这里比的是"逐字相同"而不是子串:整段重发的正文会与某条 hydrated 回复完全一致, + // 而正常流式送来的是碎片(碎片几乎必然是某条历史回复的子串,用子串判定会误伤)。 + if ( + !transcriptMessageId && + messages.some( + (message) => + message.type === 'assistant_text' && + message.transcriptMessageId && + message.content.trim() === trimmedContent, + ) + ) { + return messages + } const canMergeIntoLast = last?.type === 'assistant_text' && @@ -2179,7 +2194,7 @@ export const useChatStore = create((set, get) => ({ if (receivedLiveDelta && get().sessions[sessionId]?.chatState !== 'idle') ensureElapsedTimer() break - case 'thinking': + case 'thinking': { if (get().sessions[sessionId]?.suppressNextTaskNotificationResponse) { consumePendingDelta(sessionId) update(() => ({ @@ -2189,11 +2204,28 @@ export const useChatStore = create((set, get) => ({ })) break } + // 重放/空块都不该冒出一个新的「已思考」气泡,也不该把会话拖回 thinking 态 + // 或者启动计时器 —— 那正是"打开一个早就结束的会话,它自己开始输出"的观感。 + let skippedThinkingBlock = false update((s) => { const pendingText = `${s.streamingText}${consumePendingDelta(sessionId)}` const base = pendingText.trim() ? appendAssistantTextMessage(s.messages, pendingText, Date.now()) : s.messages + // 服务端两个 thinking 发射点都做了非空过滤,但 `&& delta.thinking` 是真值 + // 判断,纯空白仍能漏过来,落到下面就是一个点开什么都没有的空壳气泡。 + if (!msg.text.trim()) { + skippedThinkingBlock = true + return { messages: base, streamingText: '' } + } + // 真正的重放源已在服务端按 uuid 挡掉(conversationService.isReplayedSdkMessage)。 + // 这里再兜一道:thinking 没有 transcriptMessageId 之类的身份,任何漏网的 + // 重放都只能靠"整块内容与已有 thinking 逐字相同"来认。流式 delta 是碎片, + // 不会命中;命中的必然是被整块重发的同一段思考。 + if (base.some((message) => message.type === 'thinking' && message.content === msg.text)) { + skippedThinkingBlock = true + return { messages: base, streamingText: '' } + } const lastIndex = findStreamMergeTargetIndex(base) const last = lastIndex >= 0 ? base[lastIndex] : undefined if (last && last.type === 'thinking') { @@ -2216,8 +2248,9 @@ export const useChatStore = create((set, get) => ({ streamingResponseChars: s.streamingResponseChars + msg.text.length, } }) - ensureElapsedTimer() + if (!skippedThinkingBlock) ensureElapsedTimer() break + } case 'tool_use_complete': { clearPendingToolInputDelta(sessionId) @@ -2226,26 +2259,37 @@ export const useChatStore = create((set, get) => ({ const toolUseId = msg.toolUseId || session?.activeToolUseId || '' const parentToolUseId = msg.parentToolUseId ?? getPendingToolParentUseId(sessionId, toolUseId) rememberPendingToolParentUseId(sessionId, toolUseId, parentToolUseId) - update((s) => ({ - messages: toolUseId - ? upsertToolUseMessage(s.messages, toolUseId, (existing) => ({ - id: existing?.id ?? nextId(), - type: 'tool_use', - toolName, - toolUseId, - input: msg.input, - timestamp: existing?.timestamp ?? Date.now(), - parentToolUseId, - isPending: false, - })) - : [...s.messages, { - id: nextId(), type: 'tool_use', toolName, - toolUseId, - input: msg.input, timestamp: Date.now(), parentToolUseId, - isPending: false, - }], - activeToolUseId: null, activeToolName: null, activeThinkingId: null, streamingToolInput: '', - })) + update((s) => { + // 流式路径上,工具块的 content_start 已经把待定正文冲刷成一条消息了 + // (见 case 'content_start' 里 blockType !== 'text' 的分支)。但整块兜底 + // 路径只发 tool_use_complete、不发 content_start —— 不在这里补一次冲刷, + // 工具调用前后的两段正文就会跨消息粘成一条。 + const pendingText = `${s.streamingText}${consumePendingDelta(sessionId)}` + const base = pendingText.trim() + ? appendAssistantTextMessage(s.messages, pendingText, Date.now()) + : s.messages + return { + messages: toolUseId + ? upsertToolUseMessage(base, toolUseId, (existing) => ({ + id: existing?.id ?? nextId(), + type: 'tool_use', + toolName, + toolUseId, + input: msg.input, + timestamp: existing?.timestamp ?? Date.now(), + parentToolUseId, + isPending: false, + })) + : [...base, { + id: nextId(), type: 'tool_use', toolName, + toolUseId, + input: msg.input, timestamp: Date.now(), parentToolUseId, + isPending: false, + }], + streamingText: '', + activeToolUseId: null, activeToolName: null, activeThinkingId: null, streamingToolInput: '', + } + }) if (toolName === 'TodoWrite' && Array.isArray((msg.input as any)?.todos)) { useCLITaskStore.getState().setTasksFromTodos((msg.input as any).todos, sessionId) } else if (TASK_TOOL_NAMES.has(toolName)) { diff --git a/src/server/__tests__/conversation-service.test.ts b/src/server/__tests__/conversation-service.test.ts index fd3e730a..b14c29b7 100644 --- a/src/server/__tests__/conversation-service.test.ts +++ b/src/server/__tests__/conversation-service.test.ts @@ -1266,6 +1266,139 @@ describe('ConversationService', () => { expect(completionObserved).toBe(true) }) + // CLI 的 WebSocketTransport 每次重连成功都会把整个发送缓冲区重放一遍,并假定 + // 「The server deduplicates by UUID」。以前 server 没实现这个契约:笔记本睡醒后 + // CLI 重连,一整轮早已结束的对话会被重新推上来,前端当成实时输出再渲染一遍 + // (表现为满屏「已思考」)。 + test('drops SDK messages replayed by the CLI after a reconnect', () => { + const service = new ConversationService() as any + const forwarded: any[] = [] + service.sessions.set('replay-dedupe', { + outputCallbacks: [(message: any) => forwarded.push(message)], + seenSdkMessageUuids: new Set(), + sdkMessages: [], + initMessage: null, + pendingPermissionRequests: new Map(), + }) + + const turn = [ + { type: 'assistant', uuid: 'uuid-thinking-1', message: { content: [{ type: 'thinking', thinking: 'step one' }] } }, + { type: 'assistant', uuid: 'uuid-thinking-2', message: { content: [{ type: 'thinking', thinking: 'step two' }] } }, + { type: 'result', uuid: 'uuid-result', subtype: 'success', is_error: false }, + ] + const payload = turn.map((msg) => JSON.stringify(msg)).join('\n') + + service.handleSdkPayload('replay-dedupe', payload) + expect(forwarded).toHaveLength(3) + + // 重连后 CLI 把同一批消息从缓冲区头部重放 —— 一条都不该再转发出去。 + service.handleSdkPayload('replay-dedupe', payload) + service.handleSdkPayload('replay-dedupe', payload) + expect(forwarded).toHaveLength(3) + }) + + // 真机日志(cli-diagnostics.jsonl.58975)显示重放的 858 条里绝大多数是 stream_event: + // 桌面端固定传 --include-partial-messages,CLI 为每个 thinking_delta 单独产一条 + // stream_event 并现铸 uuid(QueryEngine.ts:846),所以重放到达前端时是 delta 碎片, + // 而不是整块 thinking。渲染侧的逐字比对挡不住碎片,只有这里的 uuid 判重挡得住。 + test('drops replayed partial-message stream events, not just whole assistant blocks', () => { + const service = new ConversationService() as any + const forwarded: any[] = [] + service.sessions.set('replay-stream-events', { + outputCallbacks: [(message: any) => forwarded.push(message)], + seenSdkMessageUuids: new Set(), + sdkMessages: [], + initMessage: null, + pendingPermissionRequests: new Map(), + }) + + const deltas = ['Start: get_', 'app_state to find ', 'the search box.'] + const payload = deltas + .map((thinking, index) => + JSON.stringify({ + type: 'stream_event', + uuid: `uuid-stream-${index}`, + event: { + type: 'content_block_delta', + index: 0, + delta: { type: 'thinking_delta', thinking }, + }, + }), + ) + .join('\n') + + service.handleSdkPayload('replay-stream-events', payload) + expect(forwarded).toHaveLength(deltas.length) + + service.handleSdkPayload('replay-stream-events', payload) + expect(forwarded).toHaveLength(deltas.length) + }) + + test('keeps forwarding SDK messages that carry no uuid', () => { + const service = new ConversationService() as any + const forwarded: any[] = [] + service.sessions.set('no-uuid', { + outputCallbacks: [(message: any) => forwarded.push(message)], + seenSdkMessageUuids: new Set(), + sdkMessages: [], + initMessage: null, + pendingPermissionRequests: new Map(), + }) + + // 没有 uuid 的消息不会进 CLI 的重放缓冲,所以也不该被判重丢弃。 + const payload = JSON.stringify({ type: 'result', subtype: 'success', is_error: false }) + service.handleSdkPayload('no-uuid', payload) + service.handleSdkPayload('no-uuid', payload) + + expect(forwarded).toHaveLength(2) + }) + + test('tolerates sessions created without the replay-dedupe bookkeeping', () => { + const service = new ConversationService() as any + const forwarded: any[] = [] + // 故意不带 seenSdkMessageUuids,模拟别处构造出来的会话对象。 + service.sessions.set('legacy-shape', { + outputCallbacks: [(message: any) => forwarded.push(message)], + sdkMessages: [], + initMessage: null, + pendingPermissionRequests: new Map(), + }) + + const payload = JSON.stringify({ + type: 'assistant', + uuid: 'uuid-legacy', + message: { content: [{ type: 'thinking', thinking: 'hello' }] }, + }) + service.handleSdkPayload('legacy-shape', payload) + service.handleSdkPayload('legacy-shape', payload) + + expect(forwarded).toHaveLength(1) + }) + + test('remembers enough uuids to cover a full CLI replay buffer', () => { + const service = new ConversationService() as any + const forwarded: any[] = [] + service.sessions.set('buffer-span', { + outputCallbacks: [(message: any) => forwarded.push(message)], + seenSdkMessageUuids: new Set(), + sdkMessages: [], + initMessage: null, + pendingPermissionRequests: new Map(), + }) + + // CLI 侧缓冲上限是 1000 条,整个缓冲区被重放时每一条都必须还认得出来。 + const CLI_REPLAY_BUFFER_SIZE = 1000 + const payload = Array.from({ length: CLI_REPLAY_BUFFER_SIZE }, (_unused, index) => + JSON.stringify({ type: 'assistant', uuid: `uuid-${index}`, message: { content: [] } }), + ).join('\n') + + service.handleSdkPayload('buffer-span', payload) + expect(forwarded).toHaveLength(CLI_REPLAY_BUFFER_SIZE) + + service.handleSdkPayload('buffer-span', payload) + expect(forwarded).toHaveLength(CLI_REPLAY_BUFFER_SIZE) + }) + test('removes an exited CLI session even when one output callback throws', async () => { const service = new ConversationService() as any const sessionId = 'exit-callback-isolation' diff --git a/src/server/services/conversationService.ts b/src/server/services/conversationService.ts index 7b9197d6..75ef046d 100644 --- a/src/server/services/conversationService.ts +++ b/src/server/services/conversationService.ts @@ -68,6 +68,13 @@ export const MAX_CAPTURED_SDK_MESSAGE_BYTES = 64 * 1024 export const MAX_CAPTURED_SDK_TOTAL_BYTES = 512 * 1024 const MAX_CAPTURED_SDK_DIAGNOSTIC_TEXT_BYTES = 4 * 1024 const CONTROL_READY_POLL_MS = 50 +/** + * 记住多少条已处理的 SDK 消息 uuid,用来挡掉 CLI 重连时的重放。 + * CLI 侧重放缓冲是 DEFAULT_MAX_BUFFER_SIZE = 1000 条 + * (src/cli/transports/WebSocketTransport.ts),这里留一倍余量, + * 保证整个缓冲区被重放时每一条都还认得出来。 + */ +const MAX_SEEN_SDK_MESSAGE_UUIDS = 2_000 const AUTO_MEMORY_DIRNAME = 'memory' export const DESKTOP_CLI_GRACEFUL_SHUTDOWN_TIMEOUT_MS = 6_000 @@ -186,6 +193,11 @@ type SessionProcess = { stdoutLines: string[] stderrLines: string[] outputDrain: Promise + /** + * UUID 的 SDK 消息一旦处理过就记在这里,用于挡掉 CLI 重连时的整轮重放。 + * 详见 handleSdkPayload 里的说明。插入顺序即淘汰顺序(Set 保序)。 + */ + seenSdkMessageUuids: Set sdkMessages: any[] sdkMessageBytes?: number initMessage: any | null @@ -443,6 +455,7 @@ export class ConversationService { networkDerivedFirstTokenTimeout: networkRuntimeMetadata.firstTokenTimeoutDerived, sdkToken: this.getSdkTokenFromUrl(sdkUrl), sdkSocket: null, + seenSdkMessageUuids: new Set(), sdkAttached, resolveSdkAttached, pendingOutbound: [], @@ -955,6 +968,35 @@ export class ConversationService { } } + /** + * CLI 的 WebSocketTransport 在每次重连成功后会把它的整个发送缓冲区重放一遍, + * 并且明确假定「The server deduplicates by UUID」 + * (src/cli/transports/WebSocketTransport.ts:204)。这个契约以前没有实现: + * 笔记本睡眠导致连接断开后(该 transport 有专门的睡眠检测,会无限重置重连预算), + * CLI 重连时会把最多 1000 条**早已完成**的消息重新推上来,server 原样转发给前端, + * 前端便把一整轮结束很久的对话当成实时输出重新渲染一遍 —— 表现为满屏「已思考」。 + * + * 只有带 uuid 的消息才会进入 CLI 的重放缓冲(同文件 write()),所以这里也只按 + * uuid 判重;没有 uuid 的消息(如 control_request)本就不会被重放,照常处理。 + */ + private isReplayedSdkMessage(session: SessionProcess, msg: any): boolean { + const uuid = typeof msg?.uuid === 'string' ? msg.uuid : '' + if (!uuid) return false + + // 会话对象并非只有 startSession 一条构造路径,缺字段时按空集合起步而不是抛错。 + const seen = session.seenSdkMessageUuids ?? new Set() + session.seenSdkMessageUuids = seen + if (seen.has(uuid)) return true + + if (seen.size >= MAX_SEEN_SDK_MESSAGE_UUIDS) { + // Set 保持插入顺序,最早进来的就是最该淘汰的。 + const oldest = seen.values().next().value + if (oldest !== undefined) seen.delete(oldest) + } + seen.add(uuid) + return false + } + handleSdkPayload(sessionId: string, rawPayload: string): void { const session = this.sessions.get(sessionId) if (!session) return @@ -967,6 +1009,7 @@ export class ConversationService { for (const line of lines) { try { const msg = JSON.parse(line) + if (this.isReplayedSdkMessage(session, msg)) continue this.retainSdkMessage(session, msg, Buffer.byteLength(line, 'utf-8')) const sdkError = this.extractSdkErrorEvent(msg) if (sdkError) {