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.
This commit is contained in:
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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([])
|
||||
|
||||
|
||||
@@ -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<SessionResumeResult['inflight']
|
||||
const text = inflight.assistant ?? ''
|
||||
const corrections = inflight.corrections ?? []
|
||||
const offsets = inflight.correction_offsets
|
||||
|
||||
if (!corrections.length) {
|
||||
return [text]
|
||||
}
|
||||
|
||||
// Python's len() counts Unicode code points, unlike JS string offsets.
|
||||
const characters = Array.from(text)
|
||||
|
||||
@@ -192,8 +198,7 @@ export function reconcilePersistedLiveTurn(
|
||||
messages: ChatMessage[],
|
||||
previous: ChatMessage[],
|
||||
rows: SessionMessage[],
|
||||
projection: Pick<SessionResumeResult, 'inflight' | 'queued' | 'session_id'>,
|
||||
reconcileHistory: (messages: ChatMessage[], previous: ChatMessage[]) => ChatMessage[]
|
||||
projection: Pick<SessionResumeResult, 'inflight' | 'queued' | 'session_id'>
|
||||
): 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
|
||||
)
|
||||
|
||||
@@ -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<string, ChatMessage>()
|
||||
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 &&
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -8,6 +8,7 @@ export {
|
||||
completeOpenTimelineParts,
|
||||
dedupeRepeatedTextInParts,
|
||||
mergeFinalAssistantText,
|
||||
normalizeWs,
|
||||
reasoningPart,
|
||||
renderMediaTags,
|
||||
textPart
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user