fix(api): passing body into workflow

This commit is contained in:
Emir Karabeg
2025-04-01 13:53:32 -07:00
parent 7dde4191c7
commit 24b487d1fa
4 changed files with 275 additions and 60 deletions
+38 -3
View File
@@ -37,6 +37,21 @@ async function executeWorkflow(workflow: any, requestId: string, input?: any) {
throw new Error('Workflow is already running')
}
// Log input to help debug
logger.info(`[${requestId}] Executing workflow with input:`,
input ? JSON.stringify(input, null, 2) : 'No input provided');
// Validate and structure input for maximum compatibility
let processedInput = input;
if (input && typeof input === 'object') {
// Ensure input is properly structured for the starter block
if (input.input === undefined) {
// If input is not already nested, structure it properly
processedInput = { input: input };
logger.info(`[${requestId}] Restructured input for workflow:`, JSON.stringify(processedInput, null, 2));
}
}
try {
runningExecutions.add(workflowId)
logger.info(`[${requestId}] Starting workflow execution: ${workflowId}`)
@@ -176,7 +191,7 @@ async function executeWorkflow(workflow: any, requestId: string, input?: any) {
serializedWorkflow,
processedBlockStates,
decryptedEnvVars,
input,
processedInput,
workflowVariables
)
@@ -262,9 +277,29 @@ export async function POST(request: NextRequest, { params }: { params: Promise<{
}
const bodyText = await request.text()
const body = bodyText ? JSON.parse(bodyText) : {}
logger.info(`[${requestId}] Raw request body:`, bodyText)
let body = {}
if (bodyText && bodyText.trim()) {
try {
body = JSON.parse(bodyText)
logger.info(`[${requestId}] Parsed request body:`, JSON.stringify(body, null, 2))
} catch (error) {
logger.error(`[${requestId}] Failed to parse request body:`, error)
return createErrorResponse('Invalid JSON in request body', 400, 'INVALID_JSON')
}
} else {
logger.info(`[${requestId}] No request body provided`)
}
const result = await executeWorkflow(validation.workflow, requestId, body)
// Don't double-nest the input if it's already structured
const hasContent = Object.keys(body).length > 0;
const input = hasContent ? { input: body } : {};
logger.info(`[${requestId}] Input passed to workflow:`, JSON.stringify(input, null, 2))
// Execute workflow with the structured input
const result = await executeWorkflow(validation.workflow, requestId, input)
return createSuccessResponse(result)
} catch (error: any) {
logger.error(`[${requestId}] Error executing workflow: ${id}`, error)
+27 -2
View File
@@ -819,8 +819,34 @@ export class ApiBlockHandler implements BlockHandler {
}
try {
let processedInputs = { ...inputs };
// Handle body specifically to ensure it's properly processed for API requests
if (processedInputs.body !== undefined) {
// If body is a string that looks like JSON, parse it
if (typeof processedInputs.body === 'string') {
try {
// Trim whitespace before checking for JSON pattern
const trimmedBody = processedInputs.body.trim();
if (trimmedBody.startsWith('{') || trimmedBody.startsWith('[')) {
processedInputs.body = JSON.parse(trimmedBody);
logger.info('[ApiBlockHandler] Parsed JSON body:', JSON.stringify(processedInputs.body, null, 2));
}
} catch (e) {
logger.info('[ApiBlockHandler] Failed to parse body as JSON, using as string:', e);
// Keep as string if parsing fails
}
} else if (processedInputs.body === null) {
// Convert null to undefined for consistency with API expectations
processedInputs.body = undefined;
}
}
// Ensure the final processed body is logged
logger.info('[ApiBlockHandler] Final processed request body:', JSON.stringify(processedInputs.body, null, 2));
const result = await executeTool(block.config.tool, {
...inputs,
...processedInputs,
_context: { workflowId: context.workflowId },
})
@@ -932,7 +958,6 @@ export class FunctionBlockHandler implements BlockHandler {
inputs: Record<string, any>,
context: ExecutionContext
): Promise<BlockOutput> {
// Prepare code for execution
const codeContent = Array.isArray(inputs.code)
? inputs.code.map((c: { content: string }) => c.content).join('\n')
: inputs.code
+89 -20
View File
@@ -43,8 +43,14 @@ export class Executor {
private workflowVariables: Record<string, any> = {}
) {
this.validateWorkflow()
this.workflowInput = workflowInput || {}
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.resolver = new InputResolver(
workflow,
@@ -341,6 +347,7 @@ export class Executor {
// Initialize the starter block with the workflow input
try {
const blockParams = starterBlock.config.params
/* Commenting out input format handling
const inputFormat = blockParams?.inputFormat
// If input format is defined, structure the input according to the schema
@@ -352,7 +359,14 @@ export class Executor {
for (const field of inputFormat) {
if (field.name && field.type) {
// Get the field value from workflow input if available
const inputValue = this.workflowInput?.[field.name]
// First try to access via input.field, then directly from field
// 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
logger.info(`[Executor] Processing input field ${field.name} (${field.type}):`,
inputValue !== undefined ? JSON.stringify(inputValue) : 'undefined');
// Convert the value to the appropriate type
let typedValue = inputValue
@@ -378,13 +392,27 @@ export class Executor {
}
}
// Check if we managed to process any fields - if not, use the raw input
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
// Use the structured input if we processed fields, otherwise use raw input
const finalInput = hasProcessedFields ? structuredInput : rawInputData;
// Initialize the starter block with structured input
// Ensure both input and direct fields are available
const starterOutput = {
response: {
input: structuredInput,
...structuredInput, // Add input fields directly at response level too
input: finalInput,
...finalInput, // Add input fields directly at response level too
},
}
logger.info(`[Executor] Starter output:`, JSON.stringify(starterOutput, null, 2));
context.blockStates.set(starterBlock.id, {
output: starterOutput,
@@ -392,23 +420,31 @@ export class Executor {
executionTime: 0,
})
} else {
*/
// No input format defined or not an array,
// check if we're receiving input from API call
// Handle API call - prioritize using the input as-is
if (this.workflowInput && typeof this.workflowInput === 'object') {
// For API calls, use the raw input but make it accessible at both paths
// For API calls, extract input from the nested structure if it exists
const inputData = this.workflowInput.input !== undefined
? this.workflowInput.input // Use the nested input data
: this.workflowInput; // Fallback to direct input
// Create starter output with both formats for maximum compatibility
const starterOutput = {
response: {
input: this.workflowInput,
...this.workflowInput, // Make fields directly accessible at response level
input: inputData,
...inputData, // Make fields directly accessible at response level
},
}
logger.info(`Using API input type: ${typeof this.workflowInput}`, {
isArray: Array.isArray(this.workflowInput),
keys: Object.keys(this.workflowInput),
rawInput: JSON.stringify(this.workflowInput),
inputEmpty: Object.keys(this.workflowInput).length === 0,
})
logger.info(`[Executor] API input for starter block:`, {
type: typeof inputData,
isArray: Array.isArray(inputData),
keys: Object.keys(inputData),
inputData: JSON.stringify(inputData, null, 2),
});
logger.info(`[Executor] Final starter block output:`, JSON.stringify(starterOutput, null, 2));
context.blockStates.set(starterBlock.id, {
output: starterOutput,
@@ -422,6 +458,8 @@ export class Executor {
input: this.workflowInput,
},
}
logger.info(`[Executor] Simple starter output:`, JSON.stringify(starterOutput, null, 2));
context.blockStates.set(starterBlock.id, {
output: starterOutput,
@@ -429,16 +467,24 @@ export class Executor {
executionTime: 0,
})
}
}
//} // End of inputFormat conditional
} catch (e) {
logger.warn('Error processing starter block input format:', e)
// Fallback to raw input with both paths accessible
// Ensure we handle both input formats
const inputData = this.workflowInput.input !== undefined
? this.workflowInput.input // Use nested input if available
: this.workflowInput; // Fallback to direct input
const starterOutput = {
response: {
input: this.workflowInput,
...this.workflowInput, // Add input fields directly at response level too
input: inputData,
...inputData, // Add input fields directly at response level too
},
}
logger.info(`[Executor] Fallback starter output:`, JSON.stringify(starterOutput, null, 2));
context.blockStates.set(starterBlock.id, {
output: starterOutput,
@@ -446,10 +492,12 @@ export class Executor {
executionTime: 0,
})
}
// Mark the starter block as executed and add its connections to the active path
// Ensure the starter block is in the active execution path
context.activeExecutionPath.add(starterBlock.id);
// Mark the starter block as executed
context.executedBlocks.add(starterBlock.id)
// Add all blocks connected to the starter to the active execution path
const connectedToStarter = this.workflow.connections
.filter((conn) => conn.source === starterBlock.id)
.map((conn) => conn.target)
@@ -634,25 +682,46 @@ export class Executor {
throw new Error(`Cannot execute disabled block: ${block.metadata?.name || block.id}`)
}
const inputs = this.resolver.resolveInputs(block, context)
// Check if this block needs the starter block's output
// This is especially relevant for API, function, and conditions that might reference <start.response.input>
const starterBlock = this.workflow.blocks.find((b) => b.metadata?.id === 'starter')
if (starterBlock) {
const starterState = context.blockStates.get(starterBlock.id)
if (!starterState) {
logger.warn(`Starter block state not found when executing ${block.metadata?.name || blockId}. This may cause reference errors.`)
} else {
logger.debug(`Starter block state available for ${block.metadata?.name || blockId}:`,
JSON.stringify(starterState.output || {}, null, 2));
}
}
// Resolve inputs (which will look up references to other blocks including starter)
const inputs = this.resolver.resolveInputs(block, context)
logger.debug(`Resolved inputs for ${block.metadata?.name || blockId}:`,
JSON.stringify(inputs, null, 2));
// Find the appropriate handler
const handler = this.blockHandlers.find((h) => h.canHandle(block))
if (!handler) {
throw new Error(`No handler found for block type: ${block.metadata?.id}`)
}
// Execute the block
const startTime = performance.now()
const rawOutput = await handler.execute(block, inputs, context)
const executionTime = performance.now() - startTime
// Normalize the output
const output = this.normalizeBlockOutput(rawOutput, block)
// Update the context with the execution result
context.blockStates.set(blockId, {
output,
executed: true,
executionTime,
})
// Update the execution log
blockLog.success = true
blockLog.output = output
blockLog.durationMs = Math.round(executionTime)
+121 -35
View File
@@ -1,6 +1,9 @@
import { SerializedBlock, SerializedWorkflow } from '@/serializer/types'
import { LoopManager } from './loops'
import { ExecutionContext } from './types'
import { createLogger } from '@/lib/logs/console-logger'
const logger = createLogger('InputResolver')
/**
* Resolves input values for blocks by handling references and variable substitution.
@@ -51,8 +54,6 @@ export class InputResolver {
resolveInputs(block: SerializedBlock, context: ExecutionContext): Record<string, any> {
const inputs = { ...block.config.params }
const result: Record<string, any> = {}
const isFunctionBlock = block.metadata?.id === 'function'
// Process each input parameter
for (const [key, value] of Object.entries(inputs)) {
// Skip null or undefined values
@@ -78,16 +79,39 @@ export class InputResolver {
// Resolve environment variables
resolvedValue = this.resolveEnvVariables(resolvedValue, isApiKey)
// Special handling for different block types
const isFunctionBlock = block.metadata?.id === 'function'
const isApiBlock = block.metadata?.id === 'api'
// For function blocks, we need special handling for code input
if (isFunctionBlock && key === 'code') {
// For code input in function blocks, we don't want to parse JSON
result[key] = resolvedValue
logger.debug(`[resolveInputs] Function block code input preserved as string`);
}
// For API blocks, handle body input specially
else if (isApiBlock && key === 'body') {
try {
// If it's JSON-looking, preserve its structure
if (resolvedValue.trim().startsWith('{') || resolvedValue.trim().startsWith('[')) {
result[key] = JSON.parse(resolvedValue)
logger.debug(`[resolveInputs] API block body parsed as JSON object`);
} else {
result[key] = resolvedValue
logger.debug(`[resolveInputs] API block body preserved as string`);
}
} catch {
// If parsing fails, keep as string
result[key] = resolvedValue
logger.debug(`[resolveInputs] API block body JSON parsing failed, keeping as string`);
}
}
// For other inputs, try to convert JSON strings to objects
else {
try {
if (resolvedValue.startsWith('{') || resolvedValue.startsWith('[')) {
result[key] = JSON.parse(resolvedValue)
logger.debug(`[resolveInputs] Parsed JSON value for ${key}`);
} else {
result[key] = resolvedValue
}
@@ -272,6 +296,101 @@ export class InputResolver {
const path = match.slice(1, -1)
const [blockRef, ...pathParts] = path.split('.')
// Log the reference being processed
logger.debug(`[resolveBlockReferences] Processing block reference: ${match}`, {
blockRef,
pathParts,
currentBlock: currentBlock.id,
currentBlockType: currentBlock.metadata?.id
});
// Special case for "start" references
// This allows users to reference the starter block using <start.response.type.input>
// regardless of the actual name of the starter block
if (blockRef.toLowerCase() === 'start') {
// Find the starter block
const starterBlock = this.workflow.blocks.find((block) => block.metadata?.id === 'starter')
if (starterBlock) {
logger.debug(`[resolveBlockReferences] Found starter block with ID: ${starterBlock.id}`);
const blockState = context.blockStates.get(starterBlock.id)
if (blockState) {
logger.debug(`[resolveBlockReferences] Starter block state:`, JSON.stringify(blockState, null, 2));
// Navigate through the path parts
let replacementValue: any = blockState.output
// Log the initial output value from the starter block
logger.debug(`[resolveBlockReferences] Initial starter output:`, JSON.stringify(replacementValue, null, 2));
for (const part of pathParts) {
logger.debug(`[resolveBlockReferences] Navigating path part: ${part}`, {
currentValue: typeof replacementValue === 'object' ?
JSON.stringify(replacementValue) : replacementValue
});
if (!replacementValue || typeof replacementValue !== 'object') {
logger.warn(`[resolveBlockReferences] Invalid path "${part}" - replacementValue is not an object:`, replacementValue);
throw new Error(`Invalid path "${part}" in "${path}" for starter block.`)
}
replacementValue = replacementValue[part]
if (replacementValue === undefined) {
logger.warn(`[resolveBlockReferences] No value found at path "${part}" in starter block.`);
throw new Error(`No value found at path "${path}" in starter block.`)
}
}
// Format the value based on block type and path
let formattedValue: string;
// Special handling for all blocks referencing starter input
if (blockRef.toLowerCase() === 'start' && pathParts.join('.').includes('input')) {
const blockType = currentBlock.metadata?.id;
// Format based on which block is consuming this value
if (typeof replacementValue === 'object' && replacementValue !== null) {
// For function blocks, preserve the object structure for code usage
if (blockType === 'function') {
logger.debug(`[resolveBlockReferences] Special handling for function input:`,
JSON.stringify(replacementValue, null, 2));
formattedValue = JSON.stringify(replacementValue);
}
// For API blocks, handle body special case
else if (blockType === 'api') {
logger.debug(`[resolveBlockReferences] Special handling for API input:`,
JSON.stringify(replacementValue, null, 2));
formattedValue = JSON.stringify(replacementValue);
}
// For condition blocks, ensure proper formatting
else if (blockType === 'condition') {
logger.debug(`[resolveBlockReferences] Special handling for condition input:`,
JSON.stringify(replacementValue, null, 2));
formattedValue = this.stringifyForCondition(replacementValue);
}
// For all other blocks, stringify objects
else {
formattedValue = JSON.stringify(replacementValue);
}
} else {
// For primitive values
formattedValue = String(replacementValue);
}
} else {
// Standard handling for non-input references
formattedValue = typeof replacementValue === 'object'
? JSON.stringify(replacementValue)
: String(replacementValue);
}
logger.debug(`[resolveBlockReferences] Resolved value:`, formattedValue);
resolvedValue = resolvedValue.replace(match, formattedValue)
continue
}
}
}
// Special case for "loop" references - allows accessing loop properties
if (blockRef.toLowerCase() === 'loop') {
// Find which loop this block belongs to
@@ -383,39 +502,6 @@ export class InputResolver {
}
}
// Special case for "start" references
// This allows users to reference the starter block using <start.response.type.input>
// regardless of the actual name of the starter block
if (blockRef.toLowerCase() === 'start') {
// Find the starter block
const starterBlock = this.workflow.blocks.find((block) => block.metadata?.id === 'starter')
if (starterBlock) {
const blockState = context.blockStates.get(starterBlock.id)
if (blockState) {
// Navigate through the path parts
let replacementValue: any = blockState.output
for (const part of pathParts) {
if (!replacementValue || typeof replacementValue !== 'object') {
throw new Error(`Invalid path "${part}" in "${path}" for starter block.`)
}
replacementValue = replacementValue[part]
if (replacementValue === undefined) {
throw new Error(`No value found at path "${path}" in starter block.`)
}
}
// Format the value
const formattedValue =
typeof replacementValue === 'object'
? JSON.stringify(replacementValue)
: String(replacementValue)
resolvedValue = resolvedValue.replace(match, formattedValue)
continue
}
}
}
// Standard block reference resolution
let sourceBlock = this.blockById.get(blockRef)
if (!sourceBlock) {