From 5aff954e155eb40ddbe3ee14cc22ff2a193ea09a Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Tue, 22 Sep 2026 12:10:42 +0530 Subject: [PATCH] refactor(desktop): dedupe the live-turn reconciler seams - `reconcilePersistedLiveTurn` called `reconcileAuthoritativeChatMessages` back through a callback that was the caller itself; the projection-less branch is now `reconcileDurableHistory` in utils and imported directly. - The activation path retried the live-turn reconcile inside its fallback with the same rows; `null` does not depend on `previous`, so the retry was dead work. The fallback now runs the durable path only. - `isLiveTailReplyId` (spoken-reply) replaces the third and fourth copies of the `assistant-stream-`/`inflight-assistant-` prefix predicate; `normalizeWs` is exported from chat-messages instead of re-declared in coverage.ts and the journal; `isCommittedRow` names the "not a projection, not a recovery" predicate the journal repeated three times. - Loop-invariant `lastPreviousUser` hoisted; `renderedText(raw)` computed once per row instead of per part; `snapshotIntervals` skips the code-point split when there are no corrections; `persistInFlightTurnState` checks for a recovered row before building the recoverable tail on every idle commit. --- .../hooks/use-session-actions/index.ts | 35 ++++++------------- .../live-turn-remainder.test.ts | 2 +- .../persisted-live-turn.ts | 17 +++++---- .../hooks/use-session-actions/utils.ts | 32 +++++++++++------ .../desktop/src/lib/chat-messages/coverage.ts | 3 +- apps/desktop/src/lib/chat-messages/index.ts | 1 + apps/desktop/src/lib/inflight-turn-journal.ts | 33 +++++++++-------- 7 files changed, 64 insertions(+), 59 deletions(-) diff --git a/apps/desktop/src/app/session/hooks/use-session-actions/index.ts b/apps/desktop/src/app/session/hooks/use-session-actions/index.ts index fc195e3a45..2f68b5b666 100644 --- a/apps/desktop/src/app/session/hooks/use-session-actions/index.ts +++ b/apps/desktop/src/app/session/hooks/use-session-actions/index.ts @@ -166,7 +166,7 @@ import { patchSessionWorkspace, preserveEquivalentTranscript, preserveLocalPendingTurnMessages, - reconcileResumeMessages, + reconcileDurableHistory, removeRepresentedLocalLiveProjection, resolveResumedBusy, resolveSessionProfile, @@ -266,27 +266,17 @@ function reconcileAuthoritativeChatMessages( sourceRows?: SessionMessage[] ): ChatMessage[] { if (liveProjection && sourceRows) { - const reconciled = reconcilePersistedLiveTurn( - authoritativeMessages, - previousMessages, - sourceRows, - liveProjection, - (messages, previous) => reconcileAuthoritativeChatMessages(messages, previous) - ) + const reconciled = reconcilePersistedLiveTurn(authoritativeMessages, previousMessages, sourceRows, liveProjection) if (reconciled) { return reconciled } } - const withLiveProjection = liveProjection - ? appendLiveSessionProjection(authoritativeMessages, liveProjection) - : authoritativeMessages - - const reconciled = reconcileResumeMessages(withLiveProjection, previousMessages) - const withPendingTurn = preserveLocalPendingTurnMessages(reconciled, previousMessages) - - return preserveLocalAssistantErrors(withPendingTurn, previousMessages) + return reconcileDurableHistory( + liveProjection ? appendLiveSessionProjection(authoritativeMessages, liveProjection) : authoritativeMessages, + previousMessages + ) } function reconcileAuthoritativeMessages( @@ -1528,19 +1518,14 @@ export function useSessionActions({ persistedMessages, sessionStateByRuntimeIdRef.current.get(cachedRuntimeId)?.messages ?? previousMessages, persisted.messages, - liveProjection, - (messages, previous) => reconcileAuthoritativeChatMessages(messages, previous) + liveProjection ) + // `null` does not depend on `previous`; retrying the live-turn + // reconcile inside the fallback would return `null` again. reconciledCurrentLiveTurn = currentLiveTurn !== null activatedMessages = - currentLiveTurn ?? - reconcileAuthoritativeChatMessages( - persistedMessages, - previousMessages, - liveProjection, - persisted.messages - ) + currentLiveTurn ?? reconcileAuthoritativeChatMessages(persistedMessages, previousMessages, liveProjection) } } diff --git a/apps/desktop/src/app/session/hooks/use-session-actions/live-turn-remainder.test.ts b/apps/desktop/src/app/session/hooks/use-session-actions/live-turn-remainder.test.ts index 50b5854c42..edb5fecd93 100644 --- a/apps/desktop/src/app/session/hooks/use-session-actions/live-turn-remainder.test.ts +++ b/apps/desktop/src/app/session/hooks/use-session-actions/live-turn-remainder.test.ts @@ -113,7 +113,7 @@ it('pairs only the queue projection, preserving equal corrections and different } const reconcile = (previous: ChatMessage[]) => - reconcilePersistedLiveTurn(toChatMessages(rows), previous, rows, projection, messages => messages)! + reconcilePersistedLiveTurn(toChatMessages(rows), previous, rows, projection)! let current = reconcile([]) diff --git a/apps/desktop/src/app/session/hooks/use-session-actions/persisted-live-turn.ts b/apps/desktop/src/app/session/hooks/use-session-actions/persisted-live-turn.ts index 2a149f1b8f..695df1b8a2 100644 --- a/apps/desktop/src/app/session/hooks/use-session-actions/persisted-live-turn.ts +++ b/apps/desktop/src/app/session/hooks/use-session-actions/persisted-live-turn.ts @@ -5,6 +5,7 @@ import { parseErrorSurface } from '@/lib/error-surface' import type { SessionMessage, SessionResumeResult } from '@/types/hermes' import { mergeLiveAssistantRun } from './live-turn-remainder' +import { reconcileDurableHistory } from './utils' const rowId = (row: SessionMessage) => row.row_id ?? row.id const userText = (text: string) => textWithoutReferenceLines(text).trim() @@ -68,12 +69,12 @@ function candidateTurn( // MEDIA rendering and folding change lengths. Only source-row provenance, // not equal prose elsewhere in the transcript, proves display coverage. + const rendered = renderedText(raw).trim() + if ( id === undefined || !turnMessages.some(message => - message.parts.some( - part => part.type === 'text' && part.sourceRowId === id && part.text.trim() === renderedText(raw).trim() - ) + message.parts.some(part => part.type === 'text' && part.sourceRowId === id && part.text.trim() === rendered) ) ) { return null @@ -162,6 +163,11 @@ function snapshotIntervals(inflight: NonNullable, - reconcileHistory: (messages: ChatMessage[], previous: ChatMessage[]) => ChatMessage[] + projection: Pick ): ChatMessage[] | null { const inflight = projection.inflight @@ -233,7 +238,7 @@ export function reconcilePersistedLiveTurn( const snapshots = snapshotIntervals(inflight) const corrections = inflight.corrections ?? [] - const result = reconcileHistory( + const result = reconcileDurableHistory( messages.slice(0, turn.start + 1), localStart >= 0 ? previous.slice(0, localStart + 1) : previous ) diff --git a/apps/desktop/src/app/session/hooks/use-session-actions/utils.ts b/apps/desktop/src/app/session/hooks/use-session-actions/utils.ts index 99bcaa2b75..caaf87c1ba 100644 --- a/apps/desktop/src/app/session/hooks/use-session-actions/utils.ts +++ b/apps/desktop/src/app/session/hooks/use-session-actions/utils.ts @@ -1,11 +1,19 @@ import { resolveSessionRpcOwner } from '@/app/contrib/wiring-routing' import { textWithoutReferenceLines } from '@/components/assistant-ui/reference-kinds' import { getSession } from '@/hermes' -import { assistantTextPart, type ChatMessage, chatMessageText, textPart, toChatMessages } from '@/lib/chat-messages' +import { + assistantTextPart, + type ChatMessage, + chatMessageText, + preserveLocalAssistantErrors, + textPart, + toChatMessages +} from '@/lib/chat-messages' import { normalizePersonalityValue } from '@/lib/chat-runtime' import { embeddedImageUrls, textWithoutEmbeddedImages } from '@/lib/embedded-images' import { parseErrorSurface } from '@/lib/error-surface' import { isMessagingSource, normalizeSessionSource } from '@/lib/session-source' +import { isLiveTailReplyId } from '@/lib/spoken-reply' import { reconcileApprovalModeForProfile } from '@/store/approval-mode' import { requestDesktopOnboardingForCredentialWarning } from '@/store/onboarding' import { $activeGatewayProfile, $profiles, normalizeProfileKey } from '@/store/profile' @@ -78,12 +86,7 @@ function hasStructuralParts(message: ChatMessage): boolean { * turn — as opposed to a committed transcript row. */ function isLiveTailRow(message: ChatMessage): boolean { - return ( - message.pending === true || - message.id.startsWith('assistant-stream-') || - message.id.startsWith('inflight-assistant-') || - message.interim === true - ) + return message.pending === true || isLiveTailReplyId(message.id) || message.interim === true } /** @@ -334,6 +337,15 @@ export function preserveEquivalentTranscript(current: ChatMessage[], next: ChatM return chatMessageArraysEquivalent(current, next) ? current : next } +/** Durable history against the local view: role-ordinal pairing, then the + * local pending turn and local assistant errors the DB cannot know about. */ +export function reconcileDurableHistory(messages: ChatMessage[], previous: ChatMessage[]): ChatMessage[] { + const reconciled = reconcileResumeMessages(messages, previous) + const withPendingTurn = preserveLocalPendingTurnMessages(reconciled, previous) + + return preserveLocalAssistantErrors(withPendingTurn, previous) +} + export function reconcileResumeMessages(nextMessages: ChatMessage[], previousMessages: ChatMessage[]): ChatMessage[] { if (!previousMessages.length) { return nextMessages @@ -580,8 +592,7 @@ function durableFoldCoversLiveResponse(messages: ChatMessage[], live: ChatMessag return messages.slice(lastUser + 1).some(message => { if ( message.role !== 'assistant' || - message.id.startsWith('assistant-stream-') || - message.id.startsWith('inflight-assistant-') + isLiveTailReplyId(message.id) ) { return false } @@ -655,6 +666,7 @@ export function preserveLocalPendingTurnMessages( // Authoritative id → richer local pending row. Replacing (not appending) // avoids painting both the empty inflight shell and the full stream bubble. const replacements = new Map() + const lastPreviousUser = previousMessages.findLastIndex(row => row.role === 'user' && !isGatewaySystemMarker(row)) for (const message of previousMessages) { if (isGatewaySystemMarker(message)) { @@ -794,8 +806,6 @@ export function preserveLocalPendingTurnMessages( } } - const lastPreviousUser = previousMessages.findLastIndex(row => row.role === 'user' && !isGatewaySystemMarker(row)) - if ( isPendingAssistant && previousMessages.indexOf(message) > lastPreviousUser && diff --git a/apps/desktop/src/lib/chat-messages/coverage.ts b/apps/desktop/src/lib/chat-messages/coverage.ts index fd777464f4..fa7efc0b2e 100644 --- a/apps/desktop/src/lib/chat-messages/coverage.ts +++ b/apps/desktop/src/lib/chat-messages/coverage.ts @@ -1,7 +1,6 @@ +import { normalizeWs as normalizedText } from './parts' import type { ChatMessage, ChatMessagePart } from './types' -const normalizedText = (value: string) => value.replace(/\s+/g, ' ').trim() - function sameOccurrencePart(stored: ChatMessagePart, local: ChatMessagePart): boolean { if (stored.type === 'tool-call' && local.type === 'tool-call') { return Boolean(stored.toolCallId) && stored.toolCallId === local.toolCallId diff --git a/apps/desktop/src/lib/chat-messages/index.ts b/apps/desktop/src/lib/chat-messages/index.ts index 7777224c46..e806f020b5 100644 --- a/apps/desktop/src/lib/chat-messages/index.ts +++ b/apps/desktop/src/lib/chat-messages/index.ts @@ -8,6 +8,7 @@ export { completeOpenTimelineParts, dedupeRepeatedTextInParts, mergeFinalAssistantText, + normalizeWs, reasoningPart, renderMediaTags, textPart diff --git a/apps/desktop/src/lib/inflight-turn-journal.ts b/apps/desktop/src/lib/inflight-turn-journal.ts index 96ca50e23c..55e51892bf 100644 --- a/apps/desktop/src/lib/inflight-turn-journal.ts +++ b/apps/desktop/src/lib/inflight-turn-journal.ts @@ -1,5 +1,11 @@ -import { type ChatMessage, type ChatMessagePart, chatMessageText } from '@/lib/chat-messages' +import { + type ChatMessage, + type ChatMessagePart, + chatMessageText, + normalizeWs as normalizedText +} from '@/lib/chat-messages' import { withoutCoveredAssistantPrefix } from '@/lib/chat-messages/coverage' +import { isLiveTailReplyId } from '@/lib/spoken-reply' /** * Crash-survivable in-flight turn journal. @@ -523,10 +529,6 @@ function cloneMessages(messages: ChatMessage[]): ChatMessage[] { } } -function normalizedText(value: string): string { - return value.replace(/\s+/g, ' ').trim() -} - function attachmentSignature(message: ChatMessage): string { return (message.attachmentRefs ?? []).join('\n') } @@ -556,11 +558,13 @@ function assistantHasRecoverableContent(message: ChatMessage): boolean { /** A live-turn projection row (backend `inflight` via appendLiveSessionProjection, * or a still-streaming local bubble) — as opposed to a completed transcript row. */ function isLiveProjectionRow(message: ChatMessage): boolean { - return ( - Boolean(message.pending) || - message.id.startsWith('assistant-stream-') || - message.id.startsWith('inflight-assistant-') - ) + return Boolean(message.pending) || isLiveTailReplyId(message.id) +} + +/** A row the transcript actually holds: neither a live projection nor a + * journal recovery that is still waiting for its durable counterpart. */ +function isCommittedRow(message: ChatMessage): boolean { + return !isLiveProjectionRow(message) && !message.recovered } /** Visible tail of the running turn: the streaming assistant row (plus any @@ -716,7 +720,7 @@ function journalTailAlreadyCommitted(tailAssistants: ChatMessage[], baseMessages const lastTurnStart = baseMessages.findLastIndex(message => message.role === 'user' && !message.hidden) return recoverable.every(message => baseMessages.some((base, index) => - base.role === 'assistant' && !base.hidden && !isLiveProjectionRow(base) && !base.recovered && + base.role === 'assistant' && !base.hidden && isCommittedRow(base) && (base.id === message.id || (message.rowId === undefined ? index > lastTurnStart : base.rowId !== undefined && base.rowId === message.rowId)) && base.error === message.error && @@ -806,7 +810,7 @@ export function mergeInFlightMessages( const afterUser = baseMessages.slice(matchingUserIndex + 1, end) const completedReply = afterUser.find( - message => assistantHasRecoverableContent(message) && !isLiveProjectionRow(message) && !message.recovered && + message => assistantHasRecoverableContent(message) && isCommittedRow(message) && (message.durableComplete === true || (message.durableComplete === undefined && !message.interim && !message.parts.some(part => part.type === 'tool-call'))) ) @@ -817,7 +821,7 @@ export function mergeInFlightMessages( } tailAssistants = withoutCoveredAssistantPrefix( - afterUser.filter(message => !isLiveProjectionRow(message) && !message.recovered), tailAssistants + afterUser.filter(isCommittedRow), tailAssistants ) lastJournalRow = tailAssistants.findLast(assistantHasRecoverableContent) ?? null @@ -961,7 +965,8 @@ export function persistInFlightTurnState(state: JournalableSessionState): void { } if (!state.busy && !state.awaitingResponse && !state.streamId && - !recoverableTail(state.messages, null).some(message => message.recovered)) { + !(state.messages.some(message => message.recovered) && + recoverableTail(state.messages, null).some(message => message.recovered))) { clearInFlightTurnJournal(storedSessionId) return