From 2af611e938864ab3ae6d2d22729c8b434937fdb6 Mon Sep 17 00:00:00 2001 From: Thomas Trompette Date: Mon, 1 Sep 2025 16:00:26 +0200 Subject: [PATCH] [Iterator] Implement Backend for Iterator step (#14145) Manually tested by starting to implement the frontend: - loopStepIds contains a create record action that takes one item at a time - items to iterate are a list of string that will be used as record title https://github.com/user-attachments/assets/728310d9-4728-422f-a324-3783da24e517 --- .../workflow-schema.workspace-service.ts | 9 + ...workflow-version-step.workspace-service.ts | 14 + .../factories/workflow-action.factory.ts | 4 + .../types/workflow-action-input.ts | 6 + .../types/workflow-action-output.type.ts | 1 + .../utils/can-execute-step.util.ts | 24 +- .../utils/find-step-or-throw.util.ts | 24 ++ .../utils/step-has-been-started.util.ts | 12 + .../ai-agent/ai-agent.workflow-action.ts | 13 +- .../code/code.workflow-action.ts | 13 +- .../filter/filter.workflow-action.ts | 13 +- .../form/form.workflow-action.ts | 13 +- .../is-workflow-iterator-action.guard.ts | 10 + .../iterator/iterator-action.module.ts | 11 + .../iterator/iterator.workflow-action.ts | 161 ++++++++ .../workflow-iterator-action-settings.type.ts | 12 + .../types/workflow-iterator-result.type.ts | 5 + .../can-execute-iterator-step.util.spec.ts | 376 ++++++++++++++++++ .../get-all-step-ids-in-loop.util.spec.ts | 228 +++++++++++ .../utils/can-execute-iterator-step.util.ts | 54 +++ .../utils/get-all-step-ids-in-loop.util.ts | 91 +++++ .../create-record.workflow-action.ts | 27 +- .../delete-record.workflow-action.ts | 12 +- .../find-records.workflow-action.ts | 13 +- .../update-record.workflow-action.ts | 13 +- .../types/workflow-action-settings.type.ts | 4 +- .../types/workflow-action.type.ts | 10 +- .../workflow-executor.module.ts | 2 + ...orkflow-executor.workspace-service.spec.ts | 6 +- .../workflow-executor.workspace-service.ts | 145 +++++-- .../workflow-run.workspace-service.ts | 32 ++ 31 files changed, 1232 insertions(+), 126 deletions(-) create mode 100644 packages/twenty-server/src/modules/workflow/workflow-executor/utils/find-step-or-throw.util.ts create mode 100644 packages/twenty-server/src/modules/workflow/workflow-executor/utils/step-has-been-started.util.ts create mode 100644 packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/guards/is-workflow-iterator-action.guard.ts create mode 100644 packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/iterator-action.module.ts create mode 100644 packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/iterator.workflow-action.ts create mode 100644 packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/types/workflow-iterator-action-settings.type.ts create mode 100644 packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/types/workflow-iterator-result.type.ts create mode 100644 packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/utils/__tests__/can-execute-iterator-step.util.spec.ts create mode 100644 packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/utils/__tests__/get-all-step-ids-in-loop.util.spec.ts create mode 100644 packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/utils/can-execute-iterator-step.util.ts create mode 100644 packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/utils/get-all-step-ids-in-loop.util.ts diff --git a/packages/twenty-server/src/modules/workflow/workflow-builder/workflow-schema/workflow-schema.workspace-service.ts b/packages/twenty-server/src/modules/workflow/workflow-builder/workflow-schema/workflow-schema.workspace-service.ts index a26f023371..e46754042b 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-builder/workflow-schema/workflow-schema.workspace-service.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-builder/workflow-schema/workflow-schema.workspace-service.ts @@ -76,6 +76,15 @@ export class WorkflowSchemaWorkspaceService { formFieldMetadataItems: step.settings.input, workspaceId, }); + case WorkflowActionType.ITERATOR: { + return { + nextItemToProcess: { + isLeaf: true, + type: 'unknown', + value: generateFakeValue('unknown'), + }, + }; + } case WorkflowActionType.CODE: // StepOutput schema is computed on serverlessFunction draft execution default: return {}; diff --git a/packages/twenty-server/src/modules/workflow/workflow-builder/workflow-version-step/workflow-version-step.workspace-service.ts b/packages/twenty-server/src/modules/workflow/workflow-builder/workflow-version-step/workflow-version-step.workspace-service.ts index 1c2108f7ea..5912eae2b0 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-builder/workflow-version-step/workflow-version-step.workspace-service.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-builder/workflow-version-step/workflow-version-step.workspace-service.ts @@ -654,6 +654,20 @@ export class WorkflowVersionStepWorkspaceService { }, }; } + case WorkflowActionType.ITERATOR: { + return { + ...baseStep, + name: 'Iterator', + type: WorkflowActionType.ITERATOR, + settings: { + ...BASE_STEP_DEFINITION, + input: { + items: [], + initialLoopStepIds: [], + }, + }, + }; + } default: throw new WorkflowVersionStepException( `WorkflowActionType '${type}' unknown`, diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/factories/workflow-action.factory.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/factories/workflow-action.factory.ts index a1de36f689..ae4c14ecbe 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-executor/factories/workflow-action.factory.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/factories/workflow-action.factory.ts @@ -10,6 +10,7 @@ import { AiAgentWorkflowAction } from 'src/modules/workflow/workflow-executor/wo import { CodeWorkflowAction } from 'src/modules/workflow/workflow-executor/workflow-actions/code/code.workflow-action'; import { FilterWorkflowAction } from 'src/modules/workflow/workflow-executor/workflow-actions/filter/filter.workflow-action'; import { FormWorkflowAction } from 'src/modules/workflow/workflow-executor/workflow-actions/form/form.workflow-action'; +import { IteratorWorkflowAction } from 'src/modules/workflow/workflow-executor/workflow-actions/iterator/iterator.workflow-action'; import { CreateRecordWorkflowAction } from 'src/modules/workflow/workflow-executor/workflow-actions/record-crud/create-record.workflow-action'; import { DeleteRecordWorkflowAction } from 'src/modules/workflow/workflow-executor/workflow-actions/record-crud/delete-record.workflow-action'; import { FindRecordsWorkflowAction } from 'src/modules/workflow/workflow-executor/workflow-actions/record-crud/find-records.workflow-action'; @@ -27,6 +28,7 @@ export class WorkflowActionFactory { private readonly findRecordsWorkflowAction: FindRecordsWorkflowAction, private readonly formWorkflowAction: FormWorkflowAction, private readonly filterWorkflowAction: FilterWorkflowAction, + private readonly iteratorWorkflowAction: IteratorWorkflowAction, private readonly toolExecutorWorkflowAction: ToolExecutorWorkflowAction, private readonly aiAgentWorkflowAction: AiAgentWorkflowAction, ) {} @@ -49,6 +51,8 @@ export class WorkflowActionFactory { return this.formWorkflowAction; case WorkflowActionType.FILTER: return this.filterWorkflowAction; + case WorkflowActionType.ITERATOR: + return this.iteratorWorkflowAction; case WorkflowActionType.HTTP_REQUEST: return this.toolExecutorWorkflowAction; case WorkflowActionType.AI_AGENT: diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/types/workflow-action-input.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/types/workflow-action-input.ts index 912833c574..02bcc1f4d8 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-executor/types/workflow-action-input.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/types/workflow-action-input.ts @@ -1,7 +1,13 @@ import { type WorkflowAction } from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type'; +export type WorkflowRunInfo = { + workflowRunId: string; + workspaceId: string; +}; + export type WorkflowActionInput = { currentStepId: string; steps: WorkflowAction[]; context: Record; + runInfo: WorkflowRunInfo; }; diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/types/workflow-action-output.type.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/types/workflow-action-output.type.ts index 0e25ec18fb..6e5f16f9e3 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-executor/types/workflow-action-output.type.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/types/workflow-action-output.type.ts @@ -3,4 +3,5 @@ export type WorkflowActionOutput = { error?: string; pendingEvent?: boolean; shouldEndWorkflowRun?: boolean; + shouldRemainRunning?: boolean; }; diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/utils/can-execute-step.util.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/utils/can-execute-step.util.ts index a248a136dd..15cdcbee44 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-executor/utils/can-execute-step.util.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/utils/can-execute-step.util.ts @@ -1,8 +1,11 @@ import { isDefined } from 'twenty-shared/utils'; import { StepStatus, type WorkflowRunStepInfos } from 'twenty-shared/workflow'; -import { type WorkflowAction } from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type'; import { WorkflowRunStatus } from 'src/modules/workflow/common/standard-objects/workflow-run.workspace-entity'; +import { stepHasBeenStarted } from 'src/modules/workflow/workflow-executor/utils/step-has-been-started.util'; +import { isWorkflowIteratorAction } from 'src/modules/workflow/workflow-executor/workflow-actions/iterator/guards/is-workflow-iterator-action.guard'; +import { canExecuteIteratorStep } from 'src/modules/workflow/workflow-executor/workflow-actions/iterator/utils/can-execute-iterator-step.util'; +import { type WorkflowAction } from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type'; export const canExecuteStep = ({ stepId, @@ -19,10 +22,21 @@ export const canExecuteStep = ({ return false; } - if ( - isDefined(stepInfos[stepId]?.status) && - stepInfos[stepId].status !== StepStatus.NOT_STARTED - ) { + const step = steps.find((step) => step.id === stepId); + + if (!step) { + return false; + } + + if (isWorkflowIteratorAction(step)) { + return canExecuteIteratorStep({ + step, + steps, + stepInfos, + }); + } + + if (stepHasBeenStarted(stepId, stepInfos)) { return false; } diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/utils/find-step-or-throw.util.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/utils/find-step-or-throw.util.ts new file mode 100644 index 0000000000..f27a25175c --- /dev/null +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/utils/find-step-or-throw.util.ts @@ -0,0 +1,24 @@ +import { + WorkflowStepExecutorException, + WorkflowStepExecutorExceptionCode, +} from 'src/modules/workflow/workflow-executor/exceptions/workflow-step-executor.exception'; +import { type WorkflowAction } from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type'; + +export const findStepOrThrow = ({ + stepId, + steps, +}: { + stepId: string; + steps: WorkflowAction[]; +}) => { + const step = steps.find((step) => step.id === stepId); + + if (!step) { + throw new WorkflowStepExecutorException( + 'Step not found', + WorkflowStepExecutorExceptionCode.STEP_NOT_FOUND, + ); + } + + return step; +}; diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/utils/step-has-been-started.util.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/utils/step-has-been-started.util.ts new file mode 100644 index 0000000000..552fdf207b --- /dev/null +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/utils/step-has-been-started.util.ts @@ -0,0 +1,12 @@ +import { isDefined } from 'twenty-shared/utils'; +import { StepStatus, type WorkflowRunStepInfos } from 'twenty-shared/workflow'; + +export const stepHasBeenStarted = ( + stepId: string, + stepInfos: WorkflowRunStepInfos, +) => { + return ( + isDefined(stepInfos[stepId]?.status) && + stepInfos[stepId].status !== StepStatus.NOT_STARTED + ); +}; diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/ai-agent/ai-agent.workflow-action.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/ai-agent/ai-agent.workflow-action.ts index 962c042c67..0178709eaa 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/ai-agent/ai-agent.workflow-action.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/ai-agent/ai-agent.workflow-action.ts @@ -18,6 +18,7 @@ import { } from 'src/modules/workflow/workflow-executor/exceptions/workflow-step-executor.exception'; import { type WorkflowActionInput } from 'src/modules/workflow/workflow-executor/types/workflow-action-input'; import { type WorkflowActionOutput } from 'src/modules/workflow/workflow-executor/types/workflow-action-output.type'; +import { findStepOrThrow } from 'src/modules/workflow/workflow-executor/utils/find-step-or-throw.util'; import { AiAgentExecutorService } from 'src/modules/workflow/workflow-executor/workflow-actions/ai-agent/services/ai-agent-executor.service'; import { isWorkflowAiAgentAction } from './guards/is-workflow-ai-agent-action.guard'; @@ -36,14 +37,10 @@ export class AiAgentWorkflowAction implements WorkflowAction { steps, context, }: WorkflowActionInput): Promise { - const step = steps.find((step) => step.id === currentStepId); - - if (!step) { - throw new WorkflowStepExecutorException( - 'Step not found', - WorkflowStepExecutorExceptionCode.STEP_NOT_FOUND, - ); - } + const step = findStepOrThrow({ + stepId: currentStepId, + steps, + }); if (!isWorkflowAiAgentAction(step)) { throw new WorkflowStepExecutorException( diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/code/code.workflow-action.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/code/code.workflow-action.ts index a06222e6a3..dc82e12854 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/code/code.workflow-action.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/code/code.workflow-action.ts @@ -12,6 +12,7 @@ import { } from 'src/modules/workflow/workflow-executor/exceptions/workflow-step-executor.exception'; import { type WorkflowActionInput } from 'src/modules/workflow/workflow-executor/types/workflow-action-input'; import { type WorkflowActionOutput } from 'src/modules/workflow/workflow-executor/types/workflow-action-output.type'; +import { findStepOrThrow } from 'src/modules/workflow/workflow-executor/utils/find-step-or-throw.util'; import { isWorkflowCodeAction } from 'src/modules/workflow/workflow-executor/workflow-actions/code/guards/is-workflow-code-action.guard'; import { type WorkflowCodeActionInput } from 'src/modules/workflow/workflow-executor/workflow-actions/code/types/workflow-code-action-input.type'; @@ -27,14 +28,10 @@ export class CodeWorkflowAction implements WorkflowAction { steps, context, }: WorkflowActionInput): Promise { - const step = steps.find((step) => step.id === currentStepId); - - if (!step) { - throw new WorkflowStepExecutorException( - 'Step not found', - WorkflowStepExecutorExceptionCode.STEP_NOT_FOUND, - ); - } + const step = findStepOrThrow({ + stepId: currentStepId, + steps, + }); if (!isWorkflowCodeAction(step)) { throw new WorkflowStepExecutorException( diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/filter/filter.workflow-action.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/filter/filter.workflow-action.ts index 41f9cda3d6..2f404cb3aa 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/filter/filter.workflow-action.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/filter/filter.workflow-action.ts @@ -10,6 +10,7 @@ import { } from 'src/modules/workflow/workflow-executor/exceptions/workflow-step-executor.exception'; import { type WorkflowActionInput } from 'src/modules/workflow/workflow-executor/types/workflow-action-input'; import { type WorkflowActionOutput } from 'src/modules/workflow/workflow-executor/types/workflow-action-output.type'; +import { findStepOrThrow } from 'src/modules/workflow/workflow-executor/utils/find-step-or-throw.util'; import { isWorkflowFilterAction } from 'src/modules/workflow/workflow-executor/workflow-actions/filter/guards/is-workflow-filter-action.guard'; import { evaluateFilterConditions } from 'src/modules/workflow/workflow-executor/workflow-actions/filter/utils/evaluate-filter-conditions.util'; @@ -18,14 +19,10 @@ export class FilterWorkflowAction implements WorkflowAction { async execute(input: WorkflowActionInput): Promise { const { currentStepId, steps, context } = input; - const step = steps.find((step) => step.id === currentStepId); - - if (!step) { - throw new WorkflowStepExecutorException( - 'Step not found', - WorkflowStepExecutorExceptionCode.STEP_NOT_FOUND, - ); - } + const step = findStepOrThrow({ + stepId: currentStepId, + steps, + }); if (!isWorkflowFilterAction(step)) { throw new WorkflowStepExecutorException( diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/form/form.workflow-action.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/form/form.workflow-action.ts index e184fafe03..ef449d8891 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/form/form.workflow-action.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/form/form.workflow-action.ts @@ -8,6 +8,7 @@ import { } from 'src/modules/workflow/workflow-executor/exceptions/workflow-step-executor.exception'; import { type WorkflowActionInput } from 'src/modules/workflow/workflow-executor/types/workflow-action-input'; import { type WorkflowActionOutput } from 'src/modules/workflow/workflow-executor/types/workflow-action-output.type'; +import { findStepOrThrow } from 'src/modules/workflow/workflow-executor/utils/find-step-or-throw.util'; import { isWorkflowFormAction } from 'src/modules/workflow/workflow-executor/workflow-actions/form/guards/is-workflow-form-action.guard'; @Injectable() @@ -16,14 +17,10 @@ export class FormWorkflowAction implements WorkflowAction { currentStepId, steps, }: WorkflowActionInput): Promise { - const step = steps.find((step) => step.id === currentStepId); - - if (!step) { - throw new WorkflowStepExecutorException( - 'Step not found', - WorkflowStepExecutorExceptionCode.STEP_NOT_FOUND, - ); - } + const step = findStepOrThrow({ + stepId: currentStepId, + steps, + }); if (!isWorkflowFormAction(step)) { throw new WorkflowStepExecutorException( diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/guards/is-workflow-iterator-action.guard.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/guards/is-workflow-iterator-action.guard.ts new file mode 100644 index 0000000000..a1f8684692 --- /dev/null +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/guards/is-workflow-iterator-action.guard.ts @@ -0,0 +1,10 @@ +import { + type WorkflowAction, + WorkflowActionType, + type WorkflowIteratorAction, +} from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type'; + +export const isWorkflowIteratorAction = ( + action: WorkflowAction, +): action is WorkflowIteratorAction => + action.type === WorkflowActionType.ITERATOR; diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/iterator-action.module.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/iterator-action.module.ts new file mode 100644 index 0000000000..5e6883f2fb --- /dev/null +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/iterator-action.module.ts @@ -0,0 +1,11 @@ +import { Module } from '@nestjs/common'; + +import { IteratorWorkflowAction } from 'src/modules/workflow/workflow-executor/workflow-actions/iterator/iterator.workflow-action'; +import { WorkflowRunModule } from 'src/modules/workflow/workflow-runner/workflow-run/workflow-run.module'; + +@Module({ + imports: [WorkflowRunModule], + providers: [IteratorWorkflowAction], + exports: [IteratorWorkflowAction], +}) +export class IteratorActionModule {} diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/iterator.workflow-action.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/iterator.workflow-action.ts new file mode 100644 index 0000000000..3861ba81b6 --- /dev/null +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/iterator.workflow-action.ts @@ -0,0 +1,161 @@ +import { Injectable } from '@nestjs/common'; + +import { isString } from '@sniptt/guards'; +import { isDefined, resolveInput } from 'twenty-shared/utils'; +import { StepStatus, WorkflowRunStepInfo } from 'twenty-shared/workflow'; + +import { WorkflowAction as WorkflowActionInterface } from 'src/modules/workflow/workflow-executor/interfaces/workflow-action.interface'; + +import { + WorkflowStepExecutorException, + WorkflowStepExecutorExceptionCode, +} from 'src/modules/workflow/workflow-executor/exceptions/workflow-step-executor.exception'; +import { type WorkflowActionInput } from 'src/modules/workflow/workflow-executor/types/workflow-action-input'; +import { type WorkflowActionOutput } from 'src/modules/workflow/workflow-executor/types/workflow-action-output.type'; +import { findStepOrThrow } from 'src/modules/workflow/workflow-executor/utils/find-step-or-throw.util'; +import { isWorkflowIteratorAction } from 'src/modules/workflow/workflow-executor/workflow-actions/iterator/guards/is-workflow-iterator-action.guard'; +import { type WorkflowIteratorActionInput } from 'src/modules/workflow/workflow-executor/workflow-actions/iterator/types/workflow-iterator-action-settings.type'; +import { WorkflowIteratorResult } from 'src/modules/workflow/workflow-executor/workflow-actions/iterator/types/workflow-iterator-result.type'; +import { getAllStepIdsInLoop } from 'src/modules/workflow/workflow-executor/workflow-actions/iterator/utils/get-all-step-ids-in-loop.util'; +import { type WorkflowAction } from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type'; +import { WorkflowRunWorkspaceService } from 'src/modules/workflow/workflow-runner/workflow-run/workflow-run.workspace-service'; + +const MAX_ITERATIONS = 10000; + +@Injectable() +export class IteratorWorkflowAction implements WorkflowActionInterface { + constructor( + private readonly workflowRunWorkspaceService: WorkflowRunWorkspaceService, + ) {} + + async execute(input: WorkflowActionInput): Promise { + const { currentStepId: iteratorStepId, steps, context, runInfo } = input; + + const step = findStepOrThrow({ + stepId: iteratorStepId, + steps, + }); + + if (!isWorkflowIteratorAction(step)) { + throw new WorkflowStepExecutorException( + 'Step is not an iterator action', + WorkflowStepExecutorExceptionCode.INVALID_STEP_TYPE, + ); + } + + const iteratorInput = resolveInput( + step.settings.input, + context, + ) as WorkflowIteratorActionInput; + + const { items, initialLoopStepIds } = iteratorInput; + + // TODO: remove once the UI is implemented + const parsedInitialLoopStepIds = isString(initialLoopStepIds) + ? JSON.parse(initialLoopStepIds) + : initialLoopStepIds; + + const parsedItems = isString(items) ? JSON.parse(items) : items; + + if (!Array.isArray(parsedItems)) { + throw new WorkflowStepExecutorException( + 'Iterator input items must be an array', + WorkflowStepExecutorExceptionCode.INTERNAL_ERROR, + ); + } + + if (parsedInitialLoopStepIds.length === 0 || parsedItems.length === 0) { + return { + result: { + currentItemIndex: 0, + currentItem: undefined, + hasProcessedAllItems: true, + } satisfies WorkflowIteratorResult, + }; + } + + const workflowRun = + await this.workflowRunWorkspaceService.getWorkflowRunOrFail({ + workflowRunId: runInfo.workflowRunId, + workspaceId: runInfo.workspaceId, + }); + + const stepInfos = workflowRun.state.stepInfos; + const existingIteratorStepResult = stepInfos[iteratorStepId] + ?.result as WorkflowIteratorResult; + + const currentItemIndex = isDefined(existingIteratorStepResult) + ? existingIteratorStepResult.currentItemIndex + 1 + : 0; + + if (currentItemIndex >= MAX_ITERATIONS) { + throw new WorkflowStepExecutorException( + 'Iterator has reached the maximum number of iterations', + WorkflowStepExecutorExceptionCode.INTERNAL_ERROR, + ); + } + + const hasProcessedAllItems = currentItemIndex >= parsedItems.length; + + const nextIteratorStepInfoResult: WorkflowIteratorResult = { + currentItemIndex, + currentItem: + currentItemIndex < parsedItems.length + ? parsedItems[currentItemIndex] + : undefined, + hasProcessedAllItems, + }; + + if (!hasProcessedAllItems) { + await this.resetStepsInLoop({ + iteratorStepId, + initialLoopStepIds: parsedInitialLoopStepIds, + workflowRunId: runInfo.workflowRunId, + workspaceId: runInfo.workspaceId, + steps, + }); + } + + return { + result: nextIteratorStepInfoResult, + shouldRemainRunning: !hasProcessedAllItems, + }; + } + + private async resetStepsInLoop({ + iteratorStepId, + initialLoopStepIds, + workflowRunId, + workspaceId, + steps, + }: { + iteratorStepId: string; + initialLoopStepIds: string[]; + workflowRunId: string; + workspaceId: string; + steps: WorkflowAction[]; + }) { + const stepIdsToReset = getAllStepIdsInLoop({ + iteratorStepId, + initialLoopStepIds, + steps, + }); + + await this.workflowRunWorkspaceService.updateWorkflowRunStepInfos({ + stepInfos: stepIdsToReset.reduce( + (acc, stepId) => { + acc[stepId] = { + status: StepStatus.NOT_STARTED, + result: undefined, + error: undefined, + }; + + return acc; + }, + {} as Record, + ), + workflowRunId, + workspaceId, + }); + } +} diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/types/workflow-iterator-action-settings.type.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/types/workflow-iterator-action-settings.type.ts new file mode 100644 index 0000000000..b5ea25755e --- /dev/null +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/types/workflow-iterator-action-settings.type.ts @@ -0,0 +1,12 @@ +import { type BaseWorkflowActionSettings } from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action-settings.type'; + +export type WorkflowIteratorActionInput = { + // eslint-disable-next-line @typescript-eslint/no-explicit-any + items?: Array | string; + // testing purpose, should never be a string nor undefined + initialLoopStepIds?: string[] | string; +}; + +export type WorkflowIteratorActionSettings = BaseWorkflowActionSettings & { + input: WorkflowIteratorActionInput; +}; diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/types/workflow-iterator-result.type.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/types/workflow-iterator-result.type.ts new file mode 100644 index 0000000000..badb5e8f49 --- /dev/null +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/types/workflow-iterator-result.type.ts @@ -0,0 +1,5 @@ +export type WorkflowIteratorResult = { + currentItemIndex: number; + currentItem?: unknown; + hasProcessedAllItems: boolean; +}; diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/utils/__tests__/can-execute-iterator-step.util.spec.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/utils/__tests__/can-execute-iterator-step.util.spec.ts new file mode 100644 index 0000000000..0ebf2a6ea0 --- /dev/null +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/utils/__tests__/can-execute-iterator-step.util.spec.ts @@ -0,0 +1,376 @@ +import { StepStatus } from 'twenty-shared/workflow'; + +import { canExecuteIteratorStep } from 'src/modules/workflow/workflow-executor/workflow-actions/iterator/utils/can-execute-iterator-step.util'; +import { + type WorkflowAction, + type WorkflowIteratorAction, + WorkflowActionType, +} from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type'; + +// Mock the getAllStepIdsInLoop utility +jest.mock( + 'src/modules/workflow/workflow-executor/workflow-actions/iterator/utils/get-all-step-ids-in-loop.util', + () => ({ + getAllStepIdsInLoop: jest.fn(), + }), +); + +const { getAllStepIdsInLoop } = jest.requireMock( + 'src/modules/workflow/workflow-executor/workflow-actions/iterator/utils/get-all-step-ids-in-loop.util', +); + +describe('canExecuteIteratorStep', () => { + const createMockIteratorStep = ( + id: string, + initialLoopStepIds: string[] | string = [], + ): WorkflowIteratorAction => ({ + id, + name: 'Iterator Step', + type: WorkflowActionType.ITERATOR, + settings: { + input: { + initialLoopStepIds, + items: [], + }, + errorHandlingOptions: { + continueOnFailure: { value: false }, + retryOnFailure: { value: false }, + }, + outputSchema: {}, + }, + valid: true, + nextStepIds: [], + }); + + const createMockStep = ( + id: string, + nextStepIds: string[] = [], + ): WorkflowAction => ({ + id, + name: 'Mock Step', + type: WorkflowActionType.CODE, + settings: { + input: { + serverlessFunctionId: 'mock-function-id', + serverlessFunctionVersion: 'mock-function-version', + serverlessFunctionInput: {}, + }, + errorHandlingOptions: { + continueOnFailure: { value: false }, + retryOnFailure: { value: false }, + }, + outputSchema: {}, + }, + valid: true, + nextStepIds, + }); + + beforeEach(() => { + jest.clearAllMocks(); + }); + + describe('when the step has not been started', () => { + it('should return true when all parent steps are successful', () => { + const iteratorStep = createMockIteratorStep('iterator-1', [ + 'step-1', + 'step-2', + ]); + const steps = [ + createMockStep('step-1', ['iterator-1']), + createMockStep('step-2', ['iterator-1']), + createMockStep('step-3', ['step-4']), + iteratorStep, + ]; + const stepInfos = { + 'step-1': { status: StepStatus.SUCCESS }, + 'step-2': { status: StepStatus.SUCCESS }, + 'step-3': { status: StepStatus.SUCCESS }, + }; + + // Mock getAllStepIdsInLoop to return the loop step IDs + getAllStepIdsInLoop.mockReturnValue(['step-1', 'step-2']); + + const result = canExecuteIteratorStep({ + step: iteratorStep, + steps, + stepInfos, + }); + + expect(result).toBe(true); + expect(getAllStepIdsInLoop).toHaveBeenCalledWith({ + iteratorStepId: 'iterator-1', + initialLoopStepIds: ['step-1', 'step-2'], + steps, + }); + }); + + it('should return false when some parent steps not in loop have failed', () => { + const iteratorStep = createMockIteratorStep('iterator-1', [ + 'step-1', + 'step-2', + ]); + const steps = [ + createMockStep('step-1', ['iterator-1']), + createMockStep('step-2', ['iterator-1']), + iteratorStep, + ]; + const stepInfos = { + 'step-1': { status: StepStatus.FAILED }, + 'step-2': { status: StepStatus.SUCCESS }, + }; + + getAllStepIdsInLoop.mockReturnValue(['step-2']); + + const result = canExecuteIteratorStep({ + step: iteratorStep, + steps, + stepInfos, + }); + + expect(result).toBe(false); + }); + + it('should return false when some parent steps are not started', () => { + const iteratorStep = createMockIteratorStep('iterator-1', [ + 'step-1', + 'step-2', + ]); + const steps = [ + createMockStep('step-1', ['iterator-1']), + createMockStep('step-2', ['iterator-1']), + iteratorStep, + ]; + const stepInfos = { + 'step-1': { status: StepStatus.SUCCESS }, + 'step-2': { status: StepStatus.NOT_STARTED }, + }; + + getAllStepIdsInLoop.mockReturnValue(['step-1']); + + const result = canExecuteIteratorStep({ + step: iteratorStep, + steps, + stepInfos, + }); + + expect(result).toBe(false); + }); + + it('should return true even if loop step is not started', () => { + const iteratorStep = createMockIteratorStep('iterator-1', ['step-1']); + const steps = [ + createMockStep('step-1', ['iterator-1']), // In loop + createMockStep('step-2', ['iterator-1']), // Not in loop + iteratorStep, + ]; + const stepInfos = { + 'step-1': { status: StepStatus.NOT_STARTED }, // This shouldn't affect the result + 'step-2': { status: StepStatus.SUCCESS }, + }; + + // Only step-1 is in the loop + getAllStepIdsInLoop.mockReturnValue(['step-1']); + + const result = canExecuteIteratorStep({ + step: iteratorStep, + steps, + stepInfos, + }); + + expect(result).toBe(true); + }); + + it('should return true when there are no parent steps targeting the iterator', () => { + const iteratorStep = createMockIteratorStep('iterator-1', []); + const steps = [ + createMockStep('step-1', ['step-2']), + createMockStep('step-2', []), + iteratorStep, + ]; + const stepInfos = { + 'step-1': { status: StepStatus.SUCCESS }, + 'step-2': { status: StepStatus.SUCCESS }, + }; + + getAllStepIdsInLoop.mockReturnValue([]); + + const result = canExecuteIteratorStep({ + step: iteratorStep, + steps, + stepInfos, + }); + + expect(result).toBe(true); + }); + + it('should handle undefined steps gracefully', () => { + const iteratorStep = createMockIteratorStep('iterator-1', ['step-1']); + const steps = [ + createMockStep('step-1', ['iterator-1']), + undefined as unknown as WorkflowAction, + iteratorStep, + ]; + const stepInfos = { + 'step-1': { status: StepStatus.SUCCESS }, + }; + + getAllStepIdsInLoop.mockReturnValue(['step-1']); + + const result = canExecuteIteratorStep({ + step: iteratorStep, + steps, + stepInfos, + }); + + expect(result).toBe(true); + }); + + it('should work correctly with multiple steps targeting the same iterator', () => { + const iteratorStep = createMockIteratorStep('iterator-1', [ + 'step-1', + 'step-2', + 'step-3', + ]); + const steps = [ + createMockStep('step-1', ['iterator-1']), + createMockStep('step-2', ['iterator-1']), + createMockStep('step-3', ['iterator-1']), + createMockStep('step-4', ['step-5']), // Not targeting iterator + iteratorStep, + ]; + const stepInfos = { + 'step-1': { status: StepStatus.SUCCESS }, + 'step-2': { status: StepStatus.SUCCESS }, + 'step-3': { status: StepStatus.SUCCESS }, + 'step-4': { status: StepStatus.FAILED }, + }; + + getAllStepIdsInLoop.mockReturnValue(['step-1', 'step-2', 'step-3']); + + const result = canExecuteIteratorStep({ + step: iteratorStep, + steps, + stepInfos, + }); + + expect(result).toBe(true); + }); + }); + + describe('when the step has been started', () => { + it('should return true if all the steps targeting the iterator have been successful', () => { + const iteratorStep = createMockIteratorStep('iterator-1', ['step-1']); + const steps = [ + createMockStep('step-1', ['iterator-1']), + createMockStep('step-2', ['iterator-1']), + iteratorStep, + ]; + const stepInfos = { + 'iterator-1': { status: StepStatus.RUNNING }, // Iterator has been started + 'step-1': { status: StepStatus.SUCCESS }, + 'step-2': { status: StepStatus.SUCCESS }, + }; + + getAllStepIdsInLoop.mockReturnValue(['step-1']); + + const result = canExecuteIteratorStep({ + step: iteratorStep, + steps, + stepInfos, + }); + + expect(result).toBe(true); + }); + + it('should return false if some of the steps targeting the iterator have failed', () => { + const iteratorStep = createMockIteratorStep('iterator-1', ['step-1']); + const steps = [ + createMockStep('step-1', ['iterator-1']), + createMockStep('step-2', ['iterator-1']), + iteratorStep, + ]; + const stepInfos = { + 'iterator-1': { status: StepStatus.RUNNING }, // Iterator has been started + 'step-1': { status: StepStatus.SUCCESS }, + 'step-2': { status: StepStatus.FAILED }, // This step failed + }; + + getAllStepIdsInLoop.mockReturnValue(['step-1']); + + const result = canExecuteIteratorStep({ + step: iteratorStep, + steps, + stepInfos, + }); + + expect(result).toBe(false); + }); + + it('should return false if some of the steps targeting the iterator are still running', () => { + const iteratorStep = createMockIteratorStep('iterator-1', ['step-1']); + const steps = [ + createMockStep('step-1', ['iterator-1']), + createMockStep('step-2', ['iterator-1']), + iteratorStep, + ]; + const stepInfos = { + 'iterator-1': { status: StepStatus.RUNNING }, // Iterator has been started + 'step-1': { status: StepStatus.SUCCESS }, + 'step-2': { status: StepStatus.RUNNING }, // This step is still running + }; + + getAllStepIdsInLoop.mockReturnValue(['step-1']); + + const result = canExecuteIteratorStep({ + step: iteratorStep, + steps, + stepInfos, + }); + + expect(result).toBe(false); + }); + + it('should return true when there are no steps targeting the iterator', () => { + const iteratorStep = createMockIteratorStep('iterator-1', []); + const steps = [createMockStep('step-1', ['step-2']), iteratorStep]; + const stepInfos = { + 'iterator-1': { status: StepStatus.RUNNING }, // Iterator has been started + 'step-1': { status: StepStatus.SUCCESS }, + }; + + getAllStepIdsInLoop.mockReturnValue([]); + + const result = canExecuteIteratorStep({ + step: iteratorStep, + steps, + stepInfos, + }); + + expect(result).toBe(true); + }); + + it('should check all steps targeting iterator including loop steps when iterator has been started', () => { + const iteratorStep = createMockIteratorStep('iterator-1', ['step-1']); + const steps = [ + createMockStep('step-1', ['iterator-1']), // In loop and targeting iterator + createMockStep('step-2', ['iterator-1']), // Not in loop but targeting iterator + iteratorStep, + ]; + const stepInfos = { + 'iterator-1': { status: StepStatus.SUCCESS }, // Iterator has been started + 'step-1': { status: StepStatus.SUCCESS }, // Loop step successful + 'step-2': { status: StepStatus.SUCCESS }, // Non-loop step successful + }; + + getAllStepIdsInLoop.mockReturnValue(['step-1']); + + const result = canExecuteIteratorStep({ + step: iteratorStep, + steps, + stepInfos, + }); + + expect(result).toBe(true); + }); + }); +}); diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/utils/__tests__/get-all-step-ids-in-loop.util.spec.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/utils/__tests__/get-all-step-ids-in-loop.util.spec.ts new file mode 100644 index 0000000000..27dd4d7214 --- /dev/null +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/utils/__tests__/get-all-step-ids-in-loop.util.spec.ts @@ -0,0 +1,228 @@ +import { type WorkflowCodeActionSettings } from 'src/modules/workflow/workflow-executor/workflow-actions/code/types/workflow-code-action-settings.type'; +import { type WorkflowIteratorActionInput } from 'src/modules/workflow/workflow-executor/workflow-actions/iterator/types/workflow-iterator-action-settings.type'; +import { getAllStepIdsInLoop } from 'src/modules/workflow/workflow-executor/workflow-actions/iterator/utils/get-all-step-ids-in-loop.util'; +import { + type WorkflowAction, + WorkflowActionType, + type WorkflowCodeAction, + type WorkflowIteratorAction, +} from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type'; + +describe('getAllStepIdsInLoop', () => { + const createCodeMockStep = ( + id: string, + nextStepIds: string[], + ): WorkflowCodeAction => ({ + id, + name: `Step ${id}`, + type: WorkflowActionType.CODE, + valid: true, + nextStepIds, + settings: { + input: {}, + outputSchema: {}, + errorHandlingOptions: { + continueOnFailure: { value: false }, + retryOnFailure: { value: false }, + }, + } as WorkflowCodeActionSettings, + }); + + const createIteratorMockStep = ( + id: string, + nextStepIds: string[], + initialLoopStepIds: string[], + ): WorkflowIteratorAction => ({ + id, + name: `Step ${id}`, + type: WorkflowActionType.ITERATOR, + valid: true, + nextStepIds, + settings: { + input: (initialLoopStepIds + ? ({ initialLoopStepIds } as WorkflowIteratorActionInput) + : {}) as WorkflowIteratorActionInput, + outputSchema: {}, + errorHandlingOptions: { + continueOnFailure: { value: false }, + retryOnFailure: { value: false }, + }, + }, + }); + + describe('simple loop scenarios', () => { + it('should return all step IDs in a simple linear loop', () => { + const steps: WorkflowAction[] = [ + createIteratorMockStep('iterator1', ['step2'], []), + createCodeMockStep('step2', ['step3']), + createCodeMockStep('step3', ['step4']), + createCodeMockStep('step4', ['iterator1']), // loops back + ]; + + const result = getAllStepIdsInLoop({ + iteratorStepId: 'iterator1', + initialLoopStepIds: ['step2'], + steps, + }); + + expect(result).toEqual(['step2', 'step3', 'step4']); + }); + + it('should handle loop with branching paths that converge', () => { + const steps: WorkflowAction[] = [ + createIteratorMockStep('iterator1', ['step2'], []), + createCodeMockStep('step2', ['step3', 'step4']), + createCodeMockStep('step3', ['step5']), + createCodeMockStep('step4', ['step5']), + createCodeMockStep('step5', ['iterator1']), // loops back + ]; + + const result = getAllStepIdsInLoop({ + iteratorStepId: 'iterator1', + initialLoopStepIds: ['step2'], + steps, + }); + + expect(result).toEqual(['step2', 'step3', 'step5', 'step4']); + }); + + it('should handle loop with branching paths that converge to the iterator', () => { + const steps: WorkflowAction[] = [ + createIteratorMockStep('iterator1', ['step2'], []), + createCodeMockStep('step2', ['step3', 'step4']), + createCodeMockStep('step3', ['iterator1']), + createCodeMockStep('step4', ['iterator1']), + ]; + + const result = getAllStepIdsInLoop({ + iteratorStepId: 'iterator1', + initialLoopStepIds: ['step2'], + steps, + }); + + expect(result).toEqual(['step2', 'step3', 'step4']); + }); + + it('should handle multiple entry points to the loop', () => { + const steps: WorkflowAction[] = [ + createIteratorMockStep('iterator1', ['step2', 'step3'], []), + createCodeMockStep('step2', ['step4']), + createCodeMockStep('step3', ['step4']), + createCodeMockStep('step4', ['iterator1']), // loops back + ]; + + const result = getAllStepIdsInLoop({ + iteratorStepId: 'iterator1', + initialLoopStepIds: ['step2', 'step3'], + steps, + }); + + expect(result).toEqual(['step2', 'step4', 'step3']); + }); + }); + + describe('nested iterator scenarios', () => { + it('should handle a nested iterator within a loop', () => { + const steps: WorkflowAction[] = [ + createIteratorMockStep('iterator1', ['step2'], []), + createCodeMockStep('step2', ['nested_iterator']), + createIteratorMockStep('nested_iterator', ['step5'], ['step3']), + createCodeMockStep('step3', ['step4']), + createCodeMockStep('step4', ['nested_iterator']), // loops back to nested iterator + createCodeMockStep('step5', ['iterator1']), // loops back to main iterator + ]; + + const result = getAllStepIdsInLoop({ + iteratorStepId: 'iterator1', + initialLoopStepIds: ['step2'], + steps, + }); + + expect(result).toEqual([ + 'step2', + 'nested_iterator', + 'step3', + 'step4', + 'step5', + ]); + }); + + it('should handle multiple levels of nested iterators', () => { + const steps: WorkflowAction[] = [ + createIteratorMockStep('iterator1', ['step2'], []), + createCodeMockStep('step2', ['nested_iterator1']), + createIteratorMockStep('nested_iterator1', ['step6'], ['step3']), + createCodeMockStep('step3', ['nested_iterator2']), + createIteratorMockStep('nested_iterator2', ['step5'], ['step4']), + createCodeMockStep('step4', ['nested_iterator2']), // loops back to nested iterator2 + createCodeMockStep('step5', ['nested_iterator1']), // loops back to nested iterator1 + createCodeMockStep('step6', ['iterator1']), // loops back to main iterator + ]; + + const result = getAllStepIdsInLoop({ + iteratorStepId: 'iterator1', + initialLoopStepIds: ['step2'], + steps, + }); + + expect(result).toEqual([ + 'step2', + 'nested_iterator1', + 'step3', + 'nested_iterator2', + 'step4', + 'step5', + 'step6', + ]); + }); + }); + + describe('edge cases', () => { + it('should handle empty initial loop step IDs', () => { + const steps: WorkflowAction[] = [ + createIteratorMockStep('iterator1', ['step2'], []), + ]; + + const result = getAllStepIdsInLoop({ + iteratorStepId: 'iterator1', + initialLoopStepIds: [], + steps, + }); + + expect(result).toEqual([]); + }); + + it('should handle steps with no nextStepIds', () => { + const steps: WorkflowAction[] = [ + createIteratorMockStep('iterator1', ['step2'], []), + createCodeMockStep('step2', []), // no nextStepIds + ]; + + const result = getAllStepIdsInLoop({ + iteratorStepId: 'iterator1', + initialLoopStepIds: ['step2'], + steps, + }); + + expect(result).toEqual(['step2']); + }); + + it('should prevent infinite loops with circular references', () => { + const steps: WorkflowAction[] = [ + createIteratorMockStep('iterator1', ['step2'], []), + createCodeMockStep('step2', ['step3']), + createCodeMockStep('step3', ['step4']), + createCodeMockStep('step4', ['step2']), // circular reference + ]; + + const result = getAllStepIdsInLoop({ + iteratorStepId: 'iterator1', + initialLoopStepIds: ['step2'], + steps, + }); + + // Should still include all steps but not get stuck in infinite loop + expect(result).toEqual(['step2', 'step3', 'step4']); + }); + }); +}); diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/utils/can-execute-iterator-step.util.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/utils/can-execute-iterator-step.util.ts new file mode 100644 index 0000000000..2b2567c708 --- /dev/null +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/utils/can-execute-iterator-step.util.ts @@ -0,0 +1,54 @@ +import { isString } from 'class-validator'; +import { isDefined } from 'twenty-shared/utils'; +import { StepStatus, type WorkflowRunStepInfos } from 'twenty-shared/workflow'; + +import { stepHasBeenStarted } from 'src/modules/workflow/workflow-executor/utils/step-has-been-started.util'; +import { getAllStepIdsInLoop } from 'src/modules/workflow/workflow-executor/workflow-actions/iterator/utils/get-all-step-ids-in-loop.util'; +import { + type WorkflowAction, + type WorkflowIteratorAction, +} from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type'; + +export const canExecuteIteratorStep = ({ + step, + steps, + stepInfos, +}: { + step: WorkflowIteratorAction; + steps: WorkflowAction[]; + stepInfos: WorkflowRunStepInfos; +}) => { + const stepsTargetingIterator = steps.filter( + (parentStep) => + isDefined(parentStep) && parentStep.nextStepIds?.includes(step.id), + ); + + // If the step has been started, we need to check if all the steps targeting the iterator have been successful. + if (stepHasBeenStarted(step.id, stepInfos)) { + return stepsTargetingIterator.every( + (step) => stepInfos[step.id]?.status === StepStatus.SUCCESS, + ); + } else { + // On the first iteration, the loop steps will not have been started yet. + // TODO: remove parsing once the UI is implemented + const parsedInitialLoopStepIds = isString( + step.settings.input.initialLoopStepIds, + ) + ? JSON.parse(step.settings.input.initialLoopStepIds) + : step.settings.input.initialLoopStepIds; + + const stepIdsInLoop = getAllStepIdsInLoop({ + iteratorStepId: step.id, + initialLoopStepIds: parsedInitialLoopStepIds, + steps, + }); + + const parentSteps = stepsTargetingIterator.filter( + (step) => !stepIdsInLoop.includes(step.id), + ); + + return parentSteps.every( + (step) => stepInfos[step.id]?.status === StepStatus.SUCCESS, + ); + } +}; diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/utils/get-all-step-ids-in-loop.util.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/utils/get-all-step-ids-in-loop.util.ts new file mode 100644 index 0000000000..f584ce07d7 --- /dev/null +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/utils/get-all-step-ids-in-loop.util.ts @@ -0,0 +1,91 @@ +import { isWorkflowIteratorAction } from 'src/modules/workflow/workflow-executor/workflow-actions/iterator/guards/is-workflow-iterator-action.guard'; +import { type WorkflowIteratorActionInput } from 'src/modules/workflow/workflow-executor/workflow-actions/iterator/types/workflow-iterator-action-settings.type'; +import { type WorkflowAction } from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type'; + +const traverseSteps = ({ + iteratorStepId, + stepIds, + steps, + visitedStepIds, + allStepIdsInLoop, +}: { + iteratorStepId: string; + stepIds: string[]; + steps: WorkflowAction[]; + visitedStepIds: Set; + allStepIdsInLoop: Set; +}) => { + for (const stepId of stepIds) { + if (visitedStepIds.has(stepId)) { + continue; + } + + visitedStepIds.add(stepId); + allStepIdsInLoop.add(stepId); + + const step = steps.find((s) => s.id === stepId); + + if (!step || !step.nextStepIds) { + continue; + } + + if (isWorkflowIteratorAction(step)) { + const nestedIteratorInput = step.settings + .input as WorkflowIteratorActionInput; + + if ( + nestedIteratorInput.initialLoopStepIds && + // TODO: To remove once we remove the string input for initialLoopStepIds + Array.isArray(nestedIteratorInput.initialLoopStepIds) + ) { + const nestedLoopStepIds = getAllStepIdsInLoop({ + iteratorStepId: stepId, + initialLoopStepIds: nestedIteratorInput.initialLoopStepIds, + steps, + }); + + nestedLoopStepIds.forEach((nestedStepId) => { + allStepIdsInLoop.add(nestedStepId); + }); + } + } + + const connectsBackToIterator = step.nextStepIds.includes(iteratorStepId); + + if (connectsBackToIterator) { + // We've found the end of the loop, stop traversing + continue; + } + + traverseSteps({ + iteratorStepId, + stepIds: step.nextStepIds, + steps, + visitedStepIds, + allStepIdsInLoop, + }); + } +}; + +export const getAllStepIdsInLoop = ({ + iteratorStepId, + initialLoopStepIds, + steps, +}: { + iteratorStepId: string; + initialLoopStepIds: string[]; + steps: WorkflowAction[]; +}): string[] => { + const allStepIdsInLoop = new Set(); + const visitedStepIds = new Set(); + + traverseSteps({ + iteratorStepId, + stepIds: initialLoopStepIds, + steps, + visitedStepIds, + allStepIdsInLoop, + }); + + return Array.from(allStepIdsInLoop); +}; diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/record-crud/create-record.workflow-action.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/record-crud/create-record.workflow-action.ts index b0858081b8..8790ef1136 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/record-crud/create-record.workflow-action.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/record-crud/create-record.workflow-action.ts @@ -2,8 +2,8 @@ import { Injectable } from '@nestjs/common'; import { InjectRepository } from '@nestjs/typeorm'; import { isDefined } from 'class-validator'; -import { Repository } from 'typeorm'; import { resolveInput } from 'twenty-shared/utils'; +import { Repository } from 'typeorm'; import { type WorkflowAction } from 'src/modules/workflow/workflow-executor/interfaces/workflow-action.interface'; @@ -14,17 +14,13 @@ import { ObjectMetadataEntity } from 'src/engine/metadata-modules/object-metadat import { ScopedWorkspaceContextFactory } from 'src/engine/twenty-orm/factories/scoped-workspace-context.factory'; import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager'; import { WorkflowCommonWorkspaceService } from 'src/modules/workflow/common/workspace-services/workflow-common.workspace-service'; -import { - WorkflowStepExecutorException, - WorkflowStepExecutorExceptionCode, -} from 'src/modules/workflow/workflow-executor/exceptions/workflow-step-executor.exception'; import { type WorkflowActionInput } from 'src/modules/workflow/workflow-executor/types/workflow-action-input'; import { type WorkflowActionOutput } from 'src/modules/workflow/workflow-executor/types/workflow-action-output.type'; +import { findStepOrThrow } from 'src/modules/workflow/workflow-executor/utils/find-step-or-throw.util'; import { RecordCRUDActionException, RecordCRUDActionExceptionCode, } from 'src/modules/workflow/workflow-executor/workflow-actions/record-crud/exceptions/record-crud-action.exception'; -import { isWorkflowCreateRecordAction } from 'src/modules/workflow/workflow-executor/workflow-actions/record-crud/guards/is-workflow-create-record-action.guard'; import { type WorkflowCreateRecordActionInput } from 'src/modules/workflow/workflow-executor/workflow-actions/record-crud/types/workflow-record-crud-action-input.type'; @Injectable() @@ -44,21 +40,10 @@ export class CreateRecordWorkflowAction implements WorkflowAction { steps, context, }: WorkflowActionInput): Promise { - const step = steps.find((step) => step.id === currentStepId); - - if (!step) { - throw new WorkflowStepExecutorException( - 'Step not found', - WorkflowStepExecutorExceptionCode.STEP_NOT_FOUND, - ); - } - - if (!isWorkflowCreateRecordAction(step)) { - throw new WorkflowStepExecutorException( - 'Step is not a create record action', - WorkflowStepExecutorExceptionCode.INVALID_STEP_TYPE, - ); - } + const step = findStepOrThrow({ + steps, + stepId: currentStepId, + }); const workspaceId = this.scopedWorkspaceContextFactory.create().workspaceId; diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/record-crud/delete-record.workflow-action.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/record-crud/delete-record.workflow-action.ts index 0fdb2ec3d5..a427dbb1a8 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/record-crud/delete-record.workflow-action.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/record-crud/delete-record.workflow-action.ts @@ -13,6 +13,7 @@ import { } from 'src/modules/workflow/workflow-executor/exceptions/workflow-step-executor.exception'; import { type WorkflowActionInput } from 'src/modules/workflow/workflow-executor/types/workflow-action-input'; import { type WorkflowActionOutput } from 'src/modules/workflow/workflow-executor/types/workflow-action-output.type'; +import { findStepOrThrow } from 'src/modules/workflow/workflow-executor/utils/find-step-or-throw.util'; import { RecordCRUDActionException, RecordCRUDActionExceptionCode, @@ -32,14 +33,11 @@ export class DeleteRecordWorkflowAction implements WorkflowAction { steps, context, }: WorkflowActionInput): Promise { - const step = steps.find((step) => step.id === currentStepId); + const step = findStepOrThrow({ + steps, + stepId: currentStepId, + }); - if (!step) { - throw new WorkflowStepExecutorException( - 'Step not found', - WorkflowStepExecutorExceptionCode.STEP_NOT_FOUND, - ); - } if (!isWorkflowDeleteRecordAction(step)) { throw new WorkflowStepExecutorException( 'Step is not a delete record action', diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/record-crud/find-records.workflow-action.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/record-crud/find-records.workflow-action.ts index acc12bcba0..f96eb88e0f 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/record-crud/find-records.workflow-action.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/record-crud/find-records.workflow-action.ts @@ -25,6 +25,7 @@ import { } from 'src/modules/workflow/workflow-executor/exceptions/workflow-step-executor.exception'; import { type WorkflowActionInput } from 'src/modules/workflow/workflow-executor/types/workflow-action-input'; import { type WorkflowActionOutput } from 'src/modules/workflow/workflow-executor/types/workflow-action-output.type'; +import { findStepOrThrow } from 'src/modules/workflow/workflow-executor/utils/find-step-or-throw.util'; import { RecordCRUDActionException, RecordCRUDActionExceptionCode, @@ -45,14 +46,10 @@ export class FindRecordsWorkflowAction implements WorkflowAction { steps, context, }: WorkflowActionInput): Promise { - const step = steps.find((step) => step.id === currentStepId); - - if (!step) { - throw new WorkflowStepExecutorException( - 'Step not found', - WorkflowStepExecutorExceptionCode.STEP_NOT_FOUND, - ); - } + const step = findStepOrThrow({ + steps, + stepId: currentStepId, + }); if (!isWorkflowFindRecordsAction(step)) { throw new WorkflowStepExecutorException( diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/record-crud/update-record.workflow-action.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/record-crud/update-record.workflow-action.ts index f5ef2f9fed..60071719c7 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/record-crud/update-record.workflow-action.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/record-crud/update-record.workflow-action.ts @@ -15,6 +15,7 @@ import { } from 'src/modules/workflow/workflow-executor/exceptions/workflow-step-executor.exception'; import { type WorkflowActionInput } from 'src/modules/workflow/workflow-executor/types/workflow-action-input'; import { type WorkflowActionOutput } from 'src/modules/workflow/workflow-executor/types/workflow-action-output.type'; +import { findStepOrThrow } from 'src/modules/workflow/workflow-executor/utils/find-step-or-throw.util'; import { RecordCRUDActionException, RecordCRUDActionExceptionCode, @@ -36,14 +37,10 @@ export class UpdateRecordWorkflowAction implements WorkflowAction { steps, context, }: WorkflowActionInput): Promise { - const step = steps.find((step) => step.id === currentStepId); - - if (!step) { - throw new WorkflowStepExecutorException( - 'Step not found', - WorkflowStepExecutorExceptionCode.STEP_NOT_FOUND, - ); - } + const step = findStepOrThrow({ + steps, + stepId: currentStepId, + }); if (!isWorkflowUpdateRecordAction(step)) { throw new WorkflowStepExecutorException( diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action-settings.type.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action-settings.type.ts index aa0f9879a8..2cfe77e8e8 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action-settings.type.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action-settings.type.ts @@ -4,6 +4,7 @@ import { type WorkflowCodeActionSettings } from 'src/modules/workflow/workflow-e import { type WorkflowFilterActionSettings } from 'src/modules/workflow/workflow-executor/workflow-actions/filter/types/workflow-filter-action-settings.type'; import { type WorkflowFormActionSettings } from 'src/modules/workflow/workflow-executor/workflow-actions/form/types/workflow-form-action-settings.type'; import { type WorkflowHttpRequestActionSettings } from 'src/modules/workflow/workflow-executor/workflow-actions/http-request/types/workflow-http-request-action-settings.type'; +import { type WorkflowIteratorActionSettings } from 'src/modules/workflow/workflow-executor/workflow-actions/iterator/types/workflow-iterator-action-settings.type'; import { type WorkflowSendEmailActionSettings } from 'src/modules/workflow/workflow-executor/workflow-actions/mail-sender/types/workflow-send-email-action-settings.type'; import { type WorkflowCreateRecordActionSettings, @@ -34,4 +35,5 @@ export type WorkflowActionSettings = | WorkflowFormActionSettings | WorkflowFilterActionSettings | WorkflowHttpRequestActionSettings - | WorkflowAiAgentActionSettings; + | WorkflowAiAgentActionSettings + | WorkflowIteratorActionSettings; diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type.ts index 74ccb72dc7..fe86a839d0 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type.ts @@ -3,6 +3,7 @@ import { type WorkflowCodeActionSettings } from 'src/modules/workflow/workflow-e import { type WorkflowFilterActionSettings } from 'src/modules/workflow/workflow-executor/workflow-actions/filter/types/workflow-filter-action-settings.type'; import { type WorkflowFormActionSettings } from 'src/modules/workflow/workflow-executor/workflow-actions/form/types/workflow-form-action-settings.type'; import { type WorkflowHttpRequestActionSettings } from 'src/modules/workflow/workflow-executor/workflow-actions/http-request/types/workflow-http-request-action-settings.type'; +import { type WorkflowIteratorActionSettings } from 'src/modules/workflow/workflow-executor/workflow-actions/iterator/types/workflow-iterator-action-settings.type'; import { type WorkflowSendEmailActionSettings } from 'src/modules/workflow/workflow-executor/workflow-actions/mail-sender/types/workflow-send-email-action-settings.type'; import { type WorkflowCreateRecordActionSettings, @@ -23,6 +24,7 @@ export enum WorkflowActionType { FILTER = 'FILTER', HTTP_REQUEST = 'HTTP_REQUEST', AI_AGENT = 'AI_AGENT', + ITERATOR = 'ITERATOR', } type BaseWorkflowAction = { @@ -88,6 +90,11 @@ export type WorkflowAiAgentAction = BaseWorkflowAction & { settings: WorkflowAiAgentActionSettings; }; +export type WorkflowIteratorAction = BaseWorkflowAction & { + type: WorkflowActionType.ITERATOR; + settings: WorkflowIteratorActionSettings; +}; + export type WorkflowAction = | WorkflowCodeAction | WorkflowSendEmailAction @@ -98,4 +105,5 @@ export type WorkflowAction = | WorkflowFormAction | WorkflowFilterAction | WorkflowHttpRequestAction - | WorkflowAiAgentAction; + | WorkflowAiAgentAction + | WorkflowIteratorAction; diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-executor.module.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-executor.module.ts index 1cba6cae86..4687c35db2 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-executor.module.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-executor.module.ts @@ -10,6 +10,7 @@ import { AiAgentActionModule } from 'src/modules/workflow/workflow-executor/work import { CodeActionModule } from 'src/modules/workflow/workflow-executor/workflow-actions/code/code-action.module'; import { FilterActionModule } from 'src/modules/workflow/workflow-executor/workflow-actions/filter/filter-action.module'; import { FormActionModule } from 'src/modules/workflow/workflow-executor/workflow-actions/form/form-action.module'; +import { IteratorActionModule } from 'src/modules/workflow/workflow-executor/workflow-actions/iterator/iterator-action.module'; import { RecordCRUDActionModule } from 'src/modules/workflow/workflow-executor/workflow-actions/record-crud/record-crud-action.module'; import { ToolExecutorWorkflowAction } from 'src/modules/workflow/workflow-executor/workflow-actions/tool-executor-workflow-action'; import { WorkflowExecutorWorkspaceService } from 'src/modules/workflow/workflow-executor/workspace-services/workflow-executor.workspace-service'; @@ -24,6 +25,7 @@ import { WorkflowRunModule } from 'src/modules/workflow/workflow-runner/workflow WorkflowRunModule, BillingModule, FilterActionModule, + IteratorActionModule, AiAgentActionModule, FeatureFlagModule, AiModule, diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/workspace-services/__tests__/workflow-executor.workspace-service.spec.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/workspace-services/__tests__/workflow-executor.workspace-service.spec.ts index 9a9a0053e1..053ed46316 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-executor/workspace-services/__tests__/workflow-executor.workspace-service.spec.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/workspace-services/__tests__/workflow-executor.workspace-service.spec.ts @@ -8,13 +8,13 @@ import { BillingMeterEventName } from 'src/engine/core-modules/billing/enums/bil import { BillingService } from 'src/engine/core-modules/billing/services/billing.service'; import { WorkspaceEventEmitter } from 'src/engine/workspace-event-emitter/workspace-event-emitter'; import { WorkflowActionFactory } from 'src/modules/workflow/workflow-executor/factories/workflow-action.factory'; +import { canExecuteStep } from 'src/modules/workflow/workflow-executor/utils/can-execute-step.util'; import { type WorkflowAction, WorkflowActionType, } from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type'; import { WorkflowExecutorWorkspaceService } from 'src/modules/workflow/workflow-executor/workspace-services/workflow-executor.workspace-service'; import { WorkflowRunWorkspaceService } from 'src/modules/workflow/workflow-runner/workflow-run/workflow-run.workspace-service'; -import { canExecuteStep } from 'src/modules/workflow/workflow-executor/utils/can-execute-step.util'; jest.mock( 'src/modules/workflow/workflow-executor/utils/can-execute-step.util', @@ -155,6 +155,10 @@ describe('WorkflowExecutorWorkspaceService', () => { currentStepId: 'step-1', steps: mockSteps, context: getWorkflowRunContext(mockStepInfos), + runInfo: { + workflowRunId: mockWorkflowRunId, + workspaceId: mockWorkspaceId, + }, }); expect(workspaceEventEmitter.emitCustomBatchEvent).toHaveBeenCalledWith( diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/workspace-services/workflow-executor.workspace-service.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/workspace-services/workflow-executor.workspace-service.ts index b948a6a89b..b93e415017 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-executor/workspace-services/workflow-executor.workspace-service.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/workspace-services/workflow-executor.workspace-service.ts @@ -1,5 +1,6 @@ import { Injectable } from '@nestjs/common'; +import { isString } from '@sniptt/guards'; import { isDefined } from 'twenty-shared/utils'; import { getWorkflowRunContext, @@ -24,6 +25,9 @@ import { import { canExecuteStep } from 'src/modules/workflow/workflow-executor/utils/can-execute-step.util'; import { workflowShouldFail } from 'src/modules/workflow/workflow-executor/utils/workflow-should-fail.util'; import { workflowShouldKeepRunning } from 'src/modules/workflow/workflow-executor/utils/workflow-should-keep-running.util'; +import { isWorkflowIteratorAction } from 'src/modules/workflow/workflow-executor/workflow-actions/iterator/guards/is-workflow-iterator-action.guard'; +import { WorkflowIteratorResult } from 'src/modules/workflow/workflow-executor/workflow-actions/iterator/types/workflow-iterator-result.type'; +import { WorkflowAction } from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type'; import { WorkflowRunWorkspaceService } from 'src/modules/workflow/workflow-runner/workflow-run/workflow-run.workspace-service'; @Injectable() @@ -106,6 +110,7 @@ export class WorkflowExecutorWorkspaceService { await this.workflowRunWorkspaceService.updateWorkflowRunStepInfo({ stepId, stepInfo: { + ...stepInfos[stepId], status: StepStatus.RUNNING, }, workflowRunId, @@ -117,6 +122,10 @@ export class WorkflowExecutorWorkspaceService { currentStepId: stepId, steps, context: getWorkflowRunContext(stepInfos), + runInfo: { + workflowRunId, + workspaceId, + }, }); } catch (error) { actionOutput = { @@ -129,60 +138,33 @@ export class WorkflowExecutorWorkspaceService { }; } - const isPendingEvent = actionOutput.pendingEvent; - - const isSuccess = isDefined(actionOutput.result); - const isError = isDefined(actionOutput.error); - const isStopped = actionOutput.shouldEndWorkflowRun; - if (!isError) { this.sendWorkflowNodeRunEvent(workspaceId); } - let stepInfo: WorkflowRunStepInfo; - - if (isPendingEvent) { - stepInfo = { - status: StepStatus.PENDING, - }; - } else if (isStopped) { - stepInfo = { - status: StepStatus.STOPPED, - result: actionOutput?.result, - }; - } else if (isSuccess) { - stepInfo = { - status: StepStatus.SUCCESS, - result: actionOutput?.result, - }; - } else { - stepInfo = { - status: StepStatus.FAILED, - error: actionOutput?.error, - }; - } - - await this.workflowRunWorkspaceService.updateWorkflowRunStepInfo({ + const { shouldProcessNextSteps } = await this.processStepExecutionResult({ + actionOutput, stepId, - stepInfo, workflowRunId, workspaceId, }); - if ( - isSuccess && - !isStopped && - isDefined(stepToExecute.nextStepIds) && - stepToExecute.nextStepIds.length > 0 - ) { - await this.executeFromSteps({ - stepIds: stepToExecute.nextStepIds, - workflowRunId, - workspaceId, - shouldComputeWorkflowRunStatus: false, + if (shouldProcessNextSteps) { + const nextStepIdsToExecute = await this.getNextStepIdsToExecute({ + executedStep: stepToExecute, + executedStepResult: actionOutput, }); + + if (isDefined(nextStepIdsToExecute) && nextStepIdsToExecute.length > 0) { + await this.executeFromSteps({ + stepIds: nextStepIdsToExecute, + workflowRunId, + workspaceId, + shouldComputeWorkflowRunStatus: false, + }); + } } } @@ -247,4 +229,83 @@ export class WorkflowExecutorWorkspaceService { )) ); } + + private async getNextStepIdsToExecute({ + executedStep, + executedStepResult, + }: { + executedStep: WorkflowAction; + executedStepResult: WorkflowActionOutput; + }) { + const isIteratorStep = isWorkflowIteratorAction(executedStep); + + if (isIteratorStep) { + const iteratorStepResult = + executedStepResult.result as WorkflowIteratorResult; + + if (!iteratorStepResult.hasProcessedAllItems) { + return isString(executedStep.settings.input.initialLoopStepIds) + ? JSON.parse(executedStep.settings.input.initialLoopStepIds) + : executedStep.settings.input.initialLoopStepIds; + } + } + + return executedStep.nextStepIds; + } + + private async processStepExecutionResult({ + actionOutput, + stepId, + workflowRunId, + workspaceId, + }: { + actionOutput: WorkflowActionOutput; + stepId: string; + workflowRunId: string; + workspaceId: string; + }): Promise<{ shouldProcessNextSteps: boolean }> { + const isPendingEvent = actionOutput.pendingEvent; + const isSuccess = isDefined(actionOutput.result); + const isStopped = actionOutput.shouldEndWorkflowRun; + const isNotFinished = actionOutput.shouldRemainRunning; + + let stepInfo: WorkflowRunStepInfo; + + if (isPendingEvent) { + stepInfo = { + status: StepStatus.PENDING, + }; + } else if (isStopped) { + stepInfo = { + status: StepStatus.STOPPED, + result: actionOutput?.result, + }; + } else if (isNotFinished) { + stepInfo = { + status: StepStatus.RUNNING, + result: actionOutput?.result, + }; + } else if (isSuccess) { + stepInfo = { + status: StepStatus.SUCCESS, + result: actionOutput?.result, + }; + } else { + stepInfo = { + status: StepStatus.FAILED, + error: actionOutput?.error, + }; + } + + await this.workflowRunWorkspaceService.updateWorkflowRunStepInfo({ + stepId, + stepInfo, + workflowRunId, + workspaceId, + }); + + return { + shouldProcessNextSteps: isSuccess && !isStopped, + }; + } } diff --git a/packages/twenty-server/src/modules/workflow/workflow-runner/workflow-run/workflow-run.workspace-service.ts b/packages/twenty-server/src/modules/workflow/workflow-runner/workflow-run/workflow-run.workspace-service.ts index 9a41143cd3..c3d18240e1 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-runner/workflow-run/workflow-run.workspace-service.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-runner/workflow-run/workflow-run.workspace-service.ts @@ -254,6 +254,38 @@ export class WorkflowRunWorkspaceService { await this.updateWorkflowRun({ workflowRunId, workspaceId, partialUpdate }); } + @WithLock('workflowRunId') + async updateWorkflowRunStepInfos({ + stepInfos, + workflowRunId, + workspaceId, + }: { + stepInfos: Record; + workflowRunId: string; + workspaceId: string; + }) { + const workflowRunToUpdate = await this.getWorkflowRunOrFail({ + workflowRunId, + workspaceId, + }); + + const partialUpdate = { + state: { + ...workflowRunToUpdate.state, + stepInfos: { + ...workflowRunToUpdate.state?.stepInfos, + ...stepInfos, + }, + }, + }; + + await this.updateWorkflowRun({ + workflowRunId, + workspaceId, + partialUpdate, + }); + } + @WithLock('workflowRunId') async updateWorkflowRunStep({ workflowRunId,