diff --git a/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/__tests__/iterator-action.workflow-action.spec.ts b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/__tests__/iterator-action.workflow-action.spec.ts new file mode 100644 index 00000000000..442b2e96b39 --- /dev/null +++ b/packages/twenty-server/src/modules/workflow/workflow-executor/workflow-actions/iterator/__tests__/iterator-action.workflow-action.spec.ts @@ -0,0 +1,324 @@ +import { Test, type TestingModule } from '@nestjs/testing'; + +import { StepStatus } from 'twenty-shared/workflow'; + +import { WorkflowStepExecutorException } from 'src/modules/workflow/workflow-executor/exceptions/workflow-step-executor.exception'; +import { IteratorWorkflowAction } from 'src/modules/workflow/workflow-executor/workflow-actions/iterator/iterator.workflow-action'; +import { type WorkflowIteratorActionSettings } from 'src/modules/workflow/workflow-executor/workflow-actions/iterator/types/workflow-iterator-action-settings.type'; +import { type WorkflowActionSettings } from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action-settings.type'; +import { + type WorkflowAction, + WorkflowActionType, +} 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'; + +describe('IteratorWorkflowAction', () => { + let service: IteratorWorkflowAction; + let workflowRunWorkspaceService: jest.Mocked; + + const mockWorkflowRunId = 'workflow-run-1'; + const mockWorkspaceId = 'workspace-1'; + const mockIteratorStepId = 'iterator-step-1'; + const mockInitialLoopStepIds = ['step-1', 'step-2']; + + const baseSettings: WorkflowActionSettings = { + outputSchema: {}, + errorHandlingOptions: { + retryOnFailure: { value: false }, + continueOnFailure: { value: false }, + }, + input: { + items: [], + initialLoopStepIds: mockInitialLoopStepIds, + }, + }; + + const iteratorSettings: WorkflowIteratorActionSettings = { + ...baseSettings, + input: { + items: [], + initialLoopStepIds: mockInitialLoopStepIds, + }, + }; + + const createEmptyAction = (): WorkflowAction => ({ + id: mockIteratorStepId, + type: WorkflowActionType.EMPTY, + name: 'Empty Step', + valid: true, + settings: baseSettings, + }); + + const createIteratorAction = (items: any): WorkflowAction => ({ + id: mockIteratorStepId, + type: WorkflowActionType.ITERATOR, + name: 'Iterator Step', + valid: true, + settings: { + ...iteratorSettings, + input: { + items, + initialLoopStepIds: mockInitialLoopStepIds, + }, + }, + }); + + beforeEach(async () => { + workflowRunWorkspaceService = { + getWorkflowRunOrFail: jest.fn(), + updateWorkflowRunStepInfos: jest.fn(), + } as any; + + const module: TestingModule = await Test.createTestingModule({ + providers: [ + IteratorWorkflowAction, + { + provide: WorkflowRunWorkspaceService, + useValue: workflowRunWorkspaceService, + }, + ], + }).compile(); + + service = module.get(IteratorWorkflowAction); + }); + + describe('execute', () => { + it('should throw error if step is not an iterator action', async () => { + const input = { + currentStepId: mockIteratorStepId, + steps: [createEmptyAction()], + context: {}, + runInfo: { + workflowRunId: mockWorkflowRunId, + workspaceId: mockWorkspaceId, + }, + }; + + await expect(service.execute(input)).rejects.toThrow( + WorkflowStepExecutorException, + ); + }); + + it('should throw error if items is not an array', async () => { + const input = { + currentStepId: mockIteratorStepId, + steps: [createIteratorAction('not-an-array')], + context: {}, + runInfo: { + workflowRunId: mockWorkflowRunId, + workspaceId: mockWorkspaceId, + }, + }; + + await expect(service.execute(input)).rejects.toThrow(); + }); + + it('should return early if no items to process', async () => { + const input = { + currentStepId: mockIteratorStepId, + steps: [createIteratorAction([])], + context: {}, + runInfo: { + workflowRunId: mockWorkflowRunId, + workspaceId: mockWorkspaceId, + }, + }; + + const result = await service.execute(input); + + expect(result).toEqual({ + result: { + currentItemIndex: 0, + currentItem: undefined, + hasProcessedAllItems: true, + }, + }); + }); + + it('should process first item correctly', async () => { + const items = ['item1', 'item2']; + const input = { + currentStepId: mockIteratorStepId, + steps: [createIteratorAction(items)], + context: {}, + runInfo: { + workflowRunId: mockWorkflowRunId, + workspaceId: mockWorkspaceId, + }, + }; + + workflowRunWorkspaceService.getWorkflowRunOrFail.mockResolvedValue({ + state: { + stepInfos: {}, + }, + } as any); + + const result = await service.execute(input); + + expect(result).toEqual({ + result: { + currentItemIndex: 0, + currentItem: 'item1', + hasProcessedAllItems: false, + }, + shouldRemainRunning: true, + }); + }); + + it('should process next item and reset loop steps', async () => { + const items = ['item1', 'item2']; + const input = { + currentStepId: mockIteratorStepId, + steps: [createIteratorAction(items)], + context: {}, + runInfo: { + workflowRunId: mockWorkflowRunId, + workspaceId: mockWorkspaceId, + }, + }; + + const mockStepInfo = { + state: { + stepInfos: { + [mockIteratorStepId]: { + result: { + currentItemIndex: 0, + currentItem: 'item1', + hasProcessedAllItems: false, + }, + status: StepStatus.SUCCESS, + }, + }, + }, + } as any; + + workflowRunWorkspaceService.getWorkflowRunOrFail + .mockResolvedValueOnce(mockStepInfo) + .mockResolvedValueOnce(mockStepInfo); + + const result = await service.execute(input); + + expect(result).toEqual({ + result: { + currentItemIndex: 1, + currentItem: 'item2', + hasProcessedAllItems: false, + }, + shouldRemainRunning: true, + }); + + expect( + workflowRunWorkspaceService.updateWorkflowRunStepInfos, + ).toHaveBeenCalledWith({ + stepInfos: expect.any(Object), + workflowRunId: mockWorkflowRunId, + workspaceId: mockWorkspaceId, + }); + }); + + it('should complete iteration when all items are processed', async () => { + const items = ['item1']; + const input = { + currentStepId: mockIteratorStepId, + steps: [createIteratorAction(items)], + context: {}, + runInfo: { + workflowRunId: mockWorkflowRunId, + workspaceId: mockWorkspaceId, + }, + }; + + const mockStepInfo = { + state: { + stepInfos: { + [mockIteratorStepId]: { + result: { + currentItemIndex: 0, + currentItem: 'item1', + hasProcessedAllItems: false, + }, + status: StepStatus.SUCCESS, + }, + }, + }, + } as any; + + workflowRunWorkspaceService.getWorkflowRunOrFail + .mockResolvedValueOnce(mockStepInfo) + .mockResolvedValueOnce(mockStepInfo); + + const result = await service.execute(input); + + expect(result).toEqual({ + result: { + currentItemIndex: 1, + currentItem: undefined, + hasProcessedAllItems: true, + }, + shouldRemainRunning: false, + }); + }); + + it('should throw error when max iterations is reached', async () => { + const items = Array(10001).fill('item'); + const input = { + currentStepId: mockIteratorStepId, + steps: [createIteratorAction(items)], + context: {}, + runInfo: { + workflowRunId: mockWorkflowRunId, + workspaceId: mockWorkspaceId, + }, + }; + + workflowRunWorkspaceService.getWorkflowRunOrFail.mockResolvedValue({ + state: { + stepInfos: { + [mockIteratorStepId]: { + result: { + currentItemIndex: 10000, + currentItem: 'item', + hasProcessedAllItems: false, + }, + status: StepStatus.SUCCESS, + }, + }, + }, + } as any); + + await expect(service.execute(input)).rejects.toThrow( + WorkflowStepExecutorException, + ); + }); + + it('should handle JSON string input for items', async () => { + const items = JSON.stringify(['item1', 'item2']); + const input = { + currentStepId: mockIteratorStepId, + steps: [createIteratorAction(items)], + context: {}, + runInfo: { + workflowRunId: mockWorkflowRunId, + workspaceId: mockWorkspaceId, + }, + }; + + workflowRunWorkspaceService.getWorkflowRunOrFail.mockResolvedValue({ + state: { + stepInfos: {}, + }, + } as any); + + const result = await service.execute(input); + + expect(result).toEqual({ + result: { + currentItemIndex: 0, + currentItem: 'item1', + hasProcessedAllItems: false, + }, + shouldRemainRunning: true, + }); + }); + }); +}); 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 index f8c5a8340f4..cd7ca8fe2f5 100644 --- 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 @@ -105,10 +105,11 @@ export class IteratorWorkflowAction implements WorkflowActionInterface { hasProcessedAllItems, }; - if (!hasProcessedAllItems && currentItemIndex > 0) { + if (currentItemIndex > 0) { await this.resetStepsInLoop({ iteratorStepId, initialLoopStepIds, + hasProcessedAllItems, workflowRunId: runInfo.workflowRunId, workspaceId: runInfo.workspaceId, steps, @@ -124,21 +125,19 @@ export class IteratorWorkflowAction implements WorkflowActionInterface { private async resetStepsInLoop({ iteratorStepId, initialLoopStepIds, + hasProcessedAllItems, workflowRunId, workspaceId, steps, }: { iteratorStepId: string; initialLoopStepIds: string[]; + hasProcessedAllItems: boolean; workflowRunId: string; workspaceId: string; steps: WorkflowAction[]; }) { - const stepIdsToReset = getAllStepIdsInLoop({ - iteratorStepId, - initialLoopStepIds, - steps, - }); + let stepInfosToUpdate: Record = {}; const workflowRunToUpdate = await this.workflowRunWorkspaceService.getWorkflowRunOrFail({ @@ -147,7 +146,56 @@ export class IteratorWorkflowAction implements WorkflowActionInterface { }); const stepInfos = workflowRunToUpdate.state.stepInfos; - const subStepsInfos = stepIdsToReset.reduce( + + if (!hasProcessedAllItems) { + const subStepsInfos = await this.buildSubStepInfosReset({ + iteratorStepId, + initialLoopStepIds, + stepInfos, + steps, + }); + + stepInfosToUpdate = { + ...stepInfosToUpdate, + ...subStepsInfos, + }; + } + + const iteratorStepInfo = await this.buildIteratorStepInfoReset({ + iteratorStepId, + iteratorStepInfo: stepInfos[iteratorStepId], + }); + + stepInfosToUpdate = { + ...stepInfosToUpdate, + ...iteratorStepInfo, + }; + + await this.workflowRunWorkspaceService.updateWorkflowRunStepInfos({ + stepInfos: stepInfosToUpdate, + workflowRunId, + workspaceId, + }); + } + + private async buildSubStepInfosReset({ + iteratorStepId, + initialLoopStepIds, + stepInfos, + steps, + }: { + iteratorStepId: string; + initialLoopStepIds: string[]; + stepInfos: Record; + steps: WorkflowAction[]; + }) { + const stepIdsToReset = getAllStepIdsInLoop({ + iteratorStepId, + initialLoopStepIds, + steps, + }); + + return stepIdsToReset.reduce( (acc, stepId) => { acc[stepId] = { status: StepStatus.NOT_STARTED, @@ -167,11 +215,16 @@ export class IteratorWorkflowAction implements WorkflowActionInterface { }, {} as Record, ); + } - const iteratorStepInfo = stepInfos[iteratorStepId]; - - const stepInfosToUpdate = { - ...subStepsInfos, + private async buildIteratorStepInfoReset({ + iteratorStepId, + iteratorStepInfo, + }: { + iteratorStepId: string; + iteratorStepInfo: WorkflowRunStepInfo; + }) { + return { [iteratorStepId]: { ...iteratorStepInfo, result: undefined, @@ -186,11 +239,5 @@ export class IteratorWorkflowAction implements WorkflowActionInterface { ], }, }; - - await this.workflowRunWorkspaceService.updateWorkflowRunStepInfos({ - stepInfos: stepInfosToUpdate, - workflowRunId, - workspaceId, - }); } }