diff --git a/desktop/src/stores/chatStore.test.ts b/desktop/src/stores/chatStore.test.ts index fbc63835..951e6bc8 100644 --- a/desktop/src/stores/chatStore.test.ts +++ b/desktop/src/stores/chatStore.test.ts @@ -17087,3 +17087,184 @@ describe('chatStore closed-session edit recovery', () => { expect(sendMock).not.toHaveBeenCalled() }) }) + +describe('chatStore streamed turn against its durable copy (#1476)', () => { + const sessionId = 'streamed-turn-durable-copy' + const THINKING = 'This is a large change (72 files, 14k insertions).' + const REPLY = 'Context gathered. Setting up the loop.' + + beforeEach(() => { + sendMock.mockReset() + vi.mocked(sessionsApi.getFullHistory).mockReset() + vi.mocked(sessionsApi.getFullHistory).mockResolvedValue({ messages: [] }) + useChatStore.setState({ ...initialState, sessions: {} }) + }) + + afterEach(() => { + const timer = useChatStore.getState().sessions[sessionId]?.elapsedTimer + if (timer) clearInterval(timer) + useChatStore.setState({ ...initialState, sessions: {} }) + }) + + const send = (message: ServerMessage) => useChatStore.getState().handleServerMessage(sessionId, message) + const timeline = () => (useChatStore.getState().sessions[sessionId]?.messages ?? []).map((message) => + message.type === 'thinking' || message.type === 'assistant_text' || message.type === 'user_text' + ? `${message.type}: ${message.content}` + : message.type) + + // DeepSeek-style tokenization streams the gap between words as its own delta. + async function streamTurn() { + send({ type: 'thinking', text: 'This is a large change (72 files,' }) + send({ type: 'thinking', text: ' ' }) + send({ type: 'thinking', text: '14k insertions).' }) + send({ type: 'content_start', blockType: 'tool_use', toolName: 'Bash', toolUseId: 'toolu_status' }) + send({ type: 'tool_use_complete', toolName: 'Bash', toolUseId: 'toolu_status', input: { command: 'git status' } }) + send({ type: 'tool_result', toolUseId: 'toolu_status', content: 'clean', isError: false }) + send({ type: 'content_start', blockType: 'text' }) + send({ type: 'content_delta', text: REPLY }) + await new Promise((resolve) => setTimeout(resolve, 80)) + send({ type: 'message_complete', usage: { input_tokens: 1, output_tokens: 1 } }) + } + + // The CLI persists each block when it completes, after its first live delta. + function durableTurn(): MessageEntry[] { + const at = (offset: number) => new Date(Date.now() + offset).toISOString() + return [ + { id: 'durable-user', type: 'user', timestamp: at(-60_000), content: 'review it' }, + { id: 'durable-thinking', type: 'assistant', timestamp: at(1_000), content: [{ type: 'thinking', thinking: THINKING }] }, + { id: 'durable-tool', type: 'assistant', timestamp: at(2_000), content: [{ type: 'tool_use', id: 'toolu_status', name: 'Bash', input: { command: 'git status' } }] }, + { id: 'durable-result', type: 'user', timestamp: at(3_000), content: [{ type: 'tool_result', tool_use_id: 'toolu_status', content: 'clean' }] }, + { id: 'durable-reply', type: 'assistant', timestamp: at(4_000), content: [{ type: 'text', text: REPLY }] }, + ] + } + + const expectedTimeline = [ + 'user_text: review it', + `thinking: ${THINKING}`, + 'tool_use', + 'tool_result', + `assistant_text: ${REPLY}`, + ] + + function touchLiveMessages() { + useChatStore.setState((state) => ({ + sessions: { + ...state.sessions, + [sessionId]: { ...state.sessions[sessionId]!, messages: [...state.sessions[sessionId]!.messages] }, + }, + })) + } + + it('keeps whitespace-only thinking deltas that continue a streaming block', () => { + useChatStore.setState({ sessions: { [sessionId]: makeSession({ chatState: 'thinking' }) } }) + + send({ type: 'thinking', text: '\n\n' }) + send({ type: 'thinking', text: 'first' }) + send({ type: 'thinking', text: '\n\n' }) + send({ type: 'thinking', text: 'second' }) + send({ type: 'thinking', text: ' ', complete: true }) + + expect(useChatStore.getState().sessions[sessionId]?.messages).toMatchObject([ + { type: 'thinking', content: 'first\n\nsecond' }, + ]) + }) + + it('does not duplicate the streamed turn when completion runs a cold history load', async () => { + useChatStore.setState({ sessions: { [sessionId]: makeSession({ + chatState: 'thinking', + historyHydrated: false, + historyStatus: 'error', + messages: [{ id: 'live-user', type: 'user_text', content: 'review it', timestamp: Date.now() - 60_000 }], + }) } }) + let resolveHistory!: (value: { messages: MessageEntry[] }) => void + vi.mocked(sessionsApi.getFullHistory).mockReturnValueOnce(new Promise((resolve) => { + resolveHistory = resolve + })) + + await streamTurn() + // Live activity (for example a background agent row) lands while REST is + // in flight, so the cold merge cannot hydrate transcript ids first. + touchLiveMessages() + resolveHistory({ messages: durableTurn() }) + await vi.waitFor(() => expect(useChatStore.getState().sessions[sessionId]?.historyHydrated).toBe(true)) + + expect(timeline()).toEqual(expectedTimeline) + }) + + it('does not duplicate the streamed turn when a bounded reload merges it', async () => { + useChatStore.setState({ sessions: { [sessionId]: makeSession({ + chatState: 'thinking', + historyHydrated: true, + historyStatus: 'ready', + messages: [{ id: 'live-user', type: 'user_text', content: 'review it', timestamp: Date.now() - 60_000, transcriptMessageId: 'durable-user' }], + }) } }) + vi.mocked(sessionsApi.getFullHistory).mockResolvedValueOnce({ messages: durableTurn() }) + await streamTurn() + await vi.waitFor(() => expect(sessionsApi.getFullHistory).toHaveBeenCalledTimes(1)) + await new Promise((resolve) => setTimeout(resolve, 0)) + + let resolveHistory!: (value: SessionHistoryPage) => void + vi.mocked(sessionsApi.getFullHistory).mockReturnValueOnce(new Promise((resolve) => { + resolveHistory = resolve + })) + const session = useChatStore.getState().sessions[sessionId]! + const reload = useChatStore.getState().reloadHistory(sessionId, { + messages: session.messages, + backgroundAgentTasks: session.backgroundAgentTasks, + }) + await vi.waitFor(() => expect(sessionsApi.getFullHistory).toHaveBeenCalledTimes(2)) + // An oversized tool result was omitted, so the page is not authoritative. + resolveHistory({ + messages: durableTurn(), + page: { + historyComplete: false, + omittedOversizedEntries: 1, + contentTruncated: false, + nextCursor: null, + hasMore: false, + }, + } as SessionHistoryPage) + await reload + + expect(timeline()).toEqual(expectedTimeline) + }) + + it('keeps an identical reply to a newer live prompt beside the durable one', async () => { + const now = Date.now() + useChatStore.setState({ sessions: { [sessionId]: makeSession({ + chatState: 'idle', + historyHydrated: false, + historyStatus: 'error', + messages: [ + { id: 'live-tool', type: 'tool_use', toolName: 'Bash', toolUseId: 'toolu_status', input: {}, timestamp: now - 3_000 }, + { id: 'live-result', type: 'tool_result', toolUseId: 'toolu_status', content: 'clean', isError: false, timestamp: now - 2_000 }, + { id: 'live-again', type: 'user_text', content: 'again', timestamp: now - 1_000 }, + { id: 'live-done', type: 'assistant_text', content: 'Done.', timestamp: now }, + ], + }) } }) + let resolveHistory!: (value: { messages: MessageEntry[] }) => void + vi.mocked(sessionsApi.getFullHistory).mockReturnValueOnce(new Promise((resolve) => { + resolveHistory = resolve + })) + + const load = useChatStore.getState().loadHistory(sessionId) + touchLiveMessages() + const at = (offset: number) => new Date(now + offset).toISOString() + resolveHistory({ messages: [ + { id: 'durable-user', type: 'user', timestamp: at(-5_000), content: 'review it' }, + { id: 'durable-tool', type: 'assistant', timestamp: at(-3_000), content: [{ type: 'tool_use', id: 'toolu_status', name: 'Bash', input: {} }] }, + { id: 'durable-result', type: 'user', timestamp: at(-2_000), content: [{ type: 'tool_result', tool_use_id: 'toolu_status', content: 'clean' }] }, + { id: 'durable-done', type: 'assistant', timestamp: at(-1_500), content: [{ type: 'text', text: 'Done.' }] }, + ] }) + await load + + expect(timeline()).toEqual([ + 'user_text: review it', + 'tool_use', + 'tool_result', + 'assistant_text: Done.', + 'user_text: again', + 'assistant_text: Done.', + ]) + }) +}) diff --git a/desktop/src/stores/chatStore.ts b/desktop/src/stores/chatStore.ts index 18a03b22..50f02ef7 100644 --- a/desktop/src/stores/chatStore.ts +++ b/desktop/src/stores/chatStore.ts @@ -1921,6 +1921,88 @@ function overlayLiveHistoryMessage( return restoredMessage } +type CoalescibleProseMessage = Extract + +function isCoalescibleProse(message: UIMessage): message is CoalescibleProseMessage { + return message.type === 'thinking' || + message.type === 'assistant_text' || + message.type === 'user_text' +} + +function sameRestoredProse(live: CoalescibleProseMessage, restored: UIMessage): boolean { + if (restored.type !== live.type) return false + // Live thinking concatenates consecutive streamed blocks while history joins + // them with blank lines, so only the non-whitespace text identifies it. + if (live.type === 'thinking') { + return live.content.replace(/\s+/g, '') === restored.content.replace(/\s+/g, '') + } + return live.content.trim() === restored.content.trim() +} + +/** + * Thinking has no durable identity, and live text only carries a transcript id + * once an unchanged cache was hydrated. Without this pass a cold merge keeps + * the live copy beside its durable copy (issue #1476). A live row is the same + * output only when both lists put it in the same gap between shared rows and + * the same number of user turns away from that shared row; equal prose with no + * shared anchor may still be a genuine newer reply and stays separate. + */ +function coalesceAnchoredLiveProse( + merged: UIMessage[], + liveIndexes: number[], + restoredCount: number, +): void { + const claimed = new Set(liveIndexes.filter((index) => index < restoredCount)) + const countUserRows = (rows: UIMessage[]) => + rows.reduce((count, row) => count + (row.type === 'user_text' ? 1 : 0), 0) + const liveRowsBetween = (from: number, to: number) => + liveIndexes.slice(from, to).map((index) => merged[index]!) + + let previousAnchorOffset = -1 + for (let offset = 0; offset < liveIndexes.length; offset++) { + const index = liveIndexes[offset]! + if (index < restoredCount) { + previousAnchorOffset = offset + continue + } + const liveMessage = merged[index]! + if (!isCoalescibleProse(liveMessage)) continue + + let nextAnchorOffset = offset + 1 + while ( + nextAnchorOffset < liveIndexes.length && + liveIndexes[nextAnchorOffset]! >= restoredCount + ) nextAnchorOffset++ + const hasNextAnchor = nextAnchorOffset < liveIndexes.length + if (previousAnchorOffset < 0 && !hasNextAnchor) continue + + const lowerBound = previousAnchorOffset < 0 ? -1 : liveIndexes[previousAnchorOffset]! + const upperBound = hasNextAnchor ? liveIndexes[nextAnchorOffset]! : restoredCount + if (upperBound <= lowerBound) continue + + // Count user turns from the preceding shared row when there is one, + // otherwise back from the following shared row. + const liveTurnDistance = previousAnchorOffset >= 0 + ? countUserRows(liveRowsBetween(previousAnchorOffset + 1, offset)) + : countUserRows(liveRowsBetween(offset + 1, nextAnchorOffset)) + let matchedIndex: number | undefined + for (let candidate = lowerBound + 1; candidate < upperBound; candidate++) { + if (claimed.has(candidate) || !sameRestoredProse(liveMessage, merged[candidate]!)) continue + const restoredTurnDistance = previousAnchorOffset >= 0 + ? countUserRows(merged.slice(lowerBound + 1, candidate)) + : countUserRows(merged.slice(candidate + 1, upperBound)) + if (restoredTurnDistance !== liveTurnDistance) continue + matchedIndex = candidate + break + } + if (matchedIndex === undefined) continue + + claimed.add(matchedIndex) + liveIndexes[offset] = matchedIndex + previousAnchorOffset = offset + } +} + function mergeColdRestoredHistoryIntoLiveMessages( restoredMessages: UIMessage[], liveMessages: UIMessage[], @@ -1988,10 +2070,12 @@ function mergeColdRestoredHistoryIntoLiveMessages( } } + const restoredCount = restoredMessages.length + coalesceAnchoredLiveProse(merged, liveIndexes, restoredCount) + // A bounded REST page is only a suffix of the transcript. Unmatched live // rows can be an older cached prefix, not just new output. Place them by // shared identities while leaving the durable page's order untouched. - const restoredCount = restoredMessages.length const earliestTimestamp = restoredMessages.reduce((earliest, message) => Number.isFinite(message.timestamp) ? Math.min(earliest, message.timestamp) : earliest, Infinity) const olderPrefix = new Set() @@ -5186,9 +5270,13 @@ export const useChatStore = create((setState, get) => { const base = pendingText.trim() ? appendAssistantTextMessage(s.messages, pendingText, Date.now()) : s.messages + const lastIndex = findStreamMergeTargetIndex(base) + const last = lastIndex >= 0 ? base[lastIndex] : undefined // 服务端两个 thinking 发射点都做了非空过滤,但 `&& delta.thinking` 是真值 // 判断,纯空白仍能漏过来,落到下面就是一个点开什么都没有的空壳气泡。 - if (!msg.text.trim()) { + // 只挡整块空白和会新开气泡的空白;正在流式的思考里,单独的空格/换行 + // 增量是正文的一部分,丢掉它会让实时内容和 transcript 对不上(#1476)。 + if (!msg.text.trim() && (msg.complete === true || last?.type !== 'thinking')) { skippedThinkingBlock = true return { messages: base, streamingText: '' } } @@ -5200,8 +5288,6 @@ export const useChatStore = create((setState, get) => { skippedThinkingBlock = true return { messages: base, streamingText: '' } } - const lastIndex = findStreamMergeTargetIndex(base) - const last = lastIndex >= 0 ? base[lastIndex] : undefined if (last && last.type === 'thinking') { const updated = [...base] updated[lastIndex] = {