Add icon to duplicate + split step service in two (#14521)

As title. Created a new step operations service.

Also added Duplicate Icon
<img width="500" height="767" alt="Capture d’écran 2025-09-16 à 09 42
19"
src="https://github.com/user-attachments/assets/4d789639-62f6-4fea-afb8-396868d7c2cb"
/>
This commit is contained in:
Thomas Trompette
2025-09-16 14:45:16 +00:00
committed by GitHub
parent e9f277fe93
commit 4b889db06c
9 changed files with 953 additions and 589 deletions
@@ -6,7 +6,9 @@ import { RightDrawerFooter } from '@/ui/layout/right-drawer/components/RightDraw
import { SelectableList } from '@/ui/layout/selectable-list/components/SelectableList';
import { useDuplicateStep } from '@/workflow/workflow-steps/hooks/useDuplicateStep';
import { useTheme } from '@emotion/react';
import { useLingui } from '@lingui/react/macro';
import { useId } from 'react';
import { IconCopyPlus } from 'twenty-ui/display';
import { Button } from 'twenty-ui/input';
import { MenuItem } from 'twenty-ui/navigation';
import { getOsControlSymbol } from 'twenty-ui/utilities';
@@ -19,6 +21,7 @@ export const WorkflowActionFooter = ({
additionalActions?: React.ReactNode[];
}) => {
const dropdownId = useId();
const { t } = useLingui();
const theme = useTheme();
const { duplicateStep } = useDuplicateStep();
const { closeDropdown } = useCloseDropdown();
@@ -50,7 +53,8 @@ export const WorkflowActionFooter = ({
closeDropdown(dropdownId);
duplicateStep({ stepId });
}}
text="Duplicate"
text={t`Duplicate node`}
LeftIcon={IconCopyPlus}
/>
</SelectableList>
</DropdownMenuItemsContainer>
@@ -91,6 +91,40 @@ export class WorkflowSchemaWorkspaceService {
}
}
async enrichOutputSchema({
step,
workspaceId,
}: {
step: WorkflowAction;
workspaceId: string;
}): Promise<WorkflowAction> {
// We don't enrich on the fly for code and HTTP request workflow actions.
// For code actions, OutputSchema is computed and updated when testing the serverless function.
// For HTTP requests and AI agent, OutputSchema is determined by the example response input
if (
[
WorkflowActionType.CODE,
WorkflowActionType.HTTP_REQUEST,
WorkflowActionType.AI_AGENT,
].includes(step.type)
) {
return step;
}
const result = { ...step };
const outputSchema = await this.computeStepOutputSchema({
step,
workspaceId,
});
result.settings = {
...result.settings,
outputSchema: outputSchema || {},
};
return result;
}
private async computeDatabaseEventTriggerOutputSchema({
eventName,
workspaceId,
@@ -0,0 +1,312 @@
import { Test, type TestingModule } from '@nestjs/testing';
import { getRepositoryToken } from '@nestjs/typeorm';
import { AgentEntity } from 'src/engine/metadata-modules/agent/agent.entity';
import { ObjectMetadataEntity } from 'src/engine/metadata-modules/object-metadata/object-metadata.entity';
import { type ServerlessFunctionEntity } from 'src/engine/metadata-modules/serverless-function/serverless-function.entity';
import { ServerlessFunctionService } from 'src/engine/metadata-modules/serverless-function/serverless-function.service';
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 { WorkflowVersionStepOperationsWorkspaceService } from 'src/modules/workflow/workflow-builder/workflow-version-step/workflow-version-step-operations.workspace-service';
import {
type WorkflowAction,
WorkflowActionType,
} from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type';
const mockWorkspaceId = 'workspace-id';
describe('WorkflowVersionStepOperationsWorkspaceService', () => {
let service: WorkflowVersionStepOperationsWorkspaceService;
let twentyORMGlobalManager: jest.Mocked<TwentyORMGlobalManager>;
let serverlessFunctionService: jest.Mocked<ServerlessFunctionService>;
let agentRepository: jest.Mocked<any>;
let objectMetadataRepository: jest.Mocked<any>;
let workflowCommonWorkspaceService: jest.Mocked<WorkflowCommonWorkspaceService>;
beforeEach(async () => {
serverlessFunctionService = {
createOneServerlessFunction: jest.fn(),
hasServerlessFunctionPublishedVersion: jest.fn(),
deleteOneServerlessFunction: jest.fn(),
duplicateServerlessFunction: jest.fn(),
createDraftFromPublishedVersion: jest.fn(),
} as unknown as jest.Mocked<ServerlessFunctionService>;
agentRepository = {
findOne: jest.fn(),
delete: jest.fn(),
};
objectMetadataRepository = {
findOne: jest.fn(),
};
workflowCommonWorkspaceService = {
getObjectMetadataItemWithFieldsMaps: jest.fn(),
} as unknown as jest.Mocked<WorkflowCommonWorkspaceService>;
twentyORMGlobalManager = {
getRepositoryForWorkspace: jest.fn(),
} as unknown as jest.Mocked<TwentyORMGlobalManager>;
const module: TestingModule = await Test.createTestingModule({
providers: [
WorkflowVersionStepOperationsWorkspaceService,
{
provide: TwentyORMGlobalManager,
useValue: twentyORMGlobalManager,
},
{
provide: ServerlessFunctionService,
useValue: serverlessFunctionService,
},
{
provide: getRepositoryToken(AgentEntity),
useValue: agentRepository,
},
{
provide: getRepositoryToken(ObjectMetadataEntity),
useValue: objectMetadataRepository,
},
{
provide: WorkflowCommonWorkspaceService,
useValue: workflowCommonWorkspaceService,
},
{
provide: ScopedWorkspaceContextFactory,
useValue: {},
},
],
}).compile();
service = module.get(WorkflowVersionStepOperationsWorkspaceService);
});
describe('runWorkflowVersionStepDeletionSideEffects', () => {
it('should delete serverless function when deleting code step', async () => {
const step = {
id: 'step-id',
name: 'Code Step',
type: WorkflowActionType.CODE,
valid: true,
nextStepIds: [],
settings: {
input: {
serverlessFunctionId: 'function-id',
serverlessFunctionVersion: 'v1',
},
outputSchema: {},
errorHandlingOptions: {
continueOnFailure: { value: false },
retryOnFailure: { value: false },
},
},
} as unknown as WorkflowAction;
serverlessFunctionService.hasServerlessFunctionPublishedVersion.mockResolvedValue(
false,
);
await service.runWorkflowVersionStepDeletionSideEffects({
step,
workspaceId: mockWorkspaceId,
});
expect(
serverlessFunctionService.deleteOneServerlessFunction,
).toHaveBeenCalledWith({
id: 'function-id',
workspaceId: mockWorkspaceId,
softDelete: false,
});
});
it('should delete agent when deleting AI agent step', async () => {
const step = {
id: 'step-id',
name: 'AI Agent Step',
type: WorkflowActionType.AI_AGENT,
valid: true,
nextStepIds: [],
settings: {
input: {
agentId: 'agent-id',
prompt: '',
},
outputSchema: {},
errorHandlingOptions: {
continueOnFailure: { value: false },
retryOnFailure: { value: false },
},
},
} as unknown as WorkflowAction;
agentRepository.findOne.mockResolvedValue({ id: 'agent-id' });
await service.runWorkflowVersionStepDeletionSideEffects({
step,
workspaceId: mockWorkspaceId,
});
expect(agentRepository.delete).toHaveBeenCalledWith({
id: 'agent-id',
workspaceId: mockWorkspaceId,
});
});
});
describe('runStepCreationSideEffectsAndBuildStep', () => {
it('should create code step with serverless function', async () => {
const mockServerlessFunction = {
id: 'new-function-id',
name: 'Test Function',
description: 'Test Description',
latestVersion: 'v1',
publishedVersions: [],
workspaceId: mockWorkspaceId,
createdAt: new Date(),
updatedAt: new Date(),
deletedAt: null,
isActive: true,
isSystem: false,
isCustom: true,
isPublic: false,
latestVersionInputSchema: {},
runtime: 'nodejs',
timeoutSeconds: 30,
layerVersion: 1,
layerArn: '',
layerName: '',
layerSize: 0,
} as unknown as ServerlessFunctionEntity;
serverlessFunctionService.createOneServerlessFunction.mockResolvedValue(
mockServerlessFunction,
);
const result = await service.runStepCreationSideEffectsAndBuildStep({
type: WorkflowActionType.CODE,
workspaceId: mockWorkspaceId,
workflowVersionId: 'workflow-version-id',
});
expect(result.type).toBe(WorkflowActionType.CODE);
const codeResult = result as unknown as {
settings: {
input: {
serverlessFunctionId: string;
serverlessFunctionVersion: string;
};
};
};
expect(codeResult.settings.input.serverlessFunctionId).toBe(
'new-function-id',
);
expect(codeResult.settings.input.serverlessFunctionVersion).toBe('draft');
});
it('should create form step', async () => {
const result = await service.runStepCreationSideEffectsAndBuildStep({
type: WorkflowActionType.FORM,
workspaceId: mockWorkspaceId,
workflowVersionId: 'workflow-version-id',
});
expect(result.type).toBe(WorkflowActionType.FORM);
expect(result.settings.input).toEqual([]);
});
});
describe('createStepForDuplicate', () => {
it('should duplicate code step with new serverless function', async () => {
const originalStep = {
id: 'original-id',
type: WorkflowActionType.CODE,
name: 'Original Step',
valid: true,
settings: {
input: {
serverlessFunctionId: 'function-id',
serverlessFunctionVersion: 'v1',
},
},
nextStepIds: ['next-step'],
} as unknown as WorkflowAction;
const mockNewServerlessFunction = {
id: 'new-function-id',
name: 'Test Function',
description: 'Test Description',
latestVersion: 'v1',
publishedVersions: [],
workspaceId: mockWorkspaceId,
createdAt: new Date(),
updatedAt: new Date(),
deletedAt: null,
isActive: true,
isSystem: false,
isCustom: true,
isPublic: false,
latestVersionInputSchema: {},
runtime: 'nodejs',
timeoutSeconds: 30,
layerVersion: 1,
layerArn: '',
layerName: '',
layerSize: 0,
} as unknown as ServerlessFunctionEntity;
serverlessFunctionService.duplicateServerlessFunction.mockResolvedValue(
mockNewServerlessFunction,
);
const result = await service.createStepForDuplicate({
step: originalStep,
workspaceId: mockWorkspaceId,
});
expect(result.id).not.toBe('original-id');
expect(result.name).toBe('Original Step (Duplicate)');
const codeResult = result as unknown as {
settings: {
input: {
serverlessFunctionId: string;
serverlessFunctionVersion: string;
};
};
};
expect(codeResult.settings.input.serverlessFunctionId).toBe(
'new-function-id',
);
expect(codeResult.settings.input.serverlessFunctionVersion).toBe('draft');
expect(result.nextStepIds).toEqual([]);
});
it('should duplicate non-code step', async () => {
const originalStep = {
id: 'original-id',
type: WorkflowActionType.FORM,
name: 'Original Step',
valid: true,
settings: {
input: [],
},
nextStepIds: ['next-step'],
} as unknown as WorkflowAction;
const result = await service.createStepForDuplicate({
step: originalStep,
workspaceId: mockWorkspaceId,
});
expect(result.id).not.toBe('original-id');
expect(result.name).toBe('Original Step (Duplicate)');
expect(result.settings).toEqual(originalStep.settings);
expect(result.nextStepIds).toEqual([]);
});
});
});
@@ -1,17 +1,12 @@
import { Test, type TestingModule } from '@nestjs/testing';
import { getRepositoryToken } from '@nestjs/typeorm';
import { TRIGGER_STEP_ID } from 'twenty-shared/workflow';
import { AgentEntity } from 'src/engine/metadata-modules/agent/agent.entity';
import { ObjectMetadataEntity } from 'src/engine/metadata-modules/object-metadata/object-metadata.entity';
import { ServerlessFunctionService } from 'src/engine/metadata-modules/serverless-function/serverless-function.service';
import { ScopedWorkspaceContextFactory } from 'src/engine/twenty-orm/factories/scoped-workspace-context.factory';
import { type WorkspaceRepository } from 'src/engine/twenty-orm/repository/workspace.repository';
import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager';
import { type WorkflowVersionWorkspaceEntity } from 'src/modules/workflow/common/standard-objects/workflow-version.workspace-entity';
import { WorkflowCommonWorkspaceService } from 'src/modules/workflow/common/workspace-services/workflow-common.workspace-service';
import { WorkflowSchemaWorkspaceService } from 'src/modules/workflow/workflow-builder/workflow-schema/workflow-schema.workspace-service';
import { WorkflowVersionStepOperationsWorkspaceService } from 'src/modules/workflow/workflow-builder/workflow-version-step/workflow-version-step-operations.workspace-service';
import { WorkflowVersionStepWorkspaceService } from 'src/modules/workflow/workflow-builder/workflow-version-step/workflow-version-step.workspace-service';
import {
type WorkflowAction,
@@ -111,26 +106,27 @@ describe('WorkflowVersionStepWorkspaceService', () => {
{
provide: WorkflowSchemaWorkspaceService,
useValue: {
computeStepOutputSchema: jest.fn(),
},
},
{ provide: ServerlessFunctionService, useValue: {} },
{
provide: getRepositoryToken(AgentEntity),
useValue: {
findOne: jest.fn(),
enrichOutputSchema: jest
.fn()
.mockImplementation((args) => args.step),
},
},
{
provide: getRepositoryToken(ObjectMetadataEntity),
provide: WorkflowVersionStepOperationsWorkspaceService,
useValue: {
findOne: jest.fn(),
runStepCreationSideEffectsAndBuildStep: jest
.fn()
.mockImplementation(({ type }) => ({
id: 'new-step-id',
type,
settings: {},
nextStepIds: [],
})),
runWorkflowVersionStepDeletionSideEffects: jest.fn(),
},
},
{ provide: WorkflowRunWorkspaceService, useValue: {} },
{ provide: WorkflowRunnerWorkspaceService, useValue: {} },
{ provide: WorkflowCommonWorkspaceService, useValue: {} },
{ provide: ScopedWorkspaceContextFactory, useValue: {} },
],
}).compile();
@@ -0,0 +1,532 @@
import { Injectable } from '@nestjs/common';
import { InjectRepository } from '@nestjs/typeorm';
import { FieldMetadataType } from 'twenty-shared/types';
import { isDefined, isValidUuid } from 'twenty-shared/utils';
import { Repository } from 'typeorm';
import { v4 } from 'uuid';
import { BASE_TYPESCRIPT_PROJECT_INPUT_SCHEMA } from 'src/engine/core-modules/serverless/drivers/constants/base-typescript-project-input-schema';
import { type WorkflowStepPositionInput } from 'src/engine/core-modules/workflow/dtos/update-workflow-step-position-input.dto';
import { AgentEntity } from 'src/engine/metadata-modules/agent/agent.entity';
import { ObjectMetadataEntity } from 'src/engine/metadata-modules/object-metadata/object-metadata.entity';
import { ServerlessFunctionService } from 'src/engine/metadata-modules/serverless-function/serverless-function.service';
import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager';
import {
WorkflowVersionStepException,
WorkflowVersionStepExceptionCode,
} from 'src/modules/workflow/common/exceptions/workflow-version-step.exception';
import { type WorkflowVersionWorkspaceEntity } from 'src/modules/workflow/common/standard-objects/workflow-version.workspace-entity';
import { WorkflowCommonWorkspaceService } from 'src/modules/workflow/common/workspace-services/workflow-common.workspace-service';
import { type BaseWorkflowActionSettings } from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action-settings.type';
import {
type WorkflowAction,
WorkflowActionType,
type WorkflowEmptyAction,
type WorkflowFormAction,
} from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type';
const BASE_STEP_DEFINITION: BaseWorkflowActionSettings = {
outputSchema: {},
errorHandlingOptions: {
continueOnFailure: {
value: false,
},
retryOnFailure: {
value: false,
},
},
};
const DUPLICATED_STEP_POSITION_OFFSET = 50;
@Injectable()
export class WorkflowVersionStepOperationsWorkspaceService {
constructor(
private readonly twentyORMGlobalManager: TwentyORMGlobalManager,
private readonly serverlessFunctionService: ServerlessFunctionService,
@InjectRepository(AgentEntity)
private readonly agentRepository: Repository<AgentEntity>,
@InjectRepository(ObjectMetadataEntity)
private readonly objectMetadataRepository: Repository<ObjectMetadataEntity>,
private readonly workflowCommonWorkspaceService: WorkflowCommonWorkspaceService,
) {}
async runWorkflowVersionStepDeletionSideEffects({
step,
workspaceId,
}: {
step: WorkflowAction;
workspaceId: string;
}) {
switch (step.type) {
case WorkflowActionType.CODE: {
if (
!(await this.serverlessFunctionService.hasServerlessFunctionPublishedVersion(
step.settings.input.serverlessFunctionId,
))
) {
await this.serverlessFunctionService.deleteOneServerlessFunction({
id: step.settings.input.serverlessFunctionId,
workspaceId,
softDelete: false,
});
}
break;
}
case WorkflowActionType.AI_AGENT: {
if (!isDefined(step.settings.input.agentId)) {
break;
}
const agent = await this.agentRepository.findOne({
where: { id: step.settings.input.agentId, workspaceId },
});
if (isDefined(agent)) {
await this.agentRepository.delete({ id: agent.id, workspaceId });
}
break;
}
}
}
async runStepCreationSideEffectsAndBuildStep({
type,
workspaceId,
position,
workflowVersionId,
}: {
type: WorkflowActionType;
workspaceId: string;
position?: WorkflowStepPositionInput;
workflowVersionId: string;
}): Promise<WorkflowAction> {
const newStepId = v4();
const baseStep = {
id: newStepId,
position,
valid: false,
nextStepIds: [],
};
switch (type) {
case WorkflowActionType.CODE: {
const newServerlessFunction =
await this.serverlessFunctionService.createOneServerlessFunction(
{
name: 'A Serverless Function Code Workflow Step',
description: '',
},
workspaceId,
);
if (!isDefined(newServerlessFunction)) {
throw new WorkflowVersionStepException(
'Fail to create Code Step',
WorkflowVersionStepExceptionCode.CODE_STEP_FAILURE,
);
}
return {
...baseStep,
name: 'Code - Serverless Function',
type: WorkflowActionType.CODE,
settings: {
...BASE_STEP_DEFINITION,
outputSchema: {
link: {
isLeaf: true,
icon: 'IconVariable',
tab: 'test',
label: 'Generate Function Output',
},
_outputSchemaType: 'LINK',
},
input: {
serverlessFunctionId: newServerlessFunction.id,
serverlessFunctionVersion: 'draft',
serverlessFunctionInput: BASE_TYPESCRIPT_PROJECT_INPUT_SCHEMA,
},
},
};
}
case WorkflowActionType.SEND_EMAIL: {
return {
...baseStep,
name: 'Send Email',
type: WorkflowActionType.SEND_EMAIL,
settings: {
...BASE_STEP_DEFINITION,
input: {
connectedAccountId: '',
email: '',
subject: '',
body: '',
},
},
};
}
case WorkflowActionType.CREATE_RECORD: {
const activeObjectMetadataItem =
await this.objectMetadataRepository.findOne({
where: { workspaceId, isActive: true, isSystem: false },
});
return {
...baseStep,
name: 'Create Record',
type: WorkflowActionType.CREATE_RECORD,
settings: {
...BASE_STEP_DEFINITION,
input: {
objectName: activeObjectMetadataItem?.nameSingular || '',
objectRecord: {},
},
},
};
}
case WorkflowActionType.UPDATE_RECORD: {
const activeObjectMetadataItem =
await this.objectMetadataRepository.findOne({
where: { workspaceId, isActive: true, isSystem: false },
});
return {
...baseStep,
name: 'Update Record',
type: WorkflowActionType.UPDATE_RECORD,
settings: {
...BASE_STEP_DEFINITION,
input: {
objectName: activeObjectMetadataItem?.nameSingular || '',
objectRecord: {},
objectRecordId: '',
fieldsToUpdate: [],
},
},
};
}
case WorkflowActionType.DELETE_RECORD: {
const activeObjectMetadataItem =
await this.objectMetadataRepository.findOne({
where: { workspaceId, isActive: true, isSystem: false },
});
return {
...baseStep,
name: 'Delete Record',
type: WorkflowActionType.DELETE_RECORD,
settings: {
...BASE_STEP_DEFINITION,
input: {
objectName: activeObjectMetadataItem?.nameSingular || '',
objectRecordId: '',
},
},
};
}
case WorkflowActionType.FIND_RECORDS: {
const activeObjectMetadataItem =
await this.objectMetadataRepository.findOne({
where: { workspaceId, isActive: true, isSystem: false },
});
return {
...baseStep,
name: 'Search Records',
type: WorkflowActionType.FIND_RECORDS,
settings: {
...BASE_STEP_DEFINITION,
input: {
objectName: activeObjectMetadataItem?.nameSingular || '',
limit: 1,
},
},
};
}
case WorkflowActionType.FORM: {
return {
...baseStep,
name: 'Form',
type: WorkflowActionType.FORM,
settings: {
...BASE_STEP_DEFINITION,
input: [],
},
};
}
case WorkflowActionType.FILTER: {
return {
...baseStep,
name: 'Filter',
type: WorkflowActionType.FILTER,
settings: {
...BASE_STEP_DEFINITION,
input: {
stepFilterGroups: [],
stepFilters: [],
},
},
};
}
case WorkflowActionType.HTTP_REQUEST: {
return {
...baseStep,
name: 'HTTP Request',
type: WorkflowActionType.HTTP_REQUEST,
settings: {
...BASE_STEP_DEFINITION,
input: {
url: '',
method: 'GET',
headers: {},
body: {},
},
},
};
}
case WorkflowActionType.AI_AGENT: {
return {
...baseStep,
name: 'AI Agent',
type: WorkflowActionType.AI_AGENT,
settings: {
...BASE_STEP_DEFINITION,
input: {
agentId: '',
prompt: '',
},
},
};
}
case WorkflowActionType.ITERATOR: {
const emptyNodeStep = await this.createEmptyNodeForIteratorStep({
iteratorStepId: baseStep.id,
workflowVersionId,
workspaceId,
});
return {
...baseStep,
name: 'Iterator',
type: WorkflowActionType.ITERATOR,
settings: {
...BASE_STEP_DEFINITION,
input: {
items: [],
initialLoopStepIds: [emptyNodeStep.id],
},
},
};
}
default:
throw new WorkflowVersionStepException(
`WorkflowActionType '${type}' unknown`,
WorkflowVersionStepExceptionCode.INVALID_REQUEST,
);
}
}
async enrichFormStepResponse({
workspaceId,
step,
response,
}: {
workspaceId: string;
step: WorkflowFormAction;
response: object;
}) {
const responseKeys = Object.keys(response);
const enrichedResponses = await Promise.all(
responseKeys.map(async (key) => {
// @ts-expect-error legacy noImplicitAny
if (!isDefined(response[key])) {
// @ts-expect-error legacy noImplicitAny
return { key, value: response[key] };
}
const field = step.settings.input.find((field) => field.name === key);
if (
field?.type === 'RECORD' &&
field?.settings?.objectName &&
// @ts-expect-error legacy noImplicitAny
isDefined(response[key].id) &&
// @ts-expect-error legacy noImplicitAny
isValidUuid(response[key].id)
) {
const objectMetadataInfo =
await this.workflowCommonWorkspaceService.getObjectMetadataItemWithFieldsMaps(
field.settings.objectName,
workspaceId,
);
const relationFieldsNames = Object.values(
objectMetadataInfo.objectMetadataItemWithFieldsMaps.fieldsById,
)
.filter((field) => field.type === FieldMetadataType.RELATION)
.map((field) => field.name);
const repository =
await this.twentyORMGlobalManager.getRepositoryForWorkspace(
workspaceId,
field.settings.objectName,
{ shouldBypassPermissionChecks: true },
);
const record = await repository.findOne({
// @ts-expect-error legacy noImplicitAny
where: { id: response[key].id },
relations: relationFieldsNames,
});
return { key, value: record };
} else {
// @ts-expect-error legacy noImplicitAny
return { key, value: response[key] };
}
}),
);
return enrichedResponses.reduce((acc, { key, value }) => {
// @ts-expect-error legacy noImplicitAny
acc[key] = value;
return acc;
}, {});
}
async createStepForDuplicate({
step,
workspaceId,
}: {
step: WorkflowAction;
workspaceId: string;
}): Promise<WorkflowAction> {
const duplicatedStepPosition = {
x: (step.position?.x ?? 0) + DUPLICATED_STEP_POSITION_OFFSET,
y: (step.position?.y ?? 0) + DUPLICATED_STEP_POSITION_OFFSET,
};
switch (step.type) {
case WorkflowActionType.CODE: {
const newServerlessFunction =
await this.serverlessFunctionService.duplicateServerlessFunction({
id: step.settings.input.serverlessFunctionId,
version: step.settings.input.serverlessFunctionVersion,
workspaceId,
});
return {
...step,
id: v4(),
name: `${step.name} (Duplicate)`,
nextStepIds: [],
position: duplicatedStepPosition,
settings: {
...step.settings,
input: {
...step.settings.input,
serverlessFunctionId: newServerlessFunction.id,
serverlessFunctionVersion: 'draft',
},
},
};
}
default: {
return {
...step,
id: v4(),
name: `${step.name} (Duplicate)`,
nextStepIds: [],
position: duplicatedStepPosition,
};
}
}
}
async createEmptyNodeForIteratorStep({
iteratorStepId,
workflowVersionId,
workspaceId,
}: {
iteratorStepId: string;
workflowVersionId: string;
workspaceId: string;
}): Promise<WorkflowAction> {
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,
);
}
const existingSteps = workflowVersion.steps ?? [];
const emptyNodeStep: WorkflowEmptyAction = {
id: v4(),
name: 'Empty Node',
type: WorkflowActionType.EMPTY,
valid: true,
nextStepIds: [iteratorStepId],
settings: {
...BASE_STEP_DEFINITION,
input: {},
},
};
await workflowVersionRepository.update(workflowVersion.id, {
steps: [...existingSteps, emptyNodeStep],
});
return emptyNodeStep;
}
async createDraftStep({
step,
workspaceId,
}: {
step: WorkflowAction;
workspaceId: string;
}): Promise<WorkflowAction> {
switch (step.type) {
case WorkflowActionType.CODE: {
await this.serverlessFunctionService.createDraftFromPublishedVersion({
id: step.settings.input.serverlessFunctionId,
version: step.settings.input.serverlessFunctionVersion,
workspaceId,
});
return {
...step,
settings: {
...step.settings,
input: {
...step.settings.input,
serverlessFunctionVersion: 'draft',
},
},
};
}
default: {
return step;
}
}
}
}
@@ -7,6 +7,7 @@ import { ObjectMetadataEntity } from 'src/engine/metadata-modules/object-metadat
import { ServerlessFunctionModule } from 'src/engine/metadata-modules/serverless-function/serverless-function.module';
import { WorkflowCommonModule } from 'src/modules/workflow/common/workflow-common.module';
import { WorkflowSchemaModule } from 'src/modules/workflow/workflow-builder/workflow-schema/workflow-schema.module';
import { WorkflowVersionStepOperationsWorkspaceService } from 'src/modules/workflow/workflow-builder/workflow-version-step/workflow-version-step-operations.workspace-service';
import { WorkflowVersionStepWorkspaceService } from 'src/modules/workflow/workflow-builder/workflow-version-step/workflow-version-step.workspace-service';
import { WorkflowRunModule } from 'src/modules/workflow/workflow-runner/workflow-run/workflow-run.module';
import { WorkflowRunnerModule } from 'src/modules/workflow/workflow-runner/workflow-runner.module';
@@ -20,7 +21,10 @@ import { WorkflowRunnerModule } from 'src/modules/workflow/workflow-runner/workf
WorkflowCommonModule,
NestjsQueryTypeOrmModule.forFeature([ObjectMetadataEntity, AgentEntity]),
],
providers: [WorkflowVersionStepWorkspaceService],
providers: [
WorkflowVersionStepWorkspaceService,
WorkflowVersionStepOperationsWorkspaceService,
],
exports: [WorkflowVersionStepWorkspaceService],
})
export class WorkflowVersionStepModule {}
@@ -1,20 +1,11 @@
import { Injectable } from '@nestjs/common';
import { InjectRepository } from '@nestjs/typeorm';
import { t } from '@lingui/core/macro';
import { FieldMetadataType } from 'twenty-shared/types';
import { isDefined, isValidUuid } from 'twenty-shared/utils';
import { isDefined } from 'twenty-shared/utils';
import { StepStatus, TRIGGER_STEP_ID } from 'twenty-shared/workflow';
import { Repository } from 'typeorm';
import { v4 } from 'uuid';
import { BASE_TYPESCRIPT_PROJECT_INPUT_SCHEMA } from 'src/engine/core-modules/serverless/drivers/constants/base-typescript-project-input-schema';
import { type CreateWorkflowVersionStepInput } from 'src/engine/core-modules/workflow/dtos/create-workflow-version-step-input.dto';
import { type WorkflowStepPositionInput } from 'src/engine/core-modules/workflow/dtos/update-workflow-step-position-input.dto';
import { type WorkflowVersionStepChangesDTO } from 'src/engine/core-modules/workflow/dtos/workflow-version-step-changes.dto';
import { AgentEntity } from 'src/engine/metadata-modules/agent/agent.entity';
import { ObjectMetadataEntity } from 'src/engine/metadata-modules/object-metadata/object-metadata.entity';
import { ServerlessFunctionService } from 'src/engine/metadata-modules/serverless-function/serverless-function.service';
import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager';
import {
WorkflowVersionStepException,
@@ -22,48 +13,24 @@ import {
} from 'src/modules/workflow/common/exceptions/workflow-version-step.exception';
import { type WorkflowVersionWorkspaceEntity } from 'src/modules/workflow/common/standard-objects/workflow-version.workspace-entity';
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 { WorkflowSchemaWorkspaceService } from 'src/modules/workflow/workflow-builder/workflow-schema/workflow-schema.workspace-service';
import { insertStep } from 'src/modules/workflow/workflow-builder/workflow-version-step/utils/insert-step';
import { removeStep } from 'src/modules/workflow/workflow-builder/workflow-version-step/utils/remove-step';
import { type BaseWorkflowActionSettings } from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action-settings.type';
import {
type WorkflowAction,
WorkflowActionType,
WorkflowEmptyAction,
type WorkflowFormAction,
} from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type';
import { WorkflowVersionStepOperationsWorkspaceService } from 'src/modules/workflow/workflow-builder/workflow-version-step/workflow-version-step-operations.workspace-service';
import { isWorkflowFormAction } from 'src/modules/workflow/workflow-executor/workflow-actions/form/guards/is-workflow-form-action.guard';
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';
import { WorkflowRunnerWorkspaceService } from 'src/modules/workflow/workflow-runner/workspace-services/workflow-runner.workspace-service';
const BASE_STEP_DEFINITION: BaseWorkflowActionSettings = {
outputSchema: {},
errorHandlingOptions: {
continueOnFailure: {
value: false,
},
retryOnFailure: {
value: false,
},
},
};
const DUPLICATED_STEP_POSITION_OFFSET = 50;
@Injectable()
export class WorkflowVersionStepWorkspaceService {
constructor(
private readonly twentyORMGlobalManager: TwentyORMGlobalManager,
private readonly workflowSchemaWorkspaceService: WorkflowSchemaWorkspaceService,
private readonly serverlessFunctionService: ServerlessFunctionService,
@InjectRepository(AgentEntity)
private readonly agentRepository: Repository<AgentEntity>,
@InjectRepository(ObjectMetadataEntity)
private readonly objectMetadataRepository: Repository<ObjectMetadataEntity>,
private readonly workflowRunWorkspaceService: WorkflowRunWorkspaceService,
private readonly workflowRunnerWorkspaceService: WorkflowRunnerWorkspaceService,
private readonly workflowCommonWorkspaceService: WorkflowCommonWorkspaceService,
private readonly workflowVersionStepOperationsWorkspaceService: WorkflowVersionStepOperationsWorkspaceService,
) {}
async createWorkflowVersionStep({
@@ -82,17 +49,21 @@ export class WorkflowVersionStepWorkspaceService {
parentStepConnectionOptions,
} = input;
const newStep = await this.runStepCreationSideEffectsAndBuildStep({
type: stepType,
workspaceId,
position,
workflowVersionId,
});
const newStep =
await this.workflowVersionStepOperationsWorkspaceService.runStepCreationSideEffectsAndBuildStep(
{
type: stepType,
workspaceId,
position,
workflowVersionId,
},
);
const enrichedNewStep = await this.enrichOutputSchema({
step: newStep,
workspaceId,
});
const enrichedNewStep =
await this.workflowSchemaWorkspaceService.enrichOutputSchema({
step: newStep,
workspaceId,
});
const workflowVersionRepository =
await this.twentyORMGlobalManager.getRepositoryForWorkspace<WorkflowVersionWorkspaceEntity>(
@@ -179,10 +150,11 @@ export class WorkflowVersionStepWorkspaceService {
);
}
const enrichedNewStep = await this.enrichOutputSchema({
step,
workspaceId,
});
const enrichedNewStep =
await this.workflowSchemaWorkspaceService.enrichOutputSchema({
step,
workspaceId,
});
const updatedSteps = workflowVersion.steps.map((existingStep) => {
if (existingStep.id === step.id) {
@@ -276,10 +248,12 @@ export class WorkflowVersionStepWorkspaceService {
await Promise.all(
removedSteps.map((step) =>
this.runWorkflowVersionStepDeletionSideEffects({
step,
workspaceId,
}),
this.workflowVersionStepOperationsWorkspaceService.runWorkflowVersionStepDeletionSideEffects(
{
step,
workspaceId,
},
),
),
);
@@ -332,10 +306,13 @@ export class WorkflowVersionStepWorkspaceService {
);
}
const duplicatedStep = await this.createStepForDuplicate({
step: stepToDuplicate,
workspaceId,
});
const duplicatedStep =
await this.workflowVersionStepOperationsWorkspaceService.createStepForDuplicate(
{
step: stepToDuplicate,
workspaceId,
},
);
const { updatedSteps, updatedInsertedStep, updatedTrigger } = insertStep({
existingSteps: workflowVersion.steps ?? [],
@@ -383,7 +360,7 @@ export class WorkflowVersionStepWorkspaceService {
);
}
if (step.type !== WorkflowActionType.FORM) {
if (!isWorkflowFormAction(step)) {
throw new WorkflowVersionStepException(
'Step is not a form',
WorkflowVersionStepExceptionCode.INVALID_REQUEST,
@@ -393,11 +370,14 @@ export class WorkflowVersionStepWorkspaceService {
);
}
const enrichedResponse = await this.enrichFormStepResponse({
workspaceId,
step,
response,
});
const enrichedResponse =
await this.workflowVersionStepOperationsWorkspaceService.enrichFormStepResponse(
{
workspaceId,
step,
response,
},
);
await this.workflowRunWorkspaceService.updateWorkflowRunStepInfo({
stepId,
@@ -423,509 +403,9 @@ export class WorkflowVersionStepWorkspaceService {
step: WorkflowAction;
workspaceId: string;
}): Promise<WorkflowAction> {
switch (step.type) {
case WorkflowActionType.CODE: {
await this.serverlessFunctionService.createDraftFromPublishedVersion({
id: step.settings.input.serverlessFunctionId,
version: step.settings.input.serverlessFunctionVersion,
workspaceId,
});
return {
...step,
settings: {
...step.settings,
input: {
...step.settings.input,
serverlessFunctionVersion: 'draft',
},
},
};
}
default: {
return step;
}
}
}
private async enrichOutputSchema({
step,
workspaceId,
}: {
step: WorkflowAction;
workspaceId: string;
}): Promise<WorkflowAction> {
// We don't enrich on the fly for code and HTTP request workflow actions.
// For code actions, OutputSchema is computed and updated when testing the serverless function.
// For HTTP requests and AI agent, OutputSchema is determined by the expamle response input
if (
[
WorkflowActionType.CODE,
WorkflowActionType.HTTP_REQUEST,
WorkflowActionType.AI_AGENT,
].includes(step.type)
) {
return step;
}
const result = { ...step };
const outputSchema =
await this.workflowSchemaWorkspaceService.computeStepOutputSchema({
step,
workspaceId,
});
result.settings = {
...result.settings,
outputSchema: outputSchema || {},
};
return result;
}
private async runWorkflowVersionStepDeletionSideEffects({
step,
workspaceId,
}: {
step: WorkflowAction;
workspaceId: string;
}) {
switch (step.type) {
case WorkflowActionType.CODE: {
if (
!(await this.serverlessFunctionService.hasServerlessFunctionPublishedVersion(
step.settings.input.serverlessFunctionId,
))
) {
await this.serverlessFunctionService.deleteOneServerlessFunction({
id: step.settings.input.serverlessFunctionId,
workspaceId,
softDelete: false,
});
}
break;
}
case WorkflowActionType.AI_AGENT: {
if (!isDefined(step.settings.input.agentId)) {
break;
}
const agent = await this.agentRepository.findOne({
where: { id: step.settings.input.agentId, workspaceId },
});
if (isDefined(agent)) {
await this.agentRepository.delete({ id: agent.id, workspaceId });
}
break;
}
}
}
private async runStepCreationSideEffectsAndBuildStep({
type,
workspaceId,
position,
workflowVersionId,
}: {
type: WorkflowActionType;
workspaceId: string;
position?: WorkflowStepPositionInput;
workflowVersionId: string;
}): Promise<WorkflowAction> {
const newStepId = v4();
const baseStep = {
id: newStepId,
position,
valid: false,
nextStepIds: [],
};
switch (type) {
case WorkflowActionType.CODE: {
const newServerlessFunction =
await this.serverlessFunctionService.createOneServerlessFunction(
{
name: 'A Serverless Function Code Workflow Step',
description: '',
},
workspaceId,
);
if (!isDefined(newServerlessFunction)) {
throw new WorkflowVersionStepException(
'Fail to create Code Step',
WorkflowVersionStepExceptionCode.CODE_STEP_FAILURE,
);
}
return {
...baseStep,
name: 'Code - Serverless Function',
type: WorkflowActionType.CODE,
settings: {
...BASE_STEP_DEFINITION,
outputSchema: {
link: {
isLeaf: true,
icon: 'IconVariable',
tab: 'test',
label: 'Generate Function Output',
},
_outputSchemaType: 'LINK',
},
input: {
serverlessFunctionId: newServerlessFunction.id,
serverlessFunctionVersion: 'draft',
serverlessFunctionInput: BASE_TYPESCRIPT_PROJECT_INPUT_SCHEMA,
},
},
};
}
case WorkflowActionType.SEND_EMAIL: {
return {
...baseStep,
name: 'Send Email',
type: WorkflowActionType.SEND_EMAIL,
settings: {
...BASE_STEP_DEFINITION,
input: {
connectedAccountId: '',
email: '',
subject: '',
body: '',
},
},
};
}
case WorkflowActionType.CREATE_RECORD: {
const activeObjectMetadataItem =
await this.objectMetadataRepository.findOne({
where: { workspaceId, isActive: true, isSystem: false },
});
return {
...baseStep,
name: 'Create Record',
type: WorkflowActionType.CREATE_RECORD,
settings: {
...BASE_STEP_DEFINITION,
input: {
objectName: activeObjectMetadataItem?.nameSingular || '',
objectRecord: {},
},
},
};
}
case WorkflowActionType.UPDATE_RECORD: {
const activeObjectMetadataItem =
await this.objectMetadataRepository.findOne({
where: { workspaceId, isActive: true, isSystem: false },
});
return {
...baseStep,
name: 'Update Record',
type: WorkflowActionType.UPDATE_RECORD,
settings: {
...BASE_STEP_DEFINITION,
input: {
objectName: activeObjectMetadataItem?.nameSingular || '',
objectRecord: {},
objectRecordId: '',
fieldsToUpdate: [],
},
},
};
}
case WorkflowActionType.DELETE_RECORD: {
const activeObjectMetadataItem =
await this.objectMetadataRepository.findOne({
where: { workspaceId, isActive: true, isSystem: false },
});
return {
...baseStep,
name: 'Delete Record',
type: WorkflowActionType.DELETE_RECORD,
settings: {
...BASE_STEP_DEFINITION,
input: {
objectName: activeObjectMetadataItem?.nameSingular || '',
objectRecordId: '',
},
},
};
}
case WorkflowActionType.FIND_RECORDS: {
const activeObjectMetadataItem =
await this.objectMetadataRepository.findOne({
where: { workspaceId, isActive: true, isSystem: false },
});
return {
...baseStep,
name: 'Search Records',
type: WorkflowActionType.FIND_RECORDS,
settings: {
...BASE_STEP_DEFINITION,
input: {
objectName: activeObjectMetadataItem?.nameSingular || '',
limit: 1,
},
},
};
}
case WorkflowActionType.FORM: {
return {
...baseStep,
name: 'Form',
type: WorkflowActionType.FORM,
settings: {
...BASE_STEP_DEFINITION,
input: [],
},
};
}
case WorkflowActionType.FILTER: {
return {
...baseStep,
name: 'Filter',
type: WorkflowActionType.FILTER,
settings: {
...BASE_STEP_DEFINITION,
input: {
stepFilterGroups: [],
stepFilters: [],
},
},
};
}
case WorkflowActionType.HTTP_REQUEST: {
return {
...baseStep,
name: 'HTTP Request',
type: WorkflowActionType.HTTP_REQUEST,
settings: {
...BASE_STEP_DEFINITION,
input: {
url: '',
method: 'GET',
headers: {},
body: {},
},
},
};
}
case WorkflowActionType.AI_AGENT: {
return {
...baseStep,
name: 'AI Agent',
type: WorkflowActionType.AI_AGENT,
settings: {
...BASE_STEP_DEFINITION,
input: {
agentId: '',
prompt: '',
},
},
};
}
case WorkflowActionType.ITERATOR: {
const emptyNodeStep = await this.createEmptyNodeForIteratorStep({
iteratorStepId: baseStep.id,
workflowVersionId,
workspaceId,
});
return {
...baseStep,
name: 'Iterator',
type: WorkflowActionType.ITERATOR,
settings: {
...BASE_STEP_DEFINITION,
input: {
items: [],
initialLoopStepIds: [emptyNodeStep.id],
},
},
};
}
default:
throw new WorkflowVersionStepException(
`WorkflowActionType '${type}' unknown`,
WorkflowVersionStepExceptionCode.INVALID_REQUEST,
);
}
}
private async enrichFormStepResponse({
workspaceId,
step,
response,
}: {
workspaceId: string;
step: WorkflowFormAction;
response: object;
}) {
const responseKeys = Object.keys(response);
const enrichedResponses = await Promise.all(
responseKeys.map(async (key) => {
// @ts-expect-error legacy noImplicitAny
if (!isDefined(response[key])) {
// @ts-expect-error legacy noImplicitAny
return { key, value: response[key] };
}
const field = step.settings.input.find((field) => field.name === key);
if (
field?.type === 'RECORD' &&
field?.settings?.objectName &&
// @ts-expect-error legacy noImplicitAny
isDefined(response[key].id) &&
// @ts-expect-error legacy noImplicitAny
isValidUuid(response[key].id)
) {
const objectMetadataInfo =
await this.workflowCommonWorkspaceService.getObjectMetadataItemWithFieldsMaps(
field.settings.objectName,
workspaceId,
);
const relationFieldsNames = Object.values(
objectMetadataInfo.objectMetadataItemWithFieldsMaps.fieldsById,
)
.filter((field) => field.type === FieldMetadataType.RELATION)
.map((field) => field.name);
const repository =
await this.twentyORMGlobalManager.getRepositoryForWorkspace(
workspaceId,
field.settings.objectName,
{ shouldBypassPermissionChecks: true },
);
const record = await repository.findOne({
// @ts-expect-error legacy noImplicitAny
where: { id: response[key].id },
relations: relationFieldsNames,
});
return { key, value: record };
} else {
// @ts-expect-error legacy noImplicitAny
return { key, value: response[key] };
}
}),
);
return enrichedResponses.reduce((acc, { key, value }) => {
// @ts-expect-error legacy noImplicitAny
acc[key] = value;
return acc;
}, {});
}
private async createStepForDuplicate({
step,
workspaceId,
}: {
step: WorkflowAction;
workspaceId: string;
}): Promise<WorkflowAction> {
const duplicatedStepPosition = {
x: (step.position?.x ?? 0) + DUPLICATED_STEP_POSITION_OFFSET,
y: (step.position?.y ?? 0) + DUPLICATED_STEP_POSITION_OFFSET,
};
switch (step.type) {
case WorkflowActionType.CODE: {
const newServerlessFunction =
await this.serverlessFunctionService.duplicateServerlessFunction({
id: step.settings.input.serverlessFunctionId,
version: step.settings.input.serverlessFunctionVersion,
workspaceId,
});
return {
...step,
id: v4(),
name: `${step.name} (Duplicate)`,
nextStepIds: [],
position: duplicatedStepPosition,
settings: {
...step.settings,
input: {
...step.settings.input,
serverlessFunctionId: newServerlessFunction.id,
serverlessFunctionVersion: 'draft',
},
},
};
}
default: {
return {
...step,
id: v4(),
name: `${step.name} (Duplicate)`,
nextStepIds: [],
position: duplicatedStepPosition,
};
}
}
}
private async createEmptyNodeForIteratorStep({
iteratorStepId,
workflowVersionId,
workspaceId,
}: {
iteratorStepId: string;
workflowVersionId: string;
workspaceId: string;
}): Promise<WorkflowAction> {
const workflowVersionRepository =
await this.twentyORMGlobalManager.getRepositoryForWorkspace<WorkflowVersionWorkspaceEntity>(
workspaceId,
'workflowVersion',
{ shouldBypassPermissionChecks: true },
);
const workflowVersion = await workflowVersionRepository.findOne({
where: {
id: workflowVersionId,
},
return this.workflowVersionStepOperationsWorkspaceService.createDraftStep({
step,
workspaceId,
});
if (!isDefined(workflowVersion)) {
throw new WorkflowVersionStepException(
'WorkflowVersion not found',
WorkflowVersionStepExceptionCode.NOT_FOUND,
);
}
const existingSteps = workflowVersion.steps ?? [];
const emptyNodeStep: WorkflowEmptyAction = {
id: v4(),
name: 'Empty Node',
type: WorkflowActionType.EMPTY,
valid: true,
nextStepIds: [iteratorStepId],
settings: {
...BASE_STEP_DEFINITION,
input: {},
},
};
await workflowVersionRepository.update(workflowVersion.id, {
steps: [...existingSteps, emptyNodeStep],
});
return emptyNodeStep;
}
}
@@ -73,6 +73,7 @@ export {
IconColorSwatch,
IconMessageCircle as IconComment,
IconCopy,
IconCopyPlus,
IconCreativeCommonsSa,
IconCreditCard,
IconCsv,
+1
View File
@@ -135,6 +135,7 @@ export {
IconColorSwatch,
IconComment,
IconCopy,
IconCopyPlus,
IconCreativeCommonsSa,
IconCreditCard,
IconCsv,