From b58b6532aa0de8b539e770c36bf496a005635b60 Mon Sep 17 00:00:00 2001 From: Emir Karabeg Date: Tue, 11 Mar 2025 02:35:06 -0700 Subject: [PATCH] fix: waiting for env vars to resolve before executing --- app/api/webhooks/trigger/[path]/route.ts | 84 ++++++++++++++++++++---- executor/handlers.ts | 61 ++++++++++++++++- providers/openai/index.ts | 10 +++ stores/workflows/utils.ts | 78 ++++++++++++++++++++++ 4 files changed, 217 insertions(+), 16 deletions(-) diff --git a/app/api/webhooks/trigger/[path]/route.ts b/app/api/webhooks/trigger/[path]/route.ts index c516932371..ee2dfa6452 100644 --- a/app/api/webhooks/trigger/[path]/route.ts +++ b/app/api/webhooks/trigger/[path]/route.ts @@ -3,7 +3,7 @@ import { and, eq } from 'drizzle-orm' import { v4 as uuidv4 } from 'uuid' import { persistExecutionError, persistExecutionLogs } from '@/lib/logging' import { decryptSecret } from '@/lib/utils' -import { mergeSubblockState } from '@/stores/workflows/utils' +import { mergeSubblockState, mergeSubblockStateAsync } from '@/stores/workflows/utils' import { db } from '@/db' import { environment, webhook, workflow } from '@/db/schema' import { Executor } from '@/executor' @@ -225,8 +225,10 @@ export async function POST( const state = foundWorkflow.state as any const { blocks, edges, loops } = state - // Use the same execution flow as in manual executions - const mergedStates = mergeSubblockState(blocks) + // Use the async version of mergeSubblockState to ensure all values are properly resolved + console.log(`[Webhook Debug] Merging subblock states for workflow ${foundWorkflow.id}...`) + const mergedStates = await mergeSubblockStateAsync(blocks, foundWorkflow.id) + console.log(`[Webhook Debug] Subblock states merged successfully`) // Retrieve environment variables for this user const [userEnv] = await db @@ -236,19 +238,48 @@ export async function POST( .limit(1) // Create a map of decrypted environment variables - const decryptedEnvVars: Record = {} + let decryptedEnvVars: Record = {} if (userEnv) { - for (const [key, encryptedValue] of Object.entries( - userEnv.variables as Record - )) { - try { - const { decrypted } = await decryptSecret(encryptedValue) - decryptedEnvVars[key] = decrypted - } catch (error: any) { - console.error(`Failed to decrypt ${key}:`, error) - throw new Error(`Failed to decrypt environment variable "${key}": ${error.message}`) + const decryptionPromises = Object.entries(userEnv.variables as Record).map( + async ([key, encryptedValue]) => { + try { + const { decrypted } = await decryptSecret(encryptedValue) + return [key, decrypted] as const + } catch (error: any) { + console.error(`Failed to decrypt ${key}:`, error) + throw new Error(`Failed to decrypt environment variable "${key}": ${error.message}`) + } } + ) + + const decryptedEntries = await Promise.all(decryptionPromises) + decryptedEnvVars = Object.fromEntries(decryptedEntries) + } + + // Log environment variables for debugging (masking sensitive values) + console.log('Available environment variables before execution:') + Object.keys(decryptedEnvVars).forEach((key) => { + const value = decryptedEnvVars[key] + // Show full value for OpenAI API key for debugging, mask others + if (key === 'OPENAI_API_KEY') { + console.log(` ${key}: ${value} (FULL KEY FOR DEBUGGING)`) + } else { + // Mask sensitive values like API keys, showing only first and last few characters + const maskedValue = + key.toLowerCase().includes('key') || + key.toLowerCase().includes('secret') || + key.toLowerCase().includes('token') + ? `${value.substring(0, 4)}...${value.substring(value.length - 4)}` + : value + console.log(` ${key}: ${maskedValue}`) } + }) + + // Check specifically for OpenAI API key + if (!decryptedEnvVars['OPENAI_API_KEY']) { + console.warn('WARNING: OPENAI_API_KEY is missing from environment variables!') + } else { + console.log('OPENAI_API_KEY is present in environment variables') } // Process the block states to extract values from subBlocks @@ -266,6 +297,22 @@ export async function POST( {} as Record> ) + // Log block states for debugging + console.log(`[Webhook Debug] Block states after async merging:`) + Object.entries(currentBlockStates).forEach(([blockId, state]) => { + // Check for agent blocks specifically + if (blocks[blockId]?.type === 'agent') { + console.log(`[Webhook Debug] Agent block ${blockId} state:`, JSON.stringify(state, null, 2)) + // Check for null values + const nullKeys = Object.entries(state) + .filter(([_, value]) => value === null) + .map(([key]) => key) + if (nullKeys.length > 0) { + console.warn(`[Webhook Debug] Agent block ${blockId} has null values for:`, nullKeys) + } + } + }) + // Serialize and execute the workflow const serializedWorkflow = new Serializer().serializeWorkflow(mergedStates as any, edges, loops) @@ -311,6 +358,17 @@ export async function POST( {} as Record> ) + // Remove the timeout and use the already awaited states + console.log('[Webhook Debug] Final processed block states before execution:') + Object.entries(processedBlockStates).forEach(([blockId, state]) => { + if (blocks[blockId]?.type === 'agent') { + console.log( + `[Webhook Debug] Agent block ${blockId} final state:`, + JSON.stringify(state, null, 2) + ) + } + }) + const executor = new Executor( serializedWorkflow, processedBlockStates, // Use the processed block states diff --git a/executor/handlers.ts b/executor/handlers.ts index d70cb0986f..8345f9fbb8 100644 --- a/executor/handlers.ts +++ b/executor/handlers.ts @@ -109,6 +109,49 @@ export class AgentBlockHandler implements BlockHandler { console.log(`[AgentBlockHandler Debug] Executing agent block: ${block.id}`) console.log(`[AgentBlockHandler Debug] Block inputs:`, JSON.stringify(inputs, null, 2)) + // Check for null values and try to resolve from environment variables + const nullInputs = Object.entries(inputs) + .filter(([_, value]) => value === null) + .map(([key]) => key) + + if (nullInputs.length > 0) { + console.warn(`[AgentBlockHandler Debug] Null inputs detected:`, nullInputs) + + // Check if we can resolve API key from environment variables + if (nullInputs.includes('apiKey') && context.environmentVariables) { + console.log( + `[AgentBlockHandler Debug] Attempting to resolve API key from environment variables` + ) + console.log( + `[AgentBlockHandler Debug] Available env vars:`, + Object.keys(context.environmentVariables) + ) + + // Try different possible environment variable names for OpenAI API key + const possibleEnvVarNames = ['OPENAI_API_KEY', 'openai_api_key', 'OPENAI_KEY', 'openai_key'] + for (const envVar of possibleEnvVarNames) { + if (context.environmentVariables[envVar]) { + console.log( + `[AgentBlockHandler Debug] Found API key in environment variable: ${envVar}` + ) + inputs.apiKey = context.environmentVariables[envVar] + break + } + } + + // Check if we successfully resolved the API key + if (inputs.apiKey) { + console.log( + `[AgentBlockHandler Debug] Successfully resolved API key from environment variables` + ) + } else { + console.error( + `[AgentBlockHandler Debug] Failed to resolve API key from environment variables` + ) + } + } + } + // Parse response format if provided let responseFormat: any = undefined if (inputs.responseFormat) { @@ -213,8 +256,8 @@ export class AgentBlockHandler implements BlockHandler { ).filter((t: any): t is NonNullable => t !== null) : [] - // Ensure context is properly formatted for the provider - const response = await executeProviderRequest(providerId, { + // Debug request before sending to provider + const providerRequest = { model, systemPrompt: inputs.systemPrompt, context: Array.isArray(inputs.context) @@ -225,10 +268,22 @@ export class AgentBlockHandler implements BlockHandler { tools: formattedTools.length > 0 ? formattedTools : undefined, temperature: inputs.temperature, maxTokens: inputs.maxTokens, - apiKey: inputs.apiKey, + apiKey: inputs.apiKey || context.environmentVariables?.OPENAI_API_KEY, responseFormat, + } + + console.log(`[AgentBlockHandler Debug] Provider request details:`, { + model: providerRequest.model, + hasSystemPrompt: !!providerRequest.systemPrompt, + hasContext: !!providerRequest.context, + hasTools: !!providerRequest.tools, + hasApiKey: !!providerRequest.apiKey, + apiKeyFirstChars: providerRequest.apiKey ? providerRequest.apiKey.substring(0, 4) : 'NONE', }) + // Ensure context is properly formatted for the provider + const response = await executeProviderRequest(providerId, providerRequest) + console.log(`[AgentBlockHandler Debug] Provider response:`, { content: response.content, contentType: typeof response.content, diff --git a/providers/openai/index.ts b/providers/openai/index.ts index 0370a4729d..ebd5427d40 100644 --- a/providers/openai/index.ts +++ b/providers/openai/index.ts @@ -11,10 +11,20 @@ export const openaiProvider: ProviderConfig = { defaultModel: 'gpt-4o', executeRequest: async (request: ProviderRequest): Promise => { + console.log('Full request:', request) if (!request.apiKey) { + console.error('OpenAI API key missing in request. Request details:', { + hasModel: !!request.model, + hasSystemPrompt: !!request.systemPrompt, + hasMessages: !!request.messages, + hasTools: !!request.tools, + requestKeys: Object.keys(request), + }) throw new Error('API key is required for OpenAI') } + console.log('OpenAI API key found in request (first 4 chars):', request.apiKey.substring(0, 4)) + const openai = new OpenAI({ apiKey: request.apiKey, dangerouslyAllowBrowser: true, diff --git a/stores/workflows/utils.ts b/stores/workflows/utils.ts index 8a9d048acc..f00d481e36 100644 --- a/stores/workflows/utils.ts +++ b/stores/workflows/utils.ts @@ -68,3 +68,81 @@ export function mergeSubblockState( {} as Record ) } + +/** + * Asynchronously merges workflow block states with subblock values + * Ensures all values are properly resolved before returning + * + * @param blocks - Block configurations from workflow store + * @param workflowId - ID of the workflow to merge values for + * @param blockId - Optional specific block ID to merge (merges all if not provided) + * @returns Promise resolving to merged block states with updated values + */ +export async function mergeSubblockStateAsync( + blocks: Record, + workflowId?: string, + blockId?: string +): Promise> { + const blocksToProcess = blockId ? { [blockId]: blocks[blockId] } : blocks + const subBlockStore = useSubBlockStore.getState() + + // Process blocks in parallel for better performance + const processedBlockEntries = await Promise.all( + Object.entries(blocksToProcess).map(async ([id, block]) => { + // Skip if block is undefined or doesn't have subBlocks + if (!block || !block.subBlocks) { + return [id, block] as const + } + + // Process all subblocks in parallel + const subBlockEntries = await Promise.all( + Object.entries(block.subBlocks).map(async ([subBlockId, subBlock]) => { + // Skip if subBlock is undefined + if (!subBlock) { + return [subBlockId, subBlock] as const + } + + // Get the stored value for this subblock + let storedValue = null + + // If workflowId is provided, use it to get the value + if (workflowId) { + // Try to get the value from the subblock store for this specific workflow + const workflowValues = subBlockStore.workflowValues[workflowId] + if (workflowValues && workflowValues[id]) { + storedValue = workflowValues[id][subBlockId] + } + } else { + // Fall back to the active workflow if no workflowId is provided + storedValue = subBlockStore.getValue(id, subBlockId) + } + + // Create a new subblock object with the same structure but updated value + return [ + subBlockId, + { + ...subBlock, + value: + storedValue !== undefined && storedValue !== null ? storedValue : subBlock.value, + }, + ] as const + }) + ) + + // Convert entries back to an object + const mergedSubBlocks = Object.fromEntries(subBlockEntries) as Record + + // Return the full block state with updated subBlocks + return [ + id, + { + ...block, + subBlocks: mergedSubBlocks, + }, + ] as const + }) + ) + + // Convert entries back to an object + return Object.fromEntries(processedBlockEntries) as Record +}