diff --git a/apps/web/src/app/api/ai/global/[id]/messages/route.ts b/apps/web/src/app/api/ai/global/[id]/messages/route.ts index 65c555b350..b5c1635b2a 100644 --- a/apps/web/src/app/api/ai/global/[id]/messages/route.ts +++ b/apps/web/src/app/api/ai/global/[id]/messages/route.ts @@ -1,5 +1,5 @@ import { NextResponse } from 'next/server'; -import { streamText, convertToModelMessages, stepCountIs, UIMessage } from 'ai'; +import { streamText, convertToModelMessages, stepCountIs, UIMessage, createUIMessageStream, createUIMessageStreamResponse, type LanguageModelUsage } from 'ai'; import { incrementUsage, getCurrentUsage, getUserUsageSummary } from '@/lib/subscription/usage-service'; import { createRateLimitResponse } from '@/lib/subscription/rate-limit-middleware'; import { broadcastUsageEvent } from '@/lib/websocket'; @@ -733,70 +733,92 @@ MENTION PROCESSING: // This is separate from request.signal which fires on any client disconnect const { streamId, signal: abortSignal } = createStreamAbortController({ userId }); - const result = streamText({ - model, - system: finalSystemPrompt, - messages: modelMessages, - tools: finalTools, - stopWhen: stepCountIs(100), - abortSignal, // From registry - only aborts on explicit user stop, not client disconnect - experimental_context: { - userId, - aiProvider: currentProvider, - aiModel: currentModel, - conversationId, - locationContext, - modelCapabilities: getModelCapabilities(currentModel, currentProvider) - }, - maxRetries: 20, // Increase from default 2 to 20 for better handling of rate limits - onAbort: () => { - loggers.api.info('Global Assistant Chat API: Stream aborted by user', { - userId: maskIdentifier(userId), - conversationId, - streamId, - model: currentModel, - provider: currentProvider, - }); - }, - }); - - loggers.api.debug('📡 Global Assistant Chat API: Returning stream response', {}); - // Generate server-side message ID for the AI response // This ensures client and server use the same ID, fixing the undo-after-streaming issue - // See: https://ai-sdk.dev/docs/ai-sdk-ui/chatbot-message-persistence const serverAssistantMessageId = createId(); - return result.toUIMessageStreamResponse({ - // Provide the server-generated ID to the stream response - // The client's useChat will use this ID instead of generating its own - generateMessageId: () => serverAssistantMessageId, - // Pass streamId via headers so client can call /api/ai/abort for explicit stop - headers: { - [STREAM_ID_HEADER]: streamId, + // Track usage promise for token counting + let usagePromise: Promise | undefined; + + // Use createUIMessageStream to wrap the streaming response + // This ensures server-side processing continues even if the client disconnects + const stream = createUIMessageStream({ + originalMessages: processedMessages, + execute: async ({ writer }) => { + // Send the server-generated message ID to the client at stream start + try { + writer.write({ + type: 'start', + messageId: serverAssistantMessageId, + }); + } catch { + // Client disconnected before first write - continue processing + } + + const aiResult = streamText({ + model, + system: finalSystemPrompt, + messages: modelMessages, + tools: finalTools, + stopWhen: stepCountIs(100), + abortSignal, // From registry - only aborts on explicit user stop, not client disconnect + experimental_context: { + userId, + aiProvider: currentProvider, + aiModel: currentModel, + conversationId, + locationContext, + modelCapabilities: getModelCapabilities(currentModel, currentProvider) + }, + maxRetries: 20, + onAbort: () => { + loggers.api.info('Global Assistant Chat API: Stream aborted by user', { + userId: maskIdentifier(userId), + conversationId, + streamId, + model: currentModel, + provider: currentProvider, + }); + }, + }); + + usagePromise = aiResult.totalUsage + .catch((error) => { + loggers.api.debug('Global Assistant: Failed to retrieve token usage from stream', { + error: error instanceof Error ? error.message : 'Unknown error', + }); + return undefined; + }); + + // Stream all chunks to client, continuing server-side even if client disconnects + for await (const chunk of aiResult.toUIMessageStream()) { + try { + writer.write(chunk); + } catch { + // Client disconnected - continue processing to ensure onFinish fires + } + } }, onFinish: async ({ responseMessage }) => { // Clean up abort controller from registry removeStream({ streamId }); - loggers.api.debug('🏁 Global Assistant Chat API: onFinish callback triggered for AI response', {}); + loggers.api.debug('Global Assistant Chat API: onFinish callback triggered for AI response', {}); if (responseMessage) { try { - // Use the server-generated ID that was sent to the client - // This ensures the saved message ID matches what the client has const messageId = serverAssistantMessageId; const messageContent = extractMessageContent(responseMessage); const extractedToolCalls = extractToolCalls(responseMessage); const extractedToolResults = extractToolResults(responseMessage); - - loggers.api.debug('💾 Global Assistant Chat API: Saving AI response message:', { - id: messageId, + + loggers.api.debug('Global Assistant Chat API: Saving AI response message', { + id: messageId, contentLength: messageContent.length, toolCallsCount: extractedToolCalls.length, toolResultsCount: extractedToolResults.length, }); - + await saveGlobalAssistantMessageToDatabase({ messageId, conversationId, @@ -805,7 +827,7 @@ MENTION PROCESSING: content: messageContent, toolCalls: extractedToolCalls.length > 0 ? extractedToolCalls : undefined, toolResults: extractedToolResults.length > 0 ? extractedToolResults : undefined, - uiMessage: responseMessage, // Pass complete UIMessage to preserve part ordering + uiMessage: responseMessage, }); // Update conversation lastMessageAt @@ -817,14 +839,12 @@ MENTION PROCESSING: }) .where(eq(conversations.id, conversationId)); - loggers.api.debug('✅ Global Assistant Chat API: AI response message saved to database', {}); + loggers.api.debug('Global Assistant Chat API: AI response message saved to database', {}); // Track detailed AI usage (tokens, cost, etc.) try { - // Only attempt to get usage if the stream completed successfully - const usage = await result.usage; + const usage = usagePromise ? await usagePromise : undefined; - // Only track if we actually have usage data if (usage && usage.totalTokens && usage.totalTokens > 0) { const duration = Date.now() - startTime; @@ -839,8 +859,6 @@ MENTION PROCESSING: conversationId, messageId, success: true, - - // Context tracking - actual conversation context vs billing tokens contextMessages: contextCalculation.messageIds, contextSize: contextCalculation.totalTokens, systemPromptTokens: contextCalculation.systemPromptTokens, @@ -849,72 +867,37 @@ MENTION PROCESSING: messageCount: contextCalculation.messageCount, wasTruncated: contextCalculation.wasTruncated, truncationStrategy: contextCalculation.truncationStrategy, - metadata: { toolCallsCount: extractedToolCalls.length, toolResultsCount: extractedToolResults.length, isReadOnly: readOnlyMode, } }); - - loggers.api.debug('✅ Global Assistant: AI usage tracked', { - conversationId: maskIdentifier(conversationId), - messageId: maskIdentifier(messageId), - tokens: usage.totalTokens, - }); - } else { - loggers.api.debug('â„šī¸ Global Assistant: No usage data available (stream may have been aborted)', { - conversationId: maskIdentifier(conversationId), - messageId: maskIdentifier(messageId), - }); } } catch (trackingError) { - // Log as debug, not error - this is expected when stream is aborted - loggers.api.debug('â„šī¸ Global Assistant: Could not track AI usage (stream aborted or failed)', { + loggers.api.debug('Global Assistant: Could not track AI usage (stream aborted or failed)', { conversationId: maskIdentifier(conversationId), messageId: maskIdentifier(messageId), error: trackingError instanceof Error ? trackingError.message : 'Unknown error', }); - // Don't fail the request if tracking fails } // Track usage for PageSpace providers only (rate limiting/quota tracking) const isPageSpaceProvider = currentProvider === 'pagespace'; - const maskedUserId = maskIdentifier(userId); - const maskedConversationId = maskIdentifier(conversationId); - const maskedMessageId = maskIdentifier(messageId); - - usageLogger.info('Global Assistant usage tracking decision', { - userId: maskedUserId, - provider: currentProvider, - isPageSpaceProvider, - messageId: maskedMessageId, - conversationId: maskedConversationId, - }); - if (isPageSpaceProvider) { try { - // Determine if this is pro model based on model name const isProModel = currentModel === 'glm-4.7'; const providerType = isProModel ? 'pro' : 'standard'; - usageLogger.debug('Incrementing usage for Global Assistant response', { - userId: maskedUserId, - provider: currentProvider, - providerType, - messageId: maskedMessageId, - conversationId: maskedConversationId, - }); - const usageResult = await incrementUsage(userId, providerType); usageLogger.info('Global Assistant usage incremented', { - userId: maskedUserId, + userId: maskIdentifier(userId), provider: currentProvider, providerType, - messageId: maskedMessageId, - conversationId: maskedConversationId, + messageId: maskIdentifier(messageId), + conversationId: maskIdentifier(conversationId), currentCount: usageResult.currentCount, limit: usageResult.limit, remaining: usageResult.remainingCalls, @@ -924,7 +907,6 @@ MENTION PROCESSING: // Broadcast usage event for real-time updates try { const currentUsageSummary = await getUserUsageSummary(userId); - await broadcastUsageEvent({ userId, operation: 'updated', @@ -932,43 +914,35 @@ MENTION PROCESSING: standard: currentUsageSummary.standard, pro: currentUsageSummary.pro }); - - usageLogger.debug('Global Assistant usage broadcast sent', { - userId: maskedUserId, - conversationId: maskedConversationId, - }); } catch (broadcastError) { usageLogger.error('Global Assistant usage broadcast failed', broadcastError instanceof Error ? broadcastError : undefined, { - userId: maskedUserId, - conversationId: maskedConversationId, + userId: maskIdentifier(userId), + conversationId: maskIdentifier(conversationId), }); } - } catch (usageError) { usageLogger.error('Global Assistant usage tracking failed', usageError as Error, { - userId: maskedUserId, + userId: maskIdentifier(userId), provider: currentProvider, - messageId: maskedMessageId, - conversationId: maskedConversationId, + messageId: maskIdentifier(messageId), + conversationId: maskIdentifier(conversationId), }); - - // Don't fail the request - usage tracking errors shouldn't break the chat } - } else { - usageLogger.debug('Skipping usage tracking for non-PageSpace provider', { - provider: currentProvider, - userId: maskedUserId, - messageId: maskedMessageId, - conversationId: maskedConversationId, - }); } } catch (error) { - loggers.api.error('❌ Global Assistant Chat API: Failed to save AI response message:', error as Error); + loggers.api.error('Global Assistant Chat API: Failed to save AI response message', error as Error); } } }, }); + loggers.api.debug('Global Assistant Chat API: Returning stream response', {}); + + return createUIMessageStreamResponse({ + stream, + headers: { [STREAM_ID_HEADER]: streamId }, + }); + } catch (error) { loggers.api.error('Global Assistant Chat API Error:', error as Error); diff --git a/apps/web/src/app/api/ai/page-agents/[agentId]/conversations/[conversationId]/messages/[messageId]/route.ts b/apps/web/src/app/api/ai/page-agents/[agentId]/conversations/[conversationId]/messages/[messageId]/route.ts new file mode 100644 index 0000000000..2ec154cefa --- /dev/null +++ b/apps/web/src/app/api/ai/page-agents/[agentId]/conversations/[conversationId]/messages/[messageId]/route.ts @@ -0,0 +1,194 @@ +import { NextResponse } from 'next/server'; +import { authenticateRequestWithOptions, isAuthError } from '@/lib/auth'; +import { canUserEditPage, loggers } from '@pagespace/lib/server'; +import { maskIdentifier } from '@/lib/logging/mask'; +import { + chatMessageRepository, + processMessageContentUpdate, +} from '@/lib/repositories/chat-message-repository'; +import { getActorInfo, logMessageActivity } from '@pagespace/lib/monitoring/activity-logger'; + +const AUTH_OPTIONS = { allow: ['session', 'mcp'] as const, requireCSRF: true }; + +/** + * PATCH - Edit a page agent conversation message's content + * Updates the message text and sets editedAt timestamp + */ +export async function PATCH( + request: Request, + context: { params: Promise<{ agentId: string; conversationId: string; messageId: string }> } +) { + try { + const auth = await authenticateRequestWithOptions(request, AUTH_OPTIONS); + if (isAuthError(auth)) return auth.error; + const userId = auth.userId; + + const { agentId, conversationId, messageId } = await context.params; + const { content } = await request.json(); + + // Validate content + if (!content || typeof content !== 'string') { + return NextResponse.json( + { error: 'Content is required and must be a string' }, + { status: 400 } + ); + } + + // Check if user can edit the page (agent) this message belongs to + const canEdit = await canUserEditPage(userId, agentId); + if (!canEdit) { + loggers.api.warn('Edit agent message permission denied', { + userId: maskIdentifier(userId), + messageId: maskIdentifier(messageId), + agentId: maskIdentifier(agentId), + }); + return NextResponse.json( + { error: 'You do not have permission to edit messages in this chat' }, + { status: 403 } + ); + } + + // Get the message to verify it exists, is active, and belongs to this conversation + const message = await chatMessageRepository.getMessageById(messageId); + if (!message || !message.isActive) { + return NextResponse.json({ error: 'Message not found' }, { status: 404 }); + } + + // Verify message belongs to this agent and conversation + if (message.pageId !== agentId || message.conversationId !== conversationId) { + return NextResponse.json({ error: 'Message not found in this conversation' }, { status: 404 }); + } + + // Store original content for activity logging + const originalContent = message.content; + + // Process content, preserving structured format if present + const updatedContent = processMessageContentUpdate(message.content, content); + + // Update the message content and set editedAt + await chatMessageRepository.updateMessageContent(messageId, updatedContent); + + // Log activity for audit trail + try { + const actorInfo = await getActorInfo(userId); + logMessageActivity(userId, 'message_update', { + id: messageId, + pageId: agentId, + driveId: null, + conversationType: 'ai_chat', + }, actorInfo, { + previousContent: originalContent, + newContent: updatedContent, + aiConversationId: conversationId, + }); + } catch (loggingError) { + loggers.api.error('Failed to log agent message update activity', loggingError as Error, { + messageId: maskIdentifier(messageId), + agentId: maskIdentifier(agentId), + }); + } + + loggers.api.info('Agent message edited successfully', { + userId: maskIdentifier(userId), + messageId: maskIdentifier(messageId), + agentId: maskIdentifier(agentId), + conversationId: maskIdentifier(conversationId), + }); + + return NextResponse.json({ + success: true, + message: 'Message updated successfully', + }); + } catch (error) { + loggers.api.error('Error editing agent message', error as Error); + return NextResponse.json( + { error: 'Failed to edit message' }, + { status: 500 } + ); + } +} + +/** + * DELETE - Soft delete a page agent conversation message + * Sets isActive to false to hide the message + */ +export async function DELETE( + request: Request, + context: { params: Promise<{ agentId: string; conversationId: string; messageId: string }> } +) { + try { + const auth = await authenticateRequestWithOptions(request, AUTH_OPTIONS); + if (isAuthError(auth)) return auth.error; + const userId = auth.userId; + + const { agentId, conversationId, messageId } = await context.params; + + // Check if user can edit the page (agent) this message belongs to + const canEdit = await canUserEditPage(userId, agentId); + if (!canEdit) { + loggers.api.warn('Delete agent message permission denied', { + userId: maskIdentifier(userId), + messageId: maskIdentifier(messageId), + agentId: maskIdentifier(agentId), + }); + return NextResponse.json( + { error: 'You do not have permission to delete messages in this chat' }, + { status: 403 } + ); + } + + // Get the message to verify it exists, is active, and belongs to this conversation + const message = await chatMessageRepository.getMessageById(messageId); + if (!message || !message.isActive) { + return NextResponse.json({ error: 'Message not found' }, { status: 404 }); + } + + // Verify message belongs to this agent and conversation + if (message.pageId !== agentId || message.conversationId !== conversationId) { + return NextResponse.json({ error: 'Message not found in this conversation' }, { status: 404 }); + } + + // Store content for audit trail before deletion + const deletedContent = message.content; + + // Soft delete the message + await chatMessageRepository.softDeleteMessage(messageId); + + // Log activity for audit trail + try { + const actorInfo = await getActorInfo(userId); + logMessageActivity(userId, 'message_delete', { + id: messageId, + pageId: agentId, + driveId: null, + conversationType: 'ai_chat', + }, actorInfo, { + previousContent: deletedContent, + aiConversationId: conversationId, + }); + } catch (loggingError) { + loggers.api.error('Failed to log agent message deletion activity', loggingError as Error, { + messageId: maskIdentifier(messageId), + agentId: maskIdentifier(agentId), + }); + } + + loggers.api.info('Agent message deleted successfully', { + userId: maskIdentifier(userId), + messageId: maskIdentifier(messageId), + agentId: maskIdentifier(agentId), + conversationId: maskIdentifier(conversationId), + }); + + return NextResponse.json({ + success: true, + message: 'Message deleted successfully', + }); + } catch (error) { + loggers.api.error('Error deleting agent message', error as Error); + return NextResponse.json( + { error: 'Failed to delete message' }, + { status: 500 } + ); + } +} diff --git a/apps/web/src/components/layout/middle-content/page-views/ai-page/AiChatView.tsx b/apps/web/src/components/layout/middle-content/page-views/ai-page/AiChatView.tsx index 7e9634e757..bcbd8c95d9 100644 --- a/apps/web/src/components/layout/middle-content/page-views/ai-page/AiChatView.tsx +++ b/apps/web/src/components/layout/middle-content/page-views/ai-page/AiChatView.tsx @@ -13,8 +13,7 @@ import { useParams } from 'next/navigation'; import { Button } from '@/components/ui/button'; import { Tabs, TabsContent, TabsList, TabsTrigger } from '@/components/ui/tabs'; import { Loader2, Settings, MessageSquare, History, Plus, Save } from 'lucide-react'; -import { UIMessage, DefaultChatTransport } from 'ai'; -import { useEditingStore } from '@/stores/useEditingStore'; +import { UIMessage } from 'ai'; import { useAssistantSettingsStore } from '@/stores/useAssistantSettingsStore'; import { useVoiceModeStore } from '@/stores/useVoiceModeStore'; import { buildPagePath } from '@/lib/tree/tree-utils'; @@ -25,7 +24,7 @@ import { PageAgentSettingsTab, PageAgentHistoryTab, type PageAgentSettingsTabRef import { fetchWithAuth } from '@/lib/auth/auth-fetch'; import { VoiceModeOverlay } from '@/components/ai/voice'; -import { abortActiveStream, createStreamTrackingFetch, clearActiveStreamId } from '@/lib/ai/core/client'; +import { clearActiveStreamId } from '@/lib/ai/core/client'; import { useAppStateRecovery } from '@/hooks/useAppStateRecovery'; // Shared hooks and components @@ -34,6 +33,9 @@ import { useMessageActions, useProviderSettings, useConversations, + useChatTransport, + useStreamingRegistration, + useChatStop, AgentConfig, } from '@/lib/ai/shared'; import { @@ -148,41 +150,27 @@ const AiChatView: React.FC = ({ page }) => { // ============================================ // Use conversation ID for stream tracking (falls back to page.id before conversation is created) const streamTrackingId = currentConversationId || page.id; + + const transport = useChatTransport(streamTrackingId, '/api/ai/chat'); + const chatConfig = useMemo( - () => ({ + () => !transport ? null : ({ id: page.id, messages: initialMessages, - transport: new DefaultChatTransport({ - api: '/api/ai/chat', - // Use stream tracking fetch to capture streamId from response headers - // This enables explicit abort via /api/ai/abort endpoint - fetch: createStreamTrackingFetch({ chatId: streamTrackingId }), - }), - experimental_throttle: 100, // Increased from 50ms for better performance + transport, + experimental_throttle: 100, onError: (error: Error) => { console.error('AiChatView: Chat error:', error); }, }), - // Re-create transport when conversation changes for proper stream tracking - [page.id, streamTrackingId, initialMessages] + [page.id, transport, initialMessages] ); const { messages, sendMessage, status, error, regenerate, setMessages, stop: chatStop } = - useChat(chatConfig); + useChat(chatConfig || {}); const isStreaming = status === 'submitted' || status === 'streaming'; - - // Combined stop function that calls both abort endpoint (server-side) and useChat stop (client-side) - // Use try/finally to guarantee client-side stop runs even if server abort fails - const stop = useCallback(async () => { - try { - // Call abort endpoint to stop server-side processing - await abortActiveStream({ chatId: streamTrackingId }); - } finally { - // Call useChat's stop to abort client-side fetch - chatStop(); - } - }, [streamTrackingId, chatStop]); + const stop = useChatStop(streamTrackingId, chatStop); const isLoading = !isInitialized; // ============================================ @@ -265,20 +253,11 @@ const AiChatView: React.FC = ({ page }) => { }, [page.id]); // Register streaming state with editing store - useEffect(() => { - const componentId = `ai-chat-${page.id}`; - if (status === 'submitted' || status === 'streaming') { - useEditingStore.getState().startStreaming(componentId, { - pageId: page.id, - componentName: 'AiChatView', - }); - } else { - useEditingStore.getState().endStreaming(componentId); - } - return () => { - useEditingStore.getState().endStreaming(componentId); - }; - }, [status, page.id]); + useStreamingRegistration( + `ai-chat-${page.id}`, + isStreaming, + { pageId: page.id, componentName: 'AiChatView' } + ); // Reset error visibility when new error occurs useEffect(() => { diff --git a/apps/web/src/components/layout/middle-content/page-views/dashboard/GlobalAssistantView.tsx b/apps/web/src/components/layout/middle-content/page-views/dashboard/GlobalAssistantView.tsx index 1a22df2583..ea206dbddd 100644 --- a/apps/web/src/components/layout/middle-content/page-views/dashboard/GlobalAssistantView.tsx +++ b/apps/web/src/components/layout/middle-content/page-views/dashboard/GlobalAssistantView.tsx @@ -41,7 +41,6 @@ import React, { useEffect, useState, useRef, useMemo, useCallback } from 'react'; import { useChat } from '@ai-sdk/react'; -import { DefaultChatTransport } from 'ai'; import { usePathname } from 'next/navigation'; import { Button } from '@/components/ui/button'; import { Activity, Plus, History } from 'lucide-react'; @@ -49,7 +48,6 @@ import { AiUsageMonitor, AISelector, TasksDropdown } from '@/components/ai/share import { useLayoutStore } from '@/stores/useLayoutStore'; import { useDriveStore } from '@/hooks/useDrive'; import { fetchWithAuth } from '@/lib/auth/auth-fetch'; -import { useEditingStore } from '@/stores/useEditingStore'; import { useAssistantSettingsStore } from '@/stores/useAssistantSettingsStore'; import { useGlobalChat } from '@/contexts/GlobalChatContext'; import { usePageAgentDashboardStore } from '@/stores/page-agents'; @@ -62,9 +60,12 @@ import { useMCPTools, useMessageActions, useProviderSettings, + useChatTransport, + useStreamingRegistration, + useChatStop, LocationContext, } from '@/lib/ai/shared'; -import { abortActiveStream, createStreamTrackingFetch, clearActiveStreamId } from '@/lib/ai/core/client'; +import { abortActiveStream, clearActiveStreamId } from '@/lib/ai/core/client'; import { useAppStateRecovery } from '@/hooks/useAppStateRecovery'; import { ProviderSetupCard, @@ -257,24 +258,22 @@ const GlobalAssistantView: React.FC = () => { // CHAT CONFIGURATION // ============================================ + const agentTransport = useChatTransport(agentConversationId, '/api/ai/chat'); + // Agent mode chat config const agentChatConfig = useMemo(() => { - if (!selectedAgent || !agentConversationId) return null; + if (!selectedAgent || !agentConversationId || !agentTransport) return null; + return { id: agentConversationId, messages: agentInitialMessages, - transport: new DefaultChatTransport({ - api: '/api/ai/chat', - // Use stream tracking fetch to capture streamId from response headers - // This enables explicit abort via /api/ai/abort endpoint - fetch: createStreamTrackingFetch({ chatId: agentConversationId }), - }), - experimental_throttle: 100, // Increased from 50ms for better performance + transport: agentTransport, + experimental_throttle: 100, onError: (error: Error) => { console.error('Agent Chat error:', error); }, }; - }, [selectedAgent, agentConversationId, agentInitialMessages]); + }, [selectedAgent, agentConversationId, agentTransport, agentInitialMessages]); // Global mode chat const { @@ -308,19 +307,7 @@ const GlobalAssistantView: React.FC = () => { const regenerate = selectedAgent ? agentRegenerate : globalRegenerate; const rawStop = selectedAgent ? agentStop : globalStop; const isStreaming = status === 'submitted' || status === 'streaming'; - - // Wrap stop handler to abort server-side stream before client-side stop - // This ensures the server stops processing when user clicks Stop - // Use try/finally to guarantee client-side stop runs even if server abort fails - const stop = useCallback(async () => { - try { - if (currentConversationId) { - await abortActiveStream({ chatId: currentConversationId }); - } - } finally { - rawStop(); - } - }, [currentConversationId, rawStop]); + const stop = useChatStop(currentConversationId, rawStop); // Agent mode: initialized when we have a conversationId and not loading // Global mode: use globalIsInitialized from context const agentIsInitialized = selectedAgent ? (!!agentConversationId && !agentIsLoading) : false; @@ -525,20 +512,11 @@ const GlobalAssistantView: React.FC = () => { }, [selectedAgent, agentStatus, agentStop, agentConversationId, setAgentStopStreaming]); // Register streaming state with editing store - useEffect(() => { - const componentId = `global-assistant-${currentConversationId || 'init'}`; - if (status === 'submitted' || status === 'streaming') { - useEditingStore.getState().startStreaming(componentId, { - conversationId: currentConversationId || undefined, - componentName: 'GlobalAssistantView', - }); - } else { - useEditingStore.getState().endStreaming(componentId); - } - return () => { - useEditingStore.getState().endStreaming(componentId); - }; - }, [status, currentConversationId]); + useStreamingRegistration( + `global-assistant-${currentConversationId || 'init'}`, + isStreaming, + { conversationId: currentConversationId || undefined, componentName: 'GlobalAssistantView' } + ); // Reset error visibility when new error occurs useEffect(() => { diff --git a/apps/web/src/components/layout/right-sidebar/ai-assistant/SidebarChatTab.tsx b/apps/web/src/components/layout/right-sidebar/ai-assistant/SidebarChatTab.tsx index 3f61e2a41b..91926e0e16 100644 --- a/apps/web/src/components/layout/right-sidebar/ai-assistant/SidebarChatTab.tsx +++ b/apps/web/src/components/layout/right-sidebar/ai-assistant/SidebarChatTab.tsx @@ -1,5 +1,5 @@ import React, { useEffect, useState, useRef, useMemo, useCallback } from 'react'; -import { DefaultChatTransport, UIMessage } from 'ai'; +import { UIMessage } from 'ai'; import { usePathname } from 'next/navigation'; import { Button } from '@/components/ui/button'; import { ChatInput, type ChatInputRef } from '@/components/ai/chat/input'; @@ -15,7 +15,6 @@ import { } from '@/components/ai/ui/conversation'; import { useDriveStore } from '@/hooks/useDrive'; import { fetchWithAuth, patch, del } from '@/lib/auth/auth-fetch'; -import { useEditingStore } from '@/stores/useEditingStore'; import { useAssistantSettingsStore } from '@/stores/useAssistantSettingsStore'; import { useVoiceModeStore } from '@/stores/useVoiceModeStore'; import { useGlobalChat } from '@/contexts/GlobalChatContext'; @@ -23,7 +22,8 @@ import { usePageAgentSidebarState, usePageAgentSidebarChat, type SidebarAgentInf import { usePageAgentDashboardStore } from '@/stores/page-agents'; import { toast } from 'sonner'; import { LocationContext } from '@/lib/ai/shared'; -import { abortActiveStream, createStreamTrackingFetch, clearActiveStreamId } from '@/lib/ai/core/client'; +import { abortActiveStream, clearActiveStreamId } from '@/lib/ai/core/client'; +import { useChatTransport, useStreamingRegistration } from '@/lib/ai/shared'; import { useMobileKeyboard } from '@/hooks/useMobileKeyboard'; import { useAppStateRecovery } from '@/hooks/useAppStateRecovery'; import { VoiceModeOverlay } from '@/components/ai/voice'; @@ -173,24 +173,22 @@ const SidebarChatTab: React.FC = () => { // ============================================ // Agent Chat Configuration // ============================================ + const agentTransport = useChatTransport(agentConversationId, '/api/ai/chat'); + const agentChatConfig = useMemo(() => { - if (!selectedAgent || !agentConversationId) return null; + if (!selectedAgent || !agentConversationId || !agentTransport) return null; + return { id: agentConversationId, messages: agentInitialMessages, - transport: new DefaultChatTransport({ - api: '/api/ai/chat', - // Use stream tracking fetch to capture streamId from response headers - // This enables explicit abort via /api/ai/abort endpoint - fetch: createStreamTrackingFetch({ chatId: agentConversationId }), - }), - experimental_throttle: 100, // Increased from 50ms for better performance + transport: agentTransport, + experimental_throttle: 100, onError: (error: Error) => { console.error('Sidebar Agent Chat error:', error); toast.error('Chat error. Please try again.'); }, }; - }, [selectedAgent, agentConversationId, agentInitialMessages]); + }, [selectedAgent, agentConversationId, agentTransport, agentInitialMessages]); // ============================================ // Sidebar Chat (custom hook - unified interface) @@ -427,22 +425,11 @@ const SidebarChatTab: React.FC = () => { // ============================================ // Effects: Editing Store Registration // ============================================ - useEffect(() => { - const componentId = `assistant-sidebar-${currentConversationId || 'init'}`; - - if (status === 'submitted' || status === 'streaming') { - useEditingStore.getState().startStreaming(componentId, { - conversationId: currentConversationId || undefined, - componentName: 'SidebarChatTab', - }); - } else { - useEditingStore.getState().endStreaming(componentId); - } - - return () => { - useEditingStore.getState().endStreaming(componentId); - }; - }, [status, currentConversationId]); + useStreamingRegistration( + `assistant-sidebar-${currentConversationId || 'init'}`, + status === 'submitted' || status === 'streaming', + { conversationId: currentConversationId || undefined, componentName: 'SidebarChatTab' } + ); // ============================================ // Effects: UI State diff --git a/apps/web/src/contexts/GlobalChatContext.tsx b/apps/web/src/contexts/GlobalChatContext.tsx index 3001eff90c..502dcf2f43 100644 --- a/apps/web/src/contexts/GlobalChatContext.tsx +++ b/apps/web/src/contexts/GlobalChatContext.tsx @@ -5,7 +5,7 @@ import { DefaultChatTransport, UIMessage } from 'ai'; import { fetchWithAuth } from '@/lib/auth/auth-fetch'; import { conversationState } from '@/lib/ai/core/conversation-state'; import { getAgentId, getConversationId, setConversationId } from '@/lib/url-state'; -import { createStreamTrackingFetch } from '@/lib/ai/core/client'; +import { useChatTransport } from '@/lib/ai/shared'; /** * Global Chat Context - ONLY for Global Assistant state @@ -203,35 +203,31 @@ export function GlobalChatProvider({ children }: { children: ReactNode }) { } }, [currentConversationId]); + // Stable transport that only recreates when conversation ID changes + const apiEndpoint = currentConversationId ? `/api/ai/global/${currentConversationId}/messages` : ''; + const transport = useChatTransport(currentConversationId, apiEndpoint); + // Create stable chat config // IMPORTANT: Uses initialMessages which is set by loadConversation when switching conversations. // The chatConfig only changes when: // 1. currentConversationId changes (switching conversations) // 2. initialMessages changes (set during loadConversation) - // This keeps the config stable during normal messaging to avoid confusing useChat. const chatConfig = useMemo(() => { - if (!currentConversationId) return null; - - const apiEndpoint = `/api/ai/global/${currentConversationId}/messages`; + if (!currentConversationId || !transport) return null; return { id: currentConversationId, messages: initialMessages, - transport: new DefaultChatTransport({ - api: apiEndpoint, - // Use stream tracking fetch to capture streamId from response headers - // This enables explicit abort via /api/ai/abort endpoint - fetch: createStreamTrackingFetch({ chatId: currentConversationId }), - }), - experimental_throttle: 100, // Match agent mode throttle for consistent streaming feel + transport, + experimental_throttle: 100, onError: (error: Error) => { - console.error('❌ Global Chat Error:', error); + console.error('Global Chat Error:', error); if (error.message?.includes('Unauthorized') || error.message?.includes('401')) { - console.error('🔒 Authentication failed - user may need to log in again'); + console.error('Authentication failed - user may need to log in again'); } }, }; - }, [currentConversationId, initialMessages]); + }, [currentConversationId, transport, initialMessages]); // Context value const contextValue: GlobalChatContextValue = useMemo(() => ({ diff --git a/apps/web/src/lib/ai/shared/hooks/index.ts b/apps/web/src/lib/ai/shared/hooks/index.ts index 841c1f6640..94eb41cad5 100644 --- a/apps/web/src/lib/ai/shared/hooks/index.ts +++ b/apps/web/src/lib/ai/shared/hooks/index.ts @@ -6,3 +6,6 @@ export { useMCPTools } from './useMCPTools'; export { useConversations } from './useConversations'; export { useMessageActions } from './useMessageActions'; export { useProviderSettings } from './useProviderSettings'; +export { useChatTransport } from './useChatTransport'; +export { useStreamingRegistration } from './useStreamingRegistration'; +export { useChatStop } from './useChatStop'; diff --git a/apps/web/src/lib/ai/shared/hooks/useChatStop.ts b/apps/web/src/lib/ai/shared/hooks/useChatStop.ts new file mode 100644 index 0000000000..cb366840e8 --- /dev/null +++ b/apps/web/src/lib/ai/shared/hooks/useChatStop.ts @@ -0,0 +1,23 @@ +import { useCallback } from 'react'; +import { abortActiveStream } from '@/lib/ai/core/client'; + +/** + * Returns a memoized stop function that aborts the server-side stream + * then stops the client-side fetch via useChat's stop. + * + * Uses try/finally to guarantee client-side stop runs even if server abort fails. + */ +export function useChatStop( + chatId: string | null, + chatStop: () => void +): () => Promise { + return useCallback(async () => { + try { + if (chatId) { + await abortActiveStream({ chatId }); + } + } finally { + chatStop(); + } + }, [chatId, chatStop]); +} diff --git a/apps/web/src/lib/ai/shared/hooks/useChatTransport.ts b/apps/web/src/lib/ai/shared/hooks/useChatTransport.ts new file mode 100644 index 0000000000..b06451e3fa --- /dev/null +++ b/apps/web/src/lib/ai/shared/hooks/useChatTransport.ts @@ -0,0 +1,34 @@ +import { useRef } from 'react'; +import { DefaultChatTransport, UIMessage } from 'ai'; +import { createStreamTrackingFetch } from '@/lib/ai/core/client'; + +/** + * Creates a stable DefaultChatTransport instance that only recreates when + * the conversation ID or API endpoint changes. Avoids unnecessary useChat + * state resets caused by new transport identity. + * + * Returns null when conversationId is null (no active conversation). + */ +export function useChatTransport( + conversationId: string | null, + api: string +): DefaultChatTransport | null { + const transportRef = useRef | null>(null); + const trackingIdRef = useRef(null); + const apiRef = useRef(api); + + if (!conversationId) { + return null; + } + + if (trackingIdRef.current !== conversationId || apiRef.current !== api || !transportRef.current) { + transportRef.current = new DefaultChatTransport({ + api, + fetch: createStreamTrackingFetch({ chatId: conversationId }), + }); + trackingIdRef.current = conversationId; + apiRef.current = api; + } + + return transportRef.current; +} diff --git a/apps/web/src/lib/ai/shared/hooks/useMessageActions.ts b/apps/web/src/lib/ai/shared/hooks/useMessageActions.ts index abc75a90a0..3991c75368 100644 --- a/apps/web/src/lib/ai/shared/hooks/useMessageActions.ts +++ b/apps/web/src/lib/ai/shared/hooks/useMessageActions.ts @@ -68,49 +68,57 @@ export function useMessageActions({ async (messageId: string, newContent: string) => { if (!conversationId) return; + // Optimistically update the message in local state first + // This ensures the UI reflects the edit immediately, even if the API call fails + const updatedMessages = messages.map((m) => { + if (m.id !== messageId) return m; + return { + ...m, + parts: m.parts.map((part) => + part.type === 'text' ? { ...part, text: newContent } : part + ), + }; + }); + setMessages(updatedMessages); + try { if (isAgentMode) { - // Agent mode: Use agent API await patch( `/api/ai/page-agents/${agentId}/conversations/${conversationId}/messages/${messageId}`, { content: newContent } ); - - // Refetch agent messages - const response = await fetchWithAuth( - `/api/ai/page-agents/${agentId}/conversations/${conversationId}/messages` - ); - if (response.ok) { - const data = await response.json(); - setMessages(data.messages || []); - } } else { - // Global mode: Use global API await patch( `/api/ai/global/${conversationId}/messages/${messageId}`, { content: newContent } ); + } - // Refetch messages - const response = await fetchWithAuth( - `/api/ai/global/${conversationId}/messages` - ); + onEditVersionChange?.(); + toast.success('Message updated successfully'); + + // Refetch to reconcile with server state (non-critical) + try { + const url = isAgentMode + ? `/api/ai/page-agents/${agentId}/conversations/${conversationId}/messages` + : `/api/ai/global/${conversationId}/messages`; + const response = await fetchWithAuth(url); if (response.ok) { const data = await response.json(); - const loadedMessages = Array.isArray(data) ? data : data.messages || []; - setMessages(loadedMessages); + const loaded = isAgentMode + ? data.messages || [] + : Array.isArray(data) ? data : data.messages || []; + setMessages(loaded); } + } catch { + // Refetch failed — optimistic update already applied, server has the edit } - - onEditVersionChange?.(); - toast.success('Message updated successfully'); } catch (error) { console.error('Failed to edit message:', error); - toast.error('Failed to edit message'); - throw error; + toast.error('Failed to save edit. Your local changes may not persist.'); } }, - [isAgentMode, agentId, conversationId, setMessages, onEditVersionChange] + [isAgentMode, agentId, conversationId, messages, setMessages, onEditVersionChange] ); // Delete a message @@ -118,6 +126,11 @@ export function useMessageActions({ async (messageId: string) => { if (!conversationId) return; + // Optimistically update local state first + const previousMessages = [...messages]; + const filtered = messages.filter((m) => m.id !== messageId); + setMessages(filtered); + try { if (isAgentMode) { await del( @@ -127,13 +140,11 @@ export function useMessageActions({ await del(`/api/ai/global/${conversationId}/messages/${messageId}`); } - // Optimistically update local state - const filtered = messages.filter((m) => m.id !== messageId); - setMessages(filtered); - toast.success('Message deleted'); } catch (error) { console.error('Failed to delete message:', error); + // Revert optimistic update on failure + setMessages(previousMessages); toast.error('Failed to delete message'); throw error; } diff --git a/apps/web/src/lib/ai/shared/hooks/useStreamingRegistration.ts b/apps/web/src/lib/ai/shared/hooks/useStreamingRegistration.ts new file mode 100644 index 0000000000..34a5589db3 --- /dev/null +++ b/apps/web/src/lib/ai/shared/hooks/useStreamingRegistration.ts @@ -0,0 +1,26 @@ +import { useEffect } from 'react'; +import { useEditingStore, type EditingSession } from '@/stores/useEditingStore'; + +/** + * Registers/unregisters streaming state with the editing store. + * Prevents SWR revalidation and other UI refreshes during active AI streaming. + * + * Cleans up on unmount. + */ +export function useStreamingRegistration( + id: string, + isStreaming: boolean, + metadata?: EditingSession['metadata'] +): void { + useEffect(() => { + if (isStreaming) { + useEditingStore.getState().startStreaming(id, metadata); + } else { + useEditingStore.getState().endStreaming(id); + } + return () => { + useEditingStore.getState().endStreaming(id); + }; + // eslint-disable-next-line react-hooks/exhaustive-deps -- individual metadata fields avoid re-running on object identity changes + }, [id, isStreaming, metadata?.pageId, metadata?.conversationId, metadata?.componentName]); +}