diff --git a/ui/src/api/wizard.ts b/ui/src/api/wizard.ts index 9056b212..4096e794 100644 --- a/ui/src/api/wizard.ts +++ b/ui/src/api/wizard.ts @@ -10,6 +10,52 @@ export interface WizardConversationPayload { confirmations?: unknown[] } +export interface WizardConversationErrorDetail { + code?: string + message?: string + expectedRevision?: number + currentRevision?: number + [key: string]: unknown +} + +/** + * Error from the durable Wizard conversation endpoint. + * + * Keeping the HTTP status and structured detail at this boundary prevents + * callers from treating validation/auth/server errors as revision conflicts. + */ +export class WizardConversationRequestError extends Error { + readonly status: number + readonly detail: unknown + readonly code?: string + + constructor(message: string, status: number, detail: unknown = null) { + super(message) + this.name = 'WizardConversationRequestError' + this.status = status + this.detail = detail + const structured = detail && typeof detail === 'object' && 'detail' in detail + ? (detail as { detail?: unknown }).detail + : detail + this.code = structured && typeof structured === 'object' && typeof (structured as WizardConversationErrorDetail).code === 'string' + ? (structured as WizardConversationErrorDetail).code + : undefined + } +} + +async function readWizardError(response: Response, fallback: string): Promise { + const body = await response.json().catch(() => null) + const detail = body && typeof body === 'object' && 'detail' in body + ? (body as { detail?: unknown }).detail + : body + const message = typeof detail === 'string' + ? detail + : detail && typeof detail === 'object' && typeof (detail as WizardConversationErrorDetail).message === 'string' + ? (detail as WizardConversationErrorDetail).message as string + : fallback + return new WizardConversationRequestError(message, response.status, body) +} + export interface WizardWorkflowCollectionPayload { version: 1 revision: number @@ -21,8 +67,7 @@ export async function fetchWizardConversation(workspace: string): Promise ({ detail: 'Could not load Wizard conversation' })) - throw new Error(error.detail || 'Could not load Wizard conversation') + throw await readWizardError(response, 'Could not load Wizard conversation') } return response.json() } @@ -37,12 +82,7 @@ export async function saveWizardConversation( body: JSON.stringify({ workspace, baseRevision: conversation.revision, conversation }), }) if (!response.ok) { - const error: unknown = await response.json().catch(() => null) - if (error && typeof error === 'object') { - const detail = (error as Record).detail - if (typeof detail === 'string') throw new Error(detail) - } - throw new Error('Could not save Wizard conversation') + throw await readWizardError(response, 'Could not save Wizard conversation') } return response.json() } diff --git a/ui/src/features/agent/AgentAssistantPanel.tsx b/ui/src/features/agent/AgentAssistantPanel.tsx index 66ccbf8a..519d8448 100644 --- a/ui/src/features/agent/AgentAssistantPanel.tsx +++ b/ui/src/features/agent/AgentAssistantPanel.tsx @@ -1,7 +1,7 @@ import { useEffect, useMemo, useRef, useState, type FormEvent, type KeyboardEvent as ReactKeyboardEvent } from 'react' import { createPortal } from 'react-dom' import { ArrowUp, Loader2, Maximize2, Minimize2, PanelLeftClose, Sparkles, Trash2 } from 'lucide-react' -import { fetchWizardConversation, generateLlmText, saveWizardConversation, subscribeCanonicalTaskEvents, type CanonicalTask } from '../../api/client' +import { fetchWizardConversation, generateLlmText, subscribeCanonicalTaskEvents, type CanonicalTask, type WizardConversationPayload } from '../../api/client' import { AgentAvatar, type AgentVisualState } from './AgentAvatar' import { buildAgentTurnPrompt, HOCUSPOCUS_AGENT_SYSTEM_PROMPT, type AgentConversationEntry } from './agentKnowledge' import { @@ -15,11 +15,18 @@ import { type AgentActionResult, } from './agentActions' import { applyPollToCard, cardsFromResults, tabForExecutionTarget, type WizardExecutionCard } from './executionCards' -import { applyRemoteWizardConversation, WIZARD_WELCOME_TEXT } from './wizardConversationSync' +import { + applyRemoteWizardConversation, + isWizardConversationWriteCurrent, + normalizeRemoteWizardMessages, + shouldFollowWizardWorkspace, + WIZARD_WELCOME_TEXT, +} from './wizardConversationSync' import { AgentMarkdown } from './AgentMarkdown' import { defaultWizardWorkflowRuntime, type WizardWorkflowPendingInput, type WizardWorkflowRecord } from './wizardWorkflowRuntime' import { ensureRhythmic3dWorkflowRegistered } from './rhythmic3dWorkflow' import { defaultApplicationAdapters } from './applicationAdapters' +import { enqueueWizardConversationSave, persistQueuedWizardConversation, rebaseStaleWizardConversationHydration, rebaseWizardConversationAfterSave, resolveWizardConversationHydration } from './wizardConversationPersistence' import i18n, { useUiTranslation } from '../../i18n' export { AgentAvatar, type AgentVisualState } from './AgentAvatar' @@ -138,6 +145,13 @@ function writeMessages(workspace: string, messages: AgentMessage[]): void { } } +function taskExecutionState(status: CanonicalTask['status']): 'completed' | 'failed' | 'queued' | 'running' { + if (status === 'completed') return 'completed' + if (status === 'failed' || status === 'cancelled') return 'failed' + if (status === 'queued' || status === 'waiting_resource') return 'queued' + return 'running' +} + export function AgentAssistantPanel({ workspace, tasks, onClose, embedded = false }: AgentAssistantPanelProps) { const { t } = useUiTranslation('wizard') const { t: tCommon } = useUiTranslation('common') @@ -150,12 +164,15 @@ export function AgentAssistantPanel({ workspace, tasks, onClose, embedded = fals const [busyMessage, setBusyMessage] = useState('') const [expanded, setExpanded] = useState(false) const [errorCardId, setErrorCardId] = useState(null) + const [conversationSaveError, setConversationSaveError] = useState(null) const [activeWorkflow, setActiveWorkflow] = useState(null) const [pendingInput, setPendingInput] = useState(null) const endRef = useRef(null) const mountedRef = useRef(true) - const conversationRevisionRef = useRef(0) + const conversationSnapshotsRef = useRef>(new Map()) + const conversationSaveChainRef = useRef>(Promise.resolve()) const skipNextConversationSaveRef = useRef(false) + const conversationClearBasesRef = useRef>(new Map()) const conversationWorkspaceRef = useRef(conversationWorkspace) conversationWorkspaceRef.current = conversationWorkspace const messagesRef = useRef(messages) @@ -173,7 +190,7 @@ export function AgentAssistantPanel({ workspace, tasks, onClose, embedded = fals let active = true ensureRhythmic3dWorkflowRegistered(defaultApplicationAdapters) const unsubscribe = defaultWizardWorkflowRuntime.subscribe(({ workflow, card }) => { - if (!active || workflow.workspace !== workspace) return + if (!active || !isWizardConversationWriteCurrent(conversationWorkspace, workflow.workspace)) return setActiveWorkflow(workflow) setPendingInput(workflow.state === 'awaiting_input' ? workflow.pendingInput : null) setMessages(current => { @@ -197,10 +214,10 @@ export function AgentAssistantPanel({ workspace, tasks, onClose, embedded = fals }].slice(-40) }) }) - void defaultWizardWorkflowRuntime.open(workspace).catch(() => { + void defaultWizardWorkflowRuntime.open(conversationWorkspace).catch(() => { // Existing immediate actions remain available if workflow storage is offline. }) - const closeEvents = subscribeCanonicalTaskEvents(workspace, event => { + const closeEvents = subscribeCanonicalTaskEvents(conversationWorkspace, event => { void defaultWizardWorkflowRuntime.handleTaskEvent(event).catch(() => { // The checkpoint stays recoverable; a reconnect replays the same event. }) @@ -210,12 +227,12 @@ export function AgentAssistantPanel({ workspace, tasks, onClose, embedded = fals unsubscribe() closeEvents() } - }, [workspace]) + }, [conversationWorkspace]) useEffect(() => { setActiveWorkflow(null) setPendingInput(null) - }, [workspace]) + }, [conversationWorkspace]) useEffect(() => { writeMessages(conversationWorkspace, messages) @@ -226,79 +243,120 @@ export function AgentAssistantPanel({ workspace, tasks, onClose, embedded = fals return } const cards = messages.flatMap(message => message.cards || []) - void saveWizardConversation(conversationWorkspace, { + const capturedConversation: WizardConversationPayload = { version: 1, - revision: conversationRevisionRef.current, + revision: conversationSnapshotsRef.current.get(conversationWorkspace)?.revision || 0, messages, executions: cards, - }).then(saved => { - conversationRevisionRef.current = saved.revision - }).catch(async () => { - // A second tab may have advanced the CAS revision. Re-read and merge by - // message id; the resulting state triggers one save against the current - // backend revision. Local storage remains the fallback if this fails. - try { - const current = await fetchWizardConversation(conversationWorkspace) - if (!mountedRef.current || conversationWorkspaceRef.current !== conversationWorkspace) return - const choice = applyRemoteWizardConversation({ - localMessages: messagesRef.current, - localRevision: conversationRevisionRef.current, - remoteMessages: current.messages, - remoteRevision: current.revision || 0, - remoteExecutions: current.executions, - }) - conversationRevisionRef.current = choice.revision - skipNextConversationSaveRef.current = choice.source === 'remote' - setMessages([...choice.messages] as AgentMessage[]) - } catch { - // Local storage still holds the turn while the backend is unavailable. - } - }) + } + const queuedClearBase = conversationClearBasesRef.current.get(conversationWorkspace) + const queuedWrite = { + workspace: conversationWorkspace, + captured: capturedConversation, + // A clear is a three-way delete relative to the exact conversation the + // user saw, even if an earlier queued save advances the canonical state. + base: queuedClearBase ?? conversationSnapshotsRef.current.get(conversationWorkspace), + honorLocalDeletes: Boolean(queuedClearBase), + } + conversationSaveChainRef.current = enqueueWizardConversationSave( + conversationSaveChainRef.current, + async () => { + try { + const saved = await persistQueuedWizardConversation( + queuedWrite, + conversationSnapshotsRef.current, + ) + if (queuedClearBase && conversationClearBasesRef.current.get(conversationWorkspace) === queuedClearBase) { + conversationClearBasesRef.current.delete(conversationWorkspace) + } + if (!mountedRef.current || !isWizardConversationWriteCurrent(conversationWorkspaceRef.current, conversationWorkspace)) return + setConversationSaveError(null) + const visibleMessages = messagesRef.current + const pendingClearBase = conversationClearBasesRef.current.get(conversationWorkspace) + ?? (queuedWrite.honorLocalDeletes ? queuedWrite.base : undefined) + const rebased = rebaseWizardConversationAfterSave({ + ...queuedWrite.captured, + revision: saved.conversation.revision, + messages: visibleMessages, + executions: visibleMessages.flatMap(message => message.cards || []), + }, queuedWrite.captured, saved.conversation, pendingClearBase) + if (saved.merged || rebased.needsPersist) { + skipNextConversationSaveRef.current = !rebased.needsPersist + setMessages(normalizeRemoteWizardMessages( + rebased.conversation.messages, + rebased.conversation.executions, + ) as AgentMessage[]) + } + } catch (error) { + if (!mountedRef.current || !isWizardConversationWriteCurrent(conversationWorkspaceRef.current, conversationWorkspace)) return + setConversationSaveError(error instanceof Error ? error.message : String(error)) + } + }, + ) }, [conversationWorkspace, hydratedWorkspace, messages]) useEffect(() => { + if (workspace !== conversationWorkspace) return let cancelled = false - void fetchWizardConversation(workspace).then(payload => { - if (cancelled) return + void fetchWizardConversation(conversationWorkspace).then(payload => { + if (cancelled || !isWizardConversationWriteCurrent(conversationWorkspaceRef.current, conversationWorkspace)) return + const knownSnapshot = conversationSnapshotsRef.current.get(conversationWorkspace) + const hydration = resolveWizardConversationHydration(knownSnapshot, payload) + conversationSnapshotsRef.current.set(conversationWorkspace, hydration.snapshot) + const clearBase = conversationClearBasesRef.current.get(conversationWorkspace) + if (!hydration.applyToVisibleState || clearBase || knownSnapshot) { + const visibleMessages = messagesRef.current + const rebased = rebaseStaleWizardConversationHydration({ + ...payload, + messages: visibleMessages, + executions: visibleMessages.flatMap(message => message.cards || []), + }, clearBase ?? knownSnapshot ?? payload, hydration.snapshot, { + honorLocalDeletes: Boolean(clearBase), + }) + skipNextConversationSaveRef.current = !rebased.needsPersist + setMessages(normalizeRemoteWizardMessages( + rebased.conversation.messages, + rebased.conversation.executions, + ) as AgentMessage[]) + setHydratedWorkspace(conversationWorkspace) + return + } + const canonicalSnapshot = hydration.snapshot const choice = applyRemoteWizardConversation({ localMessages: messagesRef.current, - localRevision: conversationRevisionRef.current, - remoteMessages: payload.messages, - remoteRevision: payload.revision || 0, - remoteExecutions: payload.executions, + // A known snapshot is handled by the three-way branch above. + localRevision: 0, + remoteMessages: canonicalSnapshot.messages, + remoteRevision: canonicalSnapshot.revision || 0, + remoteExecutions: canonicalSnapshot.executions, }) - conversationRevisionRef.current = choice.revision skipNextConversationSaveRef.current = choice.source === 'remote' // A local choice may still adopt the backend's newer CAS revision and // merge remote-only messages. Use a fresh array so the persistence // effect retries the canonical save with that revision. setMessages([...choice.messages] as AgentMessage[]) - setHydratedWorkspace(workspace) + setHydratedWorkspace(conversationWorkspace) }).catch(() => { // Fall back to the local cache already loaded for this workspace. - if (!cancelled) setHydratedWorkspace(workspace) + if (!cancelled && isWizardConversationWriteCurrent(conversationWorkspaceRef.current, conversationWorkspace)) { + setHydratedWorkspace(conversationWorkspace) + } }) return () => { cancelled = true } - }, [workspace]) + }, [conversationWorkspace, workspace]) useEffect(() => { - if (workspace === conversationWorkspace) return - conversationRevisionRef.current = 0 + if (!shouldFollowWizardWorkspace({ activeWorkspace: workspace, conversationWorkspace, busy })) return skipNextConversationSaveRef.current = false + setConversationSaveError(null) setHydratedWorkspace(null) - if (busy) { - // A Wizard action changed workspace while this turn was executing. - // Keep the visible turn alive and persist it in the destination so its - // real action result is not lost when the footer updates. - setConversationWorkspace(workspace) - return - } setMessages(readMessages(workspace)) setConversationWorkspace(workspace) setState('idle') }, [busy, conversationWorkspace, workspace]) useEffect(() => { + if (workspace !== conversationWorkspace) return setMessages(current => current.map(message => { if (!message.cards?.length) return message let changed = false @@ -308,10 +366,7 @@ export function AgentAssistantPanel({ workspace, tasks, onClose, embedded = fals || (card.pipelineId && item.pipeline_id === card.pipelineId) )) if (!task) return card - const state = task.status === 'completed' ? 'completed' - : task.status === 'failed' || task.status === 'cancelled' ? 'failed' - : task.status === 'queued' || task.status === 'waiting_resource' ? 'queued' - : 'running' + const state = taskExecutionState(task.status) const outputNames = task.result_refs?.length ? task.result_refs : card.outputNames if (state === card.state && (task.message || card.message) === card.message && outputNames === card.outputNames) { return card @@ -327,7 +382,7 @@ export function AgentAssistantPanel({ workspace, tasks, onClose, embedded = fals }) return changed ? { ...message, cards } : message })) - }, [tasks]) + }, [conversationWorkspace, tasks, workspace]) useEffect(() => { const closeOnEscape = (event: KeyboardEvent) => { @@ -345,6 +400,17 @@ export function AgentAssistantPanel({ workspace, tasks, onClose, embedded = fals const clearConversation = () => { const next = [welcomeMessage()] + const snapshot = conversationSnapshotsRef.current.get(conversationWorkspace) + conversationClearBasesRef.current.set(conversationWorkspace, { + ...snapshot, + version: 1, + revision: snapshot?.revision || 0, + messages: messagesRef.current, + executions: messagesRef.current.flatMap(message => message.cards || []), + }) + // An explicit user mutation must never consume a skip reserved for an + // earlier canonical hydration render. + skipNextConversationSaveRef.current = false setMessages(next) setState('idle') } @@ -552,6 +618,11 @@ export function AgentAssistantPanel({ workspace, tasks, onClose, embedded = fals ))} + {conversationSaveError && ( +

+ {t('conversationSaveError', { message: conversationSaveError })} +

+ )} {busy && (
diff --git a/ui/src/features/agent/wizardConversationPersistence.ts b/ui/src/features/agent/wizardConversationPersistence.ts new file mode 100644 index 00000000..c3ed96c4 --- /dev/null +++ b/ui/src/features/agent/wizardConversationPersistence.ts @@ -0,0 +1,316 @@ +import { + fetchWizardConversation, + saveWizardConversation, + WizardConversationRequestError, + type WizardConversationPayload, +} from '../../api/wizard' +import { + mergeWizardMessages, + normalizeRemoteWizardMessages, +} from './wizardConversationSync' + +export interface WizardConversationTransport { + fetch: (workspace: string) => Promise + save: (workspace: string, conversation: WizardConversationPayload) => Promise +} + +export interface WizardConversationSaveResult { + conversation: WizardConversationPayload + merged: boolean +} + +export type WizardConversationSnapshotStore = Map + +export interface QueuedWizardConversationWrite { + workspace: string + captured: WizardConversationPayload + base?: WizardConversationPayload + /** Only explicit user mutations such as Clear may delete values absent locally. */ + honorLocalDeletes?: boolean +} + +export function resolveWizardConversationHydration( + known: WizardConversationPayload | undefined, + incoming: WizardConversationPayload, +): { snapshot: WizardConversationPayload; applyToVisibleState: boolean } { + if (known && known.revision >= incoming.revision) { + return { snapshot: known, applyToVisibleState: false } + } + return { snapshot: incoming, applyToVisibleState: true } +} + +/** + * Rebase browser-visible edits over a hydration response that lost a race to + * a confirmed save. The stale response is the three-way base, so confirmed + * turns added after it are retained while local edits still win. Missing + * browser-cache values are not treated as deletes unless the caller records + * an explicit user clear. + */ +export function rebaseStaleWizardConversationHydration( + visible: WizardConversationPayload, + stale: WizardConversationPayload, + confirmed: WizardConversationPayload, + options: { honorLocalDeletes?: boolean } = {}, +): { conversation: WizardConversationPayload; needsPersist: boolean } { + const honorLocalDeletes = options.honorLocalDeletes ?? false + const conversation: WizardConversationPayload = { + version: 1, + revision: confirmed.revision, + messages: mergeQueuedValues(visible.messages, stale.messages, confirmed.messages, honorLocalDeletes), + executions: mergeQueuedValues(visible.executions, stale.executions, confirmed.executions, honorLocalDeletes), + requestedActions: mergeQueuedValues( + visible.requestedActions, + stale.requestedActions, + confirmed.requestedActions, + honorLocalDeletes, + ), + executedActions: mergeQueuedValues( + visible.executedActions, + stale.executedActions, + confirmed.executedActions, + honorLocalDeletes, + ), + confirmations: mergeQueuedValues( + visible.confirmations, + stale.confirmations, + confirmed.confirmations, + honorLocalDeletes, + ), + } + const semanticContent = (value: WizardConversationPayload) => ({ + messages: value.messages, + executions: value.executions, + requestedActions: value.requestedActions ?? [], + executedActions: value.executedActions ?? [], + confirmations: value.confirmations ?? [], + }) + return { + conversation, + needsPersist: stableKey(semanticContent(conversation)) !== stableKey(semanticContent(confirmed)), + } +} + +/** + * Reconcile a save result with UI changes made while that save was in flight. + * A pending Clear supplies its own visible ancestor and is the only case where + * absent values represent intentional deletion. + */ +export function rebaseWizardConversationAfterSave( + visible: WizardConversationPayload, + captured: WizardConversationPayload, + confirmed: WizardConversationPayload, + pendingClearBase?: WizardConversationPayload, +): { conversation: WizardConversationPayload; needsPersist: boolean } { + return rebaseStaleWizardConversationHydration( + visible, + pendingClearBase ?? captured, + confirmed, + { honorLocalDeletes: Boolean(pendingClearBase) }, + ) +} + +const defaultTransport: WizardConversationTransport = { + fetch: fetchWizardConversation, + save: saveWizardConversation, +} + +function isEmptyStableValue(value: unknown): boolean { + return value == null || value === '' || (Array.isArray(value) && value.length === 0) +} + +function stableKey(value: unknown): string { + if (value === null || typeof value !== 'object') return JSON.stringify(value) ?? String(value) + if (Array.isArray(value)) return `[${value.map(stableKey).join(',')}]` + const record = value as Record + return `{${Object.keys(record) + .sort() + .filter(key => !isEmptyStableValue(record[key])) + .map(key => `${JSON.stringify(key)}:${stableKey(record[key])}`) + .join(',')}}` +} + +function mergeUniqueValues(remote: unknown, local: unknown): unknown[] { + const merged: unknown[] = [] + const seen = new Set() + for (const value of [ + ...(Array.isArray(remote) ? remote : []), + ...(Array.isArray(local) ? local : []), + ]) { + const key = stableKey(value) + if (seen.has(key)) continue + seen.add(key) + merged.push(value) + } + return merged.slice(-80) +} + +function valueIdentity(value: unknown): string { + if (value && typeof value === 'object') { + const record = value as Record + if (typeof record.id === 'string' && record.id) return `id:${record.id}` + if (typeof record.executionKey === 'string' && record.executionKey) return `execution:${record.executionKey}` + } + return `value:${stableKey(value)}` +} + +/** Apply local edits/deletes since base without discarding concurrent values. */ +function mergeQueuedValues( + local: unknown, + base: unknown, + canonical: unknown, + honorLocalDeletes = true, +): unknown[] { + const localValues = Array.isArray(local) ? local : [] + const baseValues = Array.isArray(base) ? base : [] + const canonicalValues = Array.isArray(canonical) ? canonical : [] + const localById = new Map(localValues.map(value => [valueIdentity(value), value])) + const baseById = new Map(baseValues.map(value => [valueIdentity(value), value])) + const merged: unknown[] = [] + const seen = new Set() + + canonicalValues.forEach(value => { + const id = valueIdentity(value) + const baseValue = baseById.get(id) + const localValue = localById.get(id) + if (honorLocalDeletes && baseById.has(id) && !localById.has(id)) return + if (localById.has(id) && (!baseById.has(id) || stableKey(localValue) !== stableKey(baseValue))) { + merged.push(localValue) + } else { + merged.push(value) + } + seen.add(id) + }) + localValues.forEach(value => { + const id = valueIdentity(value) + if (seen.has(id)) return + merged.push(value) + seen.add(id) + }) + return merged.slice(-80) +} + +export function mergeQueuedWizardConversationSnapshots( + local: WizardConversationPayload, + base: WizardConversationPayload | undefined, + canonical: WizardConversationPayload, + options: { honorLocalDeletes?: boolean } = {}, +): WizardConversationPayload { + const honorLocalDeletes = options.honorLocalDeletes ?? true + return { + version: 1, + revision: canonical.revision, + messages: mergeQueuedValues(local.messages, base?.messages, canonical.messages, honorLocalDeletes), + executions: mergeQueuedValues(local.executions, base?.executions, canonical.executions, honorLocalDeletes), + requestedActions: local.requestedActions === undefined + ? canonical.requestedActions + : mergeQueuedValues(local.requestedActions, base?.requestedActions, canonical.requestedActions, honorLocalDeletes), + executedActions: local.executedActions === undefined + ? canonical.executedActions + : mergeQueuedValues(local.executedActions, base?.executedActions, canonical.executedActions, honorLocalDeletes), + confirmations: local.confirmations === undefined + ? canonical.confirmations + : mergeQueuedValues(local.confirmations, base?.confirmations, canonical.confirmations, honorLocalDeletes), + } +} + +/** + * Build the one payload used after a CAS conflict. + * + * The server snapshot is canonical. Existing ids remain in server order and + * local-only turn ids are appended once, so repeating the merge is harmless. + */ +export function mergeWizardConversationSnapshots( + local: WizardConversationPayload, + remote: WizardConversationPayload, +): WizardConversationPayload { + const remoteMessages = normalizeRemoteWizardMessages(remote.messages, remote.executions) + const localMessages = normalizeRemoteWizardMessages(local.messages, local.executions) + return { + version: 1, + revision: remote.revision, + messages: mergeWizardMessages(localMessages, remoteMessages), + executions: mergeUniqueValues(remote.executions, local.executions), + requestedActions: mergeUniqueValues(remote.requestedActions, local.requestedActions), + executedActions: mergeUniqueValues(remote.executedActions, local.executedActions), + confirmations: mergeUniqueValues(remote.confirmations, local.confirmations), + } +} + +export function isWizardConversationConflict(error: unknown): boolean { + if (error instanceof WizardConversationRequestError) return error.status === 409 + return Boolean( + error + && typeof error === 'object' + && (error as { status?: unknown }).status === 409, + ) +} + +/** + * Save a Wizard conversation, recovering one and only one CAS conflict. + * + * A validation, auth or server error is propagated immediately. If the + * retry conflicts again, that second error is propagated to the caller rather + * than starting an unbounded refetch/save loop. + */ +export async function saveWizardConversationWithRecovery( + workspace: string, + conversation: WizardConversationPayload, + transport: WizardConversationTransport = defaultTransport, + base?: WizardConversationPayload, + options: { honorLocalDeletes?: boolean } = {}, +): Promise { + try { + return { + conversation: await transport.save(workspace, conversation), + merged: false, + } + } catch (error) { + if (!isWizardConversationConflict(error)) throw error + const remote = await transport.fetch(workspace) + const merged = base + ? mergeQueuedWizardConversationSnapshots(conversation, base, remote, options) + : mergeWizardConversationSnapshots(conversation, remote) + return { + conversation: await transport.save(workspace, merged), + merged: true, + } + } +} + +/** + * Persist one queued snapshot against the latest canonical state known for its + * workspace. Queued React effects intentionally capture their visible turn, + * but must not capture the revision/canonical payload: an earlier queued save + * may advance both before this write starts. + * + * The store is keyed by workspace so changing the visible workspace never + * rebinds or drops a write that was already accepted into the queue. + */ +export async function persistQueuedWizardConversation( + write: QueuedWizardConversationWrite, + snapshots: WizardConversationSnapshotStore, + transport: WizardConversationTransport = defaultTransport, +): Promise { + const canonical = snapshots.get(write.workspace) + const outgoing = canonical + ? mergeQueuedWizardConversationSnapshots(write.captured, write.base, canonical, { + honorLocalDeletes: write.honorLocalDeletes ?? false, + }) + : write.captured + const saved = await saveWizardConversationWithRecovery(write.workspace, outgoing, transport, write.base, { + honorLocalDeletes: write.honorLocalDeletes ?? false, + }) + snapshots.set(write.workspace, saved.conversation) + return saved +} + +/** + * Serialize browser conversation writes so every save reads the revision + * confirmed by its predecessor. A rejected write does not poison the queue. + */ +export function enqueueWizardConversationSave( + previous: Promise, + write: () => Promise, +): Promise { + return previous.catch(() => undefined).then(write) +} diff --git a/ui/src/features/agent/wizardConversationSync.ts b/ui/src/features/agent/wizardConversationSync.ts index 376f9b59..0ffe9d80 100644 --- a/ui/src/features/agent/wizardConversationSync.ts +++ b/ui/src/features/agent/wizardConversationSync.ts @@ -7,6 +7,10 @@ export interface WizardSyncMessage { createdAt: number language?: string cards?: unknown[] + executionKey?: string + jobLinks?: unknown[] + lastState?: string + error?: string } export interface WizardConversationChoice { @@ -36,6 +40,10 @@ export function normalizeRemoteWizardMessages( createdAt: typeof message.createdAt === 'number' ? message.createdAt : 0, ...(typeof message.language === 'string' && message.language ? { language: message.language } : {}), cards: Array.isArray(message.cards) && message.cards.length ? message.cards : undefined, + ...(typeof message.executionKey === 'string' && message.executionKey ? { executionKey: message.executionKey } : {}), + jobLinks: Array.isArray(message.jobLinks) && message.jobLinks.length ? message.jobLinks : undefined, + ...(typeof message.lastState === 'string' && message.lastState ? { lastState: message.lastState } : {}), + ...(typeof message.error === 'string' && message.error ? { error: message.error } : {}), }] }) if (!restored.length && Array.isArray(remoteExecutions) && remoteExecutions.length) { @@ -50,6 +58,43 @@ export function normalizeRemoteWizardMessages( return restored.slice(-40) } +/** + * Merge two snapshots without duplicating a message id. + * + * Remote order is canonical, but a local value wins for a shared id. Callers + * use this fallback only when no common ancestor is available, so preserving + * an in-browser card/workflow update is safer than silently reverting it. + * Local-only messages are appended in their existing order. + */ +export function mergeWizardMessages( + localMessages: WizardSyncMessage[], + remoteMessages: WizardSyncMessage[], +): WizardSyncMessage[] { + const localById = new Map(localMessages.map(message => [message.id, message])) + const merged: WizardSyncMessage[] = [] + const seen = new Set() + for (const message of remoteMessages) { + if (!message.id || seen.has(message.id)) continue + seen.add(message.id) + merged.push(localById.get(message.id) ?? message) + } + for (const message of localMessages) { + if (!message.id || seen.has(message.id)) continue + seen.add(message.id) + merged.push(message) + } + return merged.slice(-40) +} + +/** True when the visible client state still contains a turn absent from a saved snapshot. */ +export function hasExclusiveWizardMessages( + visibleMessages: WizardSyncMessage[], + savedMessages: WizardSyncMessage[], +): boolean { + const savedIds = new Set(savedMessages.map(message => message.id)) + return visibleMessages.some(message => Boolean(message.id) && !savedIds.has(message.id)) +} + export function isTransientWizardChat(messages: WizardSyncMessage[]): boolean { if (!messages.length) return true return !messages.some(message => ( @@ -78,15 +123,28 @@ export function applyRemoteWizardConversation(input: { const localHasExclusiveTurn = localMessages.some(message => !remoteIds.has(message.id)) && !isTransientWizardChat(localMessages) if (localHasExclusiveTurn) { - const localById = new Map(localMessages.map(message => [message.id, message])) - const merged = remoteMessages.map(message => localById.get(message.id) || message) - for (const message of localMessages) { - if (!remoteIds.has(message.id)) merged.push(message) + return { + source: 'local', + messages: mergeWizardMessages(localMessages, remoteMessages), + revision: Math.max(localRevision, remoteRevision), } - return { source: 'local', messages: merged.slice(-40), revision: Math.max(localRevision, remoteRevision) } } if (localRevision > remoteRevision) { return { source: 'local', messages: localMessages, revision: localRevision } } return { source: 'remote', messages: remoteMessages, revision: remoteRevision } } + +/** Follow the footer workspace only after the in-flight turn finishes. */ +export function shouldFollowWizardWorkspace(input: { + activeWorkspace: string + conversationWorkspace: string + busy: boolean +}): boolean { + return input.activeWorkspace !== input.conversationWorkspace && !input.busy +} + +/** Drop async conversation writes that finished after the owner changed. */ +export function isWizardConversationWriteCurrent(owner: string, current: string): boolean { + return Boolean(owner) && owner === current +} diff --git a/ui/src/i18n/locales/en/wizard.json b/ui/src/i18n/locales/en/wizard.json index 6d5d22a0..8e882ac6 100644 --- a/ui/src/i18n/locales/en/wizard.json +++ b/ui/src/i18n/locales/en/wizard.json @@ -21,6 +21,7 @@ "emptyReply": "My crystal ball did not return a usable answer.", "pendingError": "I cannot apply that answer to the blocked step yet: {{message}}", "llmError": "I could not query the LLM: {{message}}. Check Settings → Services and try again.", + "conversationSaveError": "I could not save this Wizard conversation: {{message}}. Your local copy is still available; retry after reloading.", "pendingChoice": "Choose one of these options to continue the spell: {{choices}}.", "pendingFields": "Answer the fields {{fields}}.", "needDecisionBody": "🪄 **I need a decision to continue the same spell.**\n\n{{reason}}", diff --git a/ui/src/i18n/locales/es/wizard.json b/ui/src/i18n/locales/es/wizard.json index cee3c0a9..c10ab838 100644 --- a/ui/src/i18n/locales/es/wizard.json +++ b/ui/src/i18n/locales/es/wizard.json @@ -21,6 +21,7 @@ "emptyReply": "Mi bola de cristal no ha devuelto una respuesta utilizable.", "pendingError": "No puedo aplicar todavía esa respuesta al paso bloqueado: {{message}}", "llmError": "No he podido consultar el LLM: {{message}}. Comprueba Ajustes → Servicios y vuelve a intentarlo.", + "conversationSaveError": "No he podido guardar esta conversación del mago: {{message}}. Tu copia local sigue disponible; vuelve a cargar para reintentarlo.", "pendingChoice": "Elige una de estas opciones para continuar el hechizo: {{choices}}.", "pendingFields": "Responde los campos {{fields}}.", "needDecisionBody": "🪄 **Necesito una decisión para continuar el mismo hechizo.**\n\n{{reason}}", diff --git a/ui/tests/wizardConversationPersistence.test.mjs b/ui/tests/wizardConversationPersistence.test.mjs new file mode 100644 index 00000000..3c46196b --- /dev/null +++ b/ui/tests/wizardConversationPersistence.test.mjs @@ -0,0 +1,608 @@ +import assert from 'node:assert/strict' +import test from 'node:test' + +const { + saveWizardConversationWithRecovery, + mergeWizardConversationSnapshots, + mergeQueuedWizardConversationSnapshots, + enqueueWizardConversationSave, + persistQueuedWizardConversation, + rebaseStaleWizardConversationHydration, + rebaseWizardConversationAfterSave, + resolveWizardConversationHydration, +} = await import('../src/features/agent/wizardConversationPersistence.ts') +const { + WizardConversationRequestError, + fetchWizardConversation, +} = await import('../src/api/wizard.ts') +const { + isWizardConversationWriteCurrent, + hasExclusiveWizardMessages, + mergeWizardMessages, + shouldFollowWizardWorkspace, +} = await import('../src/features/agent/wizardConversationSync.ts') + +function clone(value) { + return JSON.parse(JSON.stringify(value)) +} + +function payload(revision, ids, prefix = '') { + return { + version: 1, + revision, + messages: ids.map((id, index) => ({ + id, + role: index % 2 ? 'assistant' : 'user', + text: `${prefix}${id}`, + createdAt: index + 1, + executionKey: `${prefix}${id}-execution`, + jobLinks: [{ taskId: `${prefix}${id}-task`, pipelineId: '' }], + })), + executions: [], + } +} + +function revisionConflict(expected, current) { + return new WizardConversationRequestError( + `expected ${expected}, current ${current}`, + 409, + { + code: 'wizard_conversation_revision_conflict', + expectedRevision: expected, + currentRevision: current, + }, + ) +} + +test('two writers preserve both Wizard turns exactly once after one CAS conflict', async () => { + let canonical = payload(0, []) + const calls = [] + const transport = { + async fetch(workspace) { + calls.push({ method: 'fetch', workspace }) + return clone(canonical) + }, + async save(workspace, conversation) { + calls.push({ method: 'save', workspace, conversation: clone(conversation) }) + if (conversation.revision !== canonical.revision) { + throw revisionConflict(conversation.revision, canonical.revision) + } + canonical = { ...clone(conversation), revision: canonical.revision + 1 } + return clone(canonical) + }, + } + + const first = await saveWizardConversationWithRecovery('workspace-a', payload(0, ['a-user', 'a-assistant'], 'A-'), transport) + const second = await saveWizardConversationWithRecovery('workspace-a', payload(0, ['b-user', 'b-assistant'], 'B-'), transport) + + assert.equal(first.merged, false) + assert.equal(second.merged, true) + assert.equal(canonical.revision, 2) + assert.deepEqual(canonical.messages.map(message => message.id), [ + 'a-user', 'a-assistant', 'b-user', 'b-assistant', + ]) + assert.equal(new Set(canonical.messages.map(message => message.id)).size, 4) + assert.deepEqual(canonical.messages.map(message => message.executionKey), [ + 'A-a-user-execution', 'A-a-assistant-execution', + 'B-b-user-execution', 'B-b-assistant-execution', + ]) + assert.deepEqual(calls.map(call => `${call.method}:${call.workspace}`), [ + 'save:workspace-a', 'save:workspace-a', 'fetch:workspace-a', 'save:workspace-a', + ]) + + const repeatedMerge = mergeWizardConversationSnapshots(second.conversation, canonical) + assert.deepEqual(repeatedMerge.messages.map(message => message.id), canonical.messages.map(message => message.id)) +}) + +test('second conflict is surfaced after one recovery retry and never loops', async () => { + let saves = 0 + let fetches = 0 + const transport = { + async fetch() { + fetches += 1 + return payload(4, ['remote-user', 'remote-assistant'], 'remote-') + }, + async save(_workspace, conversation) { + saves += 1 + throw revisionConflict(conversation.revision, 4) + }, + } + + await assert.rejects( + saveWizardConversationWithRecovery('workspace-a', payload(0, ['local-user', 'local-assistant'], 'local-'), transport), + error => error instanceof WizardConversationRequestError && error.status === 409, + ) + assert.equal(saves, 2) + assert.equal(fetches, 1) +}) + +test('non-recoverable 4xx is surfaced without a refetch or retry', async () => { + let saves = 0 + let fetches = 0 + const transport = { + async fetch() { + fetches += 1 + return payload(0, []) + }, + async save() { + saves += 1 + throw new WizardConversationRequestError( + 'conversation payload is invalid', + 400, + { detail: 'conversation payload is invalid' }, + ) + }, + } + + await assert.rejects( + saveWizardConversationWithRecovery('workspace-a', payload(0, ['local-user']), transport), + error => error instanceof WizardConversationRequestError && error.status === 400, + ) + assert.equal(saves, 1) + assert.equal(fetches, 0) +}) + +test('conversation HTTP errors retain nested API detail and status', async () => { + const originalFetch = globalThis.fetch + globalThis.fetch = async () => new Response(JSON.stringify({ + detail: { + code: 'wizard_conversation_revision_conflict', + message: 'expected 2, current 3', + expectedRevision: 2, + currentRevision: 3, + }, + }), { + status: 409, + headers: { 'content-type': 'application/json' }, + }) + try { + await assert.rejects( + fetchWizardConversation('workspace-a'), + error => error instanceof WizardConversationRequestError + && error.status === 409 + && error.code === 'wizard_conversation_revision_conflict' + && error.message === 'expected 2, current 3', + ) + } finally { + globalThis.fetch = originalFetch + } +}) + +test('workspace changes never rebind an in-flight conversation write', () => { + assert.equal(shouldFollowWizardWorkspace({ + activeWorkspace: 'workspace-b', + conversationWorkspace: 'workspace-a', + busy: true, + }), false) + assert.equal(shouldFollowWizardWorkspace({ + activeWorkspace: 'workspace-b', + conversationWorkspace: 'workspace-a', + busy: false, + }), true) + assert.equal(isWizardConversationWriteCurrent('workspace-a', 'workspace-b'), false) + assert.equal(isWizardConversationWriteCurrent('workspace-a', 'workspace-a'), true) +}) + +test('a conflict response cannot suppress persistence of a newer visible turn', () => { + const saved = payload(3, ['old-user', 'old-assistant']).messages + const withNewTurn = [ + ...saved, + { id: 'new-user', role: 'user', text: 'new request', createdAt: 3 }, + ] + + assert.equal(hasExclusiveWizardMessages(withNewTurn, saved), true) + assert.equal(hasExclusiveWizardMessages(saved, saved), false) +}) + +test('conversation writes are serialized and a later write sees the confirmed revision', async () => { + let revision = 0 + let active = 0 + let maximumActive = 0 + const observed = [] + let releaseFirst + let markFirstStarted + const firstGate = new Promise(resolve => { releaseFirst = resolve }) + const firstStarted = new Promise(resolve => { markFirstStarted = resolve }) + + let chain = Promise.resolve() + chain = enqueueWizardConversationSave(chain, async () => { + active += 1 + maximumActive = Math.max(maximumActive, active) + observed.push(revision) + markFirstStarted() + await firstGate + revision = 1 + active -= 1 + }) + chain = enqueueWizardConversationSave(chain, async () => { + active += 1 + maximumActive = Math.max(maximumActive, active) + observed.push(revision) + revision = 2 + active -= 1 + }) + + await firstStarted + assert.deepEqual(observed, [0]) + releaseFirst() + await chain + assert.deepEqual(observed, [0, 1]) + assert.equal(maximumActive, 1) + assert.equal(revision, 2) +}) + +test('a stale queued snapshot merges the canonical turn saved by its predecessor', async () => { + let canonical = payload(0, []) + const snapshots = new Map() + const savedPayloads = [] + const transport = { + async fetch() { return clone(canonical) }, + async save(_workspace, conversation) { + assert.equal(conversation.revision, canonical.revision) + savedPayloads.push(clone(conversation)) + canonical = { ...clone(conversation), revision: canonical.revision + 1 } + return clone(canonical) + }, + } + + const firstVisible = payload(0, ['first-user', 'first-assistant']) + const staleSecondEffect = payload(0, ['second-user', 'second-assistant']) + let chain = Promise.resolve() + chain = enqueueWizardConversationSave(chain, () => ( + persistQueuedWizardConversation({ workspace: 'workspace-a', captured: firstVisible }, snapshots, transport).then(() => undefined) + )) + chain = enqueueWizardConversationSave(chain, () => ( + persistQueuedWizardConversation({ workspace: 'workspace-a', captured: staleSecondEffect }, snapshots, transport).then(() => undefined) + )) + await chain + + assert.deepEqual(savedPayloads[1].messages.map(message => message.id), [ + 'first-user', 'first-assistant', 'second-user', 'second-assistant', + ]) + assert.equal(savedPayloads[1].revision, 1) + assert.deepEqual(snapshots.get('workspace-a'), canonical) +}) + +test('a queued clear uses its recorded ancestor after a predecessor advances the snapshot', async () => { + const clearBase = payload(1, ['cleared-user', 'cleared-assistant']) + const canonical = payload(2, ['cleared-user', 'cleared-assistant', 'concurrent-user']) + const capturedClear = payload(1, ['welcome-after-clear']) + const snapshots = new Map([['workspace-a', clone(canonical)]]) + const transport = { + async fetch() { return clone(canonical) }, + async save(_workspace, conversation) { + assert.equal(conversation.revision, 2) + assert.deepEqual(conversation.messages.map(message => message.id), [ + 'concurrent-user', + 'welcome-after-clear', + ]) + return { ...clone(conversation), revision: 3 } + }, + } + + const saved = await persistQueuedWizardConversation({ + workspace: 'workspace-a', + captured: capturedClear, + base: clearBase, + honorLocalDeletes: true, + }, snapshots, transport) + + assert.equal(saved.conversation.revision, 3) + assert.deepEqual(snapshots.get('workspace-a'), saved.conversation) +}) + +test('a queued write persists to its captured workspace after the visible workspace changes', async () => { + const snapshots = new Map() + const savedWorkspaces = [] + let visibleWorkspace = 'workspace-a' + let releaseWrite + const gate = new Promise(resolve => { releaseWrite = resolve }) + const transport = { + async fetch() { return payload(0, []) }, + async save(workspace, conversation) { + await gate + savedWorkspaces.push(workspace) + return { ...clone(conversation), revision: 1 } + }, + } + + let chain = Promise.resolve() + chain = enqueueWizardConversationSave(chain, () => ( + persistQueuedWizardConversation({ workspace: 'workspace-a', captured: payload(0, ['a-user']) }, snapshots, transport).then(() => undefined) + )) + visibleWorkspace = 'workspace-b' + releaseWrite() + await chain + + assert.equal(visibleWorkspace, 'workspace-b') + assert.deepEqual(savedWorkspaces, ['workspace-a']) + assert.equal(snapshots.get('workspace-a').revision, 1) +}) + +test('three-way queued merge applies local edits while retaining concurrent turns', () => { + const base = payload(3, ['shared-user']) + const local = clone(base) + local.messages[0].text = 'locally edited' + const canonical = payload(4, ['shared-user', 'remote-assistant']) + canonical.messages[0].text = 'old canonical text' + + const merged = mergeQueuedWizardConversationSnapshots(local, base, canonical) + + assert.equal(merged.revision, 4) + assert.equal(merged.messages.find(message => message.id === 'shared-user').text, 'locally edited') + assert.deepEqual(merged.messages.map(message => message.id), ['shared-user', 'remote-assistant']) +}) + +test('three-way queued merge keeps concurrent additions but honors a local clear', () => { + const base = payload(2, ['old-user', 'old-assistant']) + const local = payload(2, []) + const canonical = payload(3, ['old-user', 'old-assistant', 'concurrent-user']) + + const merged = mergeQueuedWizardConversationSnapshots(local, base, canonical) + + assert.deepEqual(merged.messages.map(message => message.id), ['concurrent-user']) +}) + +test('three-way queued merge persists an updated execution card by stable id', () => { + const base = payload(1, ['user']) + base.executions = [{ id: 'card-1', state: 'running' }] + const local = clone(base) + local.executions = [{ id: 'card-1', state: 'completed' }] + const canonical = clone(base) + canonical.revision = 2 + + const merged = mergeQueuedWizardConversationSnapshots(local, base, canonical) + + assert.deepEqual(merged.executions, [{ id: 'card-1', state: 'completed' }]) +}) + +test('a CAS conflict during a queued edit keeps the edit and the concurrent turn', async () => { + const base = payload(4, ['shared-user']) + const snapshots = new Map([['workspace-a', clone(base)]]) + const local = clone(base) + local.messages[0].text = 'edited after hydration' + let canonical = payload(5, ['shared-user', 'remote-assistant']) + canonical.messages[0].text = base.messages[0].text + let saves = 0 + const transport = { + async fetch() { return clone(canonical) }, + async save(_workspace, conversation) { + saves += 1 + if (saves === 1) throw revisionConflict(conversation.revision, canonical.revision) + assert.equal(conversation.revision, 5) + canonical = { ...clone(conversation), revision: 6 } + return clone(canonical) + }, + } + + const saved = await persistQueuedWizardConversation({ + workspace: 'workspace-a', + captured: local, + base, + }, snapshots, transport) + + assert.equal(saved.merged, true) + assert.equal(saved.conversation.messages.find(message => message.id === 'shared-user').text, 'edited after hydration') + assert.deepEqual(saved.conversation.messages.map(message => message.id), ['shared-user', 'remote-assistant']) +}) + +test('a normal queued save cannot delete canonical turns outside the 40-message UI window', async () => { + const ids = Array.from({ length: 50 }, (_value, index) => `message-${index + 1}`) + const canonical = payload(7, ids) + const visibleWindow = payload(7, [...ids.slice(-40), 'new-local-message']) + const snapshots = new Map([['workspace-a', clone(canonical)]]) + const transport = { + async fetch() { return clone(canonical) }, + async save(_workspace, conversation) { + assert.equal(conversation.revision, 7) + assert.deepEqual(conversation.messages.map(message => message.id), [...ids, 'new-local-message']) + return { ...clone(conversation), revision: 8 } + }, + } + + await persistQueuedWizardConversation({ + workspace: 'workspace-a', + captured: visibleWindow, + base: canonical, + }, snapshots, transport) +}) + +test('a conflict retry keeps using the recorded clear ancestor', async () => { + const clearBase = payload(1, ['cleared-user', 'cleared-assistant']) + const capturedClear = payload(1, ['welcome-after-clear']) + const snapshotAfterEarlierWrite = payload(2, ['concurrent-before-conflict']) + const remoteAfterConflict = payload(3, [ + 'cleared-user', + 'cleared-assistant', + 'concurrent-before-conflict', + 'concurrent-during-conflict', + ]) + const snapshots = new Map([['workspace-a', clone(snapshotAfterEarlierWrite)]]) + let saves = 0 + const transport = { + async fetch() { return clone(remoteAfterConflict) }, + async save(_workspace, conversation) { + saves += 1 + if (saves === 1) throw revisionConflict(conversation.revision, remoteAfterConflict.revision) + assert.equal(conversation.revision, 3) + assert.deepEqual(conversation.messages.map(message => message.id), [ + 'concurrent-before-conflict', + 'concurrent-during-conflict', + 'welcome-after-clear', + ]) + return { ...clone(conversation), revision: 4 } + }, + } + + const saved = await persistQueuedWizardConversation({ + workspace: 'workspace-a', + captured: capturedClear, + base: clearBase, + honorLocalDeletes: true, + }, snapshots, transport) + + assert.equal(saved.merged, true) + assert.equal(saved.conversation.revision, 4) +}) + +test('a shared message id retains the newer local workflow card', () => { + const remote = [{ + id: 'assistant-workflow', role: 'assistant', text: 'Generating', createdAt: 1, + cards: [{ id: 'card-1', state: 'running' }], + }] + const local = [{ + ...remote[0], + text: 'Completed', + cards: [{ id: 'card-1', state: 'completed' }], + }] + + const merged = mergeWizardMessages(local, remote) + + assert.equal(merged[0].text, 'Completed') + assert.equal(merged[0].cards[0].state, 'completed') +}) + +test('a Clear made while an earlier save is in flight is not undone by that save', () => { + const earlierCaptured = payload(5, ['old-user', 'old-assistant']) + const confirmedEarlierSave = payload(6, [ + 'old-user', + 'old-assistant', + 'concurrent-before-clear', + ]) + const pendingClearBase = payload(5, ['old-user', 'old-assistant']) + const visibleAfterClear = payload(5, ['welcome-after-clear']) + + const rebased = rebaseWizardConversationAfterSave( + visibleAfterClear, + earlierCaptured, + confirmedEarlierSave, + pendingClearBase, + ) + + assert.equal(rebased.needsPersist, true) + assert.deepEqual(rebased.conversation.messages.map(message => message.id), [ + 'concurrent-before-clear', + 'welcome-after-clear', + ]) +}) + +test('a late hydration fetch cannot replace a newer confirmed snapshot or visible edits', () => { + const confirmed = payload(8, ['confirmed-user', 'confirmed-assistant']) + const staleFetch = payload(7, ['stale-user']) + + assert.deepEqual(resolveWizardConversationHydration(confirmed, staleFetch), { + snapshot: confirmed, + applyToVisibleState: false, + }) + assert.deepEqual(resolveWizardConversationHydration(confirmed, clone(confirmed)), { + snapshot: confirmed, + applyToVisibleState: false, + }) + assert.deepEqual(resolveWizardConversationHydration(undefined, staleFetch), { + snapshot: staleFetch, + applyToVisibleState: true, + }) + assert.deepEqual(resolveWizardConversationHydration(staleFetch, confirmed), { + snapshot: confirmed, + applyToVisibleState: true, + }) +}) + +test('stale hydration rebases visible edits without dropping confirmed turns', () => { + const stale = payload(7, ['shared-user']) + const confirmed = payload(8, ['shared-user', 'confirmed-assistant']) + const visible = clone(stale) + visible.messages[0].text = 'edited while hydration was in flight' + + const rebased = rebaseStaleWizardConversationHydration(visible, stale, confirmed) + + assert.equal(rebased.needsPersist, true) + assert.equal(rebased.conversation.revision, 8) + assert.equal(rebased.conversation.messages[0].text, 'edited while hydration was in flight') + assert.deepEqual(rebased.conversation.messages.map(message => message.id), [ + 'shared-user', + 'confirmed-assistant', + ]) +}) + +test('stale hydration restores confirmed-only turns without scheduling a redundant save', () => { + const stale = payload(7, ['shared-user']) + const confirmed = payload(8, ['shared-user', 'confirmed-assistant']) + + const rebased = rebaseStaleWizardConversationHydration(clone(stale), stale, confirmed) + + assert.equal(rebased.needsPersist, false) + assert.deepEqual(rebased.conversation.messages, confirmed.messages) +}) + +test('UI and server message shape differences are not treated as local edits', () => { + const serverMessage = (id, role, text, createdAt, extra = {}) => ({ + id, + role, + text, + createdAt, + cards: [], + executionKey: '', + jobLinks: [], + lastState: '', + error: '', + ...extra, + }) + const stale = { + version: 1, + revision: 7, + messages: [serverMessage('shared-user', 'user', 'hello', 1)], + executions: [], + } + const confirmed = { + version: 1, + revision: 8, + messages: [ + serverMessage('shared-user', 'user', 'hello', 1, { executionKey: 'server-key' }), + serverMessage('confirmed-assistant', 'assistant', 'reply', 2), + ], + executions: [], + } + const visible = { + version: 1, + revision: 7, + messages: [{ id: 'shared-user', role: 'user', text: 'hello', createdAt: 1 }], + executions: [], + } + + const rebased = rebaseStaleWizardConversationHydration(visible, stale, confirmed) + + assert.equal(rebased.needsPersist, false) + assert.equal(rebased.conversation.messages[0].executionKey, 'server-key') + assert.deepEqual(rebased.conversation.messages.map(message => message.id), [ + 'shared-user', + 'confirmed-assistant', + ]) +}) + +test('stale browser cache omissions do not delete newer canonical turns during hydration', () => { + const stale = payload(7, ['shared-user', 'stale-assistant']) + const confirmed = payload(8, ['shared-user', 'stale-assistant', 'confirmed-user']) + const olderVisibleCache = payload(5, ['shared-user']) + + const rebased = rebaseStaleWizardConversationHydration(olderVisibleCache, stale, confirmed) + + assert.equal(rebased.needsPersist, false) + assert.deepEqual(rebased.conversation.messages, confirmed.messages) +}) + +test('explicit clear deletes base turns while retaining concurrent canonical additions', () => { + const stale = payload(7, ['old-user', 'old-assistant']) + const confirmed = payload(8, ['old-user', 'old-assistant', 'concurrent-user']) + const cleared = payload(7, ['welcome-after-clear']) + + const rebased = rebaseStaleWizardConversationHydration(cleared, stale, confirmed, { + honorLocalDeletes: true, + }) + + assert.equal(rebased.needsPersist, true) + assert.deepEqual(rebased.conversation.messages.map(message => message.id), [ + 'concurrent-user', + 'welcome-after-clear', + ]) +})