Implement branch front end (#13489)

First PR to support workflow branches on frontend
This commit is contained in:
martmull
2025-08-04 16:11:46 +02:00
committed by GitHub
parent 2c87c6fc1b
commit 9a05673093
124 changed files with 3735 additions and 938 deletions
@@ -0,0 +1,559 @@
import { Test, TestingModule } from '@nestjs/testing';
import { getRepositoryToken } from '@nestjs/typeorm';
import { WorkflowVersionStepWorkspaceService } from 'src/modules/workflow/workflow-builder/workflow-step/workflow-version-step.workspace-service';
import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager';
import { WorkflowSchemaWorkspaceService } from 'src/modules/workflow/workflow-builder/workflow-schema/workflow-schema.workspace-service';
import { ServerlessFunctionService } from 'src/engine/metadata-modules/serverless-function/serverless-function.service';
import { AgentService } from 'src/engine/metadata-modules/agent/agent.service';
import { WorkflowRunWorkspaceService } from 'src/modules/workflow/workflow-runner/workflow-run/workflow-run.workspace-service';
import { WorkflowRunnerWorkspaceService } from 'src/modules/workflow/workflow-runner/workspace-services/workflow-runner.workspace-service';
import { WorkflowCommonWorkspaceService } from 'src/modules/workflow/common/workspace-services/workflow-common.workspace-service';
import { ScopedWorkspaceContextFactory } from 'src/engine/twenty-orm/factories/scoped-workspace-context.factory';
import { ObjectMetadataEntity } from 'src/engine/metadata-modules/object-metadata/object-metadata.entity';
import { WorkflowVersionWorkspaceEntity } from 'src/modules/workflow/common/standard-objects/workflow-version.workspace-entity';
import { WorkspaceRepository } from 'src/engine/twenty-orm/repository/workspace.repository';
import {
WorkflowAction,
WorkflowActionType,
} from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type';
import { WorkflowTriggerType } from 'src/modules/workflow/workflow-trigger/types/workflow-trigger.type';
type MockWorkspaceRepository = Partial<
WorkspaceRepository<WorkflowVersionWorkspaceEntity>
> & {
findOne: jest.Mock;
update: jest.Mock;
};
const mockWorkflowVersionId = 'workflow-version-id';
const mockWorkspaceId = 'workspace-id';
const mockSteps = [
{
id: 'step-1',
type: WorkflowActionType.FORM,
settings: {
errorHandlingOptions: {
continueOnFailure: { value: false },
retryOnFailure: { value: false },
},
},
nextStepIds: ['step-2'],
},
{
id: 'step-2',
type: WorkflowActionType.SEND_EMAIL,
settings: {
errorHandlingOptions: {
continueOnFailure: { value: false },
retryOnFailure: { value: false },
},
},
nextStepIds: [],
},
{
id: 'step-3',
type: WorkflowActionType.SEND_EMAIL,
settings: {
errorHandlingOptions: {
continueOnFailure: { value: false },
retryOnFailure: { value: false },
},
},
nextStepIds: [],
},
] as WorkflowAction[];
const mockTrigger = {
type: WorkflowTriggerType.MANUAL,
settings: {},
nextStepIds: ['step-1'],
};
const mockWorkflowVersion = {
id: mockWorkflowVersionId,
trigger: mockTrigger,
steps: mockSteps,
status: 'DRAFT',
} as WorkflowVersionWorkspaceEntity;
describe('WorkflowVersionStepWorkspaceService', () => {
let twentyORMGlobalManager: jest.Mocked<TwentyORMGlobalManager>;
let service: WorkflowVersionStepWorkspaceService;
let mockWorkflowVersionWorkspaceRepository: MockWorkspaceRepository;
beforeEach(async () => {
mockWorkflowVersionWorkspaceRepository = {
findOne: jest.fn(),
update: jest.fn(),
};
mockWorkflowVersionWorkspaceRepository.findOne.mockResolvedValue(
mockWorkflowVersion,
);
twentyORMGlobalManager = {
getRepositoryForWorkspace: jest
.fn()
.mockResolvedValue(mockWorkflowVersionWorkspaceRepository),
} as unknown as jest.Mocked<TwentyORMGlobalManager>;
const module: TestingModule = await Test.createTestingModule({
providers: [
WorkflowVersionStepWorkspaceService,
{
provide: TwentyORMGlobalManager,
useValue: twentyORMGlobalManager,
},
{
provide: WorkflowSchemaWorkspaceService,
useValue: {
computeStepOutputSchema: jest.fn(),
},
},
{ provide: ServerlessFunctionService, useValue: {} },
{ provide: AgentService, useValue: {} },
{
provide: getRepositoryToken(ObjectMetadataEntity, 'core'),
useValue: {
findOne: jest.fn(),
},
},
{ provide: WorkflowRunWorkspaceService, useValue: {} },
{ provide: WorkflowRunnerWorkspaceService, useValue: {} },
{ provide: WorkflowCommonWorkspaceService, useValue: {} },
{ provide: ScopedWorkspaceContextFactory, useValue: {} },
],
}).compile();
service = module.get(WorkflowVersionStepWorkspaceService);
});
describe('createWorkflowVersionStep', () => {
it('should create a step linked to trigger', async () => {
const result = await service.createWorkflowVersionStep({
input: {
stepType: WorkflowActionType.FORM,
parentStepId: 'trigger',
nextStepId: undefined,
workflowVersionId: mockWorkflowVersionId,
},
workspaceId: mockWorkspaceId,
});
expect(mockWorkflowVersionWorkspaceRepository.update).toHaveBeenCalled();
expect(result.createdStep).toBeDefined();
const createdStepId = result.createdStep?.id;
expect(result.triggerNextStepIds).toEqual(['step-1', createdStepId]);
expect(result.stepsNextStepIds).toEqual({
'step-1': ['step-2'],
'step-2': [],
'step-3': [],
});
});
it('should create a step between a trigger and a step', async () => {
const result = await service.createWorkflowVersionStep({
input: {
stepType: WorkflowActionType.FORM,
parentStepId: 'trigger',
nextStepId: 'step-1',
workflowVersionId: mockWorkflowVersionId,
},
workspaceId: mockWorkspaceId,
});
expect(mockWorkflowVersionWorkspaceRepository.update).toHaveBeenCalled();
expect(result.createdStep).toBeDefined();
const createdStepId = result.createdStep?.id as string;
expect(result.triggerNextStepIds).toEqual([createdStepId]);
expect(result.stepsNextStepIds).toEqual({
[createdStepId]: ['step-1'],
'step-1': ['step-2'],
'step-2': [],
'step-3': [],
});
});
it('should create a step between two steps', async () => {
const result = await service.createWorkflowVersionStep({
input: {
stepType: WorkflowActionType.FORM,
parentStepId: 'step-1',
nextStepId: 'step-2',
workflowVersionId: mockWorkflowVersionId,
},
workspaceId: mockWorkspaceId,
});
expect(mockWorkflowVersionWorkspaceRepository.update).toHaveBeenCalled();
expect(result.createdStep).toBeDefined();
const createdStepId = result.createdStep?.id as string;
expect(result.triggerNextStepIds).toEqual(['step-1']);
expect(result.stepsNextStepIds).toEqual({
'step-1': [createdStepId],
[createdStepId]: ['step-2'],
'step-2': [],
'step-3': [],
});
});
it('should create a step without parent or children', async () => {
const result = await service.createWorkflowVersionStep({
input: {
stepType: WorkflowActionType.FORM,
parentStepId: undefined,
nextStepId: undefined,
workflowVersionId: mockWorkflowVersionId,
},
workspaceId: mockWorkspaceId,
});
expect(mockWorkflowVersionWorkspaceRepository.update).toHaveBeenCalled();
expect(result.createdStep).toBeDefined();
expect(result.triggerNextStepIds).toEqual(['step-1']);
expect(result.stepsNextStepIds).toEqual({
'step-1': ['step-2'],
'step-2': [],
'step-3': [],
});
});
});
describe('createWorkflowVersionEdge', () => {
it('should throw if target does not exists', async () => {
const call = async () =>
await service.createWorkflowVersionEdge({
source: 'trigger',
target: 'not-existing-step',
workflowVersionId: mockWorkflowVersionId,
workspaceId: mockWorkspaceId,
});
await expect(call).rejects.toThrow(
`Target step 'not-existing-step' not found in workflowVersion '${mockWorkflowVersionId}'`,
);
});
describe('with source is the trigger', () => {
it('should create an edge between trigger and step-1', async () => {
const result = await service.createWorkflowVersionEdge({
source: 'trigger',
target: 'step-3',
workflowVersionId: mockWorkflowVersionId,
workspaceId: mockWorkspaceId,
});
expect(
mockWorkflowVersionWorkspaceRepository.update,
).toHaveBeenCalledWith(mockWorkflowVersionId, {
trigger: {
...mockTrigger,
nextStepIds: ['step-1', 'step-3'],
},
});
expect(result).toEqual({
triggerNextStepIds: ['step-1', 'step-3'],
stepsNextStepIds: {
'step-1': ['step-2'],
'step-2': [],
'step-3': [],
},
});
});
it('should not duplicate stepIds if edge already exists', async () => {
const result = await service.createWorkflowVersionEdge({
source: 'trigger',
target: 'step-1',
workflowVersionId: mockWorkflowVersionId,
workspaceId: mockWorkspaceId,
});
expect(
mockWorkflowVersionWorkspaceRepository.update,
).not.toHaveBeenCalled();
expect(result).toEqual({
triggerNextStepIds: ['step-1'],
stepsNextStepIds: {
'step-1': ['step-2'],
'step-2': [],
'step-3': [],
},
});
});
});
describe('with source is a step', () => {
it('should create an edge between step-2 and step-3', async () => {
const result = await service.createWorkflowVersionEdge({
source: 'step-2',
target: 'step-3',
workflowVersionId: mockWorkflowVersionId,
workspaceId: mockWorkspaceId,
});
expect(
mockWorkflowVersionWorkspaceRepository.update,
).toHaveBeenCalledWith(mockWorkflowVersionId, {
steps: mockSteps.map((step) => {
if (step.id === 'step-2') {
return {
...step,
nextStepIds: ['step-3'],
};
}
return step;
}),
});
expect(result).toEqual({
triggerNextStepIds: ['step-1'],
stepsNextStepIds: {
'step-1': ['step-2'],
'step-2': ['step-3'],
'step-3': [],
},
});
});
it('should not duplicate if edge already exist between 2 steps', async () => {
const result = await service.createWorkflowVersionEdge({
source: 'step-1',
target: 'step-2',
workflowVersionId: mockWorkflowVersionId,
workspaceId: mockWorkspaceId,
});
expect(
mockWorkflowVersionWorkspaceRepository.update,
).not.toHaveBeenCalled();
expect(result).toEqual({
triggerNextStepIds: ['step-1'],
stepsNextStepIds: {
'step-1': ['step-2'],
'step-2': [],
'step-3': [],
},
});
});
it('should throw if source step does not exists', async () => {
const call = async () =>
await service.createWorkflowVersionEdge({
source: 'not-existing-step',
target: 'step-2',
workflowVersionId: mockWorkflowVersionId,
workspaceId: mockWorkspaceId,
});
await expect(call).rejects.toThrow(
`Source step 'not-existing-step' not found in workflowVersion '${mockWorkflowVersionId}'`,
);
});
});
});
describe('deleteWorkflowVersionEdge', () => {
it('should throw if target does not exists', async () => {
const call = async () =>
await service.deleteWorkflowVersionEdge({
source: 'trigger',
target: 'not-existing-step',
workflowVersionId: mockWorkflowVersionId,
workspaceId: mockWorkspaceId,
});
await expect(call).rejects.toThrow(
`Target step 'not-existing-step' not found in workflowVersion '${mockWorkflowVersionId}'`,
);
});
describe('with source is the trigger', () => {
it('should delete an edge between trigger and step-1', async () => {
const result = await service.deleteWorkflowVersionEdge({
source: 'trigger',
target: 'step-1',
workflowVersionId: mockWorkflowVersionId,
workspaceId: mockWorkspaceId,
});
expect(
mockWorkflowVersionWorkspaceRepository.update,
).toHaveBeenCalledWith(mockWorkflowVersionId, {
trigger: {
...mockTrigger,
nextStepIds: [],
},
});
expect(result).toEqual({
triggerNextStepIds: [],
stepsNextStepIds: {
'step-1': ['step-2'],
'step-2': [],
'step-3': [],
},
});
});
it('should not delete if edge does not exists', async () => {
const result = await service.deleteWorkflowVersionEdge({
source: 'trigger',
target: 'step-2',
workflowVersionId: mockWorkflowVersionId,
workspaceId: mockWorkspaceId,
});
expect(
mockWorkflowVersionWorkspaceRepository.update,
).not.toHaveBeenCalled();
expect(result).toEqual({
triggerNextStepIds: ['step-1'],
stepsNextStepIds: {
'step-1': ['step-2'],
'step-2': [],
'step-3': [],
},
});
});
});
describe('with source is a step', () => {
it('should delete an existing edge between two steps', async () => {
const result = await service.deleteWorkflowVersionEdge({
source: 'step-1',
target: 'step-2',
workflowVersionId: mockWorkflowVersionId,
workspaceId: mockWorkspaceId,
});
expect(
mockWorkflowVersionWorkspaceRepository.update,
).toHaveBeenCalledWith(mockWorkflowVersionId, {
steps: mockSteps.map((step) => {
if (step.id === 'step-1') {
return {
...step,
nextStepIds: [],
};
}
return step;
}),
});
expect(result).toEqual({
triggerNextStepIds: ['step-1'],
stepsNextStepIds: {
'step-1': [],
'step-2': [],
'step-3': [],
},
});
});
it('should not delete if edge does not exist', async () => {
const result = await service.deleteWorkflowVersionEdge({
source: 'step-1',
target: 'step-3',
workflowVersionId: mockWorkflowVersionId,
workspaceId: mockWorkspaceId,
});
expect(
mockWorkflowVersionWorkspaceRepository.update,
).not.toHaveBeenCalledWith();
expect(result).toEqual({
triggerNextStepIds: ['step-1'],
stepsNextStepIds: {
'step-1': ['step-2'],
'step-2': [],
'step-3': [],
},
});
});
it('should throw if source step does not exists', async () => {
const call = async () =>
await service.deleteWorkflowVersionEdge({
source: 'not-existing-step',
target: 'step-2',
workflowVersionId: mockWorkflowVersionId,
workspaceId: mockWorkspaceId,
});
await expect(call).rejects.toThrow(
`Source step 'not-existing-step' not found in workflowVersion '${mockWorkflowVersionId}'`,
);
});
});
});
describe('deleteWorkflowVersionStep', () => {
it('should delete step linked to trigger', async () => {
const result = await service.deleteWorkflowVersionStep({
stepIdToDelete: 'step-1',
workflowVersionId: mockWorkflowVersionId,
workspaceId: mockWorkspaceId,
});
expect(
mockWorkflowVersionWorkspaceRepository.update,
).toHaveBeenCalledWith(mockWorkflowVersionId, {
trigger: { ...mockTrigger, nextStepIds: ['step-2'] },
steps: mockSteps.filter((step) => step.id !== 'step-1'),
});
expect(result).toEqual({
triggerNextStepIds: ['step-2'],
stepsNextStepIds: {
'step-2': [],
'step-3': [],
},
deletedStepId: 'step-1',
});
});
it('should delete trigger', async () => {
const result = await service.deleteWorkflowVersionStep({
stepIdToDelete: 'trigger',
workflowVersionId: mockWorkflowVersionId,
workspaceId: mockWorkspaceId,
});
expect(
mockWorkflowVersionWorkspaceRepository.update,
).toHaveBeenCalledWith(mockWorkflowVersionId, {
trigger: null,
});
expect(result).toEqual({
stepsNextStepIds: {
'step-1': ['step-2'],
'step-2': [],
'step-3': [],
},
deletedStepId: 'trigger',
});
});
});
});
@@ -0,0 +1,24 @@
import { computeWorkflowVersionStepChanges } from 'src/modules/workflow/workflow-builder/workflow-step/utils/compute-workflow-version-step-updates.util';
import { WorkflowTrigger } from 'src/modules/workflow/workflow-trigger/types/workflow-trigger.type';
import { WorkflowAction } from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type';
describe('computeWorkflowVersionStepChanges', () => {
it('should compute next step ids', () => {
const input = {
trigger: { nextStepIds: ['1', '2'] } as WorkflowTrigger,
steps: [
{ id: '1', nextStepIds: ['3'] },
{ id: '2', nextStepIds: ['3'] },
] as WorkflowAction[],
deletedStepId: '5',
};
const expectedResult = {
triggerNextStepIds: ['1', '2'],
stepsNextStepIds: { '1': ['3'], '2': ['3'] },
deletedStepId: '5',
};
expect(computeWorkflowVersionStepChanges(input)).toEqual(expectedResult);
});
});
@@ -3,6 +3,10 @@ import {
WorkflowAction,
WorkflowActionType,
} from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type';
import {
WorkflowTrigger,
WorkflowTriggerType,
} from 'src/modules/workflow/workflow-trigger/types/workflow-trigger.type';
describe('insertStep', () => {
const createMockAction = (
@@ -28,7 +32,15 @@ describe('insertStep', () => {
nextStepIds,
});
const createMockTrigger = (nextStepIds: string[]): WorkflowTrigger => ({
name: 'Trigger',
type: WorkflowTriggerType.MANUAL,
settings: { outputSchema: {} },
nextStepIds,
});
it('should insert a step at the end of the array when no parent or next step is specified', () => {
const existingTrigger = createMockTrigger(['1']);
const step1 = createMockAction('1');
const step2 = createMockAction('2');
const newStep = createMockAction('new');
@@ -36,6 +48,7 @@ describe('insertStep', () => {
const result = insertStep({
existingSteps: [step1, step2],
insertedStep: newStep,
existingTrigger,
});
expect(result.updatedSteps).toEqual([step1, step2, newStep]);
@@ -43,6 +56,7 @@ describe('insertStep', () => {
});
it('should update parent step nextStepIds when inserting a step between two steps', () => {
const existingTrigger = createMockTrigger(['1']);
const step1 = createMockAction('1', ['2']);
const step2 = createMockAction('2');
const newStep = createMockAction('new');
@@ -50,6 +64,7 @@ describe('insertStep', () => {
const result = insertStep({
existingSteps: [step1, step2],
insertedStep: newStep,
existingTrigger,
parentStepId: '1',
nextStepId: '2',
});
@@ -62,10 +77,12 @@ describe('insertStep', () => {
});
it('should handle inserting a step at the beginning of the workflow', () => {
const existingTrigger = createMockTrigger(['1']);
const step1 = createMockAction('1');
const newStep = createMockAction('new');
const result = insertStep({
existingTrigger,
existingSteps: [step1],
insertedStep: newStep,
parentStepId: undefined,
@@ -79,10 +96,12 @@ describe('insertStep', () => {
});
it('should handle inserting a step at the end of the workflow', () => {
const existingTrigger = createMockTrigger(['1']);
const step1 = createMockAction('1');
const newStep = createMockAction('new');
const result = insertStep({
existingTrigger,
existingSteps: [step1],
insertedStep: newStep,
parentStepId: '1',
@@ -96,12 +115,14 @@ describe('insertStep', () => {
});
it('should handle inserting a step between two steps with multiple nextStepIds', () => {
const existingTrigger = createMockTrigger(['1']);
const step1 = createMockAction('1', ['2', '3']);
const step2 = createMockAction('2');
const step3 = createMockAction('3');
const newStep = createMockAction('new');
const result = insertStep({
existingTrigger,
existingSteps: [step1, step2, step3],
insertedStep: newStep,
parentStepId: '1',
@@ -115,4 +136,24 @@ describe('insertStep', () => {
{ ...newStep, nextStepIds: ['2'] },
]);
});
it('should handle inserting after trigger', () => {
const existingTrigger = createMockTrigger(['1']);
const step1 = createMockAction('1');
const newStep = createMockAction('new');
const result = insertStep({
existingTrigger,
existingSteps: [step1],
insertedStep: newStep,
parentStepId: 'trigger',
nextStepId: undefined,
});
expect(result.updatedSteps).toEqual([step1, newStep]);
expect(result.updatedTrigger).toEqual({
...existingTrigger,
nextStepIds: ['1', 'new'],
});
});
});
@@ -3,6 +3,17 @@ import {
WorkflowAction,
WorkflowActionType,
} from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type';
import {
WorkflowTrigger,
WorkflowTriggerType,
} from 'src/modules/workflow/workflow-trigger/types/workflow-trigger.type';
const mockTrigger = {
name: 'Trigger',
type: WorkflowTriggerType.MANUAL,
settings: { outputSchema: {} },
nextStepIds: ['1'],
} as WorkflowTrigger;
describe('removeStep', () => {
const createMockAction = (
@@ -34,11 +45,13 @@ describe('removeStep', () => {
const step3 = createMockAction('3');
const result = removeStep({
existingTrigger: mockTrigger,
existingSteps: [step1, step2, step3],
stepIdToDelete: '2',
});
expect(result).toEqual([step1, step3]);
expect(result.steps).toEqual([step1, step3]);
expect(result.trigger).toEqual(mockTrigger);
});
it('should handle removing a step that has no next steps', () => {
@@ -47,11 +60,13 @@ describe('removeStep', () => {
const step3 = createMockAction('3');
const result = removeStep({
existingTrigger: mockTrigger,
existingSteps: [step1, step2, step3],
stepIdToDelete: '2',
});
expect(result).toEqual([{ ...step1, nextStepIds: [] }, step3]);
expect(result.steps).toEqual([{ ...step1, nextStepIds: [] }, step3]);
expect(result.trigger).toEqual(mockTrigger);
});
it('should update nextStepIds of parent steps to include children of removed step', () => {
@@ -60,12 +75,14 @@ describe('removeStep', () => {
const step3 = createMockAction('3');
const result = removeStep({
existingTrigger: mockTrigger,
existingSteps: [step1, step2, step3],
stepIdToDelete: '2',
stepToDeleteChildrenIds: ['3'],
});
expect(result).toEqual([{ ...step1, nextStepIds: ['3'] }, step3]);
expect(result.steps).toEqual([{ ...step1, nextStepIds: ['3'] }, step3]);
expect(result.trigger).toEqual(mockTrigger);
});
it('should handle multiple parent steps pointing to the same step', () => {
@@ -75,16 +92,18 @@ describe('removeStep', () => {
const step4 = createMockAction('4');
const result = removeStep({
existingTrigger: mockTrigger,
existingSteps: [step1, step2, step3, step4],
stepIdToDelete: '3',
stepToDeleteChildrenIds: ['4'],
});
expect(result).toEqual([
expect(result.steps).toEqual([
{ ...step1, nextStepIds: ['4'] },
{ ...step2, nextStepIds: ['4'] },
step4,
]);
expect(result.trigger).toEqual(mockTrigger);
});
it('should handle removing a step with multiple children', () => {
@@ -94,15 +113,33 @@ describe('removeStep', () => {
const step4 = createMockAction('4');
const result = removeStep({
existingTrigger: mockTrigger,
existingSteps: [step1, step2, step3, step4],
stepIdToDelete: '2',
stepToDeleteChildrenIds: ['3', '4'],
});
expect(result).toEqual([
expect(result.steps).toEqual([
{ ...step1, nextStepIds: ['3', '4'] },
step3,
step4,
]);
expect(result.trigger).toEqual(mockTrigger);
});
it('should handle removing a step linked to trigger', () => {
const step1 = createMockAction('1', ['2']);
const step2 = createMockAction('2', ['3']);
const step3 = createMockAction('3');
const result = removeStep({
existingTrigger: mockTrigger,
existingSteps: [step1, step2, step3],
stepIdToDelete: '1',
stepToDeleteChildrenIds: ['2'],
});
expect(result.steps).toEqual([step2, step3]);
expect(result.trigger).toEqual({ ...mockTrigger, nextStepIds: ['2'] });
});
});
@@ -0,0 +1,24 @@
import { WorkflowVersionStepChangesDTO } from 'src/engine/core-modules/workflow/dtos/workflow-version-step-changes.dto';
import { WorkflowTrigger } from 'src/modules/workflow/workflow-trigger/types/workflow-trigger.type';
import { WorkflowAction } from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type';
export const computeWorkflowVersionStepChanges = ({
trigger,
steps,
createdStep,
deletedStepId,
}: {
trigger: WorkflowTrigger | null;
steps: WorkflowAction[] | null;
createdStep?: WorkflowAction;
deletedStepId?: string;
}): WorkflowVersionStepChangesDTO => {
return {
triggerNextStepIds: trigger?.nextStepIds,
stepsNextStepIds: Object.fromEntries(
(steps || []).map((step) => [step.id, step.nextStepIds]),
),
createdStep,
deletedStepId,
};
};
@@ -1,32 +1,60 @@
import { WorkflowAction } from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type';
import { WorkflowTrigger } from 'src/modules/workflow/workflow-trigger/types/workflow-trigger.type';
export const insertStep = ({
existingSteps,
existingTrigger,
insertedStep,
parentStepId,
nextStepId,
}: {
existingSteps: WorkflowAction[];
existingTrigger: WorkflowTrigger | null;
insertedStep: WorkflowAction;
parentStepId?: string;
nextStepId?: string;
}): { updatedSteps: WorkflowAction[]; updatedInsertedStep: WorkflowAction } => {
const updatedExistingSteps = existingSteps.map((existingStep) => {
if (existingStep.id === parentStepId) {
return {
...existingStep,
nextStepIds: [
...new Set([
...(existingStep.nextStepIds?.filter((id) => id !== nextStepId) ||
[]),
insertedStep.id,
]),
],
};
}): {
updatedSteps: WorkflowAction[];
updatedInsertedStep: WorkflowAction;
updatedTrigger: WorkflowTrigger | null;
} => {
let updatedTrigger = existingTrigger;
let updatedExistingSteps = existingSteps;
if (parentStepId === 'trigger') {
if (!existingTrigger) {
throw new Error('Cannot insert step from undefined trigger');
}
return existingStep;
});
updatedTrigger = {
...existingTrigger,
nextStepIds: [
...new Set([
...(existingTrigger.nextStepIds?.filter((id) => id !== nextStepId) ||
[]),
insertedStep.id,
]),
],
};
} else {
updatedExistingSteps = existingSteps.map((existingStep) => {
if (existingStep.id === parentStepId) {
return {
...existingStep,
nextStepIds: [
...new Set([
...(existingStep.nextStepIds?.filter((id) => id !== nextStepId) ||
[]),
insertedStep.id,
]),
],
};
}
return existingStep;
});
}
const updatedInsertedStep = {
...insertedStep,
@@ -35,6 +63,7 @@ export const insertStep = ({
return {
updatedSteps: [...updatedExistingSteps, updatedInsertedStep],
updatedTrigger,
updatedInsertedStep,
};
};
@@ -1,30 +1,72 @@
import { isDefined } from 'twenty-shared/utils';
import { WorkflowAction } from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type';
import { WorkflowTrigger } from 'src/modules/workflow/workflow-trigger/types/workflow-trigger.type';
const computeUpdatedNextStepIds = ({
existingNextStepIds,
stepIdToDelete,
stepToDeleteChildrenIds,
}: {
existingNextStepIds: string[];
stepIdToDelete: string;
stepToDeleteChildrenIds?: string[];
}): string[] => {
const filteredNextStepIds = isDefined(existingNextStepIds)
? existingNextStepIds.filter((id) => id !== stepIdToDelete)
: [];
return [
...new Set([
...filteredNextStepIds,
// We automatically link parent and child steps together
...(stepToDeleteChildrenIds || []),
]),
];
};
export const removeStep = ({
existingTrigger,
existingSteps,
stepIdToDelete,
stepToDeleteChildrenIds,
}: {
existingTrigger: WorkflowTrigger | null;
existingSteps: WorkflowAction[];
stepIdToDelete: string;
stepToDeleteChildrenIds?: string[];
}): WorkflowAction[] => {
return existingSteps
}): { steps: WorkflowAction[]; trigger: WorkflowTrigger | null } => {
const updatedSteps = existingSteps
.filter((step) => step.id !== stepIdToDelete)
.map((step) => {
if (step.nextStepIds?.includes(stepIdToDelete)) {
return {
...step,
nextStepIds: [
...new Set([
...step.nextStepIds.filter((id) => id !== stepIdToDelete),
// We automatically link parent and child steps together
...(stepToDeleteChildrenIds || []),
]),
],
nextStepIds: computeUpdatedNextStepIds({
existingNextStepIds: step.nextStepIds,
stepIdToDelete,
stepToDeleteChildrenIds,
}),
};
}
return step;
});
let updatedTrigger = existingTrigger;
if (isDefined(existingTrigger)) {
if (existingTrigger.nextStepIds?.includes(stepIdToDelete)) {
updatedTrigger = {
...existingTrigger,
nextStepIds: computeUpdatedNextStepIds({
existingNextStepIds: existingTrigger.nextStepIds,
stepIdToDelete,
stepToDeleteChildrenIds,
}),
};
}
}
return { trigger: updatedTrigger, steps: updatedSteps };
};
@@ -10,8 +10,6 @@ import { v4 } from 'uuid';
import { BASE_TYPESCRIPT_PROJECT_INPUT_SCHEMA } from 'src/engine/core-modules/serverless/drivers/constants/base-typescript-project-input-schema';
import { CreateWorkflowVersionStepInput } from 'src/engine/core-modules/workflow/dtos/create-workflow-version-step-input.dto';
import { WorkflowActionDTO } from 'src/engine/core-modules/workflow/dtos/workflow-step.dto';
import { AgentChatService } from 'src/engine/metadata-modules/agent/agent-chat.service';
import { AgentService } from 'src/engine/metadata-modules/agent/agent.service';
import { ObjectMetadataEntity } from 'src/engine/metadata-modules/object-metadata/object-metadata.entity';
import { ServerlessFunctionService } from 'src/engine/metadata-modules/serverless-function/serverless-function.service';
@@ -35,8 +33,9 @@ import {
} 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';
import { WorkflowRunnerWorkspaceService } from 'src/modules/workflow/workflow-runner/workspace-services/workflow-runner.workspace-service';
const TRIGGER_STEP_ID = 'trigger';
import { WorkflowStepPositionInput } from 'src/engine/core-modules/workflow/dtos/update-workflow-step-position-input.dto';
import { WorkflowVersionStepChangesDTO } from 'src/engine/core-modules/workflow/dtos/workflow-version-step-changes.dto';
import { computeWorkflowVersionStepChanges } from 'src/modules/workflow/workflow-builder/workflow-step/utils/compute-workflow-version-step-updates.util';
const BASE_STEP_DEFINITION: BaseWorkflowActionSettings = {
outputSchema: {},
@@ -61,7 +60,6 @@ export class WorkflowVersionStepWorkspaceService {
private readonly objectMetadataRepository: Repository<ObjectMetadataEntity>,
private readonly workflowRunWorkspaceService: WorkflowRunWorkspaceService,
private readonly workflowRunnerWorkspaceService: WorkflowRunnerWorkspaceService,
private readonly agentChatService: AgentChatService,
private readonly workflowCommonWorkspaceService: WorkflowCommonWorkspaceService,
private readonly scopedWorkspaceContextFactory: ScopedWorkspaceContextFactory,
) {}
@@ -72,17 +70,21 @@ export class WorkflowVersionStepWorkspaceService {
}: {
workspaceId: string;
input: CreateWorkflowVersionStepInput;
}): Promise<WorkflowActionDTO> {
const { workflowVersionId, stepType, parentStepId, nextStepId } = input;
}): Promise<WorkflowVersionStepChangesDTO> {
const { workflowVersionId, stepType, parentStepId, nextStepId, position } =
input;
const newStep = await this.getStepDefaultDefinition({
type: stepType,
workspaceId,
position,
});
const enrichedNewStep = await this.enrichOutputSchema({
step: newStep,
workspaceId,
});
const workflowVersionRepository =
await this.twentyORMGlobalManager.getRepositoryForWorkspace<WorkflowVersionWorkspaceEntity>(
workspaceId,
@@ -107,18 +109,26 @@ export class WorkflowVersionStepWorkspaceService {
const existingSteps = workflowVersion.steps || [];
const { updatedSteps, updatedInsertedStep } = insertStep({
const existingTrigger = workflowVersion.trigger;
const { updatedSteps, updatedInsertedStep, updatedTrigger } = insertStep({
existingSteps,
existingTrigger,
insertedStep: enrichedNewStep,
parentStepId,
nextStepId,
});
await workflowVersionRepository.update(workflowVersion.id, {
trigger: updatedTrigger,
steps: updatedSteps,
});
return updatedInsertedStep;
return computeWorkflowVersionStepChanges({
createdStep: updatedInsertedStep,
trigger: updatedTrigger,
steps: updatedSteps,
});
}
async updateWorkflowVersionStep({
@@ -187,7 +197,7 @@ export class WorkflowVersionStepWorkspaceService {
workspaceId: string;
workflowVersionId: string;
stepIdToDelete: string;
}): Promise<WorkflowActionDTO> {
}): Promise<WorkflowVersionStepChangesDTO> {
const workflowVersionRepository =
await this.twentyORMGlobalManager.getRepositoryForWorkspace<WorkflowVersionWorkspaceEntity>(
workspaceId,
@@ -217,6 +227,20 @@ export class WorkflowVersionStepWorkspaceService {
);
}
if (stepIdToDelete === 'trigger') {
await workflowVersionRepository.update(workflowVersion.id, {
trigger: null,
});
return computeWorkflowVersionStepChanges({
trigger: null,
steps: workflowVersion?.steps,
deletedStepId: stepIdToDelete,
});
}
const existingTrigger = workflowVersion.trigger;
const stepToDelete = workflowVersion.steps.find(
(step) => step.id === stepIdToDelete,
);
@@ -228,16 +252,12 @@ export class WorkflowVersionStepWorkspaceService {
);
}
const workflowVersionUpdates =
stepIdToDelete === TRIGGER_STEP_ID
? { trigger: null }
: {
steps: removeStep({
existingSteps: workflowVersion.steps,
stepIdToDelete,
stepToDeleteChildrenIds: stepToDelete.nextStepIds,
}),
};
const workflowVersionUpdates = removeStep({
existingTrigger,
existingSteps: workflowVersion.steps,
stepIdToDelete,
stepToDeleteChildrenIds: stepToDelete.nextStepIds,
});
await workflowVersionRepository.update(
workflowVersion.id,
@@ -249,7 +269,10 @@ export class WorkflowVersionStepWorkspaceService {
workspaceId,
});
return stepToDelete;
return computeWorkflowVersionStepChanges({
...workflowVersionUpdates,
deletedStepId: stepIdToDelete,
});
}
async duplicateStep({
@@ -345,6 +368,244 @@ export class WorkflowVersionStepWorkspaceService {
});
}
async createWorkflowVersionEdge({
source,
target,
workflowVersionId,
workspaceId,
}: {
source: string;
target: string;
workflowVersionId: string;
workspaceId: string;
}): Promise<WorkflowVersionStepChangesDTO> {
const workflowVersionRepository =
await this.twentyORMGlobalManager.getRepositoryForWorkspace<WorkflowVersionWorkspaceEntity>(
workspaceId,
'workflowVersion',
{ shouldBypassPermissionChecks: true },
);
const workflowVersion = await workflowVersionRepository.findOne({
where: {
id: workflowVersionId,
},
});
if (!isDefined(workflowVersion)) {
throw new WorkflowVersionStepException(
'WorkflowVersion not found',
WorkflowVersionStepExceptionCode.NOT_FOUND,
);
}
assertWorkflowVersionIsDraft(workflowVersion);
const steps = workflowVersion.steps || [];
const trigger = workflowVersion.trigger;
const isSourceTrigger = source === 'trigger';
const targetStep = steps.find((step) => step.id === target);
if (!isDefined(targetStep)) {
throw new WorkflowVersionStepException(
`Target step '${target}' not found in workflowVersion '${workflowVersionId}'`,
WorkflowVersionStepExceptionCode.NOT_FOUND,
);
}
if (isSourceTrigger) {
if (!isDefined(trigger)) {
throw new WorkflowVersionStepException(
`Trigger not found in workflowVersion '${workflowVersionId}'`,
WorkflowVersionStepExceptionCode.NOT_FOUND,
);
}
if (trigger.nextStepIds?.includes(target)) {
return computeWorkflowVersionStepChanges({
trigger,
steps,
});
}
const updatedTrigger = {
...trigger,
nextStepIds: [...(trigger.nextStepIds ?? []), target],
};
await workflowVersionRepository.update(workflowVersion.id, {
trigger: updatedTrigger,
});
return computeWorkflowVersionStepChanges({
trigger: updatedTrigger,
steps,
});
}
const sourceStep = steps.find((step) => step.id === source);
if (!isDefined(sourceStep)) {
throw new WorkflowVersionStepException(
`Source step '${source}' not found in workflowVersion '${workflowVersionId}'`,
WorkflowVersionStepExceptionCode.NOT_FOUND,
);
}
if (sourceStep.nextStepIds?.includes(target)) {
return computeWorkflowVersionStepChanges({
trigger,
steps,
});
}
const updatedSourceStep = {
...sourceStep,
nextStepIds: [...(sourceStep.nextStepIds ?? []), target],
};
const updatedSteps = steps.map((step) => {
if (step.id === source) {
return updatedSourceStep;
}
return step;
});
await workflowVersionRepository.update(workflowVersion.id, {
steps: updatedSteps,
});
return computeWorkflowVersionStepChanges({
trigger,
steps: updatedSteps,
});
}
async deleteWorkflowVersionEdge({
source,
target,
workflowVersionId,
workspaceId,
}: {
source: string;
target: string;
workflowVersionId: string;
workspaceId: string;
}): Promise<WorkflowVersionStepChangesDTO> {
const workflowVersionRepository =
await this.twentyORMGlobalManager.getRepositoryForWorkspace<WorkflowVersionWorkspaceEntity>(
workspaceId,
'workflowVersion',
{ shouldBypassPermissionChecks: true },
);
const workflowVersion = await workflowVersionRepository.findOne({
where: {
id: workflowVersionId,
},
});
if (!isDefined(workflowVersion)) {
throw new WorkflowVersionStepException(
'WorkflowVersion not found',
WorkflowVersionStepExceptionCode.NOT_FOUND,
);
}
assertWorkflowVersionIsDraft(workflowVersion);
const steps = workflowVersion.steps || [];
const trigger = workflowVersion.trigger;
const isSourceTrigger = source === 'trigger';
const targetStep = steps.find((step) => step.id === target);
if (!isDefined(targetStep)) {
throw new WorkflowVersionStepException(
`Target step '${target}' not found in workflowVersion '${workflowVersionId}'`,
WorkflowVersionStepExceptionCode.NOT_FOUND,
);
}
if (isSourceTrigger) {
if (!isDefined(trigger)) {
throw new WorkflowVersionStepException(
`Trigger not found in workflowVersion '${workflowVersionId}'`,
WorkflowVersionStepExceptionCode.NOT_FOUND,
);
}
if (!trigger.nextStepIds?.includes(target)) {
return computeWorkflowVersionStepChanges({
trigger,
steps,
});
}
const updatedTrigger = {
...trigger,
nextStepIds: trigger.nextStepIds?.filter(
(nextStepId) => nextStepId !== target,
),
};
await workflowVersionRepository.update(workflowVersion.id, {
trigger: updatedTrigger,
});
return computeWorkflowVersionStepChanges({
trigger: updatedTrigger,
steps,
});
}
const sourceStep = steps.find((step) => step.id === source);
if (!isDefined(sourceStep)) {
throw new WorkflowVersionStepException(
`Source step '${source}' not found in workflowVersion '${workflowVersionId}'`,
WorkflowVersionStepExceptionCode.NOT_FOUND,
);
}
if (!sourceStep.nextStepIds?.includes(target)) {
return computeWorkflowVersionStepChanges({
trigger,
steps,
});
}
const updatedSourceStep = {
...sourceStep,
nextStepIds: sourceStep.nextStepIds?.filter(
(nextStepId) => nextStepId !== target,
),
};
const updatedSteps = steps.map((step) => {
if (step.id === source) {
return updatedSourceStep;
}
return step;
});
await workflowVersionRepository.update(workflowVersion.id, {
steps: updatedSteps,
});
return computeWorkflowVersionStepChanges({
trigger,
steps: updatedSteps,
});
}
private async enrichOutputSchema({
step,
workspaceId,
@@ -423,12 +684,21 @@ export class WorkflowVersionStepWorkspaceService {
private async getStepDefaultDefinition({
type,
workspaceId,
position,
}: {
type: WorkflowActionType;
workspaceId: string;
position?: WorkflowStepPositionInput;
}): Promise<WorkflowAction> {
const newStepId = v4();
const baseStep = {
id: newStepId,
position,
valid: false,
nextStepIds: [],
};
switch (type) {
case WorkflowActionType.CODE: {
const newServerlessFunction =
@@ -448,10 +718,9 @@ export class WorkflowVersionStepWorkspaceService {
}
return {
id: newStepId,
...baseStep,
name: 'Code - Serverless Function',
type: WorkflowActionType.CODE,
valid: false,
settings: {
...BASE_STEP_DEFINITION,
outputSchema: {
@@ -473,10 +742,9 @@ export class WorkflowVersionStepWorkspaceService {
}
case WorkflowActionType.SEND_EMAIL: {
return {
id: newStepId,
...baseStep,
name: 'Send Email',
type: WorkflowActionType.SEND_EMAIL,
valid: false,
settings: {
...BASE_STEP_DEFINITION,
input: {
@@ -495,10 +763,9 @@ export class WorkflowVersionStepWorkspaceService {
});
return {
id: newStepId,
...baseStep,
name: 'Create Record',
type: WorkflowActionType.CREATE_RECORD,
valid: false,
settings: {
...BASE_STEP_DEFINITION,
input: {
@@ -515,10 +782,9 @@ export class WorkflowVersionStepWorkspaceService {
});
return {
id: newStepId,
...baseStep,
name: 'Update Record',
type: WorkflowActionType.UPDATE_RECORD,
valid: false,
settings: {
...BASE_STEP_DEFINITION,
input: {
@@ -537,10 +803,9 @@ export class WorkflowVersionStepWorkspaceService {
});
return {
id: newStepId,
...baseStep,
name: 'Delete Record',
type: WorkflowActionType.DELETE_RECORD,
valid: false,
settings: {
...BASE_STEP_DEFINITION,
input: {
@@ -557,10 +822,9 @@ export class WorkflowVersionStepWorkspaceService {
});
return {
id: newStepId,
...baseStep,
name: 'Search Records',
type: WorkflowActionType.FIND_RECORDS,
valid: false,
settings: {
...BASE_STEP_DEFINITION,
input: {
@@ -572,10 +836,9 @@ export class WorkflowVersionStepWorkspaceService {
}
case WorkflowActionType.FORM: {
return {
id: newStepId,
...baseStep,
name: 'Form',
type: WorkflowActionType.FORM,
valid: false,
settings: {
...BASE_STEP_DEFINITION,
input: [],
@@ -584,10 +847,9 @@ export class WorkflowVersionStepWorkspaceService {
}
case WorkflowActionType.FILTER: {
return {
id: newStepId,
...baseStep,
name: 'Filter',
type: WorkflowActionType.FILTER,
valid: false,
settings: {
...BASE_STEP_DEFINITION,
input: {
@@ -599,10 +861,9 @@ export class WorkflowVersionStepWorkspaceService {
}
case WorkflowActionType.HTTP_REQUEST: {
return {
id: newStepId,
...baseStep,
name: 'HTTP Request',
type: WorkflowActionType.HTTP_REQUEST,
valid: false,
settings: {
...BASE_STEP_DEFINITION,
input: {
@@ -616,10 +877,9 @@ export class WorkflowVersionStepWorkspaceService {
}
case WorkflowActionType.AI_AGENT: {
return {
id: newStepId,
...baseStep,
name: 'AI Agent',
type: WorkflowActionType.AI_AGENT,
valid: false,
settings: {
...BASE_STEP_DEFINITION,
input: {
@@ -17,6 +17,7 @@ import { assertWorkflowVersionIsDraft } from 'src/modules/workflow/common/utils/
import { assertWorkflowVersionTriggerIsDefined } from 'src/modules/workflow/common/utils/assert-workflow-version-trigger-is-defined.util';
import { WorkflowVersionStepWorkspaceService } from 'src/modules/workflow/workflow-builder/workflow-step/workflow-version-step.workspace-service';
import { WorkflowAction } from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type';
import { WorkflowStepPositionUpdateInput } from 'src/engine/core-modules/workflow/dtos/update-workflow-step-position-update-input.dto';
@Injectable()
export class WorkflowVersionWorkspaceService {
@@ -112,4 +113,61 @@ export class WorkflowVersionWorkspaceService {
return draftWorkflowVersion.id;
}
async updateWorkflowVersionPositions({
workflowVersionId,
positions,
workspaceId,
}: {
workflowVersionId: string;
positions: WorkflowStepPositionUpdateInput[];
workspaceId: string;
}) {
const workflowVersionRepository =
await this.twentyORMGlobalManager.getRepositoryForWorkspace<WorkflowVersionWorkspaceEntity>(
workspaceId,
'workflowVersion',
{ shouldBypassPermissionChecks: true },
);
const workflowVersion = await workflowVersionRepository.findOneOrFail({
where: {
id: workflowVersionId,
},
});
assertWorkflowVersionIsDraft(workflowVersion);
const triggerPosition = positions.find(
(position) => position.id === 'trigger',
);
const updatedTrigger =
isDefined(triggerPosition) && isDefined(workflowVersion.trigger)
? {
...workflowVersion.trigger,
position: triggerPosition.position,
}
: undefined;
const updatedSteps = workflowVersion.steps?.map((step) => {
const updatedStep = positions.find((position) => position.id === step.id);
if (updatedStep) {
return {
...step,
position: updatedStep.position,
};
}
return step;
});
const updatePayload = {
...(!isDefined(updatedTrigger) ? {} : { trigger: updatedTrigger }),
...(!isDefined(updatedSteps) ? {} : { steps: updatedSteps }),
};
await workflowVersionRepository.update(workflowVersionId, updatePayload);
}
}