diff --git a/sim/app/api/chat/[subdomain]/route.ts b/sim/app/api/chat/[subdomain]/route.ts index 8e2fb71097..54aac10250 100644 --- a/sim/app/api/chat/[subdomain]/route.ts +++ b/sim/app/api/chat/[subdomain]/route.ts @@ -1,4 +1,4 @@ -import { NextRequest } from 'next/server' +import { NextRequest, NextResponse } from 'next/server' import { eq } from 'drizzle-orm' import { createLogger } from '@/lib/logs/console-logger' import { db } from '@/db' @@ -96,6 +96,28 @@ export async function POST(request: NextRequest, { params }: { params: Promise<{ // Execute the workflow using our helper function const result = await executeWorkflowForChat(deployment.id, message) + // If the executor returned a ReadableStream, stream it directly to the client + if (result instanceof ReadableStream) { + const streamResponse = new NextResponse(result, { + status: 200, + headers: { + 'Content-Type': 'text/plain; charset=utf-8', + }, + }) + return addCorsHeaders(streamResponse, request) + } + + // Handle StreamingExecution format + if (result && typeof result === 'object' && 'stream' in result && 'execution' in result) { + const streamResponse = new NextResponse(result.stream as ReadableStream, { + status: 200, + headers: { + 'Content-Type': 'text/plain; charset=utf-8', + }, + }) + return addCorsHeaders(streamResponse, request) + } + // Format the result for the client // If result.content is an object, preserve it for structured handling // If it's text or another primitive, make sure it's accessible diff --git a/sim/app/api/chat/utils.ts b/sim/app/api/chat/utils.ts index 68b5e79e53..674b92f702 100644 --- a/sim/app/api/chat/utils.ts +++ b/sim/app/api/chat/utils.ts @@ -11,6 +11,11 @@ import { Serializer } from '@/serializer' import { mergeSubblockState } from '@/stores/workflows/utils' import { persistExecutionLogs } from '@/lib/logs/execution-logger' import { buildTraceSpans } from '@/lib/logs/trace-spans' +import { BlockLog } from '@/executor/types' + +declare global { + var __chatStreamProcessingTasks: Promise<{success: boolean, error?: any}>[] | undefined +} const logger = createLogger('ChatAuthUtils') const isDevelopment = process.env.NODE_ENV === 'development' @@ -393,16 +398,108 @@ export async function executeWorkflowForChat(chatId: string, message: string) { ) // Create and execute the workflow - mimicking use-workflow-execution.ts - const executor = new Executor( - serializedWorkflow, - processedBlockStates, - decryptedEnvVars, - { input: message }, - workflowVariables - ) - + const executor = new Executor({ + workflow: serializedWorkflow, + currentBlockStates: processedBlockStates, + envVarValues: decryptedEnvVars, + workflowInput: { input: message }, + workflowVariables, + contextExtensions: { + // Always request streaming – the executor will downgrade gracefully if unsupported + stream: true, + selectedOutputIds: outputBlockIds, + edges: edges.map((e: any) => ({ source: e.source, target: e.target })), + }, + }) + // Execute and capture the result const result = await executor.execute(workflowId) + + // If the executor returned a ReadableStream, forward it directly for streaming + if (result instanceof ReadableStream) { + return result + } + + // Handle StreamingExecution format (combined stream + execution data) + if (result && typeof result === 'object' && 'stream' in result && 'execution' in result) { + // We need to stream the response to the client while *also* capturing the full + // content so that we can persist accurate logs once streaming completes. + + // Duplicate the original stream – one copy goes to the client, the other we read + // server-side for log enrichment. + const [clientStream, loggingStream] = (result.stream as ReadableStream).tee() + + // Kick off background processing to read the stream and persist enriched logs + const processingPromise = (async () => { + try { + // The stream is only used to properly drain it and prevent memory leaks + // All the execution data is already provided from the agent handler + // through the X-Execution-Data header + await drainStream(loggingStream) + + // No need to wait for a processing promise + // The execution-logger.ts will handle token estimation + + // We can use the execution data as-is since it's already properly structured + const executionData = result.execution as any + + // Before persisting, clean up any response objects with zero tokens in agent blocks + // This prevents confusion in the console logs + if (executionData.logs && Array.isArray(executionData.logs)) { + executionData.logs.forEach((log: BlockLog) => { + if (log.blockType === 'agent' && log.output?.response) { + const response = log.output.response; + + // Check for zero tokens that will be estimated later + if (response.tokens && + (!response.tokens.completion || response.tokens.completion === 0) && + (!response.toolCalls || !response.toolCalls.list || response.toolCalls.list.length === 0)) { + + // Remove tokens from console display to avoid confusion + // They'll be properly estimated in the execution logger + delete response.tokens; + } + } + }); + } + + // Build trace spans and persist + const { traceSpans, totalDuration } = buildTraceSpans(executionData) + const enrichedResult = { + ...executionData, + traceSpans, + totalDuration, + } + + const executionId = uuidv4() + await persistExecutionLogs(workflowId, executionId, enrichedResult, 'chat') + logger.debug(`[${requestId}] Persisted execution logs for streaming chat with ID: ${executionId}`) + + return { success: true } + } catch (error) { + logger.error(`[${requestId}] Failed to persist streaming chat execution logs:`, error) + return { success: false, error } + } finally { + // Ensure the stream is properly closed even if an error occurs + try { + const controller = new AbortController() + const signal = controller.signal + controller.abort() + } catch (cleanupError) { + logger.debug(`[${requestId}] Error during stream cleanup: ${cleanupError}`) + } + } + })() + + // Register this processing promise with a global handler or tracker if needed + // This allows the background task to be monitored or waited for in testing + if (typeof global.__chatStreamProcessingTasks !== 'undefined') { + global.__chatStreamProcessingTasks.push(processingPromise as Promise<{success: boolean, error?: any}>) + } + + // Return the client-facing stream + return clientStream + } // Mark as chat execution in metadata if (result) { @@ -412,7 +509,7 @@ export async function executeWorkflowForChat(chatId: string, message: string) { } } - // Persist execution logs using the 'chat' trigger type + // Persist execution logs using the 'chat' trigger type for non-streaming results try { // Build trace spans to enrich the logs (same as in use-workflow-execution.ts) const { traceSpans, totalDuration } = buildTraceSpans(result) @@ -543,4 +640,20 @@ export async function executeWorkflowForChat(chatId: string, message: string) { type: 'workflow' } } +} + +/** + * Utility function to properly drain a stream to prevent memory leaks + */ +async function drainStream(stream: ReadableStream): Promise { + const reader = stream.getReader() + try { + while (true) { + const { done, value } = await reader.read() + if (done) break + // We don't need to do anything with the value, just drain the stream + } + } finally { + reader.releaseLock() + } } \ No newline at end of file diff --git a/sim/app/api/providers/route.ts b/sim/app/api/providers/route.ts index bf1de4b9c6..12d2afc83c 100644 --- a/sim/app/api/providers/route.ts +++ b/sim/app/api/providers/route.ts @@ -2,6 +2,7 @@ import { NextRequest, NextResponse } from 'next/server' import { createLogger } from '@/lib/logs/console-logger' import { executeProviderRequest } from '@/providers' import { getApiKey } from '@/providers/utils' +import { StreamingExecution } from '@/executor/types' const logger = createLogger('ProvidersAPI') @@ -24,6 +25,7 @@ export async function POST(request: NextRequest) { apiKey, responseFormat, workflowId, + stream, } = body let finalApiKey: string @@ -48,8 +50,90 @@ export async function POST(request: NextRequest) { apiKey: finalApiKey, responseFormat, workflowId, + stream, }) + // Check if the response is a StreamingExecution + if (response && typeof response === 'object' && 'stream' in response && 'execution' in response) { + const streamingExec = response as StreamingExecution + logger.info('Received StreamingExecution from provider') + + // Extract the stream and execution data + const stream = streamingExec.stream + const executionData = streamingExec.execution + + // Attach the execution data as a custom header + // We need to safely serialize the execution data to avoid circular references + let executionDataHeader + try { + // Create a safe version of execution data with the most important fields + const safeExecutionData = { + success: executionData.success, + output: { + response: { + // Sanitize content to remove non-ASCII characters that would cause ByteString errors + content: executionData.output?.response?.content + ? String(executionData.output.response.content).replace(/[\u0080-\uFFFF]/g, '') + : '', + model: executionData.output?.response?.model, + tokens: executionData.output?.response?.tokens || { + prompt: 0, + completion: 0, + total: 0 + }, + // Sanitize any potential Unicode characters in tool calls + toolCalls: executionData.output?.response?.toolCalls + ? sanitizeToolCalls(executionData.output.response.toolCalls) + : undefined, + providerTiming: executionData.output?.response?.providerTiming, + cost: executionData.output?.response?.cost, + } + }, + error: executionData.error, + logs: [], // Strip logs from header to avoid encoding issues + metadata: { + startTime: executionData.metadata?.startTime, + endTime: executionData.metadata?.endTime, + duration: executionData.metadata?.duration + }, + isStreaming: true, // Always mark streaming execution data as streaming + blockId: executionData.logs?.[0]?.blockId, + blockName: executionData.logs?.[0]?.blockName, + blockType: executionData.logs?.[0]?.blockType, + } + executionDataHeader = JSON.stringify(safeExecutionData) + } catch (error) { + logger.error('Failed to serialize execution data:', error) + executionDataHeader = JSON.stringify({ + success: executionData.success, + error: 'Failed to serialize full execution data' + }) + } + + // Return the stream with execution data in a header + return new Response(stream, { + headers: { + 'Content-Type': 'text/event-stream', + 'Cache-Control': 'no-cache', + 'Connection': 'keep-alive', + 'X-Execution-Data': executionDataHeader + }, + }) + } + + // Check if the response is a ReadableStream for streaming + if (response instanceof ReadableStream) { + logger.info('Streaming response from provider') + return new Response(response, { + headers: { + 'Content-Type': 'text/event-stream', + 'Cache-Control': 'no-cache', + 'Connection': 'keep-alive', + }, + }) + } + + // Return regular JSON response for non-streaming return NextResponse.json(response) } catch (error) { logger.error('Provider request failed:', error) @@ -59,3 +143,89 @@ export async function POST(request: NextRequest) { ) } } + +/** + * Helper function to sanitize tool calls to remove Unicode characters + */ +function sanitizeToolCalls(toolCalls: any) { + // If it's an object with a list property, sanitize the list + if (toolCalls && typeof toolCalls === 'object' && Array.isArray(toolCalls.list)) { + return { + ...toolCalls, + list: toolCalls.list.map(sanitizeToolCall) + } + } + + // If it's an array, sanitize each item + if (Array.isArray(toolCalls)) { + return toolCalls.map(sanitizeToolCall) + } + + return toolCalls +} + +/** + * Sanitize a single tool call to remove Unicode characters + */ +function sanitizeToolCall(toolCall: any) { + if (!toolCall || typeof toolCall !== 'object') return toolCall + + // Create a sanitized copy + const sanitized = { ...toolCall } + + // Sanitize any string fields that might contain Unicode + if (typeof sanitized.name === 'string') { + sanitized.name = sanitized.name.replace(/[\u0080-\uFFFF]/g, '') + } + + // Sanitize input/arguments + if (sanitized.input && typeof sanitized.input === 'object') { + sanitized.input = sanitizeObject(sanitized.input) + } + + if (sanitized.arguments && typeof sanitized.arguments === 'object') { + sanitized.arguments = sanitizeObject(sanitized.arguments) + } + + // Sanitize output/result + if (sanitized.output && typeof sanitized.output === 'object') { + sanitized.output = sanitizeObject(sanitized.output) + } + + if (sanitized.result && typeof sanitized.result === 'object') { + sanitized.result = sanitizeObject(sanitized.result) + } + + // Sanitize error message + if (typeof sanitized.error === 'string') { + sanitized.error = sanitized.error.replace(/[\u0080-\uFFFF]/g, '') + } + + return sanitized +} + +/** + * Recursively sanitize an object to remove Unicode characters from strings + */ +function sanitizeObject(obj: any): any { + if (!obj || typeof obj !== 'object') return obj + + // Handle arrays + if (Array.isArray(obj)) { + return obj.map(item => sanitizeObject(item)) + } + + // Handle objects + const result: any = {} + for (const [key, value] of Object.entries(obj)) { + if (typeof value === 'string') { + result[key] = value.replace(/[\u0080-\uFFFF]/g, '') + } else if (typeof value === 'object' && value !== null) { + result[key] = sanitizeObject(value) + } else { + result[key] = value + } + } + + return result +} diff --git a/sim/app/api/schedules/execute/route.ts b/sim/app/api/schedules/execute/route.ts index 18868db6b4..8ea7262ee6 100644 --- a/sim/app/api/schedules/execute/route.ts +++ b/sim/app/api/schedules/execute/route.ts @@ -305,13 +305,19 @@ export async function GET(req: NextRequest) { ) const result = await executor.execute(schedule.workflowId) + // Check if we got a StreamingExecution result (with stream + execution properties) + // For scheduled executions, we only care about the ExecutionResult part, not the stream + const executionResult = 'stream' in result && 'execution' in result + ? result.execution + : result + logger.info(`[${requestId}] Workflow execution completed: ${schedule.workflowId}`, { - success: result.success, - executionTime: result.metadata?.duration, + success: executionResult.success, + executionTime: executionResult.metadata?.duration, }) // Update workflow run counts if execution was successful - if (result.success) { + if (executionResult.success) { await updateWorkflowRunCounts(schedule.workflowId) // Track scheduled execution in user stats @@ -325,11 +331,11 @@ export async function GET(req: NextRequest) { } // Build trace spans from execution logs - const { traceSpans, totalDuration } = buildTraceSpans(result) + const { traceSpans, totalDuration } = buildTraceSpans(executionResult) // Add trace spans to the execution result const enrichedResult = { - ...result, + ...executionResult, traceSpans, totalDuration, } @@ -338,7 +344,7 @@ export async function GET(req: NextRequest) { await persistExecutionLogs(schedule.workflowId, executionId, enrichedResult, 'schedule') // Only update next_run_at if execution was successful - if (result.success) { + if (executionResult.success) { logger.info(`[${requestId}] Workflow ${schedule.workflowId} executed successfully`) // Calculate the next run time based on the schedule configuration const nextRunAt = calculateNextRunTime(schedule, blocks) diff --git a/sim/app/api/telemetry/route.ts b/sim/app/api/telemetry/route.ts index dba8ec62a8..b00410cde1 100644 --- a/sim/app/api/telemetry/route.ts +++ b/sim/app/api/telemetry/route.ts @@ -123,13 +123,6 @@ async function forwardToCollector(data: any): Promise { }] } - // Safe debug log of the payload structure without sensitive data - logger.debug('Preparing to send telemetry payload', { - endpoint, - hasAttributes: safeAttrs.length > 0, - attributeCount: safeAttrs.length - }) - // Create explicit AbortController for timeout const controller = new AbortController() const timeoutId = setTimeout(() => controller.abort(), timeout) diff --git a/sim/app/api/workflows/[id]/execute/route.ts b/sim/app/api/workflows/[id]/execute/route.ts index 6e9fb07d08..ee32cfb964 100644 --- a/sim/app/api/workflows/[id]/execute/route.ts +++ b/sim/app/api/workflows/[id]/execute/route.ts @@ -236,13 +236,19 @@ async function executeWorkflow(workflow: any, requestId: string, input?: any) { const result = await executor.execute(workflowId) + // Check if we got a StreamingExecution result (with stream + execution properties) + // For API routes, we only care about the ExecutionResult part, not the stream + const executionResult = 'stream' in result && 'execution' in result + ? result.execution + : result + logger.info(`[${requestId}] Workflow execution completed: ${workflowId}`, { - success: result.success, - executionTime: result.metadata?.duration, + success: executionResult.success, + executionTime: executionResult.metadata?.duration, }) // Update workflow run counts if execution was successful - if (result.success) { + if (executionResult.success) { await updateWorkflowRunCounts(workflowId) // Track API call in user stats @@ -256,11 +262,11 @@ async function executeWorkflow(workflow: any, requestId: string, input?: any) { } // Build trace spans from execution logs - const { traceSpans, totalDuration } = buildTraceSpans(result) + const { traceSpans, totalDuration } = buildTraceSpans(executionResult) // Add trace spans to the execution result const enrichedResult = { - ...result, + ...executionResult, traceSpans, totalDuration, } @@ -268,7 +274,7 @@ async function executeWorkflow(workflow: any, requestId: string, input?: any) { // Log each execution step and the final result await persistExecutionLogs(workflowId, executionId, enrichedResult, 'api') - return result + return executionResult } catch (error: any) { logger.error(`[${requestId}] Workflow execution failed: ${workflowId}`, error) // Log the error diff --git a/sim/app/chat/[subdomain]/components/chat-client.tsx b/sim/app/chat/[subdomain]/components/chat-client.tsx index f8cd5c9a9c..9075f52262 100644 --- a/sim/app/chat/[subdomain]/components/chat-client.tsx +++ b/sim/app/chat/[subdomain]/components/chat-client.tsx @@ -409,69 +409,104 @@ export default function ChatClient({ subdomain }: { subdomain: string }) { throw new Error('Failed to get response') } - const responseData = await response.json() - console.log('Message response:', responseData) + // Detect streaming response via content-type (text/plain) or absence of JSON content-type + const contentType = response.headers.get('Content-Type') || '' - // Handle different response formats from API - if (responseData.multipleOutputs && responseData.contents && Array.isArray(responseData.contents)) { - // For multiple outputs, create separate assistant messages for each - const assistantMessages = responseData.contents.map((content: any) => { - // Format the content appropriately - let formattedContent = content - - // Convert objects to strings for display - if (typeof formattedContent === 'object' && formattedContent !== null) { - try { - formattedContent = JSON.stringify(formattedContent) - } catch (e) { - formattedContent = 'Received structured data response' - } - } - - return { - id: crypto.randomUUID(), - content: formattedContent || "No content found", - type: 'assistant' as const, + if (contentType.includes('text/plain')) { + // Handle streaming response + const messageId = crypto.randomUUID() + + // Add placeholder message + setMessages((prev) => [ + ...prev, + { + id: messageId, + content: '', + type: 'assistant', timestamp: new Date(), - } - }) - - // Add all messages at once - setMessages((prev) => [...prev, ...assistantMessages]) - } else { - // Handle single output as before - // Extract content from the response - could be in content or output - let messageContent = responseData.output + }, + ]) - // Handle different response formats from API - if (!messageContent && responseData.content) { - // Content could be an object or a string - if (typeof responseData.content === 'object') { - // If it's an object with a text property, use that - if (responseData.content.text) { - messageContent = responseData.content.text - } else { - // Try to convert to string for display - try { - messageContent = JSON.stringify(responseData.content) - } catch (e) { - messageContent = 'Received structured data response' + // Ensure the response body exists and is a ReadableStream + const reader = response.body?.getReader() + if (reader) { + const decoder = new TextDecoder() + let done = false + while (!done) { + const { value, done: readerDone } = await reader.read() + if (value) { + const chunk = decoder.decode(value, { stream: true }) + if (chunk) { + setMessages((prev) => + prev.map((msg) => + msg.id === messageId ? { ...msg, content: msg.content + chunk } : msg + ) + ) } } - } else { - // Direct string content - messageContent = responseData.content + done = readerDone } } + } else { + // Fallback to JSON response handling + const responseData = await response.json() + console.log('Message response:', responseData) - const assistantMessage: ChatMessage = { - id: crypto.randomUUID(), - content: messageContent || "Sorry, I couldn't process your request.", - type: 'assistant', - timestamp: new Date(), + // Handle different response formats from API + if (responseData.multipleOutputs && responseData.contents && Array.isArray(responseData.contents)) { + // For multiple outputs, create separate assistant messages for each + const assistantMessages = responseData.contents.map((content: any) => { + // Format the content appropriately + let formattedContent = content + + // Convert objects to strings for display + if (typeof formattedContent === 'object' && formattedContent !== null) { + try { + formattedContent = JSON.stringify(formattedContent) + } catch (e) { + formattedContent = 'Received structured data response' + } + } + + return { + id: crypto.randomUUID(), + content: formattedContent || 'No content found', + type: 'assistant' as const, + timestamp: new Date(), + } + }) + + // Add all messages at once + setMessages((prev) => [...prev, ...assistantMessages]) + } else { + // Handle single output as before + let messageContent = responseData.output + + if (!messageContent && responseData.content) { + if (typeof responseData.content === 'object') { + if (responseData.content.text) { + messageContent = responseData.content.text + } else { + try { + messageContent = JSON.stringify(responseData.content) + } catch (e) { + messageContent = 'Received structured data response' + } + } + } else { + messageContent = responseData.content + } + } + + const assistantMessage: ChatMessage = { + id: crypto.randomUUID(), + content: messageContent || "Sorry, I couldn't process your request.", + type: 'assistant', + timestamp: new Date(), + } + + setMessages((prev) => [...prev, assistantMessage]) } - - setMessages((prev) => [...prev, assistantMessage]) } } catch (error) { console.error('Error sending message:', error) diff --git a/sim/app/w/[id]/components/panel/components/chat/chat.tsx b/sim/app/w/[id]/components/panel/components/chat/chat.tsx index 61ceb744b2..912098bbcf 100644 --- a/sim/app/w/[id]/components/panel/components/chat/chat.tsx +++ b/sim/app/w/[id]/components/panel/components/chat/chat.tsx @@ -12,6 +12,9 @@ import { useWorkflowRegistry } from '@/stores/workflows/registry/store' import { useWorkflowExecution } from '../../../../hooks/use-workflow-execution' import { ChatMessage } from './components/chat-message/chat-message' import { OutputSelect } from './components/output-select/output-select' +import { BlockLog } from '@/executor/types' +import { calculateCost } from '@/providers/utils' +import { buildTraceSpans } from '@/lib/logs/trace-spans' interface ChatProps { panelWidth: number @@ -21,8 +24,14 @@ interface ChatProps { export function Chat({ panelWidth, chatMessage, setChatMessage }: ChatProps) { const { activeWorkflowId } = useWorkflowRegistry() - const { messages, addMessage, selectedWorkflowOutputs, setSelectedWorkflowOutput } = - useChatStore() + const { + messages, + addMessage, + selectedWorkflowOutputs, + setSelectedWorkflowOutput, + appendMessageContent, + finalizeMessageStream + } = useChatStore() const { entries } = useConsoleStore() const messagesEndRef = useRef(null) @@ -93,8 +102,192 @@ export function Chat({ panelWidth, chatMessage, setChatMessage }: ChatProps) { setChatMessage('') // Execute the workflow to generate a response, passing the chat message as input - // The workflow execution will trigger block executions which will add messages to the chat via the console store - await handleRunWorkflow({ input: sentMessage }) + const result = await handleRunWorkflow({ input: sentMessage }) + + // Check if we got a streaming response + if (result && 'stream' in result && result.stream instanceof ReadableStream) { + // Generate a unique ID for the message + const messageId = crypto.randomUUID() + + // Create a content buffer to collect initial content + let initialContent = '' + let fullContent = '' // Store the complete content for updating logs later + let hasAddedMessage = false + let executionResult = (result as any).execution // Store the execution result with type assertion + + try { + // Process the stream + const reader = result.stream.getReader() + const decoder = new TextDecoder() + + console.log("Starting to read from stream") + + while (true) { + try { + const { done, value } = await reader.read() + if (done) { + console.log("Stream complete") + break + } + + // Decode and append chunk + const chunk = decoder.decode(value, { stream: true }) // Use stream option + + if (chunk) { + initialContent += chunk + fullContent += chunk + + // Only add the message to UI once we have some actual content to show + if (!hasAddedMessage && initialContent.trim().length > 0) { + // Add message with initial content - cast to any to bypass type checking for id + addMessage({ + content: initialContent, + workflowId: activeWorkflowId, + type: 'workflow', + isStreaming: true, + id: messageId + } as any) + hasAddedMessage = true + } else if (hasAddedMessage) { + // Append to existing message + appendMessageContent(messageId, chunk) + } + } + } catch (streamError) { + console.error('Error reading from stream:', streamError) + // Break the loop on error + break + } + } + + // If we never added a message (no content received), add it now + if (!hasAddedMessage && initialContent.trim().length > 0) { + addMessage({ + content: initialContent, + workflowId: activeWorkflowId, + type: 'workflow', + id: messageId + } as any) + } + + // Update logs with the full streaming content if available + if (executionResult && fullContent.trim().length > 0) { + try { + // Format the final content properly to match what's shown for manual executions + // Include all the markdown and formatting from the streamed response + const formattedContent = fullContent + + // Calculate cost based on token usage if available + let costData = undefined + + if (executionResult.output?.response?.tokens) { + const tokens = executionResult.output.response.tokens + const model = executionResult.output?.response?.model || 'gpt-4o' + const cost = calculateCost( + model, + tokens.prompt || 0, + tokens.completion || 0, + false // Don't use cached input for chat responses + ) + costData = { ...cost, model } as any + } + + // Build trace spans and total duration before persisting + const { traceSpans, totalDuration } = buildTraceSpans(executionResult as any) + + // Create a completed execution ID + const completedExecutionId = executionResult.metadata?.executionId || crypto.randomUUID() + + // Import the workflow execution hook for direct access to the workflow service + const workflowExecutionApi = await fetch(`/api/workflows/${activeWorkflowId}/log`, { + method: 'POST', + headers: { + 'Content-Type': 'application/json', + }, + body: JSON.stringify({ + executionId: completedExecutionId, + result: { + ...executionResult, + output: { + ...executionResult.output, + response: { + ...executionResult.output?.response, + content: formattedContent, + model: executionResult.output?.response?.model, + tokens: executionResult.output?.response?.tokens, + toolCalls: executionResult.output?.response?.toolCalls, + providerTiming: executionResult.output?.response?.providerTiming, + cost: costData || executionResult.output?.response?.cost, + } + }, + cost: costData, + // Update the message to include the formatted content + logs: (executionResult.logs || []).map((log: BlockLog) => { + // Check if this is the streaming block by comparing with the selected output IDs + // Selected output IDs typically include the block ID we are streaming from + const isStreamingBlock = selectedOutputs.some(outputId => + outputId === log.blockId || outputId.startsWith(`${log.blockId}_`) + ) + + if (isStreamingBlock && log.blockType === 'agent' && log.output?.response) { + return { + ...log, + output: { + ...log.output, + response: { + ...log.output.response, + content: formattedContent, + providerTiming: log.output.response.providerTiming, + cost: costData || log.output.response.cost, + } + } + } + } + return log + }), + metadata: { + ...executionResult.metadata, + source: 'chat', + completedAt: new Date().toISOString(), + isStreamingComplete: true, + cost: costData || executionResult.metadata?.cost, + providerTiming: executionResult.output?.response?.providerTiming, + }, + traceSpans: traceSpans, + totalDuration: totalDuration, + } + }), + }) + + if (!workflowExecutionApi.ok) { + console.error('Failed to log complete streaming execution') + } + } catch (logError) { + console.error('Error logging complete streaming execution:', logError) + } + } + } catch (error) { + console.error('Error processing stream:', error) + + // If there's an error and we haven't added a message yet, add an error message + if (!hasAddedMessage) { + addMessage({ + content: "Error: Failed to process the streaming response.", + workflowId: activeWorkflowId, + type: 'workflow', + id: messageId + } as any) + } else { + // Otherwise append the error to the existing message + appendMessageContent(messageId, "\n\nError: Failed to process the streaming response.") + } + } finally { + console.log("Finalizing stream") + if (hasAddedMessage) { + finalizeMessageStream(messageId) + } + } + } } // Handle key press diff --git a/sim/app/w/[id]/components/panel/components/chat/components/chat-message/chat-message.tsx b/sim/app/w/[id]/components/panel/components/chat/components/chat-message/chat-message.tsx index 775a2a6c4a..f02059836c 100644 --- a/sim/app/w/[id]/components/panel/components/chat/components/chat-message/chat-message.tsx +++ b/sim/app/w/[id]/components/panel/components/chat/components/chat-message/chat-message.tsx @@ -1,6 +1,6 @@ import { useMemo } from 'react' -import { format, formatDistanceToNow } from 'date-fns' -import { Clock, Terminal, User } from 'lucide-react' +import { formatDistanceToNow } from 'date-fns' +import { Clock } from 'lucide-react' import { JSONView } from '../../../console/components/json-view/json-view' interface ChatMessageProps { @@ -9,6 +9,7 @@ interface ChatMessageProps { content: any timestamp: string | Date type: 'user' | 'workflow' + isStreaming?: boolean } containerWidth: number } @@ -66,7 +67,7 @@ export function ChatMessage({ message, containerWidth }: ChatMessageProps) { return JSON.stringify(message.content) // Return stringified version for type safety } - return String(message.content) + return String(message.content || '') }, [message.content, isJsonObject]) return ( @@ -79,6 +80,9 @@ export function ChatMessage({ message, containerWidth }: ChatMessageProps) {
{message.type !== 'user' && Workflow} + {message.isStreaming && ( + ••• + )}
@@ -89,6 +93,9 @@ export function ChatMessage({ message, containerWidth }: ChatMessageProps) { ) : (
+ {message.isStreaming && ( + + )}
)} diff --git a/sim/app/w/[id]/hooks/use-workflow-execution.ts b/sim/app/w/[id]/hooks/use-workflow-execution.ts index 405664b8cb..a23cad4b9f 100644 --- a/sim/app/w/[id]/hooks/use-workflow-execution.ts +++ b/sim/app/w/[id]/hooks/use-workflow-execution.ts @@ -42,7 +42,7 @@ export function useWorkflowExecution() { } = useExecutionStore() const [executionResult, setExecutionResult] = useState(null) - const persistLogs = async (executionId: string, result: ExecutionResult) => { + const persistLogs = async (executionId: string, result: ExecutionResult, streamContent?: string) => { try { // Build trace spans from execution logs const { traceSpans, totalDuration } = buildTraceSpans(result) @@ -54,6 +54,26 @@ export function useWorkflowExecution() { totalDuration, } + // If this was a streaming response and we have the final content, update it + if (streamContent && result.output?.response && typeof streamContent === 'string') { + // Update the content with the final streaming content + enrichedResult.output.response.content = streamContent + + // Also update any block logs to include the content where appropriate + if (enrichedResult.logs) { + // Get the streaming block ID from metadata if available + const streamingBlockId = (result.metadata as any)?.streamingBlockId || null; + + for (const log of enrichedResult.logs) { + // Only update the specific agent block that was streamed + const isStreamingBlock = streamingBlockId && log.blockId === streamingBlockId; + if (isStreamingBlock && log.blockType === 'agent' && log.output?.response) { + log.output.response.content = streamContent + } + } + } + } + const response = await fetch(`/api/workflows/${activeWorkflowId}/log`, { method: 'POST', headers: { @@ -68,8 +88,11 @@ export function useWorkflowExecution() { if (!response.ok) { throw new Error('Failed to persist logs') } + + return executionId } catch (error) { logger.error('Error persisting logs:', { error }) + return executionId } } @@ -104,6 +127,15 @@ export function useWorkflowExecution() { const isChatExecution = activeTab === 'chat' && (workflowInput && typeof workflowInput === 'object' && 'input' in workflowInput) + // If this is a chat execution, get the selected outputs + let selectedOutputIds: string[] | undefined = undefined + if (isChatExecution && activeWorkflowId) { + // Get selected outputs from chat store + const chatStore = await import('@/stores/panel/chat/store').then(mod => mod.useChatStore) + selectedOutputIds = chatStore.getState().getSelectedWorkflowOutput(activeWorkflowId) + logger.info('Chat execution with selected outputs:', selectedOutputIds) + } + try { // Clear any existing state setDebugContext(null) @@ -147,26 +179,116 @@ export function useWorkflowExecution() { // Create serialized workflow const workflow = new Serializer().serializeWorkflow(mergedStates, edges, loops) - // Create executor and store in global state - const newExecutor = new Executor( + // Create executor options with streaming support for chat + const executorOptions: any = { + // Default executor options workflow, currentBlockStates, envVarValues, workflowInput, - workflowVariables - ) + workflowVariables, + } + + // Add streaming context for chat executions + if (isChatExecution && selectedOutputIds && selectedOutputIds.length > 0) { + executorOptions.contextExtensions = { + stream: true, + selectedOutputIds, + edges: workflow.connections.map(conn => ({ + source: conn.source, + target: conn.target + })) + } + } + + // Create executor and store in global state + const newExecutor = new Executor(executorOptions) setExecutor(newExecutor) // Execute workflow const result = await newExecutor.execute(activeWorkflowId) + // Streaming results are handled differently - they won't have a standard result + if (result instanceof ReadableStream) { + logger.info('Received streaming result from executor') + + // For streaming results, we need to handle them in the component + // that initiated the execution (chat panel) + return { + success: true, + stream: result, + } + } + + // Handle StreamingExecution format (combined stream + execution result) + if (result && typeof result === 'object' && 'stream' in result && 'execution' in result) { + logger.info('Received combined stream+execution result from executor') + + // Generate an executionId and store it in the execution metadata so that + // the chat component can persist the logs *after* the stream finishes. + const executionId = uuidv4() + + // Determine which block is streaming - typically the one that matches a selected output ID + let streamingBlockId = null; + if (selectedOutputIds && selectedOutputIds.length > 0 && result.execution.logs) { + // Find the agent block in the logs that matches one of our selected outputs + const streamingBlock = result.execution.logs.find(log => + log.blockType === 'agent' && selectedOutputIds.some(id => id === log.blockId || id.startsWith(`${log.blockId}_`)) + ); + if (streamingBlock) { + streamingBlockId = streamingBlock.blockId; + logger.info(`Identified streaming block: ${streamingBlockId}`); + } + } + + // Attach streaming / source metadata and the newly generated executionId + result.execution.metadata = { + ...(result.execution.metadata || {}), + executionId, + source: isChatExecution ? 'chat' : 'manual', + streamingBlockId, // Add the block ID to the metadata + } as any + + // Clean up any response objects with zero tokens in agent blocks to avoid confusion in console + if (result.execution.logs && Array.isArray(result.execution.logs)) { + result.execution.logs.forEach((log: any) => { + if (log.blockType === 'agent' && log.output?.response) { + const response = log.output.response; + + // Check for zero tokens that will be estimated later + if (response.tokens && + (!response.tokens.completion || response.tokens.completion === 0) && + (!response.toolCalls || !response.toolCalls.list || response.toolCalls.list.length === 0)) { + + // Remove tokens from console display to avoid confusion + // They'll be properly estimated in the execution logger + delete response.tokens; + } + } + }); + } + + // Mark the execution as streaming so that downstream code can recognise it + (result.execution as any).isStreaming = true + + // Return both the stream and the execution object so the caller (chat panel) + // can collect the full content and then persist the logs in one go. + // Also include processingPromise if available to ensure token counts are final + return { + success: true, + stream: result.stream, + execution: result.execution, + processingPromise: (result as any).processingPromise + } + } + // Add metadata about source being chat if applicable if (isChatExecution) { // Use type assertion for adding custom metadata (result as any).metadata = { ...(result.metadata || {}), source: 'chat' - }; + } } // If we're in debug mode, store the execution context for later steps @@ -204,6 +326,8 @@ export function useWorkflowExecution() { logger.error('Error persisting logs:', { error: err }) }) } + + return result } catch (error: any) { logger.error('Workflow Execution Error:', error) @@ -292,6 +416,8 @@ export function useWorkflowExecution() { persistLogs(executionId, errorResult).catch((err) => { logger.error('Error persisting logs:', { error: err }) }) + + return errorResult } }, [ activeWorkflowId, diff --git a/sim/app/w/components/sidebar/components/settings-modal/settings-modal.tsx b/sim/app/w/components/sidebar/components/settings-modal/settings-modal.tsx index e42a1d308b..82122a4456 100644 --- a/sim/app/w/components/sidebar/components/settings-modal/settings-modal.tsx +++ b/sim/app/w/components/sidebar/components/settings-modal/settings-modal.tsx @@ -158,9 +158,6 @@ export function SettingsModal({ open, onOpenChange }: SettingsModalProps) {
-
- -
{isSubscriptionEnabled && (
)} +
+ +
diff --git a/sim/executor/handlers/agent/agent-handler.test.ts b/sim/executor/handlers/agent/agent-handler.test.ts index f361ca3daf..4e5067cc50 100644 --- a/sim/executor/handlers/agent/agent-handler.test.ts +++ b/sim/executor/handlers/agent/agent-handler.test.ts @@ -4,7 +4,7 @@ import { getAllBlocks } from '@/blocks' import { getProviderFromModel, transformBlockTool } from '@/providers/utils' import { SerializedBlock, SerializedWorkflow } from '@/serializer/types' import { executeTool } from '@/tools' -import { ExecutionContext } from '../../types' +import { ExecutionContext, StreamingExecution } from '../../types' import { AgentBlockHandler } from './agent-handler' process.env.NEXT_PUBLIC_APP_URL = 'http://localhost:3000' @@ -88,10 +88,16 @@ describe('AgentBlockHandler', () => { mockIsHosted.mockReturnValue(false) // Default to non-hosted env for tests mockGetProviderFromModel.mockReturnValue('mock-provider') - // Set up fetch mock to return a successful response mockFetch.mockImplementation(() => { return Promise.resolve({ ok: true, + headers: { + get: (name: string) => { + if (name === 'Content-Type') return 'application/json' + if (name === 'X-Execution-Data') return null + return null + } + }, json: () => Promise.resolve({ content: 'Mocked response content', @@ -112,7 +118,6 @@ describe('AgentBlockHandler', () => { })) mockGetAllBlocks.mockReturnValue([]) - // Set up executeTool mock mockExecuteTool.mockImplementation((toolId, params) => { if (toolId === 'function_execute') { return Promise.resolve({ @@ -194,12 +199,9 @@ describe('AgentBlockHandler', () => { }) it('should preserve executeFunction for custom tools with different usageControl settings', async () => { - // Set up a spy for Promise.all to capture the tools array before it's serialized let capturedTools: any[] = [] - // Mock Promise.all to capture tools Promise.all = vi.fn().mockImplementation((promises: Promise[]) => { - // Store result of the original Promise.all const result = originalPromiseAll.call(Promise, promises) // When result resolves, capture the tools @@ -212,10 +214,16 @@ describe('AgentBlockHandler', () => { return result }) - // Configure response with tool calls mockFetch.mockImplementationOnce(() => { return Promise.resolve({ ok: true, + headers: { + get: (name: string) => { + if (name === 'Content-Type') return 'application/json' + if (name === 'X-Execution-Data') return null + return null + } + }, json: () => Promise.resolve({ content: 'Using tools to respond', @@ -303,34 +311,26 @@ describe('AgentBlockHandler', () => { mockGetProviderFromModel.mockReturnValue('openai') - // Execute with the tools await handler.execute(mockBlock, inputs, mockContext) - // Verify Promise.all was called (tools were processed) expect(Promise.all).toHaveBeenCalled() - // Verify that the none tool was filtered out expect(capturedTools.length).toBe(2) - // Find the tools by name const autoTool = capturedTools.find((t) => t.name === 'auto_tool') const forceTool = capturedTools.find((t) => t.name === 'force_tool') const noneTool = capturedTools.find((t) => t.name === 'none_tool') - // Verify that auto and force tools are included expect(autoTool).toBeDefined() expect(forceTool).toBeDefined() - expect(noneTool).toBeUndefined() // None tool shouldn't be included + expect(noneTool).toBeUndefined() - // Verify usageControl properties expect(autoTool.usageControl).toBe('auto') expect(forceTool.usageControl).toBe('force') - // Verify that the executeFunction property exists on both tools expect(typeof autoTool.executeFunction).toBe('function') expect(typeof forceTool.executeFunction).toBe('function') - // Test that executeFunction can be called const autoResult = await autoTool.executeFunction({ input: 'test input' }) expect(mockExecuteTool).toHaveBeenCalledWith( 'function_execute', @@ -349,15 +349,10 @@ describe('AgentBlockHandler', () => { }) ) - // Extract the request body from the fetch call to verify serialized tools const fetchCall = mockFetch.mock.calls[0] const requestBody = JSON.parse(fetchCall[1].body) - // Verify that only two tools were passed to the API expect(requestBody.tools.length).toBe(2) - - // Note: executeFunction won't be in the serialized tools since functions aren't serializable - // But we've verified above that they exist before serialization }) it('should filter out tools with usageControl set to "none"', async () => { @@ -392,17 +387,13 @@ describe('AgentBlockHandler', () => { mockGetProviderFromModel.mockReturnValue('openai') - // Execute the handler await handler.execute(mockBlock, inputs, mockContext) - // Extract the actual request from the fetch call const fetchCall = mockFetch.mock.calls[0] const requestBody = JSON.parse(fetchCall[1].body) - // Verify that only two tools were passed (the ones with auto and force settings) expect(requestBody.tools.length).toBe(2) - // Check that the filtered tools are the right ones const toolIds = requestBody.tools.map((t: any) => t.id) expect(toolIds).toContain('transformed_tool_1') expect(toolIds).toContain('transformed_tool_3') @@ -432,7 +423,6 @@ describe('AgentBlockHandler', () => { ], } - // Custom implementation to preserve the usageControl property mockTransformBlockTool.mockImplementation((tool: any) => ({ id: `transformed_${tool.id}`, name: `${tool.id}_${tool.operation}`, @@ -442,14 +432,11 @@ describe('AgentBlockHandler', () => { mockGetProviderFromModel.mockReturnValue('openai') - // Execute the handler await handler.execute(mockBlock, inputs, mockContext) - // Extract the actual request from the fetch call const fetchCall = mockFetch.mock.calls[0] const requestBody = JSON.parse(fetchCall[1].body) - // Verify that tools have the usageControl property expect(requestBody.tools[0].usageControl).toBe('auto') expect(requestBody.tools[1].usageControl).toBe('force') }) @@ -510,23 +497,18 @@ describe('AgentBlockHandler', () => { mockGetProviderFromModel.mockReturnValue('openai') - // Execute the handler await handler.execute(mockBlock, inputs, mockContext) - // Extract the actual request from the fetch call const fetchCall = mockFetch.mock.calls[0] const requestBody = JSON.parse(fetchCall[1].body) - // Verify that only two custom tools were passed (auto and force) expect(requestBody.tools.length).toBe(2) - // Check the tools by name const toolNames = requestBody.tools.map((t: any) => t.name) expect(toolNames).toContain('custom_tool_auto') expect(toolNames).toContain('custom_tool_force') expect(toolNames).not.toContain('custom_tool_none') - // Verify usageControl properties const autoTool = requestBody.tools.find((t: any) => t.name === 'custom_tool_auto') const forceTool = requestBody.tools.find((t: any) => t.name === 'custom_tool_force') @@ -535,7 +517,6 @@ describe('AgentBlockHandler', () => { }) it('should not require API key for gpt-4o on hosted version', async () => { - // Mock hosted environment mockIsHosted.mockReturnValue(true) const inputs = { @@ -544,7 +525,6 @@ describe('AgentBlockHandler', () => { context: 'User query: Hello!', temperature: 0.7, maxTokens: 100, - // No API key provided - this will be handled server-side } mockGetProviderFromModel.mockReturnValue('openai') @@ -560,10 +540,8 @@ describe('AgentBlockHandler', () => { responseFormat: undefined, } - // Execute should work even without API key await handler.execute(mockBlock, inputs, mockContext) - // Verify the proxy was called with the right parameters expect(mockFetch).toHaveBeenCalledWith(expect.any(String), expect.any(Object)) }) @@ -577,7 +555,6 @@ describe('AgentBlockHandler', () => { id: 'block_tool_1', title: 'Data Analysis Tool', operation: 'analyze', - // Assume transformBlockTool resolves this based on blocks/tools }, ], } @@ -669,18 +646,22 @@ describe('AgentBlockHandler', () => { mockGetProviderFromModel.mockReturnValue('openai') - // Process the tools to see what they'll be transformed into await handler.execute(mockBlock, inputs, mockContext) - // Verify that mockExecuteProviderRequest was called expect(mockFetch).toHaveBeenCalledWith(expect.any(String), expect.any(Object)) }) it('should handle responseFormat with valid JSON', async () => { - // Create a special mock for this test only mockFetch.mockImplementationOnce(() => { return Promise.resolve({ ok: true, + headers: { + get: (name: string) => { + if (name === 'Content-Type') return 'application/json' + if (name === 'X-Execution-Data') return null + return null + } + }, json: () => Promise.resolve({ content: '{"result": "Success", "score": 0.95}', @@ -712,10 +693,16 @@ describe('AgentBlockHandler', () => { }) it('should handle responseFormat when it is an empty string', async () => { - // Create a special mock for this test only mockFetch.mockImplementationOnce(() => { return Promise.resolve({ ok: true, + headers: { + get: (name: string) => { + if (name === 'Content-Type') return 'application/json' + if (name === 'X-Execution-Data') return null + return null + } + }, json: () => Promise.resolve({ content: 'Regular text response', @@ -773,5 +760,164 @@ describe('AgentBlockHandler', () => { 'Provider API Error' ) }) + + it('should handle streaming responses with text/event-stream content type', async () => { + const mockStreamBody = { + getReader: vi.fn().mockReturnValue({ + read: vi.fn().mockResolvedValue({ done: true, value: undefined }), + }), + } + + mockFetch.mockImplementationOnce(() => { + return Promise.resolve({ + ok: true, + headers: { + get: (name: string) => { + if (name === 'Content-Type') return 'text/event-stream' + if (name === 'X-Execution-Data') return null + return null + } + }, + body: mockStreamBody, + }) + }) + + const inputs = { + model: 'gpt-4o', + context: 'Stream this response.', + apiKey: 'test-api-key', + stream: true, + } + + mockContext.stream = true + mockContext.selectedOutputIds = [mockBlock.id] + + const result = await handler.execute(mockBlock, inputs, mockContext) + + expect(result).toHaveProperty('stream') + expect(result).toHaveProperty('execution') + + expect((result as StreamingExecution).execution).toHaveProperty('success', true) + expect((result as StreamingExecution).execution).toHaveProperty('output') + expect((result as StreamingExecution).execution.output).toHaveProperty('response') + expect((result as StreamingExecution).execution).toHaveProperty('logs') + }) + + it('should handle streaming responses with execution data in header', async () => { + const mockStreamBody = { + getReader: vi.fn().mockReturnValue({ + read: vi.fn().mockResolvedValue({ done: true, value: undefined }), + }), + } + + const mockExecutionData = { + success: true, + output: { + response: { + content: '', + model: 'mock-model', + tokens: { prompt: 10, completion: 20, total: 30 }, + } + }, + logs: [ + { blockId: 'some-id', blockType: 'agent', startedAt: new Date().toISOString(), endedAt: new Date().toISOString(), durationMs: 100, success: true } + ], + metadata: { + startTime: new Date().toISOString(), + duration: 100, + } + } + + mockFetch.mockImplementationOnce(() => { + return Promise.resolve({ + ok: true, + headers: { + get: (name: string) => { + if (name === 'Content-Type') return 'text/event-stream' + if (name === 'X-Execution-Data') return JSON.stringify(mockExecutionData) + return null + } + }, + body: mockStreamBody, + }) + }) + + const inputs = { + model: 'gpt-4o', + context: 'Stream this response with execution data.', + apiKey: 'test-api-key', + stream: true, + } + + mockContext.stream = true + mockContext.selectedOutputIds = [mockBlock.id] + + const result = await handler.execute(mockBlock, inputs, mockContext) + + expect(result).toHaveProperty('stream') + expect(result).toHaveProperty('execution') + + expect((result as StreamingExecution).execution.success).toBe(true) + expect((result as StreamingExecution).execution.output.response.model).toBe('mock-model') + const logs = (result as StreamingExecution).execution.logs + expect(logs?.length).toBe(1) + if (logs && logs.length > 0 && logs[0]) { + expect(logs[0].blockType).toBe('agent') + } + }) + + it('should handle combined stream+execution responses', async () => { + const mockStreamObj = new ReadableStream({ + start(controller) { + controller.close() + } + }) + + mockFetch.mockImplementationOnce(() => { + return Promise.resolve({ + ok: true, + headers: { + get: (name: string) => name === 'Content-Type' ? 'application/json' : null + }, + json: () => Promise.resolve({ + stream: {}, // Serialized stream placeholder + execution: { + success: true, + output: { + response: { + content: 'Test streaming content', + model: 'gpt-4o', + tokens: { prompt: 10, completion: 5, total: 15 }, + } + }, + logs: [], + metadata: { + startTime: new Date().toISOString(), + duration: 150 + } + } + }) + }) + }) + + const inputs = { + model: 'gpt-4o', + context: 'Return a combined response.', + apiKey: 'test-api-key', + stream: true, + } + + mockContext.stream = true + mockContext.selectedOutputIds = [mockBlock.id] + + const result = await handler.execute(mockBlock, inputs, mockContext) + + expect(result).toHaveProperty('stream') + expect(result).toHaveProperty('execution') + + expect((result as StreamingExecution).execution.success).toBe(true) + expect((result as StreamingExecution).execution.output.response.content).toBe('Test streaming content') + expect((result as StreamingExecution).execution.output.response.model).toBe('gpt-4o') + }) }) }) diff --git a/sim/executor/handlers/agent/agent-handler.ts b/sim/executor/handlers/agent/agent-handler.ts index ede97d7d4a..184d1bd4a2 100644 --- a/sim/executor/handlers/agent/agent-handler.ts +++ b/sim/executor/handlers/agent/agent-handler.ts @@ -5,7 +5,7 @@ import { getProviderFromModel, transformBlockTool } from '@/providers/utils' import { SerializedBlock } from '@/serializer/types' import { executeTool } from '@/tools' import { getToolAsync, getTool } from '@/tools/utils' -import { BlockHandler, ExecutionContext } from '../../types' +import { BlockHandler, ExecutionContext, StreamingExecution } from '../../types' const logger = createLogger('AgentBlockHandler') @@ -21,14 +21,9 @@ export class AgentBlockHandler implements BlockHandler { block: SerializedBlock, inputs: Record, context: ExecutionContext - ): Promise { + ): Promise { logger.info(`Executing agent block: ${block.id}`) - // Check for null values and try to resolve from environment variables - const nullInputs = Object.entries(inputs) - .filter(([_, value]) => value === null) - .map(([key]) => key) - // Parse response format if provided let responseFormat: any = undefined if (inputs.responseFormat) { @@ -155,6 +150,39 @@ export class AgentBlockHandler implements BlockHandler { ) ).filter((t: any): t is NonNullable => t !== null) : [] + + // Check if streaming is requested and this block is selected for streaming + const isBlockSelectedForOutput = context.selectedOutputIds?.some(outputId => { + // First check for direct match (if the entire outputId is the blockId) + if (outputId === block.id) { + logger.info(`Direct match found for block ${block.id} in selected outputs`) + return true + } + + // Then try parsing the blockId from the blockId_path format + const firstUnderscoreIndex = outputId.indexOf('_') + if (firstUnderscoreIndex !== -1) { + const blockId = outputId.substring(0, firstUnderscoreIndex) + const isMatch = blockId === block.id + if (isMatch) { + logger.info(`Path match found for block ${block.id} in selected outputs (from ${outputId})`) + } + return isMatch + } + return false + }) ?? false + + // Check if this block has any outgoing connections + const hasOutgoingConnections = context.edges?.some(edge => edge.source === block.id) ?? false + + // Determine if we should use streaming for this block + const shouldUseStreaming = context.stream && + isBlockSelectedForOutput && + !hasOutgoingConnections + + if (shouldUseStreaming) { + logger.info(`Block ${block.id} will use streaming response (selected for output with no outgoing connections)`) + } // Debug request before sending to provider const providerRequest = { @@ -172,6 +200,7 @@ export class AgentBlockHandler implements BlockHandler { apiKey: inputs.apiKey, responseFormat, workflowId: context.workflowId, + stream: shouldUseStreaming, } logger.info(`Provider request prepared`, { @@ -181,6 +210,9 @@ export class AgentBlockHandler implements BlockHandler { hasTools: !!providerRequest.tools, hasApiKey: !!providerRequest.apiKey, workflowId: providerRequest.workflowId, + stream: shouldUseStreaming, + isBlockSelectedForOutput, + hasOutgoingConnections, }) const baseUrl = process.env.NEXT_PUBLIC_APP_URL || '' @@ -209,7 +241,102 @@ export class AgentBlockHandler implements BlockHandler { throw new Error(errorMessage) } + // Check if we're getting a streaming response + const contentType = response.headers.get('Content-Type') + if (contentType?.includes('text/event-stream')) { + logger.info(`Received streaming response for block ${block.id}`) + + // Ensure we have a valid body stream + if (!response.body) { + throw new Error(`No response body in streaming response for block ${block.id}`) + } + + // Check if we have execution data in the header + const executionDataHeader = response.headers.get('X-Execution-Data') + if (executionDataHeader) { + try { + // Parse the execution data from the header + const executionData = JSON.parse(executionDataHeader) + + // Add block-specific data to the execution logs if needed + if (executionData && executionData.logs) { + for (const log of executionData.logs) { + if (!log.blockId) log.blockId = block.id + if (!log.blockName && block.metadata?.name) log.blockName = block.metadata.name + if (!log.blockType && block.metadata?.id) log.blockType = block.metadata.id + } + } + + // Add block metadata to the execution data if missing + if (executionData.output?.response) { + // Ensure model and block info is set + if (block.metadata?.name && !executionData.blockName) { + executionData.blockName = block.metadata.name + } + if (block.metadata?.id && !executionData.blockType) { + executionData.blockType = block.metadata.id + } + if (!executionData.blockId) { + executionData.blockId = block.id + } + + // Add explicit streaming flag to make it easier to identify streaming executions + executionData.isStreaming = true + } + + // Return both the stream and the execution data as separate properties + const streamingExecution: StreamingExecution = { + stream: response.body, + execution: executionData + } + return streamingExecution + } catch (error) { + logger.error(`Error parsing execution data header: ${error}`) + // Continue with just the stream if there's an error + } + } + + // No execution data in header, just return the stream + // Create a minimal StreamingExecution with empty execution data + const minimalExecution: StreamingExecution = { + stream: response.body, + execution: { + success: true, + output: { response: {} }, + logs: [], + metadata: { + duration: 0, + startTime: new Date().toISOString() + } + } + } + return minimalExecution + } + + // Check if we have a combined response with both stream and execution data const result = await response.json() + + if (result && typeof result === 'object' && 'stream' in result && 'execution' in result) { + logger.info(`Received combined streaming response for block ${block.id}`) + + // Get the stream as a ReadableStream (need to convert from serialized format) + const stream = new ReadableStream({ + start(controller) { + // Since stream was serialized as JSON, we need to reconstruct it + // For now, we'll just use a placeholder message + const encoder = new TextEncoder() + controller.enqueue(encoder.encode('Stream data cannot be serialized as JSON. You will need to return a proper stream.')) + controller.close() + } + }) + + // Return both in a format the executor can handle + const streamingExecution: StreamingExecution = { + stream, + execution: result.execution + } + return streamingExecution + } logger.info(`Provider response received`, { contentLength: result.content ? result.content.length : 0, diff --git a/sim/executor/index.ts b/sim/executor/index.ts index 0d8a17f9b6..e0c3adf77a 100644 --- a/sim/executor/index.ts +++ b/sim/executor/index.ts @@ -22,6 +22,7 @@ import { ExecutionContext, ExecutionResult, NormalizedBlockOutput, + StreamingExecution, } from './types' const logger = createLogger('Executor') @@ -58,31 +59,71 @@ export class Executor { private blockHandlers: BlockHandler[] private workflowInput: any private isDebugging: boolean = false + private contextExtensions: any = {} + private actualWorkflow: SerializedWorkflow constructor( - private workflow: SerializedWorkflow, + private workflowParam: SerializedWorkflow | { + workflow: SerializedWorkflow, + currentBlockStates?: Record, + envVarValues?: Record, + workflowInput?: any, + workflowVariables?: Record, + contextExtensions?: { + stream?: boolean, + selectedOutputIds?: string[], + edges?: Array<{source: string, target: string}> + } + }, private initialBlockStates: Record = {}, private environmentVariables: Record = {}, workflowInput?: any, private workflowVariables: Record = {} ) { - this.validateWorkflow() - - if (workflowInput) { - this.workflowInput = workflowInput - logger.info('[Executor] Using workflow input:', JSON.stringify(this.workflowInput, null, 2)) + // Handle new constructor format with options object + if (typeof workflowParam === 'object' && 'workflow' in workflowParam) { + const options = workflowParam + this.actualWorkflow = options.workflow + this.initialBlockStates = options.currentBlockStates || {} + this.environmentVariables = options.envVarValues || {} + this.workflowInput = options.workflowInput || {} + this.workflowVariables = options.workflowVariables || {} + + // Store context extensions for streaming and output selection + if (options.contextExtensions) { + this.contextExtensions = options.contextExtensions + + if (this.contextExtensions.stream) { + logger.info('Executor initialized with streaming enabled', { + hasSelectedOutputIds: Array.isArray(this.contextExtensions.selectedOutputIds), + selectedOutputCount: Array.isArray(this.contextExtensions.selectedOutputIds) + ? this.contextExtensions.selectedOutputIds.length + : 0, + selectedOutputIds: this.contextExtensions.selectedOutputIds || [], + }) + } + } } else { - this.workflowInput = {} + this.actualWorkflow = workflowParam + + if (workflowInput) { + this.workflowInput = workflowInput + logger.info('[Executor] Using workflow input:', JSON.stringify(this.workflowInput, null, 2)) + } else { + this.workflowInput = {} + } } - this.loopManager = new LoopManager(workflow.loops || {}) + this.validateWorkflow() + + this.loopManager = new LoopManager(this.actualWorkflow.loops || {}) this.resolver = new InputResolver( - workflow, - environmentVariables, - workflowVariables, + this.actualWorkflow, + this.environmentVariables, + this.workflowVariables, this.loopManager ) - this.pathTracker = new PathTracker(workflow) + this.pathTracker = new PathTracker(this.actualWorkflow) this.blockHandlers = [ new AgentBlockHandler(), @@ -101,9 +142,9 @@ export class Executor { * Executes the workflow and returns the result. * * @param workflowId - Unique identifier for the workflow execution - * @returns Execution result containing output, logs, and metadata + * @returns Execution result containing output, logs, and metadata, or a stream, or combined execution and stream */ - async execute(workflowId: string): Promise { + async execute(workflowId: string): Promise { const { setIsExecuting, setIsDebugging, setPendingBlocks, reset } = useExecutionStore.getState() const startTime = new Date() let finalOutput: NormalizedBlockOutput = { response: {} } @@ -111,8 +152,8 @@ export class Executor { // Track workflow execution start trackWorkflowTelemetry('workflow_execution_started', { workflowId, - blockCount: this.workflow.blocks.length, - connectionCount: this.workflow.connections.length, + blockCount: this.actualWorkflow.blocks.length, + connectionCount: this.actualWorkflow.connections.length, startTime: startTime.toISOString() }) @@ -153,7 +194,7 @@ export class Executor { pendingBlocks: nextLayer, isDebugSession: true, context: context, // Include context for resumption - workflowConnections: this.workflow.connections.map((conn) => ({ + workflowConnections: this.actualWorkflow.connections.map((conn: any) => ({ source: conn.source, target: conn.target, })), @@ -167,9 +208,162 @@ export class Executor { hasMoreLayers = false } else { const outputs = await this.executeLayer(nextLayer, context) + + // Check if we got a StreamingExecution response from any block + const streamingOutput = outputs.find(output => + typeof output === 'object' && output !== null && + 'stream' in output && 'execution' in output + ) + + if (streamingOutput) { + // This is a combined response with both stream and execution data + logger.info('Found combined stream+execution response from block') + + // Incorporate the execution data from the block into our context + const executionData = streamingOutput.execution + + // Add any logs from the execution data to our context + if (executionData.logs && Array.isArray(executionData.logs)) { + context.blockLogs.push(...executionData.logs) + } + + // Add proper console entry for the streaming block + // This ensures identical formatting between streamed and non-streamed outputs + if (executionData.output) { + const blockLog = executionData.logs?.find((log: BlockLog) => log.blockId === executionData.blockId) + const consoleStore = useConsoleStore.getState() + + // Create a complete console entry with the full output structure, not the raw streaming object + const consoleEntry = { + output: executionData.output, // Use just the output, not the whole streaming structure + durationMs: blockLog?.durationMs || executionData.metadata?.duration || 0, + startedAt: blockLog?.startedAt || executionData.metadata?.startTime || new Date().toISOString(), + endedAt: blockLog?.endedAt || executionData.metadata?.endTime || new Date().toISOString(), + workflowId: context.workflowId, + timestamp: blockLog?.startedAt || executionData.metadata?.startTime || new Date().toISOString(), + blockId: executionData.blockId, + blockName: executionData.blockName || blockLog?.blockName || 'Agent Block', + blockType: executionData.blockType || blockLog?.blockType || 'agent' + } + + // Add to console + const newEntry = consoleStore.addConsole(consoleEntry) + + // Save the entryId for potential updates when stream completes + const consoleEntryId = newEntry?.id + + // Set up a stream completion handler to update the console with final content + if (consoleEntryId && 'stream' in streamingOutput) { + // Clone the stream so we don't consume the original one + const originalStream = streamingOutput.stream + const [contentStream, returnStream] = originalStream.tee() + + // Replace the original stream with our cloned version that will be returned + streamingOutput.stream = returnStream + + // Create a reader to process the cloned stream for content collection + const reader = contentStream.getReader() + const decoder = new TextDecoder() + let fullContent = ''; + + // Process the stream in the background to collect the full content + (async () => { + try { + while (true) { + const { done, value } = await reader.read() + if (done) break + const chunk = decoder.decode(value, { stream: true }) + fullContent += chunk + } + // Once stream is complete, update the console entry with the final content + if (fullContent.length > 0 && executionData.output?.response) { + const updatedOutput = { + ...executionData.output, + response: { + ...executionData.output.response, + content: fullContent + } + } + + // Update the console UI with the final content + consoleStore.updateConsole(consoleEntryId, { output: updatedOutput }) + + // Update the execution data itself with the final content + // so that when logs are persisted, they have the complete content + executionData.output.response.content = fullContent + + // If there's a block log for this execution, update it with the final content + if (executionData.blockId) { + const blockLog = context.blockLogs.find(log => log.blockId === executionData.blockId) + if (blockLog?.output?.response) { + blockLog.output.response.content = fullContent + } + } + } + } catch (e) { + logger.error('Error processing stream for console update:', e) + } + })() + } + } + + // Build a complete execution result with our context's logs + const execution: ExecutionResult & { isStreaming: boolean } = { + success: executionData.success !== false, + output: executionData.output || { response: {} }, + error: executionData.error, + logs: context.blockLogs, + metadata: { + duration: Date.now() - startTime.getTime(), + startTime: context.metadata.startTime!, + endTime: new Date().toISOString(), + workflowConnections: this.actualWorkflow.connections.map((conn: any) => ({ + source: conn.source, + target: conn.target, + })), + }, + isStreaming: true, + } + + // Add block metadata to logs if missing + if (context.blockLogs.length > 0) { + for (const log of context.blockLogs) { + if (!log.output) log.output = { response: {} } + + // For blocks matching the streaming block, ensure we add response and content properly + if (log.blockId === executionData.blockId) { + if (!log.output.response) log.output.response = {} + + // Add the output structure, preferring direct response content if available + if (executionData.output?.response) { + // Copy all properties from executionData response + Object.assign(log.output.response, executionData.output.response) + + // For streaming, we may not have content yet, so we store a placeholder + // that will be updated when the stream completes + if (!log.output.response.content && executionData.output.response.content) { + log.output.response.content = executionData.output.response.content + } + } + } + } + } + + // Return a properly formed StreamingExecution object + return { + stream: streamingOutput.stream, + execution, + } + } if (outputs.length > 0) { - finalOutput = outputs[outputs.length - 1] + // Filter out StreamingExecution objects (already handled above) + const normalizedOutputs = outputs.filter(output => + !(typeof output === 'object' && output !== null && 'stream' in output && 'execution' in output) + ) + if (normalizedOutputs.length > 0) { + finalOutput = normalizedOutputs[normalizedOutputs.length - 1] as NormalizedBlockOutput + } } // Process loop iterations - this will activate external paths when loops complete @@ -194,7 +388,7 @@ export class Executor { trackWorkflowTelemetry('workflow_execution_completed', { workflowId, duration, - blockCount: this.workflow.blocks.length, + blockCount: this.actualWorkflow.blocks.length, executedBlockCount: context.executedBlocks.size, startTime: startTime.toISOString(), endTime: endTime.toISOString(), @@ -208,7 +402,7 @@ export class Executor { duration: duration, startTime: context.metadata.startTime!, endTime: context.metadata.endTime!, - workflowConnections: this.workflow.connections.map((conn) => ({ + workflowConnections: this.actualWorkflow.connections.map((conn: any) => ({ source: conn.source, target: conn.target, })), @@ -278,7 +472,7 @@ export class Executor { endTime: context.metadata.endTime!, pendingBlocks: [], isDebugSession: false, - workflowConnections: this.workflow.connections.map((conn) => ({ + workflowConnections: this.actualWorkflow.connections.map((conn) => ({ source: conn.source, target: conn.target, })), @@ -319,27 +513,27 @@ export class Executor { * @throws Error if workflow validation fails */ private validateWorkflow(): void { - const starterBlock = this.workflow.blocks.find((block) => block.metadata?.id === 'starter') + const starterBlock = this.actualWorkflow.blocks.find((block) => block.metadata?.id === 'starter') if (!starterBlock || !starterBlock.enabled) { throw new Error('Workflow must have an enabled starter block') } - const incomingToStarter = this.workflow.connections.filter( + const incomingToStarter = this.actualWorkflow.connections.filter( (conn) => conn.target === starterBlock.id ) if (incomingToStarter.length > 0) { throw new Error('Starter block cannot have incoming connections') } - const outgoingFromStarter = this.workflow.connections.filter( + const outgoingFromStarter = this.actualWorkflow.connections.filter( (conn) => conn.source === starterBlock.id ) if (outgoingFromStarter.length === 0) { throw new Error('Starter block must have at least one outgoing connection') } - const blockIds = new Set(this.workflow.blocks.map((block) => block.id)) - for (const conn of this.workflow.connections) { + const blockIds = new Set(this.actualWorkflow.blocks.map((block) => block.id)) + for (const conn of this.actualWorkflow.connections) { if (!blockIds.has(conn.source)) { throw new Error(`Connection references non-existent source block: ${conn.source}`) } @@ -348,7 +542,7 @@ export class Executor { } } - for (const [loopId, loop] of Object.entries(this.workflow.loops || {})) { + for (const [loopId, loop] of Object.entries(this.actualWorkflow.loops || {})) { for (const nodeId of loop.nodes) { if (!blockIds.has(nodeId)) { throw new Error(`Loop ${loopId} references non-existent block: ${nodeId}`) @@ -388,7 +582,11 @@ export class Executor { completedLoops: new Set(), executedBlocks: new Set(), activeExecutionPath: new Set(), - workflow: this.workflow, + workflow: this.actualWorkflow, + // Add streaming context from contextExtensions + stream: this.contextExtensions.stream || false, + selectedOutputIds: this.contextExtensions.selectedOutputIds || [], + edges: this.contextExtensions.edges || [], } Object.entries(this.initialBlockStates).forEach(([blockId, output]) => { @@ -400,14 +598,14 @@ export class Executor { }) // Initialize loop iterations - if (this.workflow.loops) { - for (const loopId of Object.keys(this.workflow.loops)) { + if (this.actualWorkflow.loops) { + for (const loopId of Object.keys(this.actualWorkflow.loops)) { // Start all loops at iteration 0 context.loopIterations.set(loopId, 0) } } - const starterBlock = this.workflow.blocks.find((block) => block.metadata?.id === 'starter') + const starterBlock = this.actualWorkflow.blocks.find((block) => block.metadata?.id === 'starter') if (starterBlock) { // Initialize the starter block with the workflow input try { @@ -428,10 +626,10 @@ export class Executor { // This handles both input formats: { input: { field: value } } and { field: value } const inputValue = this.workflowInput?.input?.[field.name] !== undefined ? this.workflowInput.input[field.name] // Try to get from input.field - : this.workflowInput?.[field.name]; // Fallback to direct field access + : this.workflowInput?.[field.name] // Fallback to direct field access logger.info(`[Executor] Processing input field ${field.name} (${field.type}):`, - inputValue !== undefined ? JSON.stringify(inputValue) : 'undefined'); + inputValue !== undefined ? JSON.stringify(inputValue) : 'undefined') // Convert the value to the appropriate type let typedValue = inputValue @@ -458,15 +656,15 @@ export class Executor { } // Check if we managed to process any fields - if not, use the raw input - const hasProcessedFields = Object.keys(structuredInput).length > 0; + const hasProcessedFields = Object.keys(structuredInput).length > 0 // If no fields matched the input format, extract the raw input to use instead const rawInputData = this.workflowInput?.input !== undefined ? this.workflowInput.input // Use the nested input data - : this.workflowInput; // Fallback to direct input + : this.workflowInput // Fallback to direct input // Use the structured input if we processed fields, otherwise use raw input - const finalInput = hasProcessedFields ? structuredInput : rawInputData; + const finalInput = hasProcessedFields ? structuredInput : rawInputData // Initialize the starter block with structured input // Ensure both input and direct fields are available @@ -477,7 +675,7 @@ export class Executor { }, } - logger.info(`[Executor] Starter output:`, JSON.stringify(starterOutput, null, 2)); + logger.info(`[Executor] Starter output:`, JSON.stringify(starterOutput, null, 2)) context.blockStates.set(starterBlock.id, { output: starterOutput, @@ -554,7 +752,7 @@ export class Executor { context.executedBlocks.add(starterBlock.id) // Add all blocks connected to the starter to the active execution path - const connectedToStarter = this.workflow.connections + const connectedToStarter = this.actualWorkflow.connections .filter((conn) => conn.source === starterBlock.id) .map((conn) => conn.target) @@ -577,7 +775,7 @@ export class Executor { const executedBlocks = context.executedBlocks const pendingBlocks = new Set() - for (const block of this.workflow.blocks) { + for (const block of this.actualWorkflow.blocks) { if (executedBlocks.has(block.id) || block.enabled === false) { continue } @@ -587,12 +785,12 @@ export class Executor { continue } - const incomingConnections = this.workflow.connections.filter( + const incomingConnections = this.actualWorkflow.connections.filter( (conn) => conn.target === block.id ) // Find all loops that this block is a part of - const containingLoops = Object.values(this.workflow.loops || {}).filter((loop) => + const containingLoops = Object.values(this.actualWorkflow.loops || {}).filter((loop) => loop.nodes.includes(block.id) ) @@ -605,7 +803,7 @@ export class Executor { ) // Check if there's a direct self-connection - const hasSelfConnection = this.workflow.connections.some( + const hasSelfConnection = this.actualWorkflow.connections.some( (conn) => conn.source === block.id && conn.target === block.id ) @@ -628,7 +826,7 @@ export class Executor { // Regular non-loop block handling (unchanged) const allDependenciesMet = incomingConnections.every((conn) => { const sourceExecuted = executedBlocks.has(conn.source) - const sourceBlock = this.workflow.blocks.find((b) => b.id === conn.source) + const sourceBlock = this.actualWorkflow.blocks.find((b) => b.id === conn.source) const sourceBlockState = context.blockStates.get(conn.source) const hasSourceError = sourceBlockState?.output?.error !== undefined || @@ -636,7 +834,7 @@ export class Executor { // For condition blocks, check if this is the selected path if (conn.sourceHandle?.startsWith('condition-')) { - const sourceBlock = this.workflow.blocks.find((b) => b.id === conn.source) + const sourceBlock = this.actualWorkflow.blocks.find((b) => b.id === conn.source) if (sourceBlock?.metadata?.id === 'condition') { const conditionId = conn.sourceHandle.replace('condition-', '') const selectedCondition = context.decisions.condition.get(conn.source) @@ -741,7 +939,7 @@ export class Executor { blockId: string, context: ExecutionContext ): Promise { - const block = this.workflow.blocks.find((b) => b.id === blockId) + const block = this.actualWorkflow.blocks.find((b) => b.id === blockId) if (!block) { throw new Error(`Block ${blockId} not found`) } @@ -766,7 +964,7 @@ export class Executor { // Check if this block needs the starter block's output // This is especially relevant for API, function, and conditions that might reference - const starterBlock = this.workflow.blocks.find((b) => b.metadata?.id === 'starter') + const starterBlock = this.actualWorkflow.blocks.find((b) => b.metadata?.id === 'starter') if (starterBlock) { const starterState = context.blockStates.get(starterBlock.id) if (!starterState) { @@ -831,7 +1029,6 @@ export class Executor { startedAt: blockLog.startedAt, endedAt: blockLog.endedAt, workflowId: context.workflowId, - timestamp: blockLog.startedAt, blockId: block.id, blockName: block.metadata?.name || 'Unnamed Block', blockType: block.metadata?.id || 'unknown', @@ -874,7 +1071,6 @@ export class Executor { startedAt: blockLog.startedAt, endedAt: blockLog.endedAt, workflowId: context.workflowId, - timestamp: blockLog.startedAt, blockName: block.metadata?.name || 'Unnamed Block', blockType: block.metadata?.id || 'unknown', }) @@ -949,13 +1145,13 @@ export class Executor { */ private activateErrorPath(blockId: string, context: ExecutionContext): boolean { // Skip for starter blocks which don't have error handles - const block = this.workflow.blocks.find((b) => b.id === blockId) + const block = this.actualWorkflow.blocks.find((b) => b.id === blockId) if (block?.metadata?.id === 'starter' || block?.metadata?.id === 'condition') { return false } // Look for connections from this block's error handle - const errorConnections = this.workflow.connections.filter( + const errorConnections = this.actualWorkflow.connections.filter( (conn) => conn.source === blockId && conn.sourceHandle === 'error' ) diff --git a/sim/executor/types.ts b/sim/executor/types.ts index 0a21d3056b..8c8ca580a5 100644 --- a/sim/executor/types.ts +++ b/sim/executor/types.ts @@ -100,6 +100,11 @@ export interface ExecutionContext { activeExecutionPath: Set // Set of block IDs in the current execution path workflow?: SerializedWorkflow // Reference to the workflow being executed + + // Streaming support and output selection + stream?: boolean // Whether to use streaming responses when available + selectedOutputIds?: string[] // IDs of blocks selected for streaming output + edges?: Array<{source: string, target: string}> // Workflow edge connections } /** @@ -113,6 +118,15 @@ export interface ExecutionResult { metadata?: ExecutionMetadata } +/** + * Streaming execution result combining a readable stream with execution metadata. + * This allows us to stream content to the UI while still capturing all execution logs. + */ +export interface StreamingExecution { + stream: ReadableStream // The streaming response for the UI to consume + execution: ExecutionResult & { isStreaming?: boolean } // The complete execution data for logging purposes +} + /** * Interface for a block executor component. */ @@ -151,13 +165,13 @@ export interface BlockHandler { * @param block - Block to execute * @param inputs - Resolved input parameters * @param context - Current execution context - * @returns Block execution output + * @returns Block execution output or StreamingExecution for streaming */ execute( block: SerializedBlock, inputs: Record, context: ExecutionContext - ): Promise + ): Promise } /** diff --git a/sim/lib/logs/execution-logger.ts b/sim/lib/logs/execution-logger.ts index 3b033e1562..9b2c33d047 100644 --- a/sim/lib/logs/execution-logger.ts +++ b/sim/lib/logs/execution-logger.ts @@ -6,6 +6,7 @@ import { userStats, workflow, workflowLogs } from '@/db/schema' import { ExecutionResult as ExecutorResult } from '@/executor/types' import { stripCustomToolPrefix } from '../workflows/utils' import { getCostMultiplier } from '@/lib/environment' +import { calculateCost } from '@/providers/utils' const logger = createLogger('ExecutionLogger') @@ -109,6 +110,83 @@ export async function persistExecutionLogs( hasToolCalls: !!log.output.toolCalls, hasResponse: !!log.output.response, }) + + // FIRST PASS - Check if this is a no-tool scenario with tokens data not propagated + // In some cases, the token data from the streaming callback doesn't properly get into + // the agent block response. This ensures we capture it. + if (log.output.response && + (!log.output.response.tokens?.completion || log.output.response.tokens.completion === 0) && + (!log.output.response.toolCalls || !log.output.response.toolCalls.list || log.output.response.toolCalls.list.length === 0)) { + + // Check if output response has providerTiming - this indicates it's a streaming response + if (log.output.response.providerTiming) { + logger.debug('Processing streaming response without tool calls for token extraction', { + blockId: log.blockId, + hasTokens: !!log.output.response.tokens, + hasProviderTiming: !!log.output.response.providerTiming + }); + + // Only for no-tool streaming cases, extract content length and estimate token count + const contentLength = log.output.response.content?.length || 0; + if (contentLength > 0) { + // Estimate completion tokens based on content length as a fallback + const estimatedCompletionTokens = Math.ceil(contentLength / 4); + const promptTokens = log.output.response.tokens?.prompt || 8; + + // Update the tokens object + log.output.response.tokens = { + prompt: promptTokens, + completion: estimatedCompletionTokens, + total: promptTokens + estimatedCompletionTokens + }; + + // Update cost information using the provider's cost model + const model = log.output.response.model || 'gpt-4o'; + const costInfo = calculateCost(model, promptTokens, estimatedCompletionTokens); + log.output.response.cost = { + input: costInfo.input, + output: costInfo.output, + total: costInfo.total, + pricing: costInfo.pricing + }; + + logger.debug('Updated token information for streaming no-tool response', { + blockId: log.blockId, + contentLength, + estimatedCompletionTokens, + tokens: log.output.response.tokens + }); + } + } + } + + // Special case for streaming responses from agent blocks + // This format has both stream and executionData properties + if (log.output.stream && log.output.executionData) { + logger.debug('Found streaming response with executionData', { + blockId: log.blockId, + hasExecutionData: !!log.output.executionData, + executionDataKeys: log.output.executionData ? Object.keys(log.output.executionData) : [], + }) + + // Extract the executionData and use it as our primary source of information + const executionData = log.output.executionData + + // If executionData has output with response, use that as our response + // This is especially important for streaming responses where the final content + // is set in the executionData structure by the executor + if (executionData.output?.response) { + log.output.response = executionData.output.response + logger.debug('Using response from executionData', { + responseKeys: Object.keys(log.output.response), + hasContent: !!log.output.response.content, + contentLength: log.output.response.content?.length || 0, + hasToolCalls: !!log.output.response.toolCalls, + hasTokens: !!log.output.response.tokens, + hasCost: !!log.output.response.cost, + }) + } + } // Extract tool calls and other metadata if (log.output.response) { @@ -347,7 +425,49 @@ export async function persistExecutionLogs( } }) } - // Case 5: Parse the response string for toolCalls as a last resort + // Case 5: Look in executionData.output.response for streaming responses + else if (log.output.executionData?.output?.response?.toolCalls) { + const toolCallsObj = log.output.executionData.output.response.toolCalls + const list = Array.isArray(toolCallsObj) ? toolCallsObj : (toolCallsObj.list || []) + + logger.debug('Found toolCalls in executionData output response', { + count: list.length, + }) + + // Log raw timing data for debugging + list.forEach((tc: any, idx: number) => { + logger.debug(`executionData toolCalls ${idx} raw timing data:`, { + name: stripCustomToolPrefix(tc.name), + startTime: tc.startTime, + endTime: tc.endTime, + duration: tc.duration, + timing: tc.timing, + argumentKeys: tc.arguments ? Object.keys(tc.arguments) : undefined, + }) + }) + + toolCallData = list.map((toolCall: any) => { + // Extract timing info - try various formats that providers might use + const duration = extractDuration(toolCall) + const timing = extractTimingInfo( + toolCall, + blockStartTime ? new Date(blockStartTime) : undefined, + blockEndTime ? new Date(blockEndTime) : undefined + ) + + return { + name: toolCall.name, + duration: duration, + startTime: timing.startTime, + endTime: timing.endTime, + status: toolCall.error ? 'error' : 'success', + input: toolCall.arguments || toolCall.input, + output: toolCall.result || toolCall.output, + error: toolCall.error, + } + }) + } + // Case 6: Parse the response string for toolCalls as a last resort else if (typeof log.output.response === 'string') { const match = log.output.response.match(/"toolCalls"\s*:\s*({[^}]*}|(\[.*?\]))/s) if (match) { @@ -446,7 +566,11 @@ export async function persistExecutionLogs( executionId, level: log.success ? 'info' : 'error', message: log.success - ? `Block ${log.blockName || log.blockId} (${log.blockType || 'unknown'}): ${JSON.stringify(log.output?.response || {})}` + ? `Block ${log.blockName || log.blockId} (${log.blockType || 'unknown'}): ${ + log.output?.response?.content || + log.output?.executionData?.output?.response?.content || + JSON.stringify(log.output?.response || {}) + }` : `Block ${log.blockName || log.blockId} (${log.blockType || 'unknown'}): ${log.error || 'Failed'}`, duration: log.success ? `${log.durationMs}ms` : 'NA', trigger: triggerType, @@ -513,6 +637,23 @@ export async function persistExecutionLogs( } } + // If result has a direct cost field (for streaming responses completed with calculated cost), + // use that as a safety check to ensure we have cost data + if (result.metadata && 'cost' in result.metadata && (!workflowMetadata.cost || workflowMetadata.cost.total <= 0)) { + const resultCost = (result.metadata as any).cost + workflowMetadata.cost = { + model: primaryModel, + total: typeof resultCost === 'number' ? resultCost : (resultCost?.total || 0), + input: resultCost?.input || 0, + output: resultCost?.output || 0, + tokens: { + prompt: totalPromptTokens, + completion: totalCompletionTokens, + total: totalTokens, + }, + } + } + if (userId) { try { const userStatsRecords = await db diff --git a/sim/lib/logs/trace-spans.ts b/sim/lib/logs/trace-spans.ts index 508f8ab7ad..fed6547088 100644 --- a/sim/lib/logs/trace-spans.ts +++ b/sim/lib/logs/trace-spans.ts @@ -77,48 +77,63 @@ export function buildTraceSpans(result: ExecutionResult): { }, index: number ) => { - const relativeStart = segment.startTime - segmentStartTime - - // Enhance the segment name to include model information for model segments - let enhancedName = segment.name - if (segment.type === 'model') { - const modelName = log.output.response.model || '' - - if (segment.name === 'Initial response') { - enhancedName = `Initial response${modelName ? ` (${modelName})` : ''}` - } else if (segment.name.includes('iteration')) { - // Extract the iteration number - const iterationMatch = segment.name.match(/\(iteration (\d+)\)/) - const iterationNum = iterationMatch ? iterationMatch[1] : '' - - enhancedName = `Model response${iterationNum ? ` (iteration ${iterationNum})` : ''}${modelName ? ` (${modelName})` : ''}` + // Ensure we have valid startTime and endTime + let segmentStart: number + let segmentEnd: number + + // Handle different time formats - some providers use ISO strings, some use timestamps + if (typeof segment.startTime === 'string') { + try { + segmentStart = new Date(segment.startTime).getTime() + } catch (e) { + segmentStart = segmentStartTime + (index * 1000) // Fallback offset } + } else { + segmentStart = segment.startTime } - - const segmentSpan: TraceSpan = { + + if (typeof segment.endTime === 'string') { + try { + segmentEnd = new Date(segment.endTime).getTime() + } catch (e) { + segmentEnd = segmentStart + (segment.duration || 1000) // Fallback duration + } + } else { + segmentEnd = segment.endTime + } + + // For streaming responses, make sure our timing is valid + if (isNaN(segmentStart) || isNaN(segmentEnd) || segmentEnd < segmentStart) { + // Use fallback values + segmentStart = segmentStartTime + (index * 1000) + segmentEnd = segmentStart + (segment.duration || 1000) + } + + const childSpan: TraceSpan = { id: `${spanId}-segment-${index}`, - name: enhancedName, - // Make sure we handle model and tool types, and fallback to generic 'span' for anything else - type: segment.type === 'model' || segment.type === 'tool' ? segment.type : 'span', - duration: segment.duration, - startTime: new Date(segment.startTime).toISOString(), - endTime: new Date(segment.endTime).toISOString(), + name: segment.name || `${segment.type} operation`, + startTime: new Date(segmentStart).toISOString(), + endTime: new Date(segmentEnd).toISOString(), + duration: segment.duration || (segmentEnd - segmentStart), + type: segment.type === 'model' ? 'model' : segment.type === 'tool' ? 'tool' : 'processing', status: 'success', - // Add relative timing display for segments after the first one - relativeStartMs: index === 0 ? undefined : relativeStart, - // For model segments, add token info if available - ...(segment.type === 'model' && { - tokens: index === 0 ? log.output.response.tokens?.completion : undefined, - }), + children: [], } - children.push(segmentSpan) + // Add any additional metadata + if (segment.type === 'tool' && typeof segment.name === 'string') { + // Add as a custom attribute using type assertion + (childSpan as any).toolName = segment.name + } + + children.push(childSpan) } ) - // Add all segments as children - if (!span.children) span.children = [] - span.children.push(...children) + // Only add children if we have valid spans + if (children.length > 0) { + span.children = children + } } // If no segments but we have provider timing, create a provider span else { @@ -171,17 +186,59 @@ export function buildTraceSpans(result: ExecutionResult): { } } else { // When not using provider timing at all, add tool calls if they exist - if (log.output?.response?.toolCalls?.list) { - span.toolCalls = log.output.response.toolCalls.list.map((tc: any) => ({ - name: stripCustomToolPrefix(tc.name), - duration: tc.duration || 0, - startTime: tc.startTime || log.startedAt, - endTime: tc.endTime || log.endedAt, - status: tc.error ? 'error' : 'success', - input: tc.arguments || tc.input, - output: tc.result || tc.output, - error: tc.error, - })) + // Tool calls handling for different formats: + // 1. Standard format in response.toolCalls.list + // 2. Direct toolCalls array in response + // 3. Streaming response formats with executionData + + // Check all possible paths for toolCalls + let toolCallsList = null; + + // Wrap extraction in try-catch to handle unexpected toolCalls formats + try { + if (log.output?.response?.toolCalls?.list) { + // Standard format with list property + toolCallsList = log.output.response.toolCalls.list; + } else if (Array.isArray(log.output?.response?.toolCalls)) { + // Direct array format + toolCallsList = log.output.response.toolCalls; + } else if (log.output?.executionData?.output?.response?.toolCalls) { + // Streaming format with executionData + const tcObj = log.output.executionData.output.response.toolCalls; + toolCallsList = Array.isArray(tcObj) ? tcObj : (tcObj.list || []); + } + + // Validate that toolCallsList is actually an array before processing + if (toolCallsList && !Array.isArray(toolCallsList)) { + console.warn(`toolCallsList is not an array: ${typeof toolCallsList}`); + toolCallsList = []; + } + } catch (error) { + console.error(`Error extracting toolCalls: ${error}`); + toolCallsList = []; // Set to empty array as fallback + } + + if (toolCallsList && toolCallsList.length > 0) { + span.toolCalls = toolCallsList.map((tc: any) => { + // Add null check for each tool call + if (!tc) return null; + + try { + return { + name: stripCustomToolPrefix(tc.name || 'unnamed-tool'), + duration: tc.duration || 0, + startTime: tc.startTime || log.startedAt, + endTime: tc.endTime || log.endedAt, + status: tc.error ? 'error' : 'success', + input: tc.arguments || tc.input, + output: tc.result || tc.output, + error: tc.error, + }; + } catch (tcError) { + console.error(`Error processing tool call: ${tcError}`); + return null; + } + }).filter(Boolean); // Remove any null entries from failed processing } } diff --git a/sim/lib/webhooks/utils.ts b/sim/lib/webhooks/utils.ts index 9c74fddb4a..bf5165f91e 100644 --- a/sim/lib/webhooks/utils.ts +++ b/sim/lib/webhooks/utils.ts @@ -544,11 +544,17 @@ export async function executeWorkflowFromPayload( // This is THE critical line where the workflow actually executes const result = await executor.execute(foundWorkflow.id) + // Check if we got a StreamingExecution result (with stream + execution properties) + // For webhook executions, we only care about the ExecutionResult part, not the stream + const executionResult = 'stream' in result && 'execution' in result + ? result.execution + : result + // Add direct detailed logging right after executing logger.info(`[${requestId}] EXECUTION_MONITOR: executor.execute() completed with result`, { workflowId: foundWorkflow.id, executionId: executionId, - success: result.success, + success: executionResult.success, resultType: result ? typeof result : 'undefined', timestamp: new Date().toISOString() }); @@ -557,7 +563,7 @@ export async function executeWorkflowFromPayload( const executionDuration = Date.now() - executionStartTime; logger.info(`[${requestId}] TRACE: Workflow execution completed`, { workflowId: foundWorkflow.id, - success: result.success, + success: executionResult.success, duration: `${executionDuration}ms`, actualDurationMs: executionDuration, timestamp: new Date().toISOString() @@ -565,13 +571,13 @@ export async function executeWorkflowFromPayload( logger.info(`[${requestId}] Workflow execution finished`, { executionId, - success: result.success, - durationMs: result.metadata?.duration || executionDuration, + success: executionResult.success, + durationMs: executionResult.metadata?.duration || executionDuration, actualDurationMs: executionDuration }) // Update counts and stats if successful - if (result.success) { + if (executionResult.success) { await updateWorkflowRunCounts(foundWorkflow.id) await db .update(userStats) @@ -589,8 +595,8 @@ export async function executeWorkflowFromPayload( } // Build and enrich result with trace spans - const { traceSpans, totalDuration } = buildTraceSpans(result) - const enrichedResult = { ...result, traceSpans, totalDuration } + const { traceSpans, totalDuration } = buildTraceSpans(executionResult) + const enrichedResult = { ...executionResult, traceSpans, totalDuration } // Persist logs for this execution using the standard 'webhook' trigger type await persistExecutionLogs(foundWorkflow.id, executionId, enrichedResult, 'webhook') diff --git a/sim/next.config.ts b/sim/next.config.ts index 651c7d85b3..11ab47a88b 100644 --- a/sim/next.config.ts +++ b/sim/next.config.ts @@ -84,7 +84,7 @@ const nextConfig: NextConfig = { }, { key: 'Content-Security-Policy', - value: "default-src 'self'; script-src 'self' 'unsafe-inline' 'unsafe-eval' https://*.google.com https://apis.google.com https://*.vercel-insights.com https://vercel.live https://*.vercel.live; style-src 'self' 'unsafe-inline' https://fonts.googleapis.com; img-src 'self' data: blob: https://*.googleusercontent.com https://*.google.com https://*.atlassian.com; font-src 'self' https://fonts.gstatic.com; connect-src 'self' http://localhost:11434 http://host.docker.internal:11434 https://*.googleapis.com https://*.amazonaws.com https://*.s3.amazonaws.com https://s3.*.amazonaws.com https://*.vercel-insights.com https://*.atlassian.com https://vercel.live https://*.vercel.live; frame-src https://drive.google.com https://*.google.com; frame-ancestors 'self'; form-action 'self'; base-uri 'self'; object-src 'none'", + value: "default-src 'self'; script-src 'self' 'unsafe-inline' 'unsafe-eval' https://*.google.com https://apis.google.com https://*.vercel-insights.com https://vercel.live https://*.vercel.live; style-src 'self' 'unsafe-inline' https://fonts.googleapis.com; img-src 'self' data: blob: https://*.googleusercontent.com https://*.google.com https://*.atlassian.com; font-src 'self' https://fonts.gstatic.com; connect-src 'self' http://localhost:11434 http://host.docker.internal:11434 https://*.googleapis.com https://*.amazonaws.com https://*.s3.amazonaws.com https://*.vercel-insights.com https://*.atlassian.com https://vercel.live https://*.vercel.live; frame-src https://drive.google.com https://*.google.com; frame-ancestors 'self'; form-action 'self'; base-uri 'self'; object-src 'none'", }, ], }, diff --git a/sim/providers/anthropic/index.ts b/sim/providers/anthropic/index.ts index a4649d809c..9cd0bc85fd 100644 --- a/sim/providers/anthropic/index.ts +++ b/sim/providers/anthropic/index.ts @@ -2,10 +2,33 @@ import Anthropic from '@anthropic-ai/sdk' import { createLogger } from '@/lib/logs/console-logger' import { executeTool } from '@/tools' import { ProviderConfig, ProviderRequest, ProviderResponse, TimeSegment } from '../types' +import { StreamingExecution } from '@/executor/types' import { prepareToolsWithUsageControl, trackForcedToolUsage } from '../utils' const logger = createLogger('Anthropic Provider') +/** + * Helper to wrap Anthropic streaming (async iterable of SSE events) into a browser-friendly + * ReadableStream of raw assistant text chunks. We enqueue only `content_block_delta` events + * with `delta.type === 'text_delta'`, since that contains the incremental text tokens. + */ +function createReadableStreamFromAnthropicStream(anthropicStream: AsyncIterable): ReadableStream { + return new ReadableStream({ + async start(controller) { + try { + for await (const event of anthropicStream) { + if (event.type === 'content_block_delta' && event.delta?.text) { + controller.enqueue(new TextEncoder().encode(event.delta.text)) + } + } + controller.close() + } catch (err) { + controller.error(err) + } + }, + }) +} + export const anthropicProvider: ProviderConfig = { id: 'anthropic', name: 'Anthropic', @@ -14,7 +37,7 @@ export const anthropicProvider: ProviderConfig = { models: ['claude-3-5-sonnet-20240620', 'claude-3-7-sonnet-20250219'], defaultModel: 'claude-3-7-sonnet-20250219', - executeRequest: async (request: ProviderRequest): Promise => { + executeRequest: async (request: ProviderRequest): Promise => { if (!request.apiKey) { throw new Error('API key is required for Anthropic') } @@ -233,6 +256,73 @@ ${fieldDescriptions} } } + // EARLY STREAMING: if caller requested streaming and there are no tools to execute, + // we can directly stream the completion. + if (request.stream && (!anthropicTools || anthropicTools.length === 0)) { + logger.info('Using streaming response for Anthropic request (no tools)') + + // Start execution timer for the entire provider execution + const providerStartTime = Date.now() + const providerStartTimeISO = new Date(providerStartTime).toISOString() + + // Create a streaming request + const streamResponse: any = await anthropic.messages.create({ + ...payload, + stream: true, + }) + + // Start collecting token usage + let tokenUsage = { + prompt: 0, + completion: 0, + total: 0 + } + + // Create a StreamingExecution response with a readable stream + const streamingResult = { + stream: createReadableStreamFromAnthropicStream(streamResponse), + execution: { + success: true, + output: { + response: { + content: '', // Will be filled by streaming content in chat component + model: request.model, + tokens: tokenUsage, + toolCalls: undefined, + providerTiming: { + startTime: providerStartTimeISO, + endTime: new Date().toISOString(), + duration: Date.now() - providerStartTime, + timeSegments: [{ + type: 'model', + name: 'Streaming response', + startTime: providerStartTime, + endTime: Date.now(), + duration: Date.now() - providerStartTime, + }] + }, + // Estimate token cost based on typical Claude pricing + cost: { + total: 0.0, + input: 0.0, + output: 0.0 + } + } + }, + logs: [], // No block logs for direct streaming + metadata: { + startTime: providerStartTimeISO, + endTime: new Date().toISOString(), + duration: Date.now() - providerStartTime, + }, + isStreaming: true + } + } + + // Return the streaming execution object + return streamingResult as StreamingExecution + } + // Start execution timer for the entire provider execution const providerStartTime = Date.now() const providerStartTimeISO = new Date(providerStartTime).toISOString() @@ -519,6 +609,72 @@ ${fieldDescriptions} const providerEndTimeISO = new Date(providerEndTime).toISOString() const totalDuration = providerEndTime - providerStartTime + // After all tool processing complete, if streaming was requested and we have messages, use streaming for the final response + if (request.stream && iterationCount > 0) { + logger.info('Using streaming for final Anthropic response after tool calls') + + // When streaming after tool calls with forced tools, make sure tool_choice is removed + // This prevents the API from trying to force tool usage again in the final streaming response + const streamingPayload = { + ...payload, + messages: currentMessages, + // For Anthropic, omit tool_choice entirely rather than setting it to 'none' + stream: true, + } + + // Remove the tool_choice parameter as Anthropic doesn't accept 'none' as a string value + delete streamingPayload.tool_choice + + const streamResponse: any = await anthropic.messages.create(streamingPayload) + + // Create a StreamingExecution response with all collected data + const streamingResult = { + stream: createReadableStreamFromAnthropicStream(streamResponse), + execution: { + success: true, + output: { + response: { + content: '', // Will be filled by the callback + model: request.model || 'claude-3-7-sonnet-20250219', + tokens: { + prompt: tokens.prompt, + completion: tokens.completion, + total: tokens.total, + }, + toolCalls: toolCalls.length > 0 ? { + list: toolCalls, + count: toolCalls.length + } : undefined, + providerTiming: { + startTime: providerStartTimeISO, + endTime: new Date().toISOString(), + duration: Date.now() - providerStartTime, + modelTime: modelTime, + toolsTime: toolsTime, + firstResponseTime: firstResponseTime, + iterations: iterationCount + 1, + timeSegments: timeSegments, + }, + cost: { + total: (tokens.total || 0) * 0.0001, // Estimate cost based on tokens + input: (tokens.prompt || 0) * 0.0001, + output: (tokens.completion || 0) * 0.0001 + } + } + }, + logs: [], // No block logs at provider level + metadata: { + startTime: providerStartTimeISO, + endTime: new Date().toISOString(), + duration: Date.now() - providerStartTime, + }, + isStreaming: true + } + } + + return streamingResult as StreamingExecution + } + return { content, model: request.model || 'claude-3-7-sonnet-20250219', diff --git a/sim/providers/cerebras/index.ts b/sim/providers/cerebras/index.ts index f84e28afb1..7a73a4368e 100644 --- a/sim/providers/cerebras/index.ts +++ b/sim/providers/cerebras/index.ts @@ -2,9 +2,33 @@ import { Cerebras } from '@cerebras/cerebras_cloud_sdk' import { createLogger } from '@/lib/logs/console-logger' import { executeTool } from '@/tools' import { ProviderConfig, ProviderRequest, ProviderResponse, TimeSegment } from '../types' +import { StreamingExecution } from '@/executor/types' const logger = createLogger('Cerebras Provider') +/** + * Helper to convert a Cerebras streaming response (async iterable) into a ReadableStream. + * Enqueues only the model's text delta chunks as UTF-8 encoded bytes. + */ +function createReadableStreamFromCerebrasStream(cerebrasStream: AsyncIterable): ReadableStream { + return new ReadableStream({ + async start(controller) { + try { + for await (const chunk of cerebrasStream) { + // Expecting delta content similar to OpenAI: chunk.choices[0]?.delta?.content + const content = chunk.choices?.[0]?.delta?.content || '' + if (content) { + controller.enqueue(new TextEncoder().encode(content)) + } + } + controller.close() + } catch (error) { + controller.error(error) + } + } + }) +} + export const cerebrasProvider: ProviderConfig = { id: 'cerebras', name: 'Cerebras', @@ -12,7 +36,7 @@ export const cerebrasProvider: ProviderConfig = { version: '1.0.0', models: ['cerebras/llama-3.3-70b'], defaultModel: 'cerebras/llama-3.3-70b', - executeRequest: async (request: ProviderRequest): Promise => { + executeRequest: async (request: ProviderRequest): Promise => { if (!request.apiKey) { throw new Error('API key is required for Cerebras') } @@ -106,6 +130,66 @@ export const cerebrasProvider: ProviderConfig = { } } + // EARLY STREAMING: if streaming requested and no tools to execute, stream directly + if (request.stream && (!tools || tools.length === 0)) { + logger.info('Using streaming response for Cerebras request (no tools)') + const streamResponse: any = await client.chat.completions.create({ + ...payload, + stream: true, + }) + + // Start collecting token usage + let tokenUsage = { + prompt: 0, + completion: 0, + total: 0 + } + + // Create a StreamingExecution response with a readable stream + const streamingResult = { + stream: createReadableStreamFromCerebrasStream(streamResponse), + execution: { + success: true, + output: { + response: { + content: '', // Will be filled by streaming content in chat component + model: request.model || 'cerebras/llama-3.3-70b', + tokens: tokenUsage, + toolCalls: undefined, + providerTiming: { + startTime: providerStartTimeISO, + endTime: new Date().toISOString(), + duration: Date.now() - providerStartTime, + timeSegments: [{ + type: 'model', + name: 'Streaming response', + startTime: providerStartTime, + endTime: Date.now(), + duration: Date.now() - providerStartTime, + }] + }, + // Estimate token cost + cost: { + total: 0.0, + input: 0.0, + output: 0.0 + } + } + }, + logs: [], // No block logs for direct streaming + metadata: { + startTime: providerStartTimeISO, + endTime: new Date().toISOString(), + duration: Date.now() - providerStartTime, + }, + isStreaming: true + } + } + + // Return the streaming execution object + return streamingResult as StreamingExecution + } + // Make the initial API request const initialCallTime = Date.now() @@ -348,6 +432,70 @@ export const cerebrasProvider: ProviderConfig = { const providerEndTimeISO = new Date(providerEndTime).toISOString() const totalDuration = providerEndTime - providerStartTime + // POST-TOOL-STREAMING: stream after tool calls if requested + if (request.stream && iterationCount > 0) { + logger.info('Using streaming for final Cerebras response after tool calls') + + // When streaming after tool calls with forced tools, make sure tool_choice is set to 'auto' + // This prevents the API from trying to force tool usage again in the final streaming response + const streamingPayload = { + ...payload, + messages: currentMessages, + tool_choice: 'auto', // Always use 'auto' for the streaming response after tool calls + stream: true, + } + + const streamResponse: any = await client.chat.completions.create(streamingPayload) + + // Create a StreamingExecution response with all collected data + const streamingResult = { + stream: createReadableStreamFromCerebrasStream(streamResponse), + execution: { + success: true, + output: { + response: { + content: '', // Will be filled by the callback + model: request.model || 'cerebras/llama-3.3-70b', + tokens: { + prompt: tokens.prompt, + completion: tokens.completion, + total: tokens.total, + }, + toolCalls: toolCalls.length > 0 ? { + list: toolCalls, + count: toolCalls.length + } : undefined, + providerTiming: { + startTime: providerStartTimeISO, + endTime: new Date().toISOString(), + duration: Date.now() - providerStartTime, + modelTime: modelTime, + toolsTime: toolsTime, + firstResponseTime: firstResponseTime, + iterations: iterationCount + 1, + timeSegments: timeSegments, + }, + cost: { + total: (tokens.total || 0) * 0.0001, + input: (tokens.prompt || 0) * 0.0001, + output: (tokens.completion || 0) * 0.0001 + } + } + }, + logs: [], // No block logs at provider level + metadata: { + startTime: providerStartTimeISO, + endTime: new Date().toISOString(), + duration: Date.now() - providerStartTime, + }, + isStreaming: true + } + } + + // Return the streaming execution object + return streamingResult as StreamingExecution + } + return { content, model: request.model, diff --git a/sim/providers/deepseek/index.ts b/sim/providers/deepseek/index.ts index a37281cff6..36caaf3e62 100644 --- a/sim/providers/deepseek/index.ts +++ b/sim/providers/deepseek/index.ts @@ -2,10 +2,33 @@ import OpenAI from 'openai' import { createLogger } from '@/lib/logs/console-logger' import { executeTool } from '@/tools' import { ProviderConfig, ProviderRequest, ProviderResponse, TimeSegment } from '../types' +import { StreamingExecution } from '@/executor/types' import { prepareToolsWithUsageControl, trackForcedToolUsage } from '../utils' const logger = createLogger('Deepseek Provider') +/** + * Helper function to convert a DeepSeek (OpenAI-compatible) stream to a ReadableStream + * of text chunks that can be consumed by the browser. + */ +function createReadableStreamFromDeepseekStream(deepseekStream: any): ReadableStream { + return new ReadableStream({ + async start(controller) { + try { + for await (const chunk of deepseekStream) { + const content = chunk.choices[0]?.delta?.content || '' + if (content) { + controller.enqueue(new TextEncoder().encode(content)) + } + } + controller.close() + } catch (error) { + controller.error(error) + } + } + }) +} + export const deepseekProvider: ProviderConfig = { id: 'deepseek', name: 'Deepseek', @@ -14,7 +37,7 @@ export const deepseekProvider: ProviderConfig = { models: ['deepseek-chat'], defaultModel: 'deepseek-chat', - executeRequest: async (request: ProviderRequest): Promise => { + executeRequest: async (request: ProviderRequest): Promise => { if (!request.apiKey) { throw new Error('API key is required for Deepseek') } @@ -103,6 +126,67 @@ export const deepseekProvider: ProviderConfig = { } } + // EARLY STREAMING: if streaming requested and no tools to execute, stream directly + if (request.stream && (!tools || tools.length === 0)) { + logger.info('Using streaming response for DeepSeek request (no tools)') + + const streamResponse = await deepseek.chat.completions.create({ + ...payload, + stream: true, + }) + + // Start collecting token usage + let tokenUsage = { + prompt: 0, + completion: 0, + total: 0 + } + + // Create a StreamingExecution response with a readable stream + const streamingResult = { + stream: createReadableStreamFromDeepseekStream(streamResponse), + execution: { + success: true, + output: { + response: { + content: '', // Will be filled by streaming content in chat component + model: request.model || 'deepseek-chat', + tokens: tokenUsage, + toolCalls: undefined, + providerTiming: { + startTime: providerStartTimeISO, + endTime: new Date().toISOString(), + duration: Date.now() - providerStartTime, + timeSegments: [{ + type: 'model', + name: 'Streaming response', + startTime: providerStartTime, + endTime: Date.now(), + duration: Date.now() - providerStartTime, + }] + }, + // Estimate token cost + cost: { + total: 0.0, + input: 0.0, + output: 0.0 + } + } + }, + logs: [], // No block logs for direct streaming + metadata: { + startTime: providerStartTimeISO, + endTime: new Date().toISOString(), + duration: Date.now() - providerStartTime, + }, + isStreaming: true + } + } + + // Return the streaming execution object + return streamingResult as StreamingExecution + } + // Make the initial API request const initialCallTime = Date.now() @@ -350,6 +434,70 @@ export const deepseekProvider: ProviderConfig = { const providerEndTimeISO = new Date(providerEndTime).toISOString() const totalDuration = providerEndTime - providerStartTime + // POST-TOOL STREAMING: stream final response after tool calls if requested + if (request.stream && iterationCount > 0) { + logger.info('Using streaming for final DeepSeek response after tool calls') + + // When streaming after tool calls with forced tools, make sure tool_choice is set to 'auto' + // This prevents the API from trying to force tool usage again in the final streaming response + const streamingPayload = { + ...payload, + messages: currentMessages, + tool_choice: 'auto', // Always use 'auto' for the streaming response after tool calls + stream: true, + } + + const streamResponse = await deepseek.chat.completions.create(streamingPayload) + + // Create a StreamingExecution response with all collected data + const streamingResult = { + stream: createReadableStreamFromDeepseekStream(streamResponse), + execution: { + success: true, + output: { + response: { + content: '', // Will be filled by the callback + model: request.model || 'deepseek-chat', + tokens: { + prompt: tokens.prompt, + completion: tokens.completion, + total: tokens.total, + }, + toolCalls: toolCalls.length > 0 ? { + list: toolCalls, + count: toolCalls.length + } : undefined, + providerTiming: { + startTime: providerStartTimeISO, + endTime: new Date().toISOString(), + duration: Date.now() - providerStartTime, + modelTime: modelTime, + toolsTime: toolsTime, + firstResponseTime: firstResponseTime, + iterations: iterationCount + 1, + timeSegments: timeSegments, + }, + cost: { + total: (tokens.total || 0) * 0.0001, + input: (tokens.prompt || 0) * 0.0001, + output: (tokens.completion || 0) * 0.0001 + } + } + }, + logs: [], // No block logs at provider level + metadata: { + startTime: providerStartTimeISO, + endTime: new Date().toISOString(), + duration: Date.now() - providerStartTime, + }, + isStreaming: true + } + } + + // Return the streaming execution object + return streamingResult as StreamingExecution + } + return { content, model: request.model, diff --git a/sim/providers/google/index.ts b/sim/providers/google/index.ts index 900edd1276..c8905d3282 100644 --- a/sim/providers/google/index.ts +++ b/sim/providers/google/index.ts @@ -1,9 +1,90 @@ import { createLogger } from '@/lib/logs/console-logger' import { executeTool } from '@/tools' import { ProviderConfig, ProviderRequest, ProviderResponse, TimeSegment } from '../types' +import { StreamingExecution } from '@/executor/types' const logger = createLogger('Google Provider') +/** + * Creates a ReadableStream from Google's Gemini stream response + */ +function createReadableStreamFromGeminiStream(response: Response): ReadableStream { + const reader = response.body?.getReader() + if (!reader) { + throw new Error('Failed to get reader from response body') + } + + return new ReadableStream({ + async start(controller) { + try { + let buffer = '' + + while (true) { + const { done, value } = await reader.read() + if (done) { + controller.close() + break + } + + const text = new TextDecoder().decode(value) + buffer += text + + try { + const lines = buffer.split('\n') + buffer = '' + + for (let i = 0; i < lines.length; i++) { + const line = lines[i].trim() + + if (i === lines.length - 1 && line !== '') { + buffer = line + continue + } + + if (!line) continue + + if (line.startsWith('data: ')) { + const jsonStr = line.substring(6) + + if (jsonStr === '[DONE]') continue + + try { + const data = JSON.parse(jsonStr) + const candidate = data.candidates?.[0] + if (candidate?.content?.parts) { + const content = extractTextContent(candidate) + if (content) { + controller.enqueue(new TextEncoder().encode(content)) + } + } + } catch (e) { + logger.error('Error parsing Gemini SSE JSON data', { + error: e instanceof Error ? e.message : String(e), + data: jsonStr + }) + } + } + } + } catch (e) { + logger.error('Error processing Gemini SSE stream', { + error: e instanceof Error ? e.message : String(e), + chunk: text + }) + } + } + } catch (e) { + logger.error('Error reading Google Gemini stream', { + error: e instanceof Error ? e.message : String(e) + }) + controller.error(e) + } + }, + async cancel() { + await reader.cancel() + } + }) +} + export const googleProvider: ProviderConfig = { id: 'google', name: 'Google', @@ -12,7 +93,7 @@ export const googleProvider: ProviderConfig = { models: ['gemini-2.5-pro-exp-03-25', 'gemini-2.5-flash-preview-04-17'], defaultModel: 'gemini-2.5-pro-exp-03-25', - executeRequest: async (request: ProviderRequest): Promise => { + executeRequest: async (request: ProviderRequest): Promise => { if (!request.apiKey) { throw new Error('API key is required for Google Gemini') } @@ -24,6 +105,7 @@ export const googleProvider: ProviderConfig = { hasTools: !!request.tools?.length, toolCount: request.tools?.length || 0, hasResponseFormat: !!request.responseFormat, + streaming: !!request.stream, }) // Start execution timer for the entire provider execution @@ -90,8 +172,13 @@ export const googleProvider: ProviderConfig = { // Make the API request const initialCallTime = Date.now() + // For streaming requests, add the alt=sse parameter to the URL + const endpoint = request.stream + ? `https://generativelanguage.googleapis.com/v1beta/models/${requestedModel}:generateContent?key=${request.apiKey}&alt=sse` + : `https://generativelanguage.googleapis.com/v1beta/models/${requestedModel}:generateContent?key=${request.apiKey}` + const response = await fetch( - `https://generativelanguage.googleapis.com/v1beta/models/${requestedModel}:generateContent?key=${request.apiKey}`, + endpoint, { method: 'POST', headers: { @@ -112,6 +199,64 @@ export const googleProvider: ProviderConfig = { } const firstResponseTime = Date.now() - initialCallTime + + // Handle streaming response + if (request.stream) { + logger.info('Handling Google Gemini streaming response') + + // Create a ReadableStream from the Google Gemini stream + const stream = createReadableStreamFromGeminiStream(response) + + // Create an object that combines the stream with execution metadata + const streamingExecution: StreamingExecution = { + stream, + execution: { + success: true, + output: { + response: { + content: '', + model: request.model, + tokens: { + prompt: 0, + completion: 0, + total: 0, + }, + providerTiming: { + startTime: providerStartTimeISO, + endTime: new Date().toISOString(), + duration: firstResponseTime, + modelTime: firstResponseTime, + toolsTime: 0, + firstResponseTime, + iterations: 1, + timeSegments: [{ + type: 'model', + name: 'Initial streaming response', + startTime: initialCallTime, + endTime: initialCallTime + firstResponseTime, + duration: firstResponseTime, + }], + cost: { + total: 0.0, // Initial estimate, updated as tokens are processed + input: 0.0, + output: 0.0 + } + } + } + }, + logs: [], + metadata: { + startTime: providerStartTimeISO, + endTime: new Date().toISOString(), + duration: firstResponseTime, + }, + isStreaming: true + } + } + + return streamingExecution + } + let geminiResponse = await response.json() // Check structured output format @@ -307,7 +452,105 @@ export const googleProvider: ProviderConfig = { const nextModelStartTime = Date.now() try { - // Make the next request + // Check if we should stream the final response after tool calls + if (request.stream) { + // Create a payload for the streaming response after tool calls + const streamingPayload = { + ...payload, + contents: simplifiedMessages, + tool_config: { mode: 'AUTO' }, // Always use AUTO mode for streaming after tools + } + + // Remove any forced tool configuration to prevent issues with streaming + if ('tool_config' in streamingPayload) { + streamingPayload.tool_config = { mode: 'AUTO' }; + } + + // Make the streaming request with alt=sse parameter + const streamingResponse = await fetch( + `https://generativelanguage.googleapis.com/v1beta/models/${requestedModel}:generateContent?key=${request.apiKey}&alt=sse`, + { + method: 'POST', + headers: { + 'Content-Type': 'application/json', + }, + body: JSON.stringify(streamingPayload), + } + ) + + if (!streamingResponse.ok) { + const errorBody = await streamingResponse.text() + logger.error('Error in Gemini streaming follow-up request:', { + status: streamingResponse.status, + statusText: streamingResponse.statusText, + responseBody: errorBody + }) + throw new Error(`Gemini API streaming error: ${streamingResponse.status} ${streamingResponse.statusText}`) + } + + // Create a stream from the response + const stream = createReadableStreamFromGeminiStream(streamingResponse) + + // Calculate timing information + const nextModelEndTime = Date.now() + const thisModelTime = nextModelEndTime - nextModelStartTime + modelTime += thisModelTime + + // Add to time segments + timeSegments.push({ + type: 'model', + name: 'Final streaming response after tool calls', + startTime: nextModelStartTime, + endTime: nextModelEndTime, + duration: thisModelTime, + }) + + // Return a streaming execution with tool call information + const streamingExecution: StreamingExecution = { + stream, + execution: { + success: true, + output: { + response: { + content: '', + model: request.model, + tokens, + toolCalls: toolCalls.length > 0 ? { + list: toolCalls, + count: toolCalls.length + } : undefined, + toolResults, + providerTiming: { + startTime: providerStartTimeISO, + endTime: new Date().toISOString(), + duration: Date.now() - providerStartTime, + modelTime, + toolsTime, + firstResponseTime, + iterations: iterationCount + 1, + timeSegments, + }, + cost: { + total: (tokens.total || 0) * 0.0001, // Estimate cost based on tokens + input: (tokens.prompt || 0) * 0.0001, + output: (tokens.completion || 0) * 0.0001 + } + } + }, + logs: [], + metadata: { + startTime: providerStartTimeISO, + endTime: new Date().toISOString(), + duration: Date.now() - providerStartTime, + }, + isStreaming: true + } + } + + return streamingExecution + } + + // Make the next request for non-streaming response const nextResponse = await fetch( `https://generativelanguage.googleapis.com/v1beta/models/${requestedModel}:generateContent?key=${request.apiKey}`, { diff --git a/sim/providers/groq/index.ts b/sim/providers/groq/index.ts index 59507d7df2..891eb5ba84 100644 --- a/sim/providers/groq/index.ts +++ b/sim/providers/groq/index.ts @@ -2,9 +2,31 @@ import { Groq } from 'groq-sdk' import { createLogger } from '@/lib/logs/console-logger' import { executeTool } from '@/tools' import { ProviderConfig, ProviderRequest, ProviderResponse, TimeSegment } from '../types' +import { StreamingExecution } from '@/executor/types' const logger = createLogger('Groq Provider') +/** + * Helper to wrap Groq streaming into a browser-friendly ReadableStream + * of raw assistant text chunks. + */ +function createReadableStreamFromGroqStream(groqStream: any): ReadableStream { + return new ReadableStream({ + async start(controller) { + try { + for await (const chunk of groqStream) { + if (chunk.choices[0]?.delta?.content) { + controller.enqueue(new TextEncoder().encode(chunk.choices[0].delta.content)) + } + } + controller.close() + } catch (err) { + controller.error(err) + } + }, + }) +} + export const groqProvider: ProviderConfig = { id: 'groq', name: 'Groq', @@ -17,95 +39,161 @@ export const groqProvider: ProviderConfig = { ], defaultModel: 'groq/meta-llama/llama-4-scout-17b-16e-instruct', - executeRequest: async (request: ProviderRequest): Promise => { + executeRequest: async (request: ProviderRequest): Promise => { if (!request.apiKey) { throw new Error('API key is required for Groq') } + // Create Groq client + const groq = new Groq({ apiKey: request.apiKey }) + + // Start with an empty array for all messages + const allMessages = [] + + // Add system prompt if present + if (request.systemPrompt) { + allMessages.push({ + role: 'system', + content: request.systemPrompt, + }) + } + + // Add context if present + if (request.context) { + allMessages.push({ + role: 'user', + content: request.context, + }) + } + + // Add remaining messages + if (request.messages) { + allMessages.push(...request.messages) + } + + // Transform tools to function format if provided + const tools = request.tools?.length + ? request.tools.map((tool) => ({ + type: 'function', + function: { + name: tool.id, + description: tool.description, + parameters: tool.parameters, + }, + })) + : undefined + + // Build the request payload + const payload: any = { + model: (request.model || 'groq/meta-llama/llama-4-scout-17b-16e-instruct').replace('groq/', ''), + messages: allMessages, + } + + // Add optional parameters + if (request.temperature !== undefined) payload.temperature = request.temperature + if (request.maxTokens !== undefined) payload.max_tokens = request.maxTokens + + // Add response format for structured output if specified + if (request.responseFormat) { + payload.response_format = { + type: 'json_schema', + schema: request.responseFormat.schema || request.responseFormat, + } + } + + // Handle tools and tool usage control + if (tools?.length) { + // Filter out any tools with usageControl='none', but ignore 'force' since Groq doesn't support it + const filteredTools = tools.filter((tool) => { + const toolId = tool.function?.name + const toolConfig = request.tools?.find((t) => t.id === toolId) + // Only filter out 'none', treat 'force' as 'auto' + return toolConfig?.usageControl !== 'none' + }) + + if (filteredTools?.length) { + payload.tools = filteredTools + // Always use 'auto' for Groq, regardless of the tool_choice setting + payload.tool_choice = 'auto' + + logger.info(`Groq request configuration:`, { + toolCount: filteredTools.length, + toolChoice: 'auto', // Groq always uses auto + model: request.model || 'groq/meta-llama/llama-4-scout-17b-16e-instruct', + }) + } + } + + // EARLY STREAMING: if caller requested streaming and there are no tools to execute, + // we can directly stream the completion. + if (request.stream && (!tools || tools.length === 0)) { + logger.info('Using streaming response for Groq request (no tools)') + + // Start execution timer for the entire provider execution + const providerStartTime = Date.now() + const providerStartTimeISO = new Date(providerStartTime).toISOString() + + const streamResponse = await groq.chat.completions.create({ + ...payload, + stream: true, + }) + + // Start collecting token usage + let tokenUsage = { + prompt: 0, + completion: 0, + total: 0 + } + + // Create a StreamingExecution response with a readable stream + const streamingResult = { + stream: createReadableStreamFromGroqStream(streamResponse), + execution: { + success: true, + output: { + response: { + content: '', // Will be filled by streaming content in chat component + model: request.model || 'groq/meta-llama/llama-4-scout-17b-16e-instruct', + tokens: tokenUsage, + toolCalls: undefined, + providerTiming: { + startTime: providerStartTimeISO, + endTime: new Date().toISOString(), + duration: Date.now() - providerStartTime, + timeSegments: [{ + type: 'model', + name: 'Streaming response', + startTime: providerStartTime, + endTime: Date.now(), + duration: Date.now() - providerStartTime, + }] + }, + cost: { + total: 0.0, + input: 0.0, + output: 0.0 + } + } + }, + logs: [], // No block logs for direct streaming + metadata: { + startTime: providerStartTimeISO, + endTime: new Date().toISOString(), + duration: Date.now() - providerStartTime, + }, + isStreaming: true + } + } + + // Return the streaming execution object + return streamingResult as StreamingExecution + } + // Start execution timer for the entire provider execution const providerStartTime = Date.now() const providerStartTimeISO = new Date(providerStartTime).toISOString() try { - const groq = new Groq({ apiKey: request.apiKey }) - - // Start with an empty array for all messages - const allMessages = [] - - // Add system prompt if present - if (request.systemPrompt) { - allMessages.push({ - role: 'system', - content: request.systemPrompt, - }) - } - - // Add context if present - if (request.context) { - allMessages.push({ - role: 'user', - content: request.context, - }) - } - - // Add remaining messages - if (request.messages) { - allMessages.push(...request.messages) - } - - // Transform tools to function format if provided - const tools = request.tools?.length - ? request.tools.map((tool) => ({ - type: 'function', - function: { - name: tool.id, - description: tool.description, - parameters: tool.parameters, - }, - })) - : undefined - - // Build the request payload - const payload: any = { - model: (request.model || 'groq/meta-llama/llama-4-scout-17b-16e-instruct').replace('groq/', ''), - messages: allMessages, - } - - // Add optional parameters - if (request.temperature !== undefined) payload.temperature = request.temperature - if (request.maxTokens !== undefined) payload.max_tokens = request.maxTokens - - // Add response format for structured output if specified - if (request.responseFormat) { - payload.response_format = { - type: 'json_schema', - schema: request.responseFormat.schema || request.responseFormat, - } - } - - // Handle tools and tool usage control - if (tools?.length) { - // Filter out any tools with usageControl='none', but ignore 'force' since Groq doesn't support it - const filteredTools = tools.filter((tool) => { - const toolId = tool.function?.name - const toolConfig = request.tools?.find((t) => t.id === toolId) - // Only filter out 'none', treat 'force' as 'auto' - return toolConfig?.usageControl !== 'none' - }) - - if (filteredTools?.length) { - payload.tools = filteredTools - // Always use 'auto' for Groq, regardless of the tool_choice setting - payload.tool_choice = 'auto' - - logger.info(`Groq request configuration:`, { - toolCount: filteredTools.length, - toolChoice: 'auto', // Groq always uses auto - model: request.model || 'groq/meta-llama/llama-4-scout-17b-16e-instruct', - }) - } - } - // Make the initial API request const initialCallTime = Date.now() @@ -267,6 +355,70 @@ export const groqProvider: ProviderConfig = { logger.error('Error in Groq request:', { error }) } + // After all tool processing complete, if streaming was requested and we have messages, use streaming for the final response + if (request.stream && iterationCount > 0) { + logger.info('Using streaming for final Groq response after tool calls') + + // When streaming after tool calls with forced tools, make sure tool_choice is set to 'auto' + // This prevents the API from trying to force tool usage again in the final streaming response + const streamingPayload = { + ...payload, + messages: currentMessages, + tool_choice: 'auto', // Always use 'auto' for the streaming response after tool calls + stream: true, + } + + const streamResponse = await groq.chat.completions.create(streamingPayload) + + // Create a StreamingExecution response with all collected data + const streamingResult = { + stream: createReadableStreamFromGroqStream(streamResponse), + execution: { + success: true, + output: { + response: { + content: '', // Will be filled by the callback + model: request.model || 'groq/meta-llama/llama-4-scout-17b-16e-instruct', + tokens: { + prompt: tokens.prompt, + completion: tokens.completion, + total: tokens.total, + }, + toolCalls: toolCalls.length > 0 ? { + list: toolCalls, + count: toolCalls.length + } : undefined, + providerTiming: { + startTime: providerStartTimeISO, + endTime: new Date().toISOString(), + duration: Date.now() - providerStartTime, + modelTime: modelTime, + toolsTime: toolsTime, + firstResponseTime: firstResponseTime, + iterations: iterationCount + 1, + timeSegments: timeSegments, + }, + cost: { + total: (tokens.total || 0) * 0.0001, + input: (tokens.prompt || 0) * 0.0001, + output: (tokens.completion || 0) * 0.0001 + } + } + }, + logs: [], // No block logs at provider level + metadata: { + startTime: providerStartTimeISO, + endTime: new Date().toISOString(), + duration: Date.now() - providerStartTime, + }, + isStreaming: true + } + } + + // Return the streaming execution object + return streamingResult as StreamingExecution + } + // Calculate overall timing const providerEndTime = Date.now() const providerEndTimeISO = new Date(providerEndTime).toISOString() diff --git a/sim/providers/index.ts b/sim/providers/index.ts index d5e1fec90e..60cdb1914b 100644 --- a/sim/providers/index.ts +++ b/sim/providers/index.ts @@ -2,6 +2,7 @@ import { createLogger } from '@/lib/logs/console-logger' import { supportsTemperature } from './model-capabilities' import { ProviderRequest, ProviderResponse } from './types' import { calculateCost, generateStructuredOutputInstructions, getProvider } from './utils' +import { StreamingExecution } from '@/executor/types' const logger = createLogger('Providers') @@ -18,10 +19,20 @@ function sanitizeRequest(request: ProviderRequest): ProviderRequest { return sanitizedRequest } +// Type guard for StreamingExecution +function isStreamingExecution(response: any): response is StreamingExecution { + return response && typeof response === 'object' && 'stream' in response && 'execution' in response +} + +// Type guard for ReadableStream +function isReadableStream(response: any): response is ReadableStream { + return response instanceof ReadableStream +} + export async function executeProviderRequest( providerId: string, request: ProviderRequest -): Promise { +): Promise { logger.info(`Executing request with provider: ${providerId}`, { hasResponseFormat: !!request.responseFormat, model: request.model, @@ -65,6 +76,18 @@ export async function executeProviderRequest( // Execute the request using the provider's implementation const response = await provider.executeRequest(sanitizedRequest) + // If we received a StreamingExecution or ReadableStream, just pass it through + if (isStreamingExecution(response)) { + logger.info(`Provider returned StreamingExecution`) + return response + } + + if (isReadableStream(response)) { + logger.info(`Provider returned ReadableStream`) + return response + } + + // At this point, we know we have a ProviderResponse logger.info(`Provider response received`, { contentLength: response.content ? response.content.length : 0, model: response.model, diff --git a/sim/providers/openai/index.ts b/sim/providers/openai/index.ts index d28d4b0f92..1281ab0843 100644 --- a/sim/providers/openai/index.ts +++ b/sim/providers/openai/index.ts @@ -3,9 +3,47 @@ import { createLogger } from '@/lib/logs/console-logger' import { executeTool } from '@/tools' import { ProviderConfig, ProviderRequest, ProviderResponse, TimeSegment } from '../types' import { prepareToolsWithUsageControl, trackForcedToolUsage } from '../utils' +import { StreamingExecution } from '@/executor/types' const logger = createLogger('OpenAI Provider') +/** + * Helper function to convert an OpenAI stream to a standard ReadableStream + * and collect completion metrics + */ +function createReadableStreamFromOpenAIStream(openaiStream: any, onComplete?: (content: string, usage?: any) => void): ReadableStream { + let fullContent = '' + let usageData: any = null + + return new ReadableStream({ + async start(controller) { + try { + for await (const chunk of openaiStream) { + // Check for usage data in the final chunk + if (chunk.usage) { + usageData = chunk.usage + } + + const content = chunk.choices[0]?.delta?.content || '' + if (content) { + fullContent += content + controller.enqueue(new TextEncoder().encode(content)) + } + } + + // Once stream is complete, call the completion callback with the final content and usage + if (onComplete) { + onComplete(fullContent, usageData) + } + + controller.close() + } catch (error) { + controller.error(error) + } + } + }) +} + /** * OpenAI provider configuration */ @@ -17,7 +55,7 @@ export const openaiProvider: ProviderConfig = { models: ['gpt-4o', 'o1', 'o3', 'o4-mini'], defaultModel: 'gpt-4o', - executeRequest: async (request: ProviderRequest): Promise => { + executeRequest: async (request: ProviderRequest): Promise => { logger.info('Preparing OpenAI request', { model: request.model || 'gpt-4o', hasSystemPrompt: !!request.systemPrompt, @@ -25,6 +63,7 @@ export const openaiProvider: ProviderConfig = { hasTools: !!request.tools?.length, toolCount: request.tools?.length || 0, hasResponseFormat: !!request.responseFormat, + stream: !!request.stream, }) // API key is now handled server-side before this function is called @@ -124,6 +163,96 @@ export const openaiProvider: ProviderConfig = { const providerStartTimeISO = new Date(providerStartTime).toISOString() try { + // Check if we can stream directly (no tools required) + if (request.stream && (!tools || tools.length === 0)) { + logger.info('Using streaming response for OpenAI request') + + // Create a streaming request with token usage tracking + const streamResponse = await openai.chat.completions.create({ + ...payload, + stream: true, + stream_options: { include_usage: true }, + }) + + // Start collecting token usage from the stream + let tokenUsage = { + prompt: 0, + completion: 0, + total: 0 + } + + let streamContent = '' + + // Create a StreamingExecution response with a callback to update content and tokens + const streamingResult = { + stream: createReadableStreamFromOpenAIStream(streamResponse, (content, usage) => { + // Update the execution data with the final content and token usage + streamContent = content + streamingResult.execution.output.response.content = content + + // Update the timing information with the actual completion time + const streamEndTime = Date.now() + const streamEndTimeISO = new Date(streamEndTime).toISOString() + + if (streamingResult.execution.output.response.providerTiming) { + streamingResult.execution.output.response.providerTiming.endTime = streamEndTimeISO + streamingResult.execution.output.response.providerTiming.duration = streamEndTime - providerStartTime + + // Update the time segment as well + if (streamingResult.execution.output.response.providerTiming.timeSegments?.[0]) { + streamingResult.execution.output.response.providerTiming.timeSegments[0].endTime = streamEndTime + streamingResult.execution.output.response.providerTiming.timeSegments[0].duration = streamEndTime - providerStartTime + } + } + + // Update token usage if available from the stream + if (usage) { + const newTokens = { + prompt: usage.prompt_tokens || tokenUsage.prompt, + completion: usage.completion_tokens || tokenUsage.completion, + total: usage.total_tokens || tokenUsage.total + } + + streamingResult.execution.output.response.tokens = newTokens + } + // We don't need to estimate tokens here as execution-logger.ts will handle that + }), + execution: { + success: true, + output: { + response: { + content: '', // Will be filled by the stream completion callback + model: request.model, + tokens: tokenUsage, + toolCalls: undefined, + providerTiming: { + startTime: providerStartTimeISO, + endTime: new Date().toISOString(), + duration: Date.now() - providerStartTime, + timeSegments: [{ + type: 'model', + name: 'Streaming response', + startTime: providerStartTime, + endTime: Date.now(), + duration: Date.now() - providerStartTime, + }] + } + // Cost will be calculated in execution-logger.ts + } + }, + logs: [], // No block logs for direct streaming + metadata: { + startTime: providerStartTimeISO, + endTime: new Date().toISOString(), + duration: Date.now() - providerStartTime, + } + } + } as StreamingExecution + + // Return the streaming execution object with explicit casting + return streamingResult as StreamingExecution + } + // Make the initial API request const initialCallTime = Date.now() @@ -158,6 +287,7 @@ export const openaiProvider: ProviderConfig = { const firstResponseTime = Date.now() - initialCallTime let content = currentResponse.choices[0]?.message?.content || '' + // Collect token information but don't calculate costs - that will be done in execution-logger.ts let tokens = { prompt: currentResponse.usage?.prompt_tokens || 0, completion: currentResponse.usage?.completion_tokens || 0, @@ -343,6 +473,83 @@ export const openaiProvider: ProviderConfig = { iterationCount++ } + // After all tool processing complete, if streaming was requested and we have messages, use streaming for the final response + if (request.stream && iterationCount > 0) { + logger.info('Using streaming for final response after tool calls') + + // When streaming after tool calls with forced tools, make sure tool_choice is set to 'auto' + // This prevents OpenAI API from trying to force tool usage again in the final streaming response + const streamingPayload = { + ...payload, + messages: currentMessages, + tool_choice: 'auto', // Always use 'auto' for the streaming response after tool calls + stream: true, + stream_options: { include_usage: true }, + } + + const streamResponse = await openai.chat.completions.create(streamingPayload) + + // Create the StreamingExecution object with all collected data + let streamContent = '' + + const streamingResult = { + stream: createReadableStreamFromOpenAIStream(streamResponse, (content, usage) => { + // Update the execution data with the final content and token usage + streamContent = content + streamingResult.execution.output.response.content = content + + // Update token usage if available from the stream + if (usage) { + const newTokens = { + prompt: usage.prompt_tokens || tokens.prompt, + completion: usage.completion_tokens || tokens.completion, + total: usage.total_tokens || tokens.total + } + + streamingResult.execution.output.response.tokens = newTokens + } + }), + execution: { + success: true, + output: { + response: { + content: '', // Will be filled by the callback + model: request.model, + tokens: { + prompt: tokens.prompt, + completion: tokens.completion, + total: tokens.total, + }, + toolCalls: toolCalls.length > 0 ? { + list: toolCalls, + count: toolCalls.length + } : undefined, + providerTiming: { + startTime: providerStartTimeISO, + endTime: new Date().toISOString(), + duration: Date.now() - providerStartTime, + modelTime: modelTime, + toolsTime: toolsTime, + firstResponseTime: firstResponseTime, + iterations: iterationCount + 1, + timeSegments: timeSegments, + } + // Cost will be calculated in execution-logger.ts + } + }, + logs: [], // No block logs at provider level + metadata: { + startTime: providerStartTimeISO, + endTime: new Date().toISOString(), + duration: Date.now() - providerStartTime, + } + } + } as StreamingExecution + + // Return the streaming execution object with explicit casting + return streamingResult as StreamingExecution + } + // Calculate overall timing const providerEndTime = Date.now() const providerEndTimeISO = new Date(providerEndTime).toISOString() @@ -364,6 +571,7 @@ export const openaiProvider: ProviderConfig = { iterations: iterationCount + 1, timeSegments: timeSegments, }, + // We're not calculating cost here as it will be handled in execution-logger.ts } } catch (error) { // Include timing information even for errors @@ -389,3 +597,4 @@ export const openaiProvider: ProviderConfig = { } }, } + diff --git a/sim/providers/types.ts b/sim/providers/types.ts index b749ef60f7..20c267ed2f 100644 --- a/sim/providers/types.ts +++ b/sim/providers/types.ts @@ -1,3 +1,5 @@ +import { StreamingExecution } from '@/executor/types' + export type ProviderId = | 'openai' | 'anthropic' @@ -42,7 +44,9 @@ export interface ProviderConfig { models: string[] defaultModel: string initialize?: () => Promise - executeRequest?: (request: ProviderRequest) => Promise + executeRequest: ( + request: ProviderRequest + ) => Promise | StreamingExecution> } export interface FunctionCallResponse { @@ -142,6 +146,7 @@ export interface ProviderRequest { } local_execution?: boolean workflowId?: string // Optional workflow ID for authentication context + stream?: boolean } // Map of provider IDs to their configurations diff --git a/sim/providers/utils.ts b/sim/providers/utils.ts index bd744631fa..b8610a6fa0 100644 --- a/sim/providers/utils.ts +++ b/sim/providers/utils.ts @@ -1,6 +1,6 @@ import { createLogger } from '@/lib/logs/console-logger' import { useCustomToolsStore } from '@/stores/custom-tools/store' -import { isProd, getCostMultiplier } from '@/lib/environment' +import { getCostMultiplier } from '@/lib/environment' import { anthropicProvider } from './anthropic' import { cerebrasProvider } from './cerebras' import { deepseekProvider } from './deepseek' diff --git a/sim/providers/xai/index.ts b/sim/providers/xai/index.ts index 2986840d90..6a7b8ce72e 100644 --- a/sim/providers/xai/index.ts +++ b/sim/providers/xai/index.ts @@ -2,10 +2,33 @@ import OpenAI from 'openai' import { createLogger } from '@/lib/logs/console-logger' import { executeTool } from '@/tools' import { ProviderConfig, ProviderRequest, ProviderResponse, TimeSegment } from '../types' +import { StreamingExecution } from '@/executor/types' import { prepareToolsWithUsageControl, trackForcedToolUsage } from '../utils' const logger = createLogger('XAI Provider') +/** + * Helper to wrap XAI (OpenAI-compatible) streaming into a browser-friendly + * ReadableStream of raw assistant text chunks. + */ +function createReadableStreamFromXAIStream(xaiStream: any): ReadableStream { + return new ReadableStream({ + async start(controller) { + try { + for await (const chunk of xaiStream) { + const content = chunk.choices[0]?.delta?.content || '' + if (content) { + controller.enqueue(new TextEncoder().encode(content)) + } + } + controller.close() + } catch (err) { + controller.error(err) + } + }, + }) +} + export const xAIProvider: ProviderConfig = { id: 'xai', name: 'xAI', @@ -14,108 +37,178 @@ export const xAIProvider: ProviderConfig = { models: ['grok-3-latest', 'grok-3-fast-latest'], defaultModel: 'grok-3-latest', - executeRequest: async (request: ProviderRequest): Promise => { + executeRequest: async (request: ProviderRequest): Promise => { if (!request.apiKey) { throw new Error('API key is required for xAI') } + // Initialize OpenAI client for xAI + const xai = new OpenAI({ + apiKey: request.apiKey, + baseURL: 'https://api.x.ai/v1', + }) + + // Prepare messages + const allMessages = [] + + if (request.systemPrompt) { + allMessages.push({ + role: 'system', + content: request.systemPrompt, + }) + } + + if (request.context) { + allMessages.push({ + role: 'user', + content: request.context, + }) + } + + if (request.messages) { + allMessages.push(...request.messages) + } + + // Set up tools + const tools = request.tools?.length + ? request.tools.map((tool) => ({ + type: 'function', + function: { + name: tool.id, + description: tool.description, + parameters: tool.parameters, + }, + })) + : undefined + + // Build the request payload + const payload: any = { + model: request.model || 'grok-3-latest', + messages: allMessages, + } + + if (request.temperature !== undefined) payload.temperature = request.temperature + if (request.maxTokens !== undefined) payload.max_tokens = request.maxTokens + + if (request.responseFormat) { + payload.response_format = { + type: 'json_schema', + json_schema: { + name: request.responseFormat.name || 'structured_response', + schema: request.responseFormat.schema || request.responseFormat, + strict: request.responseFormat.strict !== false, + }, + } + + if (allMessages.length > 0 && allMessages[0].role === 'system') { + allMessages[0].content = `${allMessages[0].content}\n\nYou MUST respond with a valid JSON object. DO NOT include any other text, explanations, or markdown formatting in your response - ONLY the JSON object.` + } else { + allMessages.unshift({ + role: 'system', + content: `You MUST respond with a valid JSON object. DO NOT include any other text, explanations, or markdown formatting in your response - ONLY the JSON object.`, + }) + } + } + + // Handle tools and tool usage control + let preparedTools: ReturnType | null = null + + if (tools?.length) { + preparedTools = prepareToolsWithUsageControl(tools, request.tools, logger, 'xai') + const { tools: filteredTools, toolChoice } = preparedTools + + if (filteredTools?.length && toolChoice) { + payload.tools = filteredTools + payload.tool_choice = toolChoice + + logger.info(`XAI request configuration:`, { + toolCount: filteredTools.length, + toolChoice: + typeof toolChoice === 'string' + ? toolChoice + : toolChoice.type === 'function' + ? `force:${toolChoice.function.name}` + : toolChoice.type === 'tool' + ? `force:${toolChoice.name}` + : toolChoice.type === 'any' + ? `force:${toolChoice.any?.name || 'unknown'}` + : 'unknown', + model: request.model || 'grok-3-latest', + }) + } + } + + // EARLY STREAMING: if caller requested streaming and there are no tools to execute, + // we can directly stream the completion. + if (request.stream && (!tools || tools.length === 0)) { + logger.info('Using streaming response for XAI request (no tools)') + + // Start execution timer for the entire provider execution + const providerStartTime = Date.now() + const providerStartTimeISO = new Date(providerStartTime).toISOString() + + const streamResponse = await xai.chat.completions.create({ + ...payload, + stream: true, + }) + + // Start collecting token usage + let tokenUsage = { + prompt: 0, + completion: 0, + total: 0 + } + + // Create a StreamingExecution response with a readable stream + const streamingResult = { + stream: createReadableStreamFromXAIStream(streamResponse), + execution: { + success: true, + output: { + response: { + content: '', // Will be filled by streaming content in chat component + model: request.model || 'grok-3-latest', + tokens: tokenUsage, + toolCalls: undefined, + providerTiming: { + startTime: providerStartTimeISO, + endTime: new Date().toISOString(), + duration: Date.now() - providerStartTime, + timeSegments: [{ + type: 'model', + name: 'Streaming response', + startTime: providerStartTime, + endTime: Date.now(), + duration: Date.now() - providerStartTime, + }] + }, + // Estimate token cost + cost: { + total: 0.0, + input: 0.0, + output: 0.0 + } + } + }, + logs: [], // No block logs for direct streaming + metadata: { + startTime: providerStartTimeISO, + endTime: new Date().toISOString(), + duration: Date.now() - providerStartTime, + }, + isStreaming: true + } + } + + // Return the streaming execution object + return streamingResult as StreamingExecution + } + // Start execution timer for the entire provider execution const providerStartTime = Date.now() const providerStartTimeISO = new Date(providerStartTime).toISOString() try { - const xai = new OpenAI({ - apiKey: request.apiKey, - baseURL: 'https://api.x.ai/v1', - }) - - const allMessages = [] - - if (request.systemPrompt) { - allMessages.push({ - role: 'system', - content: request.systemPrompt, - }) - } - - if (request.context) { - allMessages.push({ - role: 'user', - content: request.context, - }) - } - - if (request.messages) { - allMessages.push(...request.messages) - } - - const tools = request.tools?.length - ? request.tools.map((tool) => ({ - type: 'function', - function: { - name: tool.id, - description: tool.description, - parameters: tool.parameters, - }, - })) - : undefined - - const payload: any = { - model: request.model || 'grok-3-latest', - messages: allMessages, - } - - if (request.temperature !== undefined) payload.temperature = request.temperature - if (request.maxTokens !== undefined) payload.max_tokens = request.maxTokens - - if (request.responseFormat) { - payload.response_format = { - type: 'json_schema', - json_schema: { - name: request.responseFormat.name || 'structured_response', - schema: request.responseFormat.schema || request.responseFormat, - strict: request.responseFormat.strict !== false, - }, - } - - if (allMessages.length > 0 && allMessages[0].role === 'system') { - allMessages[0].content = `${allMessages[0].content}\n\nYou MUST respond with a valid JSON object. DO NOT include any other text, explanations, or markdown formatting in your response - ONLY the JSON object.` - } else { - allMessages.unshift({ - role: 'system', - content: `You MUST respond with a valid JSON object. DO NOT include any other text, explanations, or markdown formatting in your response - ONLY the JSON object.`, - }) - } - } - - // Handle tools and tool usage control - let preparedTools: ReturnType | null = null - - if (tools?.length) { - preparedTools = prepareToolsWithUsageControl(tools, request.tools, logger, 'xai') - const { tools: filteredTools, toolChoice } = preparedTools - - if (filteredTools?.length && toolChoice) { - payload.tools = filteredTools - payload.tool_choice = toolChoice - - logger.info(`XAI request configuration:`, { - toolCount: filteredTools.length, - toolChoice: - typeof toolChoice === 'string' - ? toolChoice - : toolChoice.type === 'function' - ? `force:${toolChoice.function.name}` - : toolChoice.type === 'tool' - ? `force:${toolChoice.name}` - : toolChoice.type === 'any' - ? `force:${toolChoice.any?.name || 'unknown'}` - : 'unknown', - model: request.model || 'grok-3-latest', - }) - } - } - // Make the initial API request const initialCallTime = Date.now() @@ -328,6 +421,70 @@ export const xAIProvider: ProviderConfig = { logger.error('Error in xAI request:', { error }) } + // After all tool processing complete, if streaming was requested and we have messages, use streaming for the final response + if (request.stream && iterationCount > 0) { + logger.info('Using streaming for final XAI response after tool calls') + + // When streaming after tool calls with forced tools, make sure tool_choice is set to 'auto' + // This prevents the API from trying to force tool usage again in the final streaming response + const streamingPayload = { + ...payload, + messages: currentMessages, + tool_choice: 'auto', // Always use 'auto' for the streaming response after tool calls + stream: true, + } + + const streamResponse = await xai.chat.completions.create(streamingPayload) + + // Create a StreamingExecution response with all collected data + const streamingResult = { + stream: createReadableStreamFromXAIStream(streamResponse), + execution: { + success: true, + output: { + response: { + content: '', // Will be filled by the callback + model: request.model || 'grok-3-latest', + tokens: { + prompt: tokens.prompt, + completion: tokens.completion, + total: tokens.total, + }, + toolCalls: toolCalls.length > 0 ? { + list: toolCalls, + count: toolCalls.length + } : undefined, + providerTiming: { + startTime: providerStartTimeISO, + endTime: new Date().toISOString(), + duration: Date.now() - providerStartTime, + modelTime: modelTime, + toolsTime: toolsTime, + firstResponseTime: firstResponseTime, + iterations: iterationCount + 1, + timeSegments: timeSegments, + }, + cost: { + total: (tokens.total || 0) * 0.0001, + input: (tokens.prompt || 0) * 0.0001, + output: (tokens.completion || 0) * 0.0001 + } + } + }, + logs: [], // No block logs at provider level + metadata: { + startTime: providerStartTimeISO, + endTime: new Date().toISOString(), + duration: Date.now() - providerStartTime, + }, + isStreaming: true + } + } + + // Return the streaming execution object + return streamingResult as StreamingExecution + } + // Calculate overall timing const providerEndTime = Date.now() const providerEndTimeISO = new Date(providerEndTime).toISOString() diff --git a/sim/stores/panel/chat/store.ts b/sim/stores/panel/chat/store.ts index 2792a902a3..90347d96a2 100644 --- a/sim/stores/panel/chat/store.ts +++ b/sim/stores/panel/chat/store.ts @@ -16,8 +16,9 @@ export const useChatStore = create()( set((state) => { const newMessage: ChatMessage = { ...message, - id: crypto.randomUUID(), - timestamp: new Date().toISOString(), + // Preserve provided id and timestamp if they exist; otherwise generate new ones + id: (message as any).id ?? crypto.randomUUID(), + timestamp: (message as any).timestamp ?? new Date().toISOString(), } // Keep only the last MAX_MESSAGES @@ -60,6 +61,38 @@ export const useChatStore = create()( getSelectedWorkflowOutput: (workflowId) => { return get().selectedWorkflowOutputs[workflowId] || [] }, + + appendMessageContent: (messageId, content) => { + set((state) => { + const newMessages = state.messages.map((message) => { + if (message.id === messageId) { + return { + ...message, + content: typeof message.content === 'string' + ? message.content + content + : (message.content ? String(message.content) + content : content), + } + } + return message + }) + + return { messages: newMessages } + }) + }, + + finalizeMessageStream: (messageId) => { + set((state) => { + const newMessages = state.messages.map((message) => { + if (message.id === messageId) { + const { isStreaming, ...rest } = message + return rest + } + return message + }) + + return { messages: newMessages } + }) + }, }), { name: 'chat-store', diff --git a/sim/stores/panel/chat/types.ts b/sim/stores/panel/chat/types.ts index 19d35edcbf..efccaa13ee 100644 --- a/sim/stores/panel/chat/types.ts +++ b/sim/stores/panel/chat/types.ts @@ -1,10 +1,11 @@ export interface ChatMessage { id: string - content: any - workflowId: string | null + content: string | any + workflowId: string type: 'user' | 'workflow' timestamp: string blockId?: string + isStreaming?: boolean } export interface OutputConfig { @@ -20,4 +21,6 @@ export interface ChatStore { getWorkflowMessages: (workflowId: string) => ChatMessage[] setSelectedWorkflowOutput: (workflowId: string, outputIds: string[]) => void getSelectedWorkflowOutput: (workflowId: string) => string[] + appendMessageContent: (messageId: string, content: string) => void + finalizeMessageStream: (messageId: string) => void } \ No newline at end of file diff --git a/sim/stores/panel/console/store.ts b/sim/stores/panel/console/store.ts index 822da2cd21..d58926c8f7 100644 --- a/sim/stores/panel/console/store.ts +++ b/sim/stores/panel/console/store.ts @@ -69,24 +69,57 @@ export const useConsoleStore = create()( entries: [], isOpen: false, - addConsole: (entry) => { + addConsole: (entry: Omit) => { set((state) => { - // Create a new entry with redacted API keys + // Determine early if this entry represents a streaming output + const isStreamingOutput = + (typeof ReadableStream !== 'undefined' && entry.output instanceof ReadableStream) || + (typeof entry.output === 'object' && entry.output && entry.output.isStreaming === true) || + (typeof entry.output === 'object' && entry.output && 'executionData' in entry.output && + typeof entry.output.executionData === 'object' && entry.output.executionData?.isStreaming === true) || + (typeof entry.output === 'object' && entry.output && 'stream' in entry.output) || + (typeof entry.output === 'object' && entry.output && + 'stream' in entry.output && 'execution' in entry.output) + + // Skip adding raw streaming objects that have both stream and executionData + if (typeof entry.output === 'object' && entry.output && + 'stream' in entry.output && 'executionData' in entry.output) { + // Don't add this entry - it will be processed by our explicit formatting code in executor/index.ts + return { entries: state.entries } + } + + // Also skip raw StreamingExecution objects (with stream and execution properties) + if (typeof entry.output === 'object' && entry.output && + 'stream' in entry.output && 'execution' in entry.output) { + // Don't add this entry to prevent duplicate console entries for streaming responses + return { entries: state.entries } + } + + // Create a new entry with redacted API keys (if not a stream) const redactedEntry = { ...entry } - // If the entry has output and it's an object, redact API keys - if (redactedEntry.output && typeof redactedEntry.output === 'object') { + // If output is a stream, we skip redaction (it's not an object we want to recurse into) + if (!isStreamingOutput && redactedEntry.output && typeof redactedEntry.output === 'object') { redactedEntry.output = redactApiKeys(redactedEntry.output) } - const newEntry: ConsoleEntry = { - ...redactedEntry, - id: crypto.randomUUID(), - timestamp: new Date().toISOString(), + // Create the new entry with ID and timestamp + const newEntry = { + ...redactedEntry, + id: crypto.randomUUID(), + timestamp: new Date().toISOString() } // Keep only the last MAX_ENTRIES - const newEntries = [newEntry, ...state.entries].slice(0, MAX_ENTRIES) + const newEntries = [ + newEntry, + ...state.entries, + ].slice(0, MAX_ENTRIES) + + // If the block produced a streaming output, skip automatic chat message creation + if (isStreamingOutput) { + return { entries: newEntries } + } // Check if this block matches a selected workflow output if (entry.workflowId && entry.blockName) { @@ -116,7 +149,12 @@ export const useConsoleStore = create()( // Format the value appropriately for display let formattedValue: string - if (specificValue === undefined) { + // For streaming responses, use empty string and set isStreaming flag + if (isStreamingOutput) { + // Skip adding a message since we'll handle streaming in workflow execution + // This prevents the "Output value not found" message for streams + continue + } else if (specificValue === undefined) { formattedValue = "Output value not found" } else if (typeof specificValue === 'object') { formattedValue = JSON.stringify(specificValue, null, 2) @@ -124,12 +162,18 @@ export const useConsoleStore = create()( formattedValue = String(specificValue) } + // Skip empty content messages (important for preventing empty entries) + if (!formattedValue || formattedValue.trim() === '') { + continue + } + // Add the specific value to chat, not the whole output chatStore.addMessage({ content: formattedValue, workflowId: entry.workflowId, type: 'workflow', blockId: entry.blockId, + isStreaming: isStreamingOutput, }) } } @@ -138,6 +182,9 @@ export const useConsoleStore = create()( return { entries: newEntries } }) + + // Return the created entry by finding it in the updated store + return get().entries[0] }, clearConsole: (workflowId: string | null) => { @@ -155,6 +202,22 @@ export const useConsoleStore = create()( toggleConsole: () => { set((state) => ({ isOpen: !state.isOpen })) }, + + updateConsole: (entryId: string, updatedData: Partial>) => { + set((state) => { + const updatedEntries = state.entries.map(entry => { + if (entry.id === entryId) { + return { + ...entry, + ...updatedData, + output: updatedData.output ? redactApiKeys(updatedData.output) : entry.output, + } + } + return entry + }) + return { entries: updatedEntries } + }) + }, }), { name: 'console-store', diff --git a/sim/stores/panel/console/types.ts b/sim/stores/panel/console/types.ts index 9f67da42fd..f1e1160c1e 100644 --- a/sim/stores/panel/console/types.ts +++ b/sim/stores/panel/console/types.ts @@ -16,8 +16,9 @@ export interface ConsoleEntry { export interface ConsoleStore { entries: ConsoleEntry[] isOpen: boolean - addConsole: (entry: Omit) => void + addConsole: (entry: Omit) => ConsoleEntry clearConsole: (workflowId: string | null) => void getWorkflowEntries: (workflowId: string) => ConsoleEntry[] toggleConsole: () => void + updateConsole: (entryId: string, updatedData: Partial>) => void }