diff --git a/apps/http-backend/src/routes/userRoutes/executionRoutes.ts b/apps/http-backend/src/routes/userRoutes/executionRoutes.ts index ed00904..667e811 100644 --- a/apps/http-backend/src/routes/userRoutes/executionRoutes.ts +++ b/apps/http-backend/src/routes/userRoutes/executionRoutes.ts @@ -35,7 +35,9 @@ execRouter.post('/node', userMiddleware, async (req: AuthRequest, res: Response) userId: req.user.sub, config: config, credentialId: nodeData.CredentialsID || config?.credentialId || "", - authType: nodeData.AvailableNode.authType + authType: nodeData.AvailableNode.authType, + nodeId: nodeData.id, + items: [] } const result = await prismaClient.$transaction(async (tx) => { // // console.log(`Execution context: ${JSON.stringify(context)}`) @@ -93,7 +95,7 @@ execRouter.post('/node', userMiddleware, async (req: AuthRequest, res: Response) data: { status: "Failed", completedAt: new Date(), - error: executionResult.output + error: executionResult.error } }) return { success: false, executionResult } diff --git a/apps/worker/src/engine/executor.ts b/apps/worker/src/engine/executor.ts index 721872e..c9f61d5 100644 --- a/apps/worker/src/engine/executor.ts +++ b/apps/worker/src/engine/executor.ts @@ -14,6 +14,11 @@ interface NodeExecutionOutput { outputData: any; } +interface QueueItem { + nodeId: string; + inputData: any; +} + // Result tracking for looped executions interface LoopExecutionResult { totalProcessed: number; @@ -25,14 +30,7 @@ interface LoopExecutionResult { results: any[]; } -/** - * Checks if inputData is spreadsheet data that should trigger looping - */ -function isSpreadsheetInput(data: any): boolean { - return data?.rows && Array.isArray(data.rows) && - data?.columns && typeof data.columns === 'object' && - data?.dataStartIndex !== undefined; -} + /** * Small delay helper for rate limiting @@ -44,287 +42,337 @@ function delay(ms: number): Promise { export async function executeWorkflow( workflowExecutionId: string ): Promise { - console.log(`workflowExecutionId is ${workflowExecutionId}`); - const data = await prismaClient.workflowExecution.findUnique({ - where: { id: workflowExecutionId }, - include: { - workflow: { - include: { - nodes: { - include: { - AvailableNode: true, - credentials: true, + try { + console.log(`workflowExecutionId is ${workflowExecutionId}`); + const data = await prismaClient.workflowExecution.findUnique({ + where: { id: workflowExecutionId }, + include: { + workflow: { + include: { + nodes: { + include: { + AvailableNode: true, + credentials: true, + }, }, - orderBy: { stage: "asc" }, }, }, + nodeExecutions: true }, - }, - }); - let currentInputData = data?.metadata; + }); - // Collect outputs from all executed nodes for variable interpolation - const executedNodeOutputs: NodeExecutionOutput[] = []; + // Collect outputs from all executed nodes for variable interpolation + const executedNodeOutputs: NodeExecutionOutput[] = []; - if (!data) { - console.log(`No workflow execution found for id ${workflowExecutionId}`); - return; - } + if (!data) { + console.log(`No workflow execution found for id ${workflowExecutionId}`); + return; + } - const update = await prismaClient.workflowExecution.update({ - where: { - id: workflowExecutionId, - }, - data: { - status: "InProgress", - }, - }); - if (!update.error) console.log("updated the workflow execution"); - - const nodes = data?.workflow.nodes; - - console.log(`Total nodes - ${nodes.length}`); - for (const node of nodes) { - console.log(`${node.name}, ${node.stage}, ${node.id}th - started Execution`); - const nodeExecution = await prismaClient.nodeExecution.create({ + const update = await prismaClient.workflowExecution.update({ + where: { + id: workflowExecutionId, + }, data: { - nodeId: node.id, - workflowExecId: workflowExecutionId, - status: "Start", - inputData: currentInputData ? currentInputData : {}, - startedAt: new Date() - } - }) - const nodeType = node.AvailableNode.type; - - // Create mutable copy of config - let nodeConfig = { ...(node.config as Record) }; - - // Build interpolation context from all previously executed nodes - const interpolationContext = buildInterpolationContext(executedNodeOutputs); - for (const out of executedNodeOutputs) { - if (out.nodeId) interpolationContext[out.nodeId] = out.outputData; - } - console.log(`[Interpolation] Before: ${JSON.stringify(interpolationContext)}`); - // Resolve any {{variable}} references in the config - console.log(`[nodeConfig] Before: ${JSON.stringify(nodeConfig)}`); - nodeConfig = resolveConfigVariables(nodeConfig, interpolationContext); - console.log(`[Interpolation] After: ${JSON.stringify(nodeConfig)}`); - - // NOTE: Removed legacy body concatenation that appended raw JSON to email body. - // Variables should be resolved via the {{interpolation}} system instead. - if (!node.CredentialsID) { + status: "InProgress", + }, + }); + if (!update.error) console.log("updated the workflow execution"); + + const nodes = data?.workflow.nodes; + + console.log(`Total nodes - ${nodes.length}`); + // for (const node of nodes) { + // console.log(`${node.name}, ${node.stage}, ${node.id}th - started Execution`); + // const nodeExecution = await prismaClient.nodeExecution.create({ + // data: { + // nodeId: node.id, + // workflowExecId: workflowExecutionId, + // status: "Start", + // inputData: currentInputData ? currentInputData : {}, + // startedAt: new Date() + // } + // }) + // const nodeType = node.AvailableNode.type; + + // // Create mutable copy of config + // let nodeConfig = { ...(node.config as Record) }; + + // // Build interpolation context from all previously executed nodes + // const interpolationContext = buildInterpolationContext(executedNodeOutputs); + // for (const out of executedNodeOutputs) { + // if (out.nodeId) interpolationContext[out.nodeId] = out.outputData; + // } + // console.log(`[Interpolation] Before: ${JSON.stringify(interpolationContext)}`); + // // Resolve any {{variable}} references in the config + // console.log(`[nodeConfig] Before: ${JSON.stringify(nodeConfig)}`); + // nodeConfig = resolveConfigVariables(nodeConfig, interpolationContext); + // console.log(`[Interpolation] After: ${JSON.stringify(nodeConfig)}`); + + // // NOTE: Removed legacy body concatenation that appended raw JSON to email body. + // // Variables should be resolved via the {{interpolation}} system instead. + // if (!node.CredentialsID) { + // await prismaClient.workflowExecution.update({ + // where: { id: workflowExecutionId }, + // data: { + // status: "Failed", + // error: "Credential id not found", + // completedAt: new Date(), + // }, + // }); + + // await prismaClient.nodeExecution.update({ + // where: { id: nodeExecution.id }, + // data: { + // status: "Failed", + // error: "Credential id not found", + // completedAt: new Date() + // } + // }) + // return; + // } + + // // Check if we need to loop (inputData is spreadsheet + config has column variables) + + + // let execute: { success: boolean; output?: any; error?: string }; + + + // if (!execute.success) { + // // Check if it's a partial loop failure (some rows succeeded) + // const isPartialFailure = execute.output?.successful > 0 && execute.output?.failed > 0; + + // await prismaClient.workflowExecution.update({ + // where: { id: workflowExecutionId }, + // data: { + // status: "Failed", + // error: execute.error, + // completedAt: new Date(), + // } + // }); + + // await prismaClient.nodeExecution.update({ + // where: { id: nodeExecution.id }, + // data: { + // status: "Failed", + // error: execute.error, + // outputData: isPartialFailure ? execute.output : undefined, + // completedAt: new Date() + // } + // }) + // return; + // } + // await prismaClient.nodeExecution.update({ + // where: { id: nodeExecution.id }, + // data: { + // completedAt: new Date(), + // outputData: execute.output, + // status: "Completed" + // } + // }) + + // // Store this node's output for variable resolution in subsequent nodes + // executedNodeOutputs.push({ + // nodeName: node.name, + // nodeId: node.id, + // outputData: execute.output + // }); + // console.log(`[Interpolation] Added ${node.name} output to context. Total nodes in context: ${executedNodeOutputs.length}`); + + // currentInputData = execute.output; + + // console.log("output: ", JSON.stringify(execute)); + // } + const firstActionNode = data.workflow.nodes.find(n => n.stage === 0) + if (!firstActionNode) { + console.log("No Trigger node found!") await prismaClient.workflowExecution.update({ where: { id: workflowExecutionId }, data: { status: "Failed", - error: "Credential id not found", completedAt: new Date(), - }, - }); - - await prismaClient.nodeExecution.update({ - where: { id: nodeExecution.id }, - data: { - status: "Failed", - error: "Credential id not found", - completedAt: new Date() + error: "No Trigger node found!" } }) - return; + return } - // Check if we need to loop (inputData is spreadsheet + config has column variables) - const shouldLoop = isSpreadsheetInput(currentInputData) && - JSON.stringify(node.config).includes('{{'); - - let execute: { success: boolean; output?: any; error?: string }; - - if (shouldLoop) { - // === AUTO-LOOP: Execute node once per data row === - console.log(`[Loop] Detected spreadsheet input for ${node.name}, looping through rows...`); - const spreadsheet = currentInputData as any; - const startIdx = spreadsheet.dataStartIndex ?? 1; - const totalRows = spreadsheet.rows.length - startIdx; - console.log(`[Loop] Processing ${totalRows} data rows (starting at index ${startIdx})`); - - const loopResult: LoopExecutionResult = { - totalProcessed: 0, - successful: 0, - failed: 0, - skipped: 0, - failures: [], - skippedRows: [], - results: [] - }; - - for (let rowIdx = startIdx; rowIdx < spreadsheet.rows.length; rowIdx++) { - loopResult.totalProcessed++; - - // Set _currentRowIndex on the spreadsheet data for column resolution - const rowContext = { ...spreadsheet, _currentRowIndex: rowIdx }; - - // Rebuild interpolation context with current row index - const loopOutputs = executedNodeOutputs.map(o => { - if (o.outputData === currentInputData) { - return { ...o, outputData: rowContext }; + const queue: QueueItem[] = [{ + nodeId: firstActionNode.id, + inputData: data?.metadata + }] + + while (queue.length > 0) { + const currentTask = queue.shift() + let currentInputData = currentTask?.inputData; + + const node = data.workflow.nodes.find(n => n.id === currentTask?.nodeId) + if (!node) { + console.log(`Failed to find node with ID ${currentTask?.nodeId}`); + await prismaClient.workflowExecution.update({ + where: { id: workflowExecutionId }, + data: { + status: "Failed", + error: `Failed to find node with ID ${currentTask?.nodeId}`, + completedAt: new Date() } - return o; - }); - const loopInterpolationCtx = buildInterpolationContext(loopOutputs); - for (const out of loopOutputs) { - if (out.nodeId) loopInterpolationCtx[out.nodeId] = out.outputData; - } - - // Re-resolve config with current row - const originalConfig = { ...(node.config as Record) }; - const resolvedRowConfig = resolveConfigVariables(originalConfig, loopInterpolationCtx); - - console.log(`[Loop] Row ${rowIdx}: resolved config = ${JSON.stringify(resolvedRowConfig)}`); + }) + return + } - // Skip rows with empty/null required fields (e.g. empty email recipient) - const skipReasons: string[] = []; - if (resolvedRowConfig.to !== undefined && (!resolvedRowConfig.to || String(resolvedRowConfig.to).trim() === '')) { - skipReasons.push('empty recipient (to)'); - } - if (resolvedRowConfig.subject !== undefined && (!resolvedRowConfig.subject || String(resolvedRowConfig.subject).trim() === '')) { - skipReasons.push('empty subject'); - } - // Check if any resolved value still has unresolved {{variables}} - for (const [key, val] of Object.entries(resolvedRowConfig)) { - if (typeof val === 'string' && val.includes('{{') && val.includes('}}')) { - skipReasons.push(`unresolved variable in ${key}: ${val}`); - } - } - if (skipReasons.length > 0) { - loopResult.skipped++; - loopResult.skippedRows.push({ row: rowIdx, reason: skipReasons.join('; ') }); - console.log(`[Loop] Row ${rowIdx} SKIPPED: ${skipReasons.join('; ')}`); - continue; + console.log(`Executing Action: ${node?.id}`) + // creating a row in db + const nodeExecution = await prismaClient.nodeExecution.create({ + data: { + nodeId: node.id, + workflowExecId: workflowExecutionId, + status: "Start", + inputData: currentInputData ? currentInputData : {}, + startedAt: new Date() } + }) - const rowCtx = { - userId: data.workflow.userId, - credentialId: node.CredentialsID, - config: resolvedRowConfig, - inputData: rowContext, - }; - - // Retry logic: up to 3 attempts per row - let rowSuccess = false; - let lastError: string | undefined; - - for (let attempt = 1; attempt <= 3; attempt++) { - try { - const rowResult = await ExecutionRegister.execute(nodeType, rowCtx); - if (rowResult.success) { - rowSuccess = true; - loopResult.successful++; - loopResult.results.push(rowResult.output); - break; - } else { - lastError = rowResult.error; - console.log(`[Loop] Row ${rowIdx} attempt ${attempt} failed: ${lastError}`); - } - } catch (err) { - lastError = err instanceof Error ? err.message : 'Unknown error'; - console.log(`[Loop] Row ${rowIdx} attempt ${attempt} threw: ${lastError}`); - } + // checking node type + const nodeType = node?.AvailableNode; + let nodeConfig = { ...(node.config as Record) }; - if (attempt < 3) { - await delay(200 * attempt); // Backoff: 200ms, 400ms - } + let itemsToProcess = []; + if (Array.isArray(currentInputData)) { + if (currentInputData.length > 0 && Array.isArray(currentInputData[0])) { + itemsToProcess = currentInputData[0] } + else + itemsToProcess = currentInputData + } + else { + itemsToProcess = currentInputData ? [{ json: currentInputData }] : [] + } - if (!rowSuccess) { - loopResult.failed++; - loopResult.failures.push({ - row: rowIdx, - error: lastError || 'Unknown error', - retries: 3 - }); - } - // Rate limiting delay between iterations - await delay(100); + // Build interpolation context from all previously executed nodes + const interpolationContext = buildInterpolationContext(executedNodeOutputs); + for (const out of executedNodeOutputs) { + if (out.nodeId) interpolationContext[out.nodeId] = out.outputData; } + console.log(`[Interpolation] Before: ${JSON.stringify(interpolationContext)}`); + // Resolve any {{variable}} references in the config + console.log(`[nodeConfig] Before: ${JSON.stringify(nodeConfig)}`); + nodeConfig = resolveConfigVariables(nodeConfig, interpolationContext, itemsToProcess[0]?.sourceRefs); + console.log(`[Interpolation] After: ${JSON.stringify(nodeConfig)}`); + + // checking node authentication with credential id + if (nodeType.requireAuth) { + if (!node.CredentialsID) { + await prismaClient.workflowExecution.update({ + where: { id: workflowExecutionId }, + data: { + status: "Failed", + error: "Credential id not found", + completedAt: new Date(), + }, + }); - console.log(`[Loop] Completed: ${loopResult.successful}/${loopResult.totalProcessed} successful, ${loopResult.failed} failed, ${loopResult.skipped} skipped`); - - // Loop completes even if some rows fail - const hasFailures = loopResult.failed > 0; - execute = { - success: !hasFailures, - output: loopResult, - error: hasFailures - ? JSON.stringify({ - summary: `${loopResult.failed}/${loopResult.totalProcessed} rows failed`, - failures: loopResult.failures + await prismaClient.nodeExecution.update({ + where: { id: nodeExecution.id }, + data: { + status: "Failed", + error: "Credential id not found", + completedAt: new Date() + } }) - : undefined - }; - } else { - // === SINGLE EXECUTION (no loop) === + return; + } + } + const context = { + nodeId: node.id, userId: data.workflow.userId, - credentialId: node.CredentialsID, + credentialId: node.CredentialsID!, config: nodeConfig, - inputData: currentInputData, - }; - console.log(`Executing with context: ${JSON.stringify(context)}`); - execute = await ExecutionRegister.execute(nodeType, context); - } - if (!execute.success) { - // Check if it's a partial loop failure (some rows succeeded) - const isPartialFailure = execute.output?.successful > 0 && execute.output?.failed > 0; + items: itemsToProcess + } + let execute: { success: boolean; output?: any; error?: string }; - await prismaClient.workflowExecution.update({ - where: { id: workflowExecutionId }, - data: { - status: "Failed", - error: execute.error, - completedAt: new Date(), - }, - }); + try { + execute = await ExecutionRegister.execute(nodeType.type, context); + } catch (err: any) { + execute = { success: false, error: err.message || "Unknown error" }; + } + if (!execute.success) { + // Check if it's a partial loop failure (some rows succeeded) + const isPartialFailure = execute.output?.successful > 0 && execute.output?.failed > 0; + + await prismaClient.workflowExecution.update({ + where: { id: workflowExecutionId }, + data: { + status: "Failed", + error: execute.error, + completedAt: new Date(), + } + }); + + await prismaClient.nodeExecution.update({ + where: { id: nodeExecution.id }, + data: { + status: "Failed", + error: execute.error, + outputData: isPartialFailure ? execute.output : undefined, + completedAt: new Date() + } + }) + return; + } await prismaClient.nodeExecution.update({ where: { id: nodeExecution.id }, data: { - status: "Failed", - error: execute.error, - outputData: isPartialFailure ? execute.output : undefined, - completedAt: new Date() + completedAt: new Date(), + outputData: execute.output, + status: "Completed" } }) - return; + + // Store this node's output for variable resolution in subsequent nodes + executedNodeOutputs.push({ + nodeName: node.name, + nodeId: node.id, + outputData: execute.output + }); + + console.log(`[Interpolation] Added ${node.name} output to context. Total nodes in context: ${executedNodeOutputs.length}`); + console.log("output: ", JSON.stringify(execute)); + + const allEdges = (data.workflow.Edges as any[]) || [] + const outgoingEdges = allEdges.filter( + (edge) => edge.source === node.id + ) + + for (const edge of outgoingEdges) { + const outputPinIndex = 0 //need to changes this after phase 5 (enables UI with 2 pins output per node) + + const branchData = execute.output?.[outputPinIndex] || []; + + queue.push({ + nodeId: edge.target, + inputData: branchData + }) + + console.log(`Pushed Node ${edge.target} into the queue!`); + } } - await prismaClient.nodeExecution.update({ - where: { id: nodeExecution.id }, + + + const updatedStatus = await prismaClient.workflowExecution.update({ + where: { id: workflowExecutionId }, data: { + status: "Completed", completedAt: new Date(), - outputData: execute.output, - status: "Completed" - } - }) - - // Store this node's output for variable resolution in subsequent nodes - executedNodeOutputs.push({ - nodeName: node.name, - nodeId: node.id, - outputData: execute.output + }, }); - console.log(`[Interpolation] Added ${node.name} output to context. Total nodes in context: ${executedNodeOutputs.length}`); + console.log(updatedStatus); - currentInputData = execute.output; - - console.log("output: ", JSON.stringify(execute)); } - const updatedStatus = await prismaClient.workflowExecution.update({ - where: { id: workflowExecutionId }, - data: { - status: "Completed", - completedAt: new Date(), - }, - }); - console.log(updatedStatus); + catch (err: any) { + //update workflow with failed message + } } diff --git a/packages/common/src/index.ts b/packages/common/src/index.ts index 5a51d84..de6008f 100644 --- a/packages/common/src/index.ts +++ b/packages/common/src/index.ts @@ -208,4 +208,15 @@ export const FilterNodeInput = z.object({ referenceKey: z.string().optional() }) +export const ExecuteItemSchema = z.object({ + json: z.record(z.string(), z.any()), + sourceRefs: z.record(z.string(), + z.object({ + wireIndex: z.number().default(0), + rowIndex: z.number() + })).optional() +}); + +export type ExecuteItem = z.infer + export type FilterNodeInput = z.infer; \ No newline at end of file diff --git a/packages/common/src/interpolation.ts b/packages/common/src/interpolation.ts index ec7e476..1673052 100644 --- a/packages/common/src/interpolation.ts +++ b/packages/common/src/interpolation.ts @@ -81,7 +81,7 @@ export function buildInterpolationContext( * @param context - The interpolation context * @returns The resolved value or the original {{variable}} if not found */ -export function resolveVariable(variable: string, context: InterpolationContext): any { +export function resolveVariable(variable: string, context: InterpolationContext, sourceRefs?: Record): any { const trimmed = variable.trim(); const dotIndex = trimmed.indexOf('.'); @@ -97,17 +97,42 @@ export function resolveVariable(variable: string, context: InterpolationContext) return `{{${variable}}}`; } - // Column-based resolution: {{google_sheet.email}} → rows[currentRow][columnIndex] - if (nodeData.columns && nodeData.columns[path] !== undefined) { - const colIndex = nodeData.columns[path]; - const rowIndex = nodeData._currentRowIndex ?? nodeData.dataStartIndex ?? 1; - const value = nodeData.rows?.[rowIndex]?.[colIndex]; + if (sourceRefs && sourceRefs[nodeName]) { + const { wireIndex, rowIndex } = sourceRefs[nodeName] + + const item = nodeData[wireIndex]?.[rowIndex]; + + const value = getNestedValue(item?.json, path) + return value !== undefined ? value : `{{${variable}}}`; } - // Fallback: standard nested path resolution (e.g., rows[1][0]) - const value = getNestedValue(nodeData, path); - return value !== undefined ? value : `{{${variable}}}`; + // // Column-based resolution: {{google_sheet.email}} → rows[currentRow][columnIndex] + // if (nodeData.columns && nodeData.columns[path] !== undefined) { + // const colIndex = nodeData.columns[path]; + // const rowIndex = nodeData._currentRowIndex ?? nodeData.dataStartIndex ?? 1; + // const value = nodeData.rows?.[rowIndex]?.[colIndex]; + // return value !== undefined ? value : `{{${variable}}}`; + // } + + // // Fallback: standard nested path resolution (e.g., rows[1][0]) + // const value = getNestedValue(nodeData, path); + + // 1. Fallback for single-item nodes (e.g., Webhook or trigger nodes at wire 0, row 0) + const fallbackItem = nodeData?.[0]?.[0]; + if (fallbackItem && typeof fallbackItem === 'object' && 'json' in fallbackItem) { + const value = getNestedValue(fallbackItem.json, path); + if (value !== undefined) return value; + } + + // 2. Direct property fallback (if raw object was passed in context) + const directValue = getNestedValue(nodeData, path); + if (directValue !== undefined) return directValue; + + // 3. If nothing matched, preserve the raw {{variable}} placeholder + return `{{${variable}}}`; + + } /** @@ -117,7 +142,7 @@ export function resolveVariable(variable: string, context: InterpolationContext) * @param context - The interpolation context * @returns String with all variables resolved */ -export function interpolateString(template: string, context: InterpolationContext): string { +export function interpolateString(template: string, context: InterpolationContext, sourceRefs?: Record): string { if (!template || typeof template !== 'string') return template; // Create a new regex instance to avoid global flag issues @@ -129,7 +154,7 @@ export function interpolateString(template: string, context: InterpolationContex const result = template.replace(regex, (match, variable) => { console.log(`[interpolateString] Found variable: "${variable}"`); console.log(`[interpolateString] MATCH variable: "${match}"`); - const resolved = resolveVariable(variable, context); + const resolved = resolveVariable(variable, context, sourceRefs); console.log(`[interpolateString] Resolved to: ${JSON.stringify(resolved)}`); // Convert non-string values to string for template replacement @@ -153,7 +178,8 @@ export function interpolateString(template: string, context: InterpolationContex */ export function resolveConfigVariables>( config: T, - context: InterpolationContext + context: InterpolationContext, + sourceRefs?: Record ): T { if (!config || typeof config !== 'object') { console.log("[---*---] config not an onject type") @@ -162,7 +188,7 @@ export function resolveConfigVariables>( // Handle arrays if (Array.isArray(config)) { - return config.map(item => resolveConfigVariables(item, context)) as unknown as T; + return config.map(item => resolveConfigVariables(item, context, sourceRefs)) as unknown as T; } // Handle objects @@ -174,12 +200,12 @@ export function resolveConfigVariables>( const exactMatch = /^\{\{([^}]+)\}\}$/.exec(trimmed); if (exactMatch) { - resolved[key] = resolveVariable(exactMatch[1]!, context); + resolved[key] = resolveVariable(exactMatch[1]!, context, sourceRefs); } else { - resolved[key] = interpolateString(value, context); + resolved[key] = interpolateString(value, context, sourceRefs); } } else if (typeof value === 'object' && value !== null) { - resolved[key] = resolveConfigVariables(value, context); + resolved[key] = resolveConfigVariables(value, context, sourceRefs); } else { resolved[key] = value; } diff --git a/packages/nodes/src/filter/filter.executor.ts b/packages/nodes/src/filter/filter.executor.ts index 780457b..c652b32 100644 --- a/packages/nodes/src/filter/filter.executor.ts +++ b/packages/nodes/src/filter/filter.executor.ts @@ -204,34 +204,32 @@ export class FilterExecutor implements NodeExecutor { error: "sourceKey is required to group datasets" } const groupResult = this.handleGroupBy(normalizedSource, sourceKey); - + const wire0 = groupResult.groupArray.map(item => ({ json: item })) return { success: true, - output: { - groupsMap: groupResult.groupMap, - groupsArray: groupResult.groupArray, - metadata: { - operation_used: operation, - total_groups: groupResult.groupArray.length, - items_processed: groupResult.total_processed, - items_without_key: groupResult.emptyCount - } + output: [wire0], + metadata: { + operation_used: operation, + total_groups: groupResult.groupArray.length, + items_processed: groupResult.total_processed, + items_without_key: groupResult.emptyCount } } default: return { success: false, error: `Unknown operation: ${operation}` }; } + const wire0 = filteredData.map(item => ({ json: item })) + const wire1 = discardedData.map(item => ({ json: item })) + return { success: true, - output: { - filteredData: filteredData, - discardedData: discardedData, - metadata: { - operation_used: operation, - items_kept: filteredData.length, - items_discard: discardedData.length - } + output: [wire0, wire1], + + metadata: { + operation_used: operation, + items_kept: filteredData.length, + items_discard: discardedData.length } } } diff --git a/packages/nodes/src/gmail/gmail.executor.ts b/packages/nodes/src/gmail/gmail.executor.ts index 4b9b8fa..b81574c 100644 --- a/packages/nodes/src/gmail/gmail.executor.ts +++ b/packages/nodes/src/gmail/gmail.executor.ts @@ -1,17 +1,19 @@ +import { ExecuteItem } from "@repo/common/zod"; import { GoogleOAuthService } from "../common/google-oauth-service.js"; import { GmailService, GmailCredentials } from "./gmail.service.js"; interface NodeExecutionContext { + nodeId: string; credentialId: string; userId: string; config?: any; authType?: string; - inputData?: any; + items: ExecuteItem[] } interface NodeExecutionResult { success: boolean; - output?: any; + output?: ExecuteItem[][]; error?: string; } @@ -57,19 +59,33 @@ class GmailExecutor { // Send email const { to, subject, body } = context.config; - const result = await this.gmailService.sendEmail(to, subject, body); + + const outputBoxes: ExecuteItem[] = []; + + for (const item of context.items) { + + const result = await this.gmailService.sendEmail(to, subject, body); + + if (!result.success) { + throw new Error(result.error) + } + + outputBoxes.push({ + json: { + ...item.json, + gmailResponse: { + status: "sent", + messageId: result.data?.id, + threadId: result.data?.threadId + } + }, + sourceRefs: item.sourceRefs + }) + } return { - success: result.success, - output: { - status: "sent", - messageId: result.data?.id, - threadId: result.data?.threadId, - to: to, - subject: subject, - summary: `Email sent to ${to} with subject "${subject}"`, - }, - error: result.error, + success: true, + output: [outputBoxes], }; } catch (error) { return { diff --git a/packages/nodes/src/google-sheets/google-sheets.executor.ts b/packages/nodes/src/google-sheets/google-sheets.executor.ts index 85ece00..1674160 100644 --- a/packages/nodes/src/google-sheets/google-sheets.executor.ts +++ b/packages/nodes/src/google-sheets/google-sheets.executor.ts @@ -1,18 +1,23 @@ import { google } from "googleapis"; import { GoogleOAuthService } from "../common/google-oauth-service.js"; import { GoogleSheetsService, GoogleSheetsCredentials } from "./google-sheets.service.js"; +import { ExecuteItem } from "@repo/common/zod"; -interface NodeExecutionContext { - credentialId: string, - userId: string, - config?: any, //sheet id / range... - authType: string, - inputData?: any // previous node data + +interface SheetAuthContext { + userId: string; + credentialId: string; + authType?: string +} +interface NodeExecutionContext extends SheetAuthContext { + nodeId: string, + config: any, //sheet id / range... + items: ExecuteItem[] } interface NodeExecutionResult { success: boolean, - output?: any, + output?: ExecuteItem[][], error?: string, authUrl?: string, requiresAuth?: boolean @@ -27,7 +32,7 @@ class GoogleSheetsNodeExecutor { this.oauthService = new GoogleOAuthService(); } - async getSheets(context: NodeExecutionContext) { + async getSheets(context: SheetAuthContext) { const init = await this.ensureSheetService(context); if ('success' in init) return init; const sheetService = this.sheetService; @@ -48,7 +53,7 @@ class GoogleSheetsNodeExecutor { } } - async ensureCredentials(context: NodeExecutionContext) { + async ensureCredentials(context: SheetAuthContext) { return this.ensureSheetService(context); } @@ -66,7 +71,7 @@ class GoogleSheetsNodeExecutor { } } - async getSheetTabs(context: NodeExecutionContext, spreadsheetId: string) { + async getSheetTabs(context: SheetAuthContext, spreadsheetId: string) { const init = await this.ensureSheetService(context); if ('success' in init) return init; const sheetService = this.sheetService; @@ -87,7 +92,7 @@ class GoogleSheetsNodeExecutor { } } - private async ensureSheetService(context: NodeExecutionContext): Promise<{ credentialId: string } | NodeExecutionResult> { + private async ensureSheetService(context: SheetAuthContext): Promise<{ credentialId: string } | NodeExecutionResult> { try { const type = context.authType ? context.authType : 'gsheet_oauth' const credentials = await this.oauthService.getCredentials(context.userId, context.credentialId, type); @@ -96,7 +101,7 @@ class GoogleSheetsNodeExecutor { return { success: false, error: 'Google Sheets authorization required', - authUrl: this.oauthService.getAuthUrl(context.userId, context.authType), + authUrl: this.oauthService.getAuthUrl(context.userId, context.authType!), requiresAuth: true }; } @@ -109,7 +114,7 @@ class GoogleSheetsNodeExecutor { return { success: false, error: 'Google account not connected.', - authUrl: this.oauthService.getAuthUrl(context.userId, context.authType), + authUrl: this.oauthService.getAuthUrl(context.userId, context.authType!), requiresAuth: true }; } @@ -164,7 +169,7 @@ class GoogleSheetsNodeExecutor { return { success: false, error: 'Google account not connected.', - authUrl: this.oauthService.getAuthUrl(context.userId, context.authType), + authUrl: this.oauthService.getAuthUrl(context.userId, context.authType!), requiresAuth: true }; } @@ -175,9 +180,9 @@ class GoogleSheetsNodeExecutor { } } - async getHeaderRow(context: NodeExecutionContext, sheetId: string, sheetName: string): Promise { + async getHeaderRow(context: SheetAuthContext, sheetId: string, sheetName: string): Promise<{ success: boolean, output?: string[], error?: string }> { const init = await this.ensureCredentials(context); - if ('success' in init) return init; + if ('success' in init) return init as any; const sheetService = this.sheetService; if (!sheetService) { @@ -208,7 +213,7 @@ class GoogleSheetsNodeExecutor { const mode = context.config.mappingMode || 'visual'; if (mode === 'bulk') { - const rawValues = context.config.bulkValues || context.inputData; + const rawValues = context.config.bulkValues || context.items.map(item => item.json); return this.normalizeValues(rawValues) } @@ -299,22 +304,50 @@ class GoogleSheetsNodeExecutor { dataRowCount = dataRows.length; } - // Build columns mapping from first row (headers) - const columns = combinedRows.length > 0 && combinedRows[0] - ? this.buildColumnsMap(combinedRows[0] as any[]) - : {}; + // // Build columns mapping from first row (headers) + // const columns = combinedRows.length > 0 && combinedRows[0] + // ? this.buildColumnsMap(combinedRows[0] as any[]) + // : {}; + + // return { + // success: true, + // output: { + // rows: combinedRows, + // columns: columns, + // dataStartIndex: 1, + // rowCount: dataRowCount, + // sheetId: spreadsheetId, + // hasHeaders: Object.keys(columns).length > 0 + // } + // } + const headers = (combinedRows.length > 0 && combinedRows[0]) ? (combinedRows[0] as string[]) : []; + + const outputBoxes: ExecuteItem[] = []; + + for (let i = 1; i < combinedRows.length; i++) { + const rowArray = combinedRows[i] + const rowObject: Record = {}; + + headers.forEach((header, index) => { + const cleanHeader = String(header).trim().toLowerCase().replace(/\s+/g, '_'); + + if (cleanHeader) { + rowObject[cleanHeader] = rowArray[index] + } + }) + outputBoxes.push({ + json: rowObject, + sourceRefs: { + [context.nodeId]: { wireIndex: 0, rowIndex: i - 1 } + } + }) + } return { success: true, - output: { - rows: combinedRows, - columns: columns, - dataStartIndex: 1, - rowCount: dataRowCount, - sheetId: spreadsheetId, - hasHeaders: Object.keys(columns).length > 0 - } - } + output: [outputBoxes] + }; + } catch (e) { return { success: false, @@ -383,17 +416,22 @@ class GoogleSheetsNodeExecutor { range: `${context.config.sheetName}!${range}`, values: values }) + + const outputBoxes = context.items.map(item => { + return { + json: { + ...item.json, + googleSheetResponse: { + operation: "Update Rows", + rowsUpdated: response.updatedRows || 1, + } + }, + sourceRefs: item.sourceRefs + } + }) return { success: true, - output: { - operation: "Update Rows", - spreadsheetId: spreadsheetId, - sheetName: context.config.sheetName, - writtenRange: response.updatedRange || range, - rowsUpdated: response.updatedRows || 1, - columnsUpdated: response.updatedColumns || 0, - cellsUpdated: response.updatedCells || 0 - } + output: [outputBoxes] } } catch (e) { return { @@ -416,17 +454,21 @@ class GoogleSheetsNodeExecutor { values: values }) console.log(`append rows: ${response}`) + const outputBoxes = context.items.map(item => { + return { + json: { + ...item.json, + googleSheetResponse: { + operation: "Append Rows", + rowUpdated: response.updates.updatedRange || response.tableRange + } + }, + sourceRefs: item.sourceRefs + } + }) return { success: true, - output: { - operation: "Append Rows", - spreadsheetId: spreadsheetId, - sheetName: context.config.sheetName, - appendedRange: response.updates?.updatedRange || response.tableRange, - rowsAdded: response.updates?.updatedRows || 1, - columnsUpdated: response.updates?.updatedColumns || 0, - cellsUpdated: response.updates?.updatedCells || 0 - } + output: [outputBoxes] } } catch (e) { @@ -446,11 +488,22 @@ class GoogleSheetsNodeExecutor { spreadsheetId: spreadsheetId, range: range }) + const outputBoxes = context.items.map(item => { + return { + json: { + ...item.json, + googleSheetResponse: { + operation: "Clear Rows", + clearedRange: response.clearedRange + } + }, + sourceRefs: item.sourceRefs + + } + }) return { success: true, - output: { - ...response - } + output: [outputBoxes] } } catch (e) { diff --git a/packages/nodes/src/registry/Execution.config.types.ts b/packages/nodes/src/registry/Execution.config.types.ts index fb6d12b..deb0f1b 100644 --- a/packages/nodes/src/registry/Execution.config.types.ts +++ b/packages/nodes/src/registry/Execution.config.types.ts @@ -1,13 +1,15 @@ +import { ExecuteItem } from "@repo/common/zod"; + export interface ExecutionContext { - nodeId?: string; + nodeId: string; userId: string; credentialId?: string; config: Record; - inputData?: any; + items: ExecuteItem[] } export interface ExecutionResult { success: boolean; - output?: any; + output?: ExecuteItem[][]; error?: string; metadata?: Record; } diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 8a56017..5b9a3ea 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -323,6 +323,9 @@ importers: packages/nodes: dependencies: + '@repo/common': + specifier: workspace:* + version: link:../common '@types/dotenv': specifier: ^8.2.3 version: 8.2.3