diff --git a/src/main/agent/acp/compatibility/adapters.ts b/src/main/agent/acp/compatibility/adapters.ts index fd92e3fb47..28a467cff8 100644 --- a/src/main/agent/acp/compatibility/adapters.ts +++ b/src/main/agent/acp/compatibility/adapters.ts @@ -20,6 +20,7 @@ import { import { createAcpPromptTerminalEvents } from '@/agent/acp/runtime/acpContentMapper' import { createState, + markStreamChanged, type DeepChatEventPublisher, type DeepChatSessionUpdatePublisher, type IoParams, @@ -159,7 +160,7 @@ export class AcpCompatibilityProjectionAdapter implements AcpCompatibilityProjec ? block.extra.permissionType : 'all' markStreamingProviderPermissionResolved(block, granted, permissionType) - state.stream.dirty = true + markStreamChanged(state.stream) this.flushIfDirty(state) } diff --git a/src/main/agent/deepchat/runtime/accumulator.ts b/src/main/agent/deepchat/runtime/accumulator.ts index bbbfe4b3c9..55ab0acf5e 100644 --- a/src/main/agent/deepchat/runtime/accumulator.ts +++ b/src/main/agent/deepchat/runtime/accumulator.ts @@ -4,7 +4,7 @@ import type { ProviderUrlSourcePayload } from '@shared/types/core/llm-events' import type { ChatMessageProviderOptions } from '@shared/types/core/chat-message' -import type { StreamState } from './types' +import { markStreamChanged, type StreamState } from './types' const MAX_VISIBLE_SEARCH_PAGES = 6 @@ -144,7 +144,7 @@ export function accumulate(state: StreamState, event: LLMCoreStreamEvent): void if (state.firstTokenTime === null) state.firstTokenTime = Date.now() const block = getCurrentBlock(state.blocks, 'content', event.provider_options) block.content += event.content - state.dirty = true + markStreamChanged(state) break } case 'reasoning': { @@ -167,12 +167,12 @@ export function accumulate(state: StreamState, event: LLMCoreStreamEvent): void } const reasoningTime = block.reasoning_time as { start: number; end: number } updateReasoningMetadata(state, reasoningTime.start, reasoningTime.end) - state.dirty = true + markStreamChanged(state) break } case 'plan': { if (finalizeTrailingPendingNarrativeBlocks(state.blocks)) { - state.dirty = true + markStreamChanged(state) } const revision = event.revision ?? (state.latestAgentPlanSnapshot?.revision ?? 0) + 1 state.latestAgentPlanSnapshot = { @@ -209,7 +209,7 @@ export function accumulate(state: StreamState, event: LLMCoreStreamEvent): void executionOwner: event.tool_call_execution_owner ?? 'deepchat', providerOptions: event.provider_options }) - state.dirty = true + markStreamChanged(state) break } case 'tool_call_chunk': { @@ -229,7 +229,7 @@ export function accumulate(state: StreamState, event: LLMCoreStreamEvent): void } } } - state.dirty = true + markStreamChanged(state) } break } @@ -265,7 +265,7 @@ export function accumulate(state: StreamState, event: LLMCoreStreamEvent): void }) } state.pendingToolCalls.delete(event.tool_call_id) - state.dirty = true + markStreamChanged(state) } break } @@ -296,12 +296,12 @@ export function accumulate(state: StreamState, event: LLMCoreStreamEvent): void } } state.blocks.push(block) - state.dirty = true + markStreamChanged(state) break } case 'provider_url_source': { if (appendProviderUrlSource(state.blocks, event.provider_url_source)) { - state.dirty = true + markStreamChanged(state) } break } @@ -318,7 +318,7 @@ export function accumulate(state: StreamState, event: LLMCoreStreamEvent): void } } state.blocks.push(block) - state.dirty = true + markStreamChanged(state) break } case 'usage': { @@ -354,7 +354,7 @@ export function accumulate(state: StreamState, event: LLMCoreStreamEvent): void if (block.status === 'pending') block.status = 'error' } state.stopReason = 'error' - state.dirty = true + markStreamChanged(state) break } default: diff --git a/src/main/agent/deepchat/runtime/deepChatLoopRunner.ts b/src/main/agent/deepchat/runtime/deepChatLoopRunner.ts index 5272bf64b9..89ddf35b4d 100644 --- a/src/main/agent/deepchat/runtime/deepChatLoopRunner.ts +++ b/src/main/agent/deepchat/runtime/deepChatLoopRunner.ts @@ -661,6 +661,7 @@ export function buildTapeViewSelection( export class DeepChatLoopRunner { private readonly toolSurfaceAdapterHistory = new ToolSurfaceAdapterHistory() + private rateLimitRevision = 0 constructor(private readonly ports: DeepChatLoopRunnerPorts) {} @@ -2726,6 +2727,7 @@ export class DeepChatLoopRunner { sessionId, messageId, updatedAt: Date.now(), + revision: ++this.rateLimitRevision, blocks: cloneBlocksForRenderer([block]) }) } @@ -2737,6 +2739,7 @@ export class DeepChatLoopRunner { sessionId, messageId, updatedAt: Date.now(), + revision: ++this.rateLimitRevision, blocks: [] }) } diff --git a/src/main/agent/deepchat/runtime/dispatch.ts b/src/main/agent/deepchat/runtime/dispatch.ts index 823506d822..cce3e77449 100644 --- a/src/main/agent/deepchat/runtime/dispatch.ts +++ b/src/main/agent/deepchat/runtime/dispatch.ts @@ -32,6 +32,7 @@ import type { ToolCallResult, ToolDispatchCollaborators } from './types' +import { markStreamChanged } from './types' import type { ChatMessage, ChatMessageProviderOptions, @@ -1253,7 +1254,7 @@ function finalizePendingNarrativeBeforeToolSettlement(state: StreamState): void } finalizeTrailingPendingNarrativeBlocks(state.blocks) - state.dirty = true + markStreamChanged(state) } function applyFinalizedToolResults(params: { @@ -1387,7 +1388,7 @@ function applyFinalizedToolResults(params: { } } - state.dirty = true + markStreamChanged(state) return interactions } @@ -1673,7 +1674,7 @@ async function reviewAutoApproveAction(params: { } if (setToolCallAutoApproveReviewing(batchToolCallBlocks, execution.completedToolCall.id, true)) { - state.dirty = true + markStreamChanged(state) rendererFlushHandle.flush() } try { @@ -1705,7 +1706,7 @@ async function reviewAutoApproveAction(params: { if ( setToolCallAutoApproveReviewing(batchToolCallBlocks, execution.completedToolCall.id, false) ) { - state.dirty = true + markStreamChanged(state) rendererFlushHandle.flush() } } @@ -1772,7 +1773,7 @@ function appendPermissionActionBlock( ...(permission.rememberable === false ? { rememberable: false } : {}) } }) - state.dirty = true + markStreamChanged(state) return { type: 'permission', origin, @@ -1830,7 +1831,7 @@ function appendQuestionActionBlock( ...extra } }) - state.dirty = true + markStreamChanged(state) return { type: 'question', origin, @@ -1896,8 +1897,8 @@ function appendSkillDraftQuestionActionBlock( ) } -function flushBlocksToRenderer(io: IoParams, blocks: AssistantMessageBlock[]): void { - const renderedBlocks = cloneBlocksForRenderer(blocks) +function flushBlocksToRenderer(io: IoParams, state: StreamState): void { + const renderedBlocks = cloneBlocksForRenderer(state.blocks) io.publishEvent('chat.stream.updated', { kind: 'snapshot', requestId: io.requestId, @@ -1906,6 +1907,7 @@ function flushBlocksToRenderer(io: IoParams, blocks: AssistantMessageBlock[]): v providerId: io.providerId, modelId: io.modelId, updatedAt: Date.now(), + revision: state.blocksRevision, blocks: renderedBlocks }) @@ -1914,10 +1916,10 @@ function flushBlocksToRenderer(io: IoParams, blocks: AssistantMessageBlock[]): v kind: 'blocks', updatedAt: Date.now(), messageId: io.messageId, - previewMarkdown: buildAssistantPreviewMarkdown(blocks), - responseMarkdown: buildAssistantResponseMarkdown(blocks), - deliverySegments: buildAssistantDeliverySegments(io.messageId, blocks), - waitingInteraction: extractWaitingInteraction(blocks, io.messageId) + previewMarkdown: buildAssistantPreviewMarkdown(state.blocks), + responseMarkdown: buildAssistantResponseMarkdown(state.blocks), + deliverySegments: buildAssistantDeliverySegments(io.messageId, state.blocks), + waitingInteraction: extractWaitingInteraction(state.blocks, io.messageId) }) } @@ -2130,7 +2132,7 @@ async function runToolCall(params: { } state.latestAgentPlanSnapshot = snapshot publishPlanUpdated(io, snapshot) - state.dirty = true + markStreamChanged(state) scheduleRendererFlush(state, rendererFlushHandle) return } @@ -2149,7 +2151,7 @@ async function runToolCall(params: { update.responseMarkdown, update.progressJson ) - state.dirty = true + markStreamChanged(state) scheduleRendererFlush(state, rendererFlushHandle) } @@ -3204,7 +3206,7 @@ export async function settleToolBatch( content: errorText }) updateToolCallBlock(batchToolCallBlocks, tc.id, errorText, true) - state.dirty = true + markStreamChanged(state) batchState.committedResultCallIds.add(tc.id) executed += 1 persistToolExecutionState(io, state, rendererFlushHandle) @@ -3534,7 +3536,7 @@ export function finalizePaused(state: StreamState, io: IoParams): void { stampGenerationTiming(state) io.messageStore.updateAssistantContent(io.messageId, state.blocks, JSON.stringify(state.metadata)) - flushBlocksToRenderer(io, state.blocks) + flushBlocksToRenderer(io, state) io.publishEvent('chat.stream.completed', { requestId: io.requestId, sessionId: io.sessionId, @@ -3547,6 +3549,7 @@ export function finalize(state: StreamState, io: IoParams): void { for (const block of state.blocks) { if (block.status === 'pending') block.status = 'success' } + markStreamChanged(state) stampPlanTerminalIfOpen(state, io, state.planTerminalReason) stampGenerationTiming(state) @@ -3556,7 +3559,7 @@ export function finalize(state: StreamState, io: IoParams): void { state.blocks, JSON.stringify(state.metadata) ) - flushBlocksToRenderer(io, state.blocks) + flushBlocksToRenderer(io, state) io.publishEvent('chat.stream.completed', { requestId: io.requestId, sessionId: io.sessionId, @@ -3568,6 +3571,7 @@ export function finalize(state: StreamState, io: IoParams): void { export function finalizeError(state: StreamState, io: IoParams, error: unknown): void { const errorMessage = error instanceof Error ? error.message : String(error) state.blocks = buildTerminalErrorBlocks(state.blocks, errorMessage) + markStreamChanged(state) stampPlanTerminalIfOpen( state, io, @@ -3577,7 +3581,7 @@ export function finalizeError(state: StreamState, io: IoParams, error: unknown): stampGenerationTiming(state) io.messageStore.setMessageError(io.messageId, state.blocks, JSON.stringify(state.metadata)) - flushBlocksToRenderer(io, state.blocks) + flushBlocksToRenderer(io, state) io.publishEvent('chat.stream.failed', { requestId: io.requestId, sessionId: io.sessionId, @@ -3596,5 +3600,5 @@ export function persistAbortExceptionPlanState(state: StreamState, io: IoParams) } io.messageStore.updateAssistantContent(io.messageId, state.blocks) - flushBlocksToRenderer(io, state.blocks) + flushBlocksToRenderer(io, state) } diff --git a/src/main/agent/deepchat/runtime/echo.ts b/src/main/agent/deepchat/runtime/echo.ts index 98db4f5480..5d8fcc1b8a 100644 --- a/src/main/agent/deepchat/runtime/echo.ts +++ b/src/main/agent/deepchat/runtime/echo.ts @@ -23,6 +23,7 @@ export function startEcho(state: StreamState, io: IoParams): EchoHandle { providerId: io.providerId, modelId: io.modelId, updatedAt: Date.now(), + revision: state.blocksRevision, blocks: renderedBlocks }) } diff --git a/src/main/agent/deepchat/runtime/process.ts b/src/main/agent/deepchat/runtime/process.ts index 67c04c428d..4a3a7a2e43 100644 --- a/src/main/agent/deepchat/runtime/process.ts +++ b/src/main/agent/deepchat/runtime/process.ts @@ -13,6 +13,7 @@ import type { StreamState, ToolCallResult } from './types' +import { markStreamChanged } from './types' import { accumulate, commitRoundUsage, finalizeTrailingPendingNarrativeBlocks } from './accumulator' import { startEcho } from './echo' import { @@ -141,6 +142,7 @@ function stripTrailingErrorBlock(state: StreamState, message: string): void { const lastBlock = state.blocks[state.blocks.length - 1] if (lastBlock?.type === 'error' && lastBlock.content === message) { state.blocks.pop() + markStreamChanged(state) } } @@ -183,15 +185,15 @@ function stampProviderAttemptIdentity( function closePreviousProviderAttemptNarrative( blocks: AssistantMessageBlock[], identity: DeepChatProviderAttemptIdentity | null -): void { - if (!identity) return +): boolean { + if (!identity) return false const last = blocks[blocks.length - 1] if ( !last || last.status !== 'pending' || (last.type !== 'content' && last.type !== 'reasoning_content') ) { - return + return false } const previousLogicalRound = last.extra?.providerLogicalRound const previousRequestSeq = last.extra?.providerRequestSeq @@ -204,9 +206,10 @@ function closePreviousProviderAttemptNarrative( previousRequestSeq === identity.requestSeq && previousPhysicalAttempt === identity.physicalAttempt) ) { - return + return false } last.status = 'success' + return true } function stampRunOutcome( @@ -250,7 +253,7 @@ function markUnexecutedToolCallsForLimit(state: StreamState): void { ...block.extra, toolCallSkippedReason: 'max_tool_calls' } - state.dirty = true + markStreamChanged(state) } } @@ -364,7 +367,7 @@ function markOtherTruncatedToolCallsIncomplete( ...block.extra, toolCallIncompleteReason: 'max_tokens' } - state.dirty = true + markStreamChanged(state) } } @@ -655,7 +658,7 @@ export function appendStreamingProviderPermissionBlock( } state.blocks.push(actionBlock) - state.dirty = true + markStreamChanged(state) return { actionBlock, @@ -971,7 +974,9 @@ export async function processStream(params: ProcessParams): Promise 0) { state.blocks = JSON.parse(JSON.stringify(initialBlocks)) as typeof state.blocks - state.dirty = normalizeInheritedUnresolvedBlocks(state.blocks) || state.dirty + if (normalizeInheritedUnresolvedBlocks(state.blocks)) { + markStreamChanged(state) + } } state.metadata.runId = run.runId const echo = startEcho(state, io) @@ -1225,7 +1230,9 @@ export async function processStream(params: ProcessParams): Promise void @@ -363,6 +364,12 @@ export function createState(): StreamState { stopReason: null, roundUsage: null, toolCallCount: 0, - dirty: false + dirty: false, + blocksRevision: 0 } } + +export function markStreamChanged(state: StreamState): void { + state.dirty = true + state.blocksRevision += 1 +} diff --git a/src/renderer/src/stores/ui/message.ts b/src/renderer/src/stores/ui/message.ts index 4e7331f1c1..796b59dc20 100644 --- a/src/renderer/src/stores/ui/message.ts +++ b/src/renderer/src/stores/ui/message.ts @@ -56,6 +56,10 @@ export const useMessageStore = defineStore('message', () => { const currentStreamRequestId = toStoreStateRef(streamStateStore, 'currentStreamRequestId') const currentStreamMessageId = toStoreStateRef(streamStateStore, 'currentStreamMessageId') const currentStreamMetadata = toStoreStateRef(streamStateStore, 'currentStreamMetadata') + const currentStreamBlocksRevision = toStoreStateRef( + streamStateStore, + 'currentStreamBlocksRevision' + ) const streamRevision = toStoreStateRef(streamStateStore, 'streamRevision') // --- State --- @@ -81,6 +85,10 @@ export const useMessageStore = defineStore('message', () => { // Stream message ids currently being hydrated into the cache as a placeholder // record (before the backend persists them). Prevents re-entrant duplicate inserts. const hydratingStreamMessageIds = new Set() + // Applied stream revisions are scoped to the stream identity (requestId): a new + // request or resume reusing the same message id restarts from revision 0 and must + // not be deduped against the previous request's revision. + const appliedStreamRevision = new Map() let latestLoadRequestId = 0 let latestHistoryRequestId = 0 let latestLoadSessionId: string | null = null @@ -581,11 +589,14 @@ export const useMessageStore = defineStore('message', () => { streamMessageId && !isEphemeralStreamMessageId(streamMessageId) ) { + appliedStreamRevision.delete(streamMessageId) applyStreamingBlocksToMessage( streamMessageId, view.sessionId, streamingBlocks.value as AssistantMessageBlock[], - currentStreamMetadata.value ?? undefined + currentStreamMetadata.value ?? undefined, + currentStreamBlocksRevision.value, + currentStreamRequestId.value ?? undefined ) } } @@ -938,6 +949,9 @@ export const useMessageStore = defineStore('message', () => { markLiveMessageViewMutation(sessionId) for (const record of changedRecords) { parsedMessageCache.delete(record.id) + // Persisted data replaces the live-folded record; drop the applied + // revision so a later snapshot re-evaluates content from scratch. + appliedStreamRevision.delete(record.id) upsertMessageRecord(record) } lastPersistedRevision.value += 1 @@ -958,6 +972,7 @@ export const useMessageStore = defineStore('message', () => { historyLoadError.value = false parsedMessageCache.clear() hydratingStreamMessageIds.clear() + appliedStreamRevision.clear() recentSessionViews.clear() messageMutationRevisions.clear() recentViewInvalidationRevisions.clear() @@ -990,17 +1005,38 @@ export const useMessageStore = defineStore('message', () => { messageId: string, conversationId: string, blocks: AssistantMessageBlock[], - metadata?: { providerId?: string; modelId?: string } + metadata?: { providerId?: string; modelId?: string }, + revision?: number, + requestId?: string ): void { if (committedSessionId.value !== conversationId) return - const serializedBlocks = JSON.stringify(blocks) - const serializedMetadata = JSON.stringify({ - ...(metadata?.providerId ? { provider: metadata.providerId } : {}), - ...(metadata?.modelId ? { model: metadata.modelId } : {}) - }) const existing = messageCache.value.get(messageId) if (existing) { if (existing.sessionId !== conversationId) return + + const lastApplied = appliedStreamRevision.get(messageId) + const lastRevision = + lastApplied && (requestId === undefined || lastApplied.requestId === requestId) + ? lastApplied.revision + : undefined + if ( + revision !== undefined && + lastRevision !== undefined && + revision <= lastRevision && + existing.status === 'pending' + ) { + cacheStreamingAssistantBlocks(existing, blocks) + return + } + + const serializedBlocks = JSON.stringify(blocks) + const serializedMetadata = JSON.stringify({ + ...(metadata?.providerId ? { provider: metadata.providerId } : {}), + ...(metadata?.modelId ? { model: metadata.modelId } : {}) + }) + if (revision !== undefined) { + appliedStreamRevision.set(messageId, { requestId: requestId ?? null, revision }) + } const nextMetadata = serializedMetadata === '{}' ? existing.metadata : serializedMetadata if ( existing.content === serializedBlocks && @@ -1023,6 +1059,11 @@ export const useMessageStore = defineStore('message', () => { return } + const serializedBlocks = JSON.stringify(blocks) + const serializedMetadata = JSON.stringify({ + ...(metadata?.providerId ? { provider: metadata.providerId } : {}), + ...(metadata?.modelId ? { model: metadata.modelId } : {}) + }) if (hydratingStreamMessageIds.has(messageId)) return hydratingStreamMessageIds.add(messageId) markLiveMessageViewMutation(conversationId) @@ -1041,6 +1082,9 @@ export const useMessageStore = defineStore('message', () => { createdAt: now, updatedAt: now } + if (revision !== undefined) { + appliedStreamRevision.set(messageId, { requestId: requestId ?? null, revision }) + } upsertMessageRecord(nextRecord) cacheStreamingAssistantBlocks(nextRecord, blocks) hydratingStreamMessageIds.delete(messageId) @@ -1052,8 +1096,24 @@ export const useMessageStore = defineStore('message', () => { sessionId: currentStreamSessionId.value, requestId: currentStreamRequestId.value }), - setStreamingState: ({ sessionId, requestId, messageId, updatedAt, blocks, metadata }) => { - streamStateStore.setStream(sessionId, blocks, messageId, metadata, requestId, updatedAt) + setStreamingState: ({ + sessionId, + requestId, + messageId, + updatedAt, + revision, + blocks, + metadata + }) => { + streamStateStore.setStream( + sessionId, + blocks, + messageId, + metadata, + requestId, + updatedAt, + revision + ) }, clearStreamingState, loadMessages, @@ -1065,6 +1125,11 @@ export const useMessageStore = defineStore('message', () => { registerStoreCleanup(messageIpcBinding.cleanup) function purgeSessionTracking(sessionId: string): void { + for (const [id, record] of messageCache.value) { + if (record.sessionId === sessionId) { + appliedStreamRevision.delete(id) + } + } recentSessionViews.delete(sessionId) messageMutationRevisions.delete(sessionId) recentViewInvalidationRevisions.delete(sessionId) diff --git a/src/renderer/src/stores/ui/messageIpc.ts b/src/renderer/src/stores/ui/messageIpc.ts index 8cfb42479b..84366e1b6d 100644 --- a/src/renderer/src/stores/ui/messageIpc.ts +++ b/src/renderer/src/stores/ui/messageIpc.ts @@ -13,6 +13,7 @@ interface BindMessageStoreIpcOptions { requestId: string messageId?: string updatedAt: number + revision: number blocks: AssistantMessageBlock[] metadata?: { providerId?: string; modelId?: string } }) => void @@ -24,7 +25,9 @@ interface BindMessageStoreIpcOptions { messageId: string, sessionId: string, blocks: AssistantMessageBlock[], - metadata?: { providerId?: string; modelId?: string } + metadata?: { providerId?: string; modelId?: string }, + revision?: number, + requestId?: string ) => void isEphemeralStreamMessageId: (messageId: string) => boolean } @@ -180,6 +183,7 @@ export function bindMessageStoreIpc(options: BindMessageStoreIpcOptions): Messag requestId: payload.requestId, messageId: streamMessageId, updatedAt: payload.updatedAt, + revision: payload.revision, blocks, metadata: { providerId: payload.providerId, @@ -192,10 +196,17 @@ export function bindMessageStoreIpc(options: BindMessageStoreIpcOptions): Messag options.applyStreamingBlocksToMessage && !options.isEphemeralStreamMessageId(streamMessageId) ) { - options.applyStreamingBlocksToMessage(streamMessageId, payload.sessionId, blocks, { - providerId: payload.providerId, - modelId: payload.modelId - }) + options.applyStreamingBlocksToMessage( + streamMessageId, + payload.sessionId, + blocks, + { + providerId: payload.providerId, + modelId: payload.modelId + }, + payload.revision, + payload.requestId + ) } }), chatClient.onStreamCompleted((payload) => { diff --git a/src/renderer/src/stores/ui/stream.ts b/src/renderer/src/stores/ui/stream.ts index 9cc7704ec0..6f458f3d5e 100644 --- a/src/renderer/src/stores/ui/stream.ts +++ b/src/renderer/src/stores/ui/stream.ts @@ -13,6 +13,7 @@ export const useStreamStateStore = defineStore('streamState', () => { const currentStreamUpdatedAt = ref(0) const currentStreamMetadata = ref<{ providerId?: string; modelId?: string } | null>(null) const streamRevision = ref(0) + const currentStreamBlocksRevision = ref(0) function setStream( sessionId: string, @@ -20,7 +21,8 @@ export const useStreamStateStore = defineStore('streamState', () => { messageId?: string, metadata?: { providerId?: string; modelId?: string }, requestId?: string, - updatedAt?: number + updatedAt?: number, + blocksRevision?: number ): void { isStreaming.value = true currentStreamSessionId.value = sessionId @@ -28,6 +30,7 @@ export const useStreamStateStore = defineStore('streamState', () => { currentStreamMessageId.value = messageId ?? null currentStreamUpdatedAt.value = updatedAt ?? 0 currentStreamMetadata.value = metadata ?? null + currentStreamBlocksRevision.value = blocksRevision ?? 0 streamingBlocks.value = blocks streamRevision.value += 1 } @@ -40,6 +43,7 @@ export const useStreamStateStore = defineStore('streamState', () => { currentStreamMessageId.value = null currentStreamUpdatedAt.value = 0 currentStreamMetadata.value = null + currentStreamBlocksRevision.value = 0 streamRevision.value += 1 } @@ -51,6 +55,7 @@ export const useStreamStateStore = defineStore('streamState', () => { currentStreamMessageId, currentStreamUpdatedAt, currentStreamMetadata, + currentStreamBlocksRevision, streamRevision, setStream, clearStreamingState diff --git a/src/shared/contracts/events/chat.events.ts b/src/shared/contracts/events/chat.events.ts index 82666c854f..8d4349588e 100644 --- a/src/shared/contracts/events/chat.events.ts +++ b/src/shared/contracts/events/chat.events.ts @@ -17,6 +17,7 @@ export const chatStreamUpdatedEvent = defineEventContract({ providerId: z.string().optional(), modelId: z.string().optional(), updatedAt: TimestampMsSchema, + revision: z.number().int().nonnegative(), blocks: z.array(AssistantMessageBlockSchema) }) }) diff --git a/test/main/agent/deepchat/runtime/accumulator.test.ts b/test/main/agent/deepchat/runtime/accumulator.test.ts index 33c87e320a..c2cd8b96d8 100644 --- a/test/main/agent/deepchat/runtime/accumulator.test.ts +++ b/test/main/agent/deepchat/runtime/accumulator.test.ts @@ -473,6 +473,31 @@ describe('accumulate', () => { expect(state.dirty).toBe(false) }) + it('bumps blocksRevision exactly when block content changes', () => { + expect(state.blocksRevision).toBe(0) + + accumulate(state, { type: 'text', content: 'first' }) + expect(state.blocksRevision).toBe(1) + + accumulate(state, { type: 'text', content: ' second' }) + expect(state.blocksRevision).toBe(2) + + state.dirty = false + accumulate(state, { + type: 'usage', + usage: { prompt_tokens: 1, completion_tokens: 1, total_tokens: 2 } + }) + accumulate(state, { type: 'stop', stop_reason: 'tool_use' }) + expect(state.blocksRevision).toBe(2) + + accumulate(state, { + type: 'tool_call_start', + tool_call_id: 'tc1', + tool_call_name: 'search' + }) + expect(state.blocksRevision).toBe(3) + }) + it('sets firstTokenTime once on first text event', () => { expect(state.firstTokenTime).toBeNull() diff --git a/test/main/agent/deepchat/runtime/echo.test.ts b/test/main/agent/deepchat/runtime/echo.test.ts index 8a46e3e06d..7dd95cc24b 100644 --- a/test/main/agent/deepchat/runtime/echo.test.ts +++ b/test/main/agent/deepchat/runtime/echo.test.ts @@ -1,6 +1,6 @@ import { describe, it, expect, vi, beforeEach, afterEach } from 'vitest' import type { StreamState, IoParams } from '@/agent/deepchat/runtime/types' -import { createState } from '@/agent/deepchat/runtime/types' +import { createState, markStreamChanged } from '@/agent/deepchat/runtime/types' vi.mock('@/events', () => ({ STREAM_EVENTS: { @@ -11,6 +11,7 @@ vi.mock('@/events', () => ({ })) import { startEcho } from '@/agent/deepchat/runtime/echo' +import { accumulate } from '@/agent/deepchat/runtime/accumulator' import { cloneBlocksForRenderer } from '@/session/clientMessageProjection' const publishDeepchatEvent = vi.fn() @@ -180,6 +181,66 @@ describe('echo', () => { echo.stop() }) + it('emits the current blocksRevision with every renderer snapshot', () => { + const echo = startEcho(state, io) + + state.blocks.push({ type: 'content', content: 'hi', status: 'pending', timestamp: Date.now() }) + markStreamChanged(state) + + echo.flush() + expect(publishDeepchatEvent).toHaveBeenCalledWith( + 'chat.stream.updated', + expect.objectContaining({ revision: 1 }) + ) + + echo.flush() + expect(publishDeepchatEvent).toHaveBeenLastCalledWith( + 'chat.stream.updated', + expect.objectContaining({ revision: 1 }) + ) + + state.blocks.push({ type: 'content', content: 'more', status: 'pending', timestamp: Date.now() }) + markStreamChanged(state) + echo.flush() + expect(publishDeepchatEvent).toHaveBeenLastCalledWith( + 'chat.stream.updated', + expect.objectContaining({ revision: 2 }) + ) + + echo.stop() + }) + + it('drives a multi-token stream end-to-end with monotonic revisions per dirty bump', () => { + const echo = startEcho(state, io) + + accumulate(state, { type: 'text', content: 'Hello ' }) + accumulate(state, { type: 'text', content: 'world' }) + echo.schedule() + vi.advanceTimersByTime(130) + let flushes = getStreamUpdatedCalls() + expect(flushes).toHaveLength(1) + expect(flushes[0]?.[1]).toMatchObject({ revision: 2 }) + + vi.advanceTimersByTime(150) + accumulate(state, { type: 'text', content: '!' }) + echo.schedule() + vi.advanceTimersByTime(130) + flushes = getStreamUpdatedCalls() + expect(flushes).toHaveLength(2) + expect(flushes[1]?.[1]).toMatchObject({ revision: 3 }) + + echo.flush() + echo.flush() + flushes = getStreamUpdatedCalls() + expect(flushes).toHaveLength(4) + expect(flushes[2]?.[1]).toMatchObject({ revision: 3 }) + expect(flushes[3]?.[1]).toMatchObject({ revision: 3 }) + + expect(state.blocksRevision).toBe(3) + + echo.stop() + }) + it('rescheduleRenderer() resets the renderer flush window from the latest interaction', () => { const echo = startEcho(state, io) diff --git a/test/main/agent/deepchat/runtime/process.test.ts b/test/main/agent/deepchat/runtime/process.test.ts index fcc4884529..320dacc1de 100644 --- a/test/main/agent/deepchat/runtime/process.test.ts +++ b/test/main/agent/deepchat/runtime/process.test.ts @@ -644,6 +644,31 @@ describe('processStream', () => { ]) }) + it('advances the stream revision when an attempt change closes the previous narrative', async () => { + let identity = { logicalRound: 1, requestSeq: 2, physicalAttempt: 1 } + const coreStream = vi.fn(async function* () { + yield { type: 'text', content: 'Partial first attempt' } as LLMCoreStreamEvent + identity = { logicalRound: 1, requestSeq: 2, physicalAttempt: 2 } + // No further content: the narrative close is the only mutation on this event. + yield { type: 'stop', stop_reason: 'complete' } as LLMCoreStreamEvent + }) as unknown as ProcessParams['coreStream'] + const params = createParams({ + coreStream, + providerAttemptIdentity: () => identity + }) + + await expect(processStream(params)).resolves.toMatchObject({ status: 'completed' }) + + const state = params.run.streamState + expect(state.blocks).toEqual([ + expect.objectContaining({ type: 'content', status: 'success' }) + ]) + // text accumulate bumps once, the narrative close must bump again, and the + // terminal finalize adds its own mark — a close that does not advance the + // revision leaves a deferred flush deduped by the renderer. + expect(state.blocksRevision).toBe(3) + }) + it('persists normalized provider search results with the assistant message', async () => { const providerReplayJson = createDeepSeekReplayJson() const resultRow = { diff --git a/test/main/events/typedEventHub.test.ts b/test/main/events/typedEventHub.test.ts index aef32a461f..3a9d2cba33 100644 --- a/test/main/events/typedEventHub.test.ts +++ b/test/main/events/typedEventHub.test.ts @@ -187,6 +187,7 @@ describe('TypedEventHub', () => { sessionId: 'run-1', messageId: 'message-1', updatedAt: 1, + revision: 0, blocks: [] }, { kind: 'run', runId: 'run-1' } @@ -201,6 +202,7 @@ describe('TypedEventHub', () => { sessionId: 'run-1', messageId: 'message-1', updatedAt: 2, + revision: 1, blocks: [] }, { kind: 'run', runId: 'run-1' } @@ -270,6 +272,7 @@ describe('SessionEventRouter', () => { sessionId: 'cli-run', messageId: 'message-1', updatedAt: 123, + revision: 0, blocks: [] }) @@ -282,6 +285,7 @@ describe('SessionEventRouter', () => { sessionId: 'cli-run', messageId: 'message-1', updatedAt: 123, + revision: 0, blocks: [] } }) @@ -305,6 +309,7 @@ describe('SessionEventRouter', () => { sessionId: 'session-1', messageId: 'message-1', updatedAt: 123, + revision: 0, blocks: [] } router.publish('chat.stream.updated', payload) @@ -336,6 +341,7 @@ describe('SessionEventRouter', () => { sessionId: 'session-1', messageId: 'message-1', updatedAt: 123, + revision: 0, blocks: [] }) diff --git a/test/main/routes/contracts.test.ts b/test/main/routes/contracts.test.ts index 7a47372492..d08ce96f34 100644 --- a/test/main/routes/contracts.test.ts +++ b/test/main/routes/contracts.test.ts @@ -2277,6 +2277,7 @@ describe('main kernel contracts', () => { providerId: 'acp', modelId: 'dimcode', updatedAt: Date.now(), + revision: 0, blocks: [ { type: 'content', diff --git a/test/renderer/stores/messageStore.test.ts b/test/renderer/stores/messageStore.test.ts index f378e21fa4..85626508a4 100644 --- a/test/renderer/stores/messageStore.test.ts +++ b/test/renderer/stores/messageStore.test.ts @@ -1601,4 +1601,170 @@ describe('messageStore', () => { const updatedBlocks = store.getAssistantMessageBlocks(store.messages.value[0]!) expect(updatedBlocks[0]).not.toBe(firstBlocks[0]) }) + + it('skips duplicate snapshots by revision and applies bumped revisions', async () => { + const { store, streamListeners } = await setupStore() + await store.loadMessages('s1') + + const emit = (revision: number, text: string, updatedAt: number) => + streamListeners.updated[0]({ + sessionId: 's1', + requestId: 'm1', + messageId: 'm1', + providerId: 'acp', + modelId: 'dimcode', + updatedAt, + revision, + blocks: [{ type: 'content', content: text, status: 'pending', timestamp: updatedAt }] + }) + + emit(1, 'hello', 1) + const firstRecord = store.messageCache.value.get('m1')! + expect(firstRecord).toBeDefined() + expect(firstRecord.content).toContain('hello') + + emit(1, 'hello-again', 1) + const afterDuplicate = store.messageCache.value.get('m1')! + expect(afterDuplicate.content).toContain('hello') + expect(afterDuplicate.content).not.toContain('hello-again') + expect(afterDuplicate.updatedAt).toBe(firstRecord.updatedAt) + + emit(2, 'hello-again', 2) + const afterAdvance = store.messageCache.value.get('m1')! + expect(afterAdvance.content).toContain('hello-again') + }) + + it('applies the first revision-zero snapshot to an existing pending record', async () => { + const { store, sessionClient, streamListeners } = await setupStore() + sessionClient.restore.mockResolvedValueOnce({ + session: { id: 's1' }, + nextCursor: null, + hasMore: false, + messages: [ + { + id: 'm1', + sessionId: 's1', + orderSeq: 1, + role: 'assistant' as const, + content: JSON.stringify([ + { type: 'content', content: 'stale-initial', status: 'pending', timestamp: 1 } + ]), + status: 'pending' as const, + isContextEdge: 0, + metadata: '{}', + traceCount: 0, + createdAt: 1, + updatedAt: 1 + } + ] + }) + await store.loadMessages('s1') + expect(store.messageCache.value.get('m1')?.content).toContain('stale-initial') + + const emit = (revision: number, text: string, updatedAt: number) => + streamListeners.updated[0]({ + sessionId: 's1', + requestId: 'm1', + messageId: 'm1', + providerId: 'acp', + modelId: 'dimcode', + updatedAt, + revision, + blocks: [{ type: 'content', content: text, status: 'pending', timestamp: updatedAt }] + }) + + // No recorded revision yet: a first revision-0 snapshot must update the record + // instead of being treated as already applied. + emit(0, 'fresh-snapshot', 2) + const updated = store.messageCache.value.get('m1')! + expect(updated.content).toContain('fresh-snapshot') + + // Once revision 0 is recorded, duplicates are still deduped. + emit(0, 'duplicate-snapshot', 3) + const afterDuplicate = store.messageCache.value.get('m1')! + expect(afterDuplicate.content).toContain('fresh-snapshot') + expect(afterDuplicate.updatedAt).toBe(updated.updatedAt) + }) + + it('applies a revision-zero snapshot when a new request reuses the message id', async () => { + const { store, streamListeners } = await setupStore() + await store.loadMessages('s1') + + const emit = (requestId: string, revision: number, text: string, updatedAt: number) => + streamListeners.updated[0]({ + sessionId: 's1', + requestId, + messageId: 'm1', + providerId: 'acp', + modelId: 'dimcode', + updatedAt, + revision, + blocks: [{ type: 'content', content: text, status: 'pending', timestamp: updatedAt }] + }) + + emit('req-a', 1, 'request-a', 1) + emit('req-a', 2, 'request-a-more', 2) + expect(store.messageCache.value.get('m1')?.content).toContain('request-a-more') + + // A new request (or a resume) reuses the same message id and restarts from + // revision 0: it must not be deduped against request A's revision. + emit('req-b', 0, 'request-b', 3) + const reused = store.messageCache.value.get('m1')! + expect(reused.content).toContain('request-b') + + // The new request's revisions are tracked independently afterwards. + emit('req-b', 0, 'request-b-duplicate', 4) + const afterDuplicate = store.messageCache.value.get('m1')! + expect(afterDuplicate.content).toContain('request-b') + expect(afterDuplicate.updatedAt).toBe(reused.updatedAt) + + emit('req-b', 1, 'request-b-more', 5) + expect(store.messageCache.value.get('m1')?.content).toContain('request-b-more') + }) + + it('drops the applied revision on persisted record arrival so a recycled stream re-folds', async () => { + const { store, streamListeners, messageListeners } = await setupStore() + await store.loadMessages('s1') + + const emit = (revision: number, text: string, updatedAt: number) => + streamListeners.updated[0]({ + sessionId: 's1', + requestId: 'm1', + messageId: 'm1', + providerId: 'acp', + modelId: 'dimcode', + updatedAt, + revision, + blocks: [{ type: 'content', content: text, status: 'pending', timestamp: updatedAt }] + }) + + emit(1, 'first', 1) + emit(2, 'second', 2) + + messageListeners[0]({ + sessionId: 's1', + messages: [ + { + id: 'm1', + sessionId: 's1', + orderSeq: 1, + role: 'assistant' as const, + content: JSON.stringify([ + { type: 'content', content: 'persisted', status: 'success', timestamp: 5 } + ]), + status: 'success' as const, + isContextEdge: 0, + metadata: '{}', + traceCount: 0, + hasNestedExecutionAudit: false, + createdAt: 1, + updatedAt: Date.now() + 10_000 + } + ] + }) + expect(store.messageCache.value.get('m1')?.content).toContain('persisted') + + emit(1, 'recycled-stream', 6) + expect(store.messageCache.value.get('m1')?.content).toContain('recycled-stream') + }) })