From f26c1770ea0c209b4cb2eca856da73924fdc7cd7 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=A8=8B=E5=BA=8F=E5=91=98=E9=98=BF=E6=B1=9F=28Relakkes?= =?UTF-8?q?=29?= Date: Fri, 25 Sep 2026 14:26:50 +0800 Subject: [PATCH] fix(desktop): prevent duplicate prompt replay after guiding a conversation --- desktop/src/stores/chatStore.test.ts | 47 ++++++++++++++ desktop/src/stores/chatStore.ts | 29 +++++++-- desktop/src/types/chat.ts | 2 +- src/server/__tests__/conversations.test.ts | 51 +++++++++++++++ src/server/__tests__/fixtures/mock-sdk-cli.ts | 65 +++++++++++++++++++ 5 files changed, 187 insertions(+), 7 deletions(-) diff --git a/desktop/src/stores/chatStore.test.ts b/desktop/src/stores/chatStore.test.ts index 3e411dfe..b7b2c898 100644 --- a/desktop/src/stores/chatStore.test.ts +++ b/desktop/src/stores/chatStore.test.ts @@ -14024,6 +14024,53 @@ describe('chatStore history mapping', () => { ]) }) + it.each([ + ['Use TypeScript', 'Keep the toolbar compact'], + ['Use TypeScript', 'Build the editor'], + ['Use TypeScript', 'Use TypeScript'], + ['Build the editor', 'Build the editor'], + ])('reconciles the delayed initial replay across guides %s / %s without resending', (...guides) => { + const initial = 'Build the editor' + useChatStore.setState({ sessions: { + [TEST_SESSION_ID]: makeSession({ chatState: 'idle' }), + } }) + useChatStore.getState().sendMessage(TEST_SESSION_ID, initial) + useChatStore.getState().handleServerMessage(TEST_SESSION_ID, { + type: 'thinking', text: 'Planning the editor', + }) + for (const content of guides) { + const id = useChatStore.getState().queueUserMessage(TEST_SESSION_ID, { + content, displayContent: content, + }) + useChatStore.getState().sendQueuedUserMessage(TEST_SESSION_ID, id) + } + const sentBeforeReplay = sendMock.mock.calls.length + for (const content of [initial, ...guides]) { + useChatStore.getState().handleServerMessage(TEST_SESSION_ID, { + type: 'user_message_replay', content, + }) + expect(useChatStore.getState().sessions[TEST_SESSION_ID]?.messages + .filter(message => message.type === 'user_text').map(message => message.content)) + .toEqual([initial, ...guides]) + } + expect(useChatStore.getState().sessions[TEST_SESSION_ID]?.messages + .filter(message => message.type === 'user_text' && message.optimisticQueued)).toEqual([]) + expect(sendMock.mock.calls).toHaveLength(sentBeforeReplay) + expect(sendMock.mock.calls.filter(([, message]) => message.type === 'user_message') + .map(([, message]) => message.content)).toEqual([initial, ...guides]) + }) + + it('does not suppress a new replay just because an older turn has the same text', () => { + const messages: UIMessage[] = [ + { id: 'old', type: 'user_text', content: 'Try again', timestamp: 1 }, + { id: 'current', type: 'user_text', content: 'Change direction', timestamp: 2 }, + { id: 'guide', type: 'user_text', content: 'Keep it small', timestamp: 3, optimisticQueued: true }, + ] + expect(appendReplayedUserMessage(messages, 'Try again', 4) + .filter(message => message.type === 'user_text').map(message => message.content)) + .toEqual(['Try again', 'Change direction', 'Keep it small', 'Try again']) + }) + it('does not duplicate a slash-command prompt when the replay normalizes extra spaces', () => { // The composer keeps the raw input (`/ego-browser␣␣https://…` — two spaces // after the command name). The CLI preserves them inside , diff --git a/desktop/src/stores/chatStore.ts b/desktop/src/stores/chatStore.ts index efe6950d..152c5b06 100644 --- a/desktop/src/stores/chatStore.ts +++ b/desktop/src/stores/chatStore.ts @@ -3285,7 +3285,7 @@ export const useChatStore = create((setState, get) => { ...(userFacingContent !== modelFacingContent ? { modelContent: modelFacingContent } : {}), attachments: isDirectAgentSession ? undefined : uiAttachments, timestamp: now, - ...(isDirectAgentSession ? { pending: true } : {}), + ...(isDirectAgentSession ? { pending: true } : { awaitingReplay: true }), }) if (!isDirectAgentSession && session.elapsedTimer) clearInterval(session.elapsedTimer) @@ -7497,8 +7497,8 @@ export function appendReplayedUserMessage( const currentTurnUserIndex = findCurrentTurnUserMessageIndex(messages, modelContent, parsed) if (currentTurnUserIndex >= 0) { const optimisticMessage = messages[currentTurnUserIndex] - if (optimisticMessage?.type === 'user_text' && optimisticMessage.optimisticQueued) { - const { optimisticQueued: _optimisticQueued, ...confirmedMessage } = optimisticMessage + if (optimisticMessage?.type === 'user_text' && (optimisticMessage.optimisticQueued || optimisticMessage.awaitingReplay)) { + const { optimisticQueued: _optimisticQueued, awaitingReplay: _awaitingReplay, ...confirmedMessage } = optimisticMessage return [ ...messages.slice(0, currentTurnUserIndex), confirmedMessage, @@ -7572,13 +7572,30 @@ function findCurrentTurnUserMessageIndex( modelContent: string, replayDisplay: RestoredUserDisplay, ): number { + // Guides are rendered before the CLI consumes them. The initial prompt's + // ACK can arrive after several guides, so match outstanding input in send + // order (including identical prompts), bounded by the current user turn. + let turnStart = 0 for (let index = messages.length - 1; index >= 0; index -= 1) { const message = messages[index] - if (message?.type !== 'user_text') { - continue + if (message?.type === 'user_text' && !message.optimisticQueued) { + turnStart = index + break } - return replayMatchesCurrentUserMessage(message, replayDisplay, modelContent) ? index : -1 } + for (let index = turnStart; index < messages.length; index += 1) { + const message = messages[index] + if ( + message?.type === 'user_text' && + (message.awaitingReplay || message.optimisticQueued) && + replayMatchesCurrentUserMessage(message, replayDisplay, modelContent) + ) return index + } + const current = messages[turnStart] + if ( + current?.type === 'user_text' && + replayMatchesCurrentUserMessage(current, replayDisplay, modelContent) + ) return turnStart return -1 } diff --git a/desktop/src/types/chat.ts b/desktop/src/types/chat.ts index 645e6fb6..7e428793 100644 --- a/desktop/src/types/chat.ts +++ b/desktop/src/types/chat.ts @@ -346,7 +346,7 @@ export type UIMessage = * the user's own prompt render identically, which is what flattened the * member transcript. */ - | { id: string; type: 'user_text'; content: string; sessionReferences?: Array<{ sessionId: string }>; collaboration?: { sourceSessionId: string; messageId?: string }; modelContent?: string; transcriptMessageId?: string; timestamp: number; attachments?: UIAttachment[]; pending?: boolean; optimisticQueued?: boolean; teammateFrom?: string } + | { id: string; type: 'user_text'; content: string; sessionReferences?: Array<{ sessionId: string }>; collaboration?: { sourceSessionId: string; messageId?: string }; modelContent?: string; transcriptMessageId?: string; timestamp: number; attachments?: UIAttachment[]; pending?: boolean; optimisticQueued?: boolean; awaitingReplay?: boolean; teammateFrom?: string } | { id: string; type: 'assistant_text'; content: string; transcriptMessageId?: string; timestamp: number; model?: string } | { id: string; type: 'thinking'; content: string; timestamp: number } | { diff --git a/src/server/__tests__/conversations.test.ts b/src/server/__tests__/conversations.test.ts index c09ddfa5..da679f06 100644 --- a/src/server/__tests__/conversations.test.ts +++ b/src/server/__tests__/conversations.test.ts @@ -2330,6 +2330,57 @@ describe('WebSocket Chat Integration', () => { expect(messages.some((m) => m.type === 'status' && m.state === 'idle')).toBe(true) }) + it('preserves delayed initial and guide replays without resubmitting user input (#1360)', async () => { + const createRes = await fetch(`${baseUrl}/api/sessions`, { + method: 'POST', + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify({ workDir: tmpDir }), + }) + expect(createRes.status).toBe(201) + const { sessionId } = await createRes.json() as { sessionId: string } + const initial = 'MOCK_GUIDE_REPLAY original request' + const guide = 'Keep the answer concise' + const messages: any[] = [] + const ws = new WebSocket(`${wsUrl}/ws/${sessionId}`) + let guideSent = false + await new Promise((resolve, reject) => { + const finish = (error?: Error) => { + clearTimeout(timeout) + ws.close() + error ? reject(error) : resolve() + } + const timeout = setTimeout(() => finish(new Error('Timed out waiting for delayed guide replay')), 10000) + ws.onerror = () => finish(new Error('Guide replay WebSocket failed')) + ws.onmessage = (event) => { + const message = JSON.parse(event.data as string) + messages.push(message) + if (message.type === 'connected') ws.send(JSON.stringify({ type: 'user_message', content: initial })) + if (message.type === 'thinking' && !guideSent) { + guideSent = true + ws.send(JSON.stringify({ type: 'user_message', content: guide })) + } + if (message.type === 'error') finish(new Error(message.message)) + if (message.type === 'message_complete') finish() + } + }) + expect(guideSent).toBe(true) + expect(messages.filter(message => message.type === 'user_message_replay').map(message => message.content)) + .toEqual([initial, guide]) + expect(messages.findIndex(message => message.type === 'thinking')) + .toBeLessThan(messages.findIndex(message => message.type === 'user_message_replay')) + const reply = `GUIDE_REPLAY_RECEIVED ${JSON.stringify([initial, guide])}` + expect(messages.filter(message => message.type === 'content_delta').map(message => message.text).join('')) + .toBe(reply) + const transcript = await sessionService.findSessionFile(sessionId) + expect(transcript).toBeTruthy() + const entries = (await fs.readFile(transcript!.filePath, 'utf8')).trim().split('\n').map(line => JSON.parse(line)) + const userEntries = entries.filter(entry => entry.type === 'user') + expect(userEntries.map(entry => entry.message.content.filter((block: any) => block.type === 'text').map((block: any) => block.text).join(' '))) + .toEqual([initial, guide]) + expect(entries.filter(entry => entry.type === 'assistant').map(entry => entry.message.content[0].text)) + .toEqual([reply]) + }, 15000) + it('should send user_message and receive streamed SDK response', async () => { const messages: any[] = [] const ws = new WebSocket(`${wsUrl}/ws/chat-test-3`) diff --git a/src/server/__tests__/fixtures/mock-sdk-cli.ts b/src/server/__tests__/fixtures/mock-sdk-cli.ts index bca4633d..c6f5d88d 100644 --- a/src/server/__tests__/fixtures/mock-sdk-cli.ts +++ b/src/server/__tests__/fixtures/mock-sdk-cli.ts @@ -41,6 +41,69 @@ const resumeUpstreamUrl = process.env.MOCK_SDK_RESUME_UPSTREAM_URL let initSent = false let firstUserExitScheduled = false let releaseReconnectStream: (() => void) | undefined +let guideReplayInitial: any | undefined +let guideReplayInputs: string[] = [] + +// Opt-in reproduction of QueryEngine's delayed initial ACK: partial thinking +// streams before the first complete assistant block, so an intervening guide +// can already be visible when the original user message is acknowledged. +async function handleGuideReplay(message: any): Promise { + const text = extractUserText(message) + if (!guideReplayInitial && !text.startsWith('MOCK_GUIDE_REPLAY')) return false + guideReplayInputs.push(text) + if (process.env.MOCK_SDK_GUIDE_REPLAY_LOG) { + await appendFile(process.env.MOCK_SDK_GUIDE_REPLAY_LOG, `${JSON.stringify({ uuid: message.uuid, text })}\n`) + } + if (!guideReplayInitial) { + guideReplayInitial = message + emit(ws, { type: 'stream_event', event: { type: 'message_start' }, session_id: sessionId }) + emit(ws, { type: 'stream_event', event: { type: 'content_block_start', index: 0, content_block: { type: 'thinking', thinking: '' } }, session_id: sessionId }) + emit(ws, { type: 'stream_event', event: { type: 'content_block_delta', index: 0, delta: { type: 'thinking_delta', thinking: 'Waiting for a guide message before acknowledging the original prompt.' } }, session_id: sessionId }) + return true + } + emit(ws, { type: 'stream_event', event: { type: 'content_block_stop', index: 0 }, session_id: sessionId }) + for (const input of [guideReplayInitial, message]) { + emit(ws, { type: 'user', message: input.message, uuid: input.uuid, isReplay: true, parent_tool_use_id: null, session_id: sessionId }) + } + const reply = `GUIDE_REPLAY_RECEIVED ${JSON.stringify(guideReplayInputs)}` + await appendGuideReplayHistory([guideReplayInitial, message], reply) + emit(ws, { type: 'stream_event', event: { type: 'content_block_start', index: 1, content_block: { type: 'text', text: '' } }, session_id: sessionId }) + emit(ws, { type: 'stream_event', event: { type: 'content_block_delta', index: 1, delta: { type: 'text_delta', text: reply } }, session_id: sessionId }) + emit(ws, { type: 'stream_event', event: { type: 'content_block_stop', index: 1 }, session_id: sessionId }) + emit(ws, { type: 'assistant', message: { role: 'assistant', content: [{ type: 'text', text: reply }] }, session_id: sessionId }) + emit(ws, { type: 'result', subtype: 'success', is_error: false, result: reply, usage: { input_tokens: 3, output_tokens: 2 }, session_id: sessionId }) + guideReplayInitial = undefined + guideReplayInputs = [] + return true +} + +async function appendGuideReplayHistory(inputs: any[], reply: string): Promise { + if (!process.env.CLAUDE_CONFIG_DIR) return + const projects = join(process.env.CLAUDE_CONFIG_DIR, 'projects') + for (const directory of await readdir(projects).catch(() => [])) { + const file = join(projects, directory, `${sessionId}.jsonl`) + try { await readFile(file) } catch { continue } + let parentUuid: string | null = null + const now = Date.now() + const records = inputs.map((input, index) => { + const uuid = input.uuid || crypto.randomUUID() + const record = { + type: 'user', uuid, parentUuid, sessionId, userType: 'external', + isSidechain: false, cwd: process.cwd(), message: input.message, + timestamp: new Date(now + index).toISOString(), + } + parentUuid = uuid + return record + }) + const assistant = { + type: 'assistant', uuid: crypto.randomUUID(), parentUuid, sessionId, + timestamp: new Date(now + inputs.length).toISOString(), + message: { id: `msg_${crypto.randomUUID()}`, type: 'message', role: 'assistant', model: 'mock-opus', content: [{ type: 'text', text: reply }] }, + } + await appendFile(file, `${[...records, assistant].map(record => JSON.stringify(record)).join('\n')}\n`) + return + } +} /** * Deterministic tool-use support. @@ -354,6 +417,8 @@ ws.addEventListener('message', (event) => { } if (parsed.type === 'user') { + sendInit() + if (await handleGuideReplay(parsed)) continue normalRunning = true try { sendInit()