diff --git a/packages/twenty-front/src/modules/workflow/workflow-diagram/components/WorkflowDiagramEffect.tsx b/packages/twenty-front/src/modules/workflow/workflow-diagram/components/WorkflowDiagramEffect.tsx index 473e6619bf..c39c5efd75 100644 --- a/packages/twenty-front/src/modules/workflow/workflow-diagram/components/WorkflowDiagramEffect.tsx +++ b/packages/twenty-front/src/modules/workflow/workflow-diagram/components/WorkflowDiagramEffect.tsx @@ -40,6 +40,10 @@ export const WorkflowDiagramEffect = ({ FeatureFlagKey.IS_WORKFLOW_FILTERING_ENABLED, ); + const isWorkflowBranchEnabled = useIsFeatureEnabled( + FeatureFlagKey.IS_WORKFLOW_BRANCH_ENABLED, + ); + const computeAndMergeNewWorkflowDiagram = useRecoilCallback( ({ snapshot, set }) => { return (currentVersion: WorkflowVersion) => { @@ -51,6 +55,7 @@ export const WorkflowDiagramEffect = ({ const nextWorkflowDiagram = getWorkflowVersionDiagram({ workflowVersion: currentVersion, isWorkflowFilteringEnabled, + isWorkflowBranchEnabled, isEditable: true, }); @@ -87,6 +92,7 @@ export const WorkflowDiagramEffect = ({ [ workflowDiagramState, isWorkflowFilteringEnabled, + isWorkflowBranchEnabled, workflowLastCreatedStepIdState, ], ); diff --git a/packages/twenty-front/src/modules/workflow/workflow-diagram/components/WorkflowVersionVisualizerEffect.tsx b/packages/twenty-front/src/modules/workflow/workflow-diagram/components/WorkflowVersionVisualizerEffect.tsx index 053b47ce96..641927230b 100644 --- a/packages/twenty-front/src/modules/workflow/workflow-diagram/components/WorkflowVersionVisualizerEffect.tsx +++ b/packages/twenty-front/src/modules/workflow/workflow-diagram/components/WorkflowVersionVisualizerEffect.tsx @@ -35,6 +35,10 @@ export const WorkflowVersionVisualizerEffect = ({ FeatureFlagKey.IS_WORKFLOW_FILTERING_ENABLED, ); + const isWorkflowBranchEnabled = useIsFeatureEnabled( + FeatureFlagKey.IS_WORKFLOW_BRANCH_ENABLED, + ); + useEffect(() => { if (!isDefined(workflowVersion)) { setFlow(undefined); @@ -67,11 +71,17 @@ export const WorkflowVersionVisualizerEffect = ({ const nextWorkflowDiagram = getWorkflowVersionDiagram({ workflowVersion, isWorkflowFilteringEnabled, + isWorkflowBranchEnabled, isEditable: false, }); setWorkflowDiagram(nextWorkflowDiagram); - }, [isWorkflowFilteringEnabled, setWorkflowDiagram, workflowVersion]); + }, [ + isWorkflowBranchEnabled, + isWorkflowFilteringEnabled, + setWorkflowDiagram, + workflowVersion, + ]); useEffect(() => { if (!isDefined(workflowVersion)) { diff --git a/packages/twenty-front/src/modules/workflow/workflow-diagram/utils/generateWorkflowDiagram.ts b/packages/twenty-front/src/modules/workflow/workflow-diagram/utils/generateWorkflowDiagram.ts index a42ddfaa82..749faea6d4 100644 --- a/packages/twenty-front/src/modules/workflow/workflow-diagram/utils/generateWorkflowDiagram.ts +++ b/packages/twenty-front/src/modules/workflow/workflow-diagram/utils/generateWorkflowDiagram.ts @@ -21,10 +21,12 @@ export const generateWorkflowDiagram = ({ trigger, steps, defaultEdgeType, + isWorkflowBranchEnabled = false, }: { trigger: WorkflowTrigger | undefined; steps: Array; defaultEdgeType: WorkflowDiagramEdgeType; + isWorkflowBranchEnabled?: boolean; }): WorkflowDiagram => { const nodes: Array = []; const edges: Array = []; @@ -34,7 +36,8 @@ export const generateWorkflowDiagram = ({ } else { nodes.push(WORKFLOW_DIAGRAM_EMPTY_TRIGGER_NODE_DEFINITION); - const triggerNextStepIds = isDefined(steps) ? getRootStepIds(steps) : []; + const triggerNextStepIds = + isDefined(steps) && !isWorkflowBranchEnabled ? getRootStepIds(steps) : []; triggerNextStepIds.forEach((stepId) => { edges.push({ diff --git a/packages/twenty-front/src/modules/workflow/workflow-diagram/utils/getWorkflowVersionDiagram.ts b/packages/twenty-front/src/modules/workflow/workflow-diagram/utils/getWorkflowVersionDiagram.ts index 62c54bd9ec..e3f95b9cb4 100644 --- a/packages/twenty-front/src/modules/workflow/workflow-diagram/utils/getWorkflowVersionDiagram.ts +++ b/packages/twenty-front/src/modules/workflow/workflow-diagram/utils/getWorkflowVersionDiagram.ts @@ -31,10 +31,12 @@ const getEdgeTypeToCreateByDefault = ({ export const getWorkflowVersionDiagram = ({ workflowVersion, isWorkflowFilteringEnabled, + isWorkflowBranchEnabled, isEditable, }: { workflowVersion: WorkflowVersion | undefined; isWorkflowFilteringEnabled: boolean; + isWorkflowBranchEnabled?: boolean; isEditable: boolean; }): WorkflowDiagram => { if (!isDefined(workflowVersion)) { @@ -48,6 +50,7 @@ export const getWorkflowVersionDiagram = ({ isWorkflowFilteringEnabled, isEditable, }), + isWorkflowBranchEnabled, }); return transformFilterNodesAsEdges({ diff --git a/packages/twenty-server/src/modules/workflow/workflow-builder/workflow-version-step/utils/__tests__/remove-step.spec.ts b/packages/twenty-server/src/modules/workflow/workflow-builder/workflow-version-step/utils/__tests__/remove-step.spec.ts index be3e91af20..f303f5b6b1 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-builder/workflow-version-step/utils/__tests__/remove-step.spec.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-builder/workflow-version-step/utils/__tests__/remove-step.spec.ts @@ -242,4 +242,15 @@ describe('removeStep', () => { }); expect(result.updatedSteps).toEqual([]); }); + + it('should remove trigger if steps are null', () => { + const result = removeStep({ + existingTrigger: { ...mockTrigger, nextStepIds: [] }, + existingSteps: null, + stepIdToDelete: 'trigger', + }); + + expect(result.updatedTrigger).toEqual(null); + expect(result.updatedSteps).toEqual([]); + }); }); diff --git a/packages/twenty-server/src/modules/workflow/workflow-builder/workflow-version-step/utils/remove-step.ts b/packages/twenty-server/src/modules/workflow/workflow-builder/workflow-version-step/utils/remove-step.ts index f166cf57a4..8a5a98ef08 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-builder/workflow-version-step/utils/remove-step.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-builder/workflow-version-step/utils/remove-step.ts @@ -35,7 +35,7 @@ const removeOneStep = ({ stepToDeleteChildrenIds, }: { existingTrigger: WorkflowTrigger | null; - existingSteps: WorkflowAction[]; + existingSteps: WorkflowAction[] | null; stepIdToDelete: string; stepToDeleteChildrenIds?: string[]; }): { @@ -43,22 +43,23 @@ const removeOneStep = ({ updatedTrigger: WorkflowTrigger | null; removedStepIds: string[]; } => { - const updatedSteps = existingSteps - .filter((step) => step.id !== stepIdToDelete) - .map((step) => { - if (step.nextStepIds?.includes(stepIdToDelete)) { - return { - ...step, - nextStepIds: computeUpdatedNextStepIds({ - existingNextStepIds: step.nextStepIds, - stepIdToRemove: stepIdToDelete, - stepToDeleteChildrenIds: stepToDeleteChildrenIds, - }), - }; - } + const updatedSteps = + existingSteps + ?.filter((step) => step.id !== stepIdToDelete) + .map((step) => { + if (step.nextStepIds?.includes(stepIdToDelete)) { + return { + ...step, + nextStepIds: computeUpdatedNextStepIds({ + existingNextStepIds: step.nextStepIds, + stepIdToRemove: stepIdToDelete, + stepToDeleteChildrenIds: stepToDeleteChildrenIds, + }), + }; + } - return step; - }); + return step; + }) ?? []; let updatedTrigger = existingTrigger; @@ -89,7 +90,7 @@ const removeRegularStep = ({ stepToDeleteChildrenIds, }: { existingTrigger: WorkflowTrigger | null; - existingSteps: WorkflowAction[]; + existingSteps: WorkflowAction[] | null; stepIdToDelete: string; stepToDeleteChildrenIds?: string[]; }): { @@ -105,7 +106,7 @@ const removeRegularStep = ({ }); for (const stepId of stepToDeleteChildrenIds ?? []) { - const step = existingSteps.find((step) => step.id === stepId); + const step = existingSteps?.find((step) => step.id === stepId); if (step?.type === WorkflowActionType.FILTER) { const { @@ -164,19 +165,18 @@ const removeTrigger = ({ existingSteps, triggerChildrenIds, }: { - existingSteps: WorkflowAction[]; + existingSteps: WorkflowAction[] | null; triggerChildrenIds?: string[]; }) => { const stepIdsToRemove = triggerChildrenIds?.filter((id) => { - const step = existingSteps.find((step) => step.id === id); + const step = existingSteps?.find((step) => step.id === id); return step?.type === WorkflowActionType.FILTER; }) ?? []; - const updatedSteps = existingSteps.filter( - (step) => !stepIdsToRemove.includes(step.id), - ); + const updatedSteps = + existingSteps?.filter((step) => !stepIdsToRemove.includes(step.id)) ?? []; return { updatedSteps, @@ -192,7 +192,7 @@ export const removeStep = ({ stepToDeleteChildrenIds, }: { existingTrigger: WorkflowTrigger | null; - existingSteps: WorkflowAction[]; + existingSteps: WorkflowAction[] | null; stepIdToDelete: string; stepToDeleteChildrenIds?: string[]; }) => { diff --git a/packages/twenty-server/src/modules/workflow/workflow-builder/workflow-version-step/workflow-version-step.workspace-service.ts b/packages/twenty-server/src/modules/workflow/workflow-builder/workflow-version-step/workflow-version-step.workspace-service.ts index 0e0ea2eb47..644d977823 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-builder/workflow-version-step/workflow-version-step.workspace-service.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-builder/workflow-version-step/workflow-version-step.workspace-service.ts @@ -218,22 +218,23 @@ export class WorkflowVersionStepWorkspaceService { assertWorkflowVersionIsDraft(workflowVersion); - if (!isDefined(workflowVersion.steps)) { + const existingTrigger = workflowVersion.trigger; + + const isDeletingTrigger = + stepIdToDelete === 'trigger' && isDefined(existingTrigger); + + if (!isDeletingTrigger && !isDefined(workflowVersion.steps)) { throw new WorkflowVersionStepException( "Can't delete step from undefined steps", WorkflowVersionStepExceptionCode.UNDEFINED, ); } - const existingTrigger = workflowVersion.trigger; - const isDeletingTrigger = - stepIdToDelete === 'trigger' && isDefined(existingTrigger); - - const stepToDelete = workflowVersion.steps.find( + const stepToDelete = workflowVersion.steps?.find( (step) => step.id === stepIdToDelete, ); - if (!isDefined(stepToDelete) && !isDeletingTrigger) { + if (!isDeletingTrigger && !isDefined(stepToDelete)) { throw new WorkflowVersionStepException( "Can't delete not existing step", WorkflowVersionStepExceptionCode.NOT_FOUND, @@ -256,9 +257,10 @@ export class WorkflowVersionStepWorkspaceService { trigger: updatedTrigger, }); - const removedSteps = workflowVersion.steps.filter((step) => - removedStepIds.includes(step.id), - ); + const removedSteps = + workflowVersion.steps?.filter((step) => + removedStepIds.includes(step.id), + ) ?? []; await Promise.all( removedSteps.map((step) => diff --git a/packages/twenty-server/src/modules/workflow/workflow-runner/jobs/run-workflow.job.ts b/packages/twenty-server/src/modules/workflow/workflow-runner/jobs/run-workflow.job.ts index a13a403e59..f6b4fde97a 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-runner/jobs/run-workflow.job.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-runner/jobs/run-workflow.job.ts @@ -20,6 +20,8 @@ import { getRootSteps } from 'src/modules/workflow/workflow-runner/utils/get-roo import { WorkflowRunQueueWorkspaceService } from 'src/modules/workflow/workflow-runner/workflow-run-queue/workspace-services/workflow-run-queue.workspace-service'; import { WorkflowRunWorkspaceService } from 'src/modules/workflow/workflow-runner/workflow-run/workflow-run.workspace-service'; import { WorkflowTriggerType } from 'src/modules/workflow/workflow-trigger/types/workflow-trigger.type'; +import { FeatureFlagKey } from 'src/engine/core-modules/feature-flag/enums/feature-flag-key.enum'; +import { FeatureFlagService } from 'src/engine/core-modules/feature-flag/services/feature-flag.service'; export type RunWorkflowJobData = { workspaceId: string; @@ -37,6 +39,7 @@ export class RunWorkflowJob { private readonly twentyConfigService: TwentyConfigService, private readonly metricsService: MetricsService, private readonly workflowRunQueueWorkspaceService: WorkflowRunQueueWorkspaceService, + private readonly featureFlagService: FeatureFlagService, ) {} @Process(RunWorkflowJob.name) @@ -98,6 +101,8 @@ export class RunWorkflowJob { ); } + await this.throttleExecution(workflowVersion.workflowId); + await this.incrementTriggerMetrics({ workflowRunId, triggerType: workflowVersion.trigger.type, @@ -108,13 +113,20 @@ export class RunWorkflowJob { workspaceId, }); - await this.throttleExecution(workflowVersion.workflowId); - const rootSteps = getRootSteps(workflowVersion.steps); + const isWorkflowBranchEnabled = + await this.featureFlagService.isFeatureEnabled( + FeatureFlagKey.IS_WORKFLOW_BRANCH_ENABLED, + workspaceId, + ); + + const stepIds = isWorkflowBranchEnabled + ? (workflowVersion.trigger.nextStepIds ?? []) + : (rootSteps.map((step) => step.id) ?? []); + await this.workflowExecutorWorkspaceService.executeFromSteps({ - stepIds: - workflowVersion.trigger.nextStepIds ?? rootSteps.map((step) => step.id), + stepIds, workflowRunId, workspaceId, }); @@ -136,10 +148,7 @@ export class RunWorkflowJob { }); if (workflowRun.status !== WorkflowRunStatus.RUNNING) { - throw new WorkflowRunException( - 'Workflow is not running', - WorkflowRunExceptionCode.WORKFLOW_RUN_INVALID, - ); + return; } const lastExecutedStep = workflowRun.state?.flow?.steps?.find( diff --git a/packages/twenty-server/src/modules/workflow/workflow-runner/workflow-runner.module.ts b/packages/twenty-server/src/modules/workflow/workflow-runner/workflow-runner.module.ts index 5803fed61e..34296f4595 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-runner/workflow-runner.module.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-runner/workflow-runner.module.ts @@ -9,6 +9,7 @@ import { RunWorkflowJob } from 'src/modules/workflow/workflow-runner/jobs/run-wo import { WorkflowRunQueueModule } from 'src/modules/workflow/workflow-runner/workflow-run-queue/workflow-run-queue.module'; import { WorkflowRunModule } from 'src/modules/workflow/workflow-runner/workflow-run/workflow-run.module'; import { WorkflowRunnerWorkspaceService } from 'src/modules/workflow/workflow-runner/workspace-services/workflow-runner.workspace-service'; +import { FeatureFlagModule } from 'src/engine/core-modules/feature-flag/feature-flag.module'; @Module({ imports: [ @@ -19,6 +20,7 @@ import { WorkflowRunnerWorkspaceService } from 'src/modules/workflow/workflow-ru WorkflowRunModule, MetricsModule, WorkflowRunQueueModule, + FeatureFlagModule, ], providers: [WorkflowRunnerWorkspaceService, RunWorkflowJob], exports: [WorkflowRunnerWorkspaceService],