Add connection options on edge creation (#14451)

Edge can now be linked to loop handle.
Renaming the existing options to connection.
This commit is contained in:
Thomas Trompette
2025-09-12 14:08:20 +00:00
committed by GitHub
parent 3270c64a96
commit 12a5ee6bc4
13 changed files with 346 additions and 45 deletions
@@ -838,6 +838,8 @@ export type CreateWebhookDto = {
export type CreateWorkflowVersionEdgeInput = {
/** Workflow version source step ID */
source: Scalars['String'];
/** Workflow version source step connection options */
sourceConnectionOptions?: InputMaybe<Scalars['JSON']>;
/** Workflow version target step ID */
target: Scalars['String'];
/** Workflow version ID */
@@ -847,10 +849,10 @@ export type CreateWorkflowVersionEdgeInput = {
export type CreateWorkflowVersionStepInput = {
/** Next step ID */
nextStepId?: InputMaybe<Scalars['UUID']>;
/** Parent step connection options */
parentStepConnectionOptions?: InputMaybe<Scalars['JSON']>;
/** Parent step ID */
parentStepId?: InputMaybe<Scalars['String']>;
/** Step creation options */
parentStepOptions?: InputMaybe<Scalars['JSON']>;
/** Step position */
position?: InputMaybe<WorkflowStepPositionInput>;
/** New step type */
@@ -802,6 +802,8 @@ export type CreateWebhookDto = {
export type CreateWorkflowVersionEdgeInput = {
/** Workflow version source step ID */
source: Scalars['String'];
/** Workflow version source step connection options */
sourceConnectionOptions?: InputMaybe<Scalars['JSON']>;
/** Workflow version target step ID */
target: Scalars['String'];
/** Workflow version ID */
@@ -811,10 +813,10 @@ export type CreateWorkflowVersionEdgeInput = {
export type CreateWorkflowVersionStepInput = {
/** Next step ID */
nextStepId?: InputMaybe<Scalars['UUID']>;
/** Parent step connection options */
parentStepConnectionOptions?: InputMaybe<Scalars['JSON']>;
/** Parent step ID */
parentStepId?: InputMaybe<Scalars['String']>;
/** Step creation options */
parentStepOptions?: InputMaybe<Scalars['JSON']>;
/** Step position */
position?: InputMaybe<WorkflowStepPositionInput>;
/** New step type */
@@ -1,5 +1,9 @@
import { Field, InputType } from '@nestjs/graphql';
import graphqlTypeJson from 'graphql-type-json';
import { WorkflowStepConnectionOptions } from 'src/modules/workflow/workflow-builder/workflow-version-step/types/WorkflowStepCreationOptions';
@InputType()
export class CreateWorkflowVersionEdgeInput {
@Field(() => String, {
@@ -19,4 +23,10 @@ export class CreateWorkflowVersionEdgeInput {
nullable: false,
})
target: string;
@Field(() => graphqlTypeJson, {
description: 'Workflow version source step connection options',
nullable: true,
})
sourceConnectionOptions?: WorkflowStepConnectionOptions;
}
@@ -4,7 +4,7 @@ import graphqlTypeJson from 'graphql-type-json';
import { UUIDScalarType } from 'src/engine/api/graphql/workspace-schema-builder/graphql-types/scalars';
import { WorkflowStepPositionInput } from 'src/engine/core-modules/workflow/dtos/update-workflow-step-position-input.dto';
import { WorkflowStepCreationOptions } from 'src/modules/workflow/workflow-builder/workflow-version-step/types/WorkflowStepCreationOptions';
import { WorkflowStepConnectionOptions } from 'src/modules/workflow/workflow-builder/workflow-version-step/types/WorkflowStepCreationOptions';
import { WorkflowActionType } from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type';
@InputType()
@@ -29,10 +29,10 @@ export class CreateWorkflowVersionStepInput {
parentStepId?: string;
@Field(() => graphqlTypeJson, {
description: 'Step creation options',
description: 'Parent step connection options',
nullable: true,
})
parentStepOptions?: WorkflowStepCreationOptions;
parentStepConnectionOptions?: WorkflowStepConnectionOptions;
@Field(() => UUIDScalarType, {
description: 'Next step ID',
@@ -2,7 +2,10 @@ import { Catch, type ExceptionFilter } from '@nestjs/common';
import { assertUnreachable } from 'twenty-shared/utils';
import { NotFoundError } from 'src/engine/core-modules/graphql/utils/graphql-errors.util';
import {
NotFoundError,
UserInputError,
} from 'src/engine/core-modules/graphql/utils/graphql-errors.util';
import {
WorkflowVersionEdgeException,
WorkflowVersionEdgeExceptionCode,
@@ -14,6 +17,8 @@ export const handleWorkflowVersionEdgeException = (
switch (exception.code) {
case WorkflowVersionEdgeExceptionCode.NOT_FOUND:
throw new NotFoundError(exception);
case WorkflowVersionEdgeExceptionCode.INVALID_REQUEST:
throw new UserInputError(exception);
default: {
assertUnreachable(exception.code);
}
@@ -36,13 +36,19 @@ export class WorkflowVersionEdgeResolver {
async createWorkflowVersionEdge(
@AuthWorkspace() { id: workspaceId }: Workspace,
@Args('input')
{ source, target, workflowVersionId }: CreateWorkflowVersionEdgeInput,
{
source,
target,
workflowVersionId,
sourceConnectionOptions,
}: CreateWorkflowVersionEdgeInput,
): Promise<WorkflowVersionStepChangesDTO> {
return this.workflowVersionEdgeWorkspaceService.createWorkflowVersionEdge({
source,
target,
workflowVersionId,
workspaceId,
sourceConnectionOptions,
});
}
@@ -4,4 +4,5 @@ export class WorkflowVersionEdgeException extends CustomException<WorkflowVersio
export enum WorkflowVersionEdgeExceptionCode {
NOT_FOUND = 'NOT_FOUND',
INVALID_REQUEST = 'INVALID_REQUEST',
}
@@ -132,6 +132,29 @@ describe('WorkflowVersionEdgeWorkspaceService', () => {
);
});
it('should throw if trigger is not found', async () => {
const mockWorkflowVersionWithoutTrigger = {
...mockWorkflowVersion,
trigger: null,
} as WorkflowVersionWorkspaceEntity;
workflowCommonWorkspaceService.getWorkflowVersionOrFail.mockResolvedValue(
mockWorkflowVersionWithoutTrigger,
);
const call = async () =>
await service.createWorkflowVersionEdge({
source: TRIGGER_STEP_ID,
target: 'step-1',
workflowVersionId: mockWorkflowVersionId,
workspaceId: mockWorkspaceId,
});
await expect(call).rejects.toThrow(
`Trigger 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({
@@ -191,6 +214,160 @@ describe('WorkflowVersionEdgeWorkspaceService', () => {
});
describe('with source is a step', () => {
describe('with iterator step', () => {
const mockIteratorStep = {
id: 'iterator-step',
type: WorkflowActionType.ITERATOR,
settings: {
errorHandlingOptions: {
continueOnFailure: { value: false },
retryOnFailure: { value: false },
},
input: {
initialLoopStepIds: ['step-1'],
},
},
nextStepIds: ['step-2'],
} as WorkflowAction;
beforeEach(() => {
const mockStepsWithIterator = [...mockSteps, mockIteratorStep];
const mockWorkflowVersionWithIterator = {
...mockWorkflowVersion,
steps: mockStepsWithIterator,
} as WorkflowVersionWorkspaceEntity;
workflowCommonWorkspaceService.getWorkflowVersionOrFail.mockResolvedValue(
mockWorkflowVersionWithIterator,
);
});
it('should add target to initialLoopStepIds when shouldInsertToLoop is true', async () => {
const result = await service.createWorkflowVersionEdge({
source: 'iterator-step',
target: 'step-3',
workflowVersionId: mockWorkflowVersionId,
workspaceId: mockWorkspaceId,
sourceConnectionOptions: {
connectedStepType: WorkflowActionType.ITERATOR,
settings: {
shouldInsertToLoop: true,
},
},
});
expect(
mockWorkflowVersionWorkspaceRepository.update,
).toHaveBeenCalledWith(mockWorkflowVersionId, {
steps: expect.arrayContaining([
expect.objectContaining({
id: 'iterator-step',
settings: expect.objectContaining({
input: expect.objectContaining({
initialLoopStepIds: ['step-1', 'step-3'],
}),
}),
}),
]),
});
expect(result).toEqual({
triggerNextStepIds: ['step-1'],
stepsNextStepIds: {
'step-1': ['step-2'],
'step-2': [],
'step-3': [],
'iterator-step': ['step-2'],
},
});
});
it('should not duplicate target in initialLoopStepIds when it already exists', async () => {
const result = await service.createWorkflowVersionEdge({
source: 'iterator-step',
target: 'step-1',
workflowVersionId: mockWorkflowVersionId,
workspaceId: mockWorkspaceId,
sourceConnectionOptions: {
connectedStepType: WorkflowActionType.ITERATOR,
settings: {
shouldInsertToLoop: true,
},
},
});
expect(
mockWorkflowVersionWorkspaceRepository.update,
).not.toHaveBeenCalled();
expect(result).toEqual({
triggerNextStepIds: ['step-1'],
stepsNextStepIds: {
'step-1': ['step-2'],
'step-2': [],
'step-3': [],
'iterator-step': ['step-2'],
},
});
});
it('should throw if source is not an iterator but connection options are for iterator', async () => {
const call = async () =>
await service.createWorkflowVersionEdge({
source: 'step-1',
target: 'step-3',
workflowVersionId: mockWorkflowVersionId,
workspaceId: mockWorkspaceId,
sourceConnectionOptions: {
connectedStepType: WorkflowActionType.ITERATOR,
settings: {
shouldInsertToLoop: true,
},
},
});
await expect(call).rejects.toThrow(
`Source step 'step-1' is not an iterator`,
);
});
it('should create normal edge when shouldInsertToLoop is false', async () => {
const result = await service.createWorkflowVersionEdge({
source: 'iterator-step',
target: 'step-3',
workflowVersionId: mockWorkflowVersionId,
workspaceId: mockWorkspaceId,
sourceConnectionOptions: {
connectedStepType: WorkflowActionType.ITERATOR,
settings: {
shouldInsertToLoop: false,
},
},
});
expect(
mockWorkflowVersionWorkspaceRepository.update,
).toHaveBeenCalledWith(mockWorkflowVersionId, {
steps: expect.arrayContaining([
expect.objectContaining({
id: 'iterator-step',
nextStepIds: ['step-2', 'step-3'],
}),
]),
});
expect(result).toEqual({
triggerNextStepIds: ['step-1'],
stepsNextStepIds: {
'step-1': ['step-2'],
'step-2': [],
'step-3': [],
'iterator-step': ['step-2', 'step-3'],
},
});
});
});
it('should create an edge between step-2 and step-3', async () => {
const result = await service.createWorkflowVersionEdge({
source: 'step-2',
@@ -14,6 +14,7 @@ import { type WorkflowVersionWorkspaceEntity } from 'src/modules/workflow/common
import { assertWorkflowVersionIsDraft } from 'src/modules/workflow/common/utils/assert-workflow-version-is-draft.util';
import { WorkflowCommonWorkspaceService } from 'src/modules/workflow/common/workspace-services/workflow-common.workspace-service';
import { computeWorkflowVersionStepChanges } from 'src/modules/workflow/workflow-builder/utils/compute-workflow-version-step-updates.util';
import { WorkflowStepConnectionOptions } from 'src/modules/workflow/workflow-builder/workflow-version-step/types/WorkflowStepCreationOptions';
import {
type WorkflowAction,
WorkflowActionType,
@@ -32,11 +33,13 @@ export class WorkflowVersionEdgeWorkspaceService {
target,
workflowVersionId,
workspaceId,
sourceConnectionOptions,
}: {
source: string;
target: string;
workflowVersionId: string;
workspaceId: string;
sourceConnectionOptions?: WorkflowStepConnectionOptions;
}): Promise<WorkflowVersionStepChangesDTO> {
const workflowVersionRepository =
await this.twentyORMGlobalManager.getRepositoryForWorkspace<WorkflowVersionWorkspaceEntity>(
@@ -81,6 +84,7 @@ export class WorkflowVersionEdgeWorkspaceService {
steps,
source,
target,
sourceConnectionOptions,
workflowVersion,
workflowVersionRepository,
});
@@ -196,6 +200,7 @@ export class WorkflowVersionEdgeWorkspaceService {
target,
workflowVersion,
workflowVersionRepository,
sourceConnectionOptions,
}: {
trigger: WorkflowTrigger | null;
steps: WorkflowAction[];
@@ -203,6 +208,7 @@ export class WorkflowVersionEdgeWorkspaceService {
target: string;
workflowVersion: WorkflowVersionWorkspaceEntity;
workflowVersionRepository: WorkspaceRepository<WorkflowVersionWorkspaceEntity>;
sourceConnectionOptions?: WorkflowStepConnectionOptions;
}): Promise<WorkflowVersionStepChangesDTO> {
const sourceStep = steps.find((step) => step.id === source);
@@ -213,17 +219,15 @@ export class WorkflowVersionEdgeWorkspaceService {
);
}
if (sourceStep.nextStepIds?.includes(target)) {
return computeWorkflowVersionStepChanges({
trigger,
steps,
});
}
const updatedSourceStep = {
...sourceStep,
nextStepIds: [...(sourceStep.nextStepIds ?? []), target],
};
const { updatedSourceStep, shouldPersist } = isDefined(
sourceConnectionOptions,
)
? this.buildUpdatedSourceStepWithConnectionOptions({
sourceStep,
target,
sourceConnectionOptions,
})
: this.buildUpdatedSourceStep({ sourceStep, target });
const updatedSteps = steps.map((step) => {
if (step.id === source) {
@@ -233,9 +237,11 @@ export class WorkflowVersionEdgeWorkspaceService {
return step;
});
await workflowVersionRepository.update(workflowVersion.id, {
steps: updatedSteps,
});
if (shouldPersist) {
await workflowVersionRepository.update(workflowVersion.id, {
steps: updatedSteps,
});
}
return computeWorkflowVersionStepChanges({
trigger,
@@ -243,6 +249,97 @@ export class WorkflowVersionEdgeWorkspaceService {
});
}
private buildUpdatedSourceStepWithConnectionOptions({
sourceStep,
target,
sourceConnectionOptions,
}: {
sourceStep: WorkflowAction;
target: string;
sourceConnectionOptions: WorkflowStepConnectionOptions;
}): {
updatedSourceStep: WorkflowAction;
shouldPersist: boolean;
} {
switch (sourceConnectionOptions.connectedStepType) {
case WorkflowActionType.ITERATOR:
if (sourceStep.type !== WorkflowActionType.ITERATOR) {
throw new WorkflowVersionEdgeException(
`Source step '${sourceStep.id}' is not an iterator`,
WorkflowVersionEdgeExceptionCode.INVALID_REQUEST,
);
}
if (sourceConnectionOptions.settings.shouldInsertToLoop) {
const currentInitialLoopStepIds =
sourceStep.settings.input.initialLoopStepIds;
if (currentInitialLoopStepIds?.includes(target)) {
return {
updatedSourceStep: sourceStep,
shouldPersist: false,
};
}
return {
updatedSourceStep: {
...sourceStep,
settings: {
...sourceStep.settings,
input: {
...sourceStep.settings.input,
initialLoopStepIds: [
...(currentInitialLoopStepIds ?? []),
target,
],
},
},
},
shouldPersist: true,
};
} else {
return this.buildUpdatedSourceStep({
sourceStep,
target,
});
}
default:
return this.buildUpdatedSourceStep({
sourceStep,
target,
});
}
}
private buildUpdatedSourceStep({
sourceStep,
target,
}: {
sourceStep: WorkflowAction;
target: string;
}): {
updatedSourceStep: WorkflowAction;
shouldPersist: boolean;
} {
if (sourceStep.nextStepIds?.includes(target)) {
return {
updatedSourceStep: sourceStep,
shouldPersist: false,
};
}
const updatedSourceStep = {
...sourceStep,
nextStepIds: [...(sourceStep.nextStepIds ?? []), target],
};
return {
updatedSourceStep,
shouldPersist: true,
};
}
private async deleteTriggerEdge({
trigger,
steps,
@@ -1,10 +1,11 @@
import { type WorkflowActionType } from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type';
type WorkflowIteratorStepCreationOptions = {
parentStepType: WorkflowActionType.ITERATOR;
type WorkflowIteratorStepConnectionOptions = {
connectedStepType: WorkflowActionType.ITERATOR;
settings: {
shouldInsertToLoop: boolean;
};
};
export type WorkflowStepCreationOptions = WorkflowIteratorStepCreationOptions;
export type WorkflowStepConnectionOptions =
WorkflowIteratorStepConnectionOptions;
@@ -187,8 +187,8 @@ describe('insertStep', () => {
existingSteps: [mockIteratorStep],
insertedStep: newStep,
parentStepId: '1',
parentStepOptions: {
parentStepType: WorkflowActionType.ITERATOR,
parentStepConnectionOptions: {
connectedStepType: WorkflowActionType.ITERATOR,
settings: {
shouldInsertToLoop: true,
},
@@ -213,8 +213,8 @@ describe('insertStep', () => {
existingSteps: [mockIteratorStep],
insertedStep: newStep,
parentStepId: '1',
parentStepOptions: {
parentStepType: WorkflowActionType.ITERATOR,
parentStepConnectionOptions: {
connectedStepType: WorkflowActionType.ITERATOR,
settings: {
shouldInsertToLoop: false,
},
@@ -5,11 +5,11 @@ import {
WorkflowVersionStepException,
WorkflowVersionStepExceptionCode,
} from 'src/modules/workflow/common/exceptions/workflow-version-step.exception';
import { type WorkflowStepCreationOptions } from 'src/modules/workflow/workflow-builder/workflow-version-step/types/WorkflowStepCreationOptions';
import { type WorkflowStepConnectionOptions } from 'src/modules/workflow/workflow-builder/workflow-version-step/types/WorkflowStepCreationOptions';
import { type WorkflowIteratorActionSettings } from 'src/modules/workflow/workflow-executor/workflow-actions/iterator/types/workflow-iterator-action-settings.type';
import {
WorkflowActionType,
type WorkflowAction,
WorkflowActionType,
} from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type';
import { type WorkflowTrigger } from 'src/modules/workflow/workflow-trigger/types/workflow-trigger.type';
@@ -19,14 +19,14 @@ export const insertStep = ({
insertedStep,
nextStepId,
parentStepId,
parentStepOptions,
parentStepConnectionOptions,
}: {
existingSteps: WorkflowAction[];
existingTrigger: WorkflowTrigger | null;
insertedStep: WorkflowAction;
nextStepId?: string;
parentStepId?: string;
parentStepOptions?: WorkflowStepCreationOptions;
parentStepConnectionOptions?: WorkflowStepConnectionOptions;
}): {
updatedSteps: WorkflowAction[];
updatedInsertedStep: WorkflowAction;
@@ -39,7 +39,7 @@ export const insertStep = ({
parentStepId,
insertedStepId: insertedStep.id,
nextStepId,
parentStepOptions,
parentStepConnectionOptions,
})
: {
updatedSteps: existingSteps,
@@ -64,24 +64,24 @@ const updateParentStep = ({
parentStepId,
insertedStepId,
nextStepId,
parentStepOptions,
parentStepConnectionOptions,
}: {
steps: WorkflowAction[];
trigger: WorkflowTrigger | null;
parentStepId: string;
insertedStepId: string;
nextStepId?: string;
parentStepOptions?: WorkflowStepCreationOptions;
parentStepConnectionOptions?: WorkflowStepConnectionOptions;
}): {
updatedSteps: WorkflowAction[];
updatedTrigger: WorkflowTrigger | null;
} => {
if (isDefined(parentStepOptions)) {
if (isDefined(parentStepConnectionOptions)) {
return updateStepsWithOptions({
steps,
parentStepId,
insertedStepId,
parentStepOptions,
parentStepConnectionOptions,
trigger,
});
} else {
@@ -160,20 +160,20 @@ const updateStepsWithOptions = ({
parentStepId,
insertedStepId,
steps,
parentStepOptions,
parentStepConnectionOptions,
trigger,
}: {
parentStepId: string;
insertedStepId: string;
steps: WorkflowAction[];
parentStepOptions: WorkflowStepCreationOptions;
parentStepConnectionOptions: WorkflowStepConnectionOptions;
trigger: WorkflowTrigger | null;
}) => {
let updatedSteps = steps;
switch (parentStepOptions.parentStepType) {
switch (parentStepConnectionOptions.connectedStepType) {
case WorkflowActionType.ITERATOR:
if (!parentStepOptions.settings.shouldInsertToLoop) {
if (!parentStepConnectionOptions.settings.shouldInsertToLoop) {
break;
}
@@ -79,7 +79,7 @@ export class WorkflowVersionStepWorkspaceService {
parentStepId,
nextStepId,
position,
parentStepOptions,
parentStepConnectionOptions,
} = input;
const newStep = await this.runStepCreationSideEffectsAndBuildStep({
@@ -126,7 +126,7 @@ export class WorkflowVersionStepWorkspaceService {
insertedStep: enrichedNewStep,
parentStepId,
nextStepId,
parentStepOptions,
parentStepConnectionOptions,
});
await workflowVersionRepository.update(workflowVersion.id, {