diff --git a/packages/twenty-server/src/engine/core-modules/workflow/services/workflow-version-core-sync.service.ts b/packages/twenty-server/src/engine/core-modules/workflow/services/workflow-version-core-sync.service.ts index a68cd40cc4..fefc02e575 100644 --- a/packages/twenty-server/src/engine/core-modules/workflow/services/workflow-version-core-sync.service.ts +++ b/packages/twenty-server/src/engine/core-modules/workflow/services/workflow-version-core-sync.service.ts @@ -100,6 +100,15 @@ export class WorkflowVersionCoreSyncService { await this.invalidateAutomatedTriggerMaps(workspaceId); } + async findCoreVersionById( + workspaceId: string, + coreWorkflowVersionId: string, + ): Promise { + return this.coreWorkflowVersionRepository.findOne(workspaceId, { + where: { id: coreWorkflowVersionId }, + }); + } + async mirrorWorkflowVersionWrite({ workspaceId, entityManager, diff --git a/packages/twenty-server/src/modules/workflow/common/workspace-services/__tests__/workflow-common.workspace-service.spec.ts b/packages/twenty-server/src/modules/workflow/common/workspace-services/__tests__/workflow-common.workspace-service.spec.ts new file mode 100644 index 0000000000..2cc88f6d2f --- /dev/null +++ b/packages/twenty-server/src/modules/workflow/common/workspace-services/__tests__/workflow-common.workspace-service.spec.ts @@ -0,0 +1,140 @@ +import { FeatureFlagKey } from 'twenty-shared/types'; + +import { type FeatureFlagService } from 'src/engine/core-modules/feature-flag/services/feature-flag.service'; +import { type WorkflowVersionCoreSyncService } from 'src/engine/core-modules/workflow/services/workflow-version-core-sync.service'; +import { type GlobalWorkspaceOrmManager } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager'; +import { WorkflowCommonWorkspaceService } from 'src/modules/workflow/common/workspace-services/workflow-common.workspace-service'; + +const WORKSPACE_ID = 'workspace-1'; +const WORKFLOW_VERSION_ID = 'workflow-version-1'; +const CORE_WORKFLOW_VERSION_ID = 'core-workflow-version-1'; + +const workspaceTrigger = { name: 'workspace trigger', type: 'MANUAL' }; +const workspaceSteps = [{ id: 'workspace-step' }]; + +const coreTrigger = { name: 'core trigger', type: 'DATABASE_EVENT' }; +const coreSteps = [{ id: 'core-step' }]; + +const buildService = ({ + isCoreReadEnabled, + coreVersion = null, + coreWorkflowVersionId = CORE_WORKFLOW_VERSION_ID, +}: { + isCoreReadEnabled: boolean; + coreVersion?: unknown; + coreWorkflowVersionId?: string | null; +}) => { + const workspaceVersion = { + id: WORKFLOW_VERSION_ID, + workflowId: 'workflow-1', + name: 'Draft', + trigger: workspaceTrigger, + steps: workspaceSteps, + status: 'DRAFT', + coreWorkflowVersionId, + }; + + const isFeatureEnabled = jest.fn().mockResolvedValue(isCoreReadEnabled); + const findCoreVersionById = jest.fn().mockResolvedValue(coreVersion); + + const globalWorkspaceOrmManager = { + executeInWorkspaceContext: (fn: () => unknown) => fn(), + getRepository: jest.fn().mockResolvedValue({ + findOne: jest.fn().mockResolvedValue(workspaceVersion), + }), + } as unknown as GlobalWorkspaceOrmManager; + + const service = new WorkflowCommonWorkspaceService( + globalWorkspaceOrmManager, + undefined as unknown as never, + undefined as unknown as never, + undefined as unknown as never, + { findCoreVersionById } as unknown as WorkflowVersionCoreSyncService, + { isFeatureEnabled } as unknown as FeatureFlagService, + ); + + return { service, isFeatureEnabled, findCoreVersionById }; +}; + +describe('WorkflowCommonWorkspaceService', () => { + describe('getWorkflowVersionOrFail core read overlay', () => { + it('returns workspace content when the core read flag is off', async () => { + const { service, isFeatureEnabled, findCoreVersionById } = buildService({ + isCoreReadEnabled: false, + }); + + const result = await service.getWorkflowVersionOrFail({ + workspaceId: WORKSPACE_ID, + workflowVersionId: WORKFLOW_VERSION_ID, + }); + + expect(isFeatureEnabled).toHaveBeenCalledWith( + FeatureFlagKey.IS_WORKFLOW_VERSION_IN_CORE_ENABLED, + WORKSPACE_ID, + ); + expect(findCoreVersionById).not.toHaveBeenCalled(); + expect(result.trigger).toEqual(workspaceTrigger); + expect(result.steps).toEqual(workspaceSteps); + expect(result.status).toBe('DRAFT'); + }); + + it('overlays trigger, steps and status from core when the flag is on and the core row exists', async () => { + const { service, findCoreVersionById } = buildService({ + isCoreReadEnabled: true, + coreVersion: { + triggers: [coreTrigger], + steps: coreSteps, + status: 'ACTIVE', + }, + }); + + const result = await service.getWorkflowVersionOrFail({ + workspaceId: WORKSPACE_ID, + workflowVersionId: WORKFLOW_VERSION_ID, + }); + + expect(findCoreVersionById).toHaveBeenCalledWith( + WORKSPACE_ID, + CORE_WORKFLOW_VERSION_ID, + ); + expect(result.trigger).toEqual(coreTrigger); + expect(result.steps).toEqual(coreSteps); + expect(result.status).toBe('ACTIVE'); + + // identity stays from the workspace row: core has neither of these + expect(result.id).toBe(WORKFLOW_VERSION_ID); + expect(result.name).toBe('Draft'); + }); + + it('falls back to workspace content when the flag is on but the core row is missing', async () => { + const { service } = buildService({ + isCoreReadEnabled: true, + coreVersion: null, + }); + + const result = await service.getWorkflowVersionOrFail({ + workspaceId: WORKSPACE_ID, + workflowVersionId: WORKFLOW_VERSION_ID, + }); + + expect(result.trigger).toEqual(workspaceTrigger); + expect(result.steps).toEqual(workspaceSteps); + expect(result.status).toBe('DRAFT'); + }); + + it('skips the core read when the version has no soft-ref', async () => { + const { service, findCoreVersionById } = buildService({ + isCoreReadEnabled: true, + coreWorkflowVersionId: null, + }); + + const result = await service.getWorkflowVersionOrFail({ + workspaceId: WORKSPACE_ID, + workflowVersionId: WORKFLOW_VERSION_ID, + }); + + expect(findCoreVersionById).not.toHaveBeenCalled(); + expect(result.trigger).toEqual(workspaceTrigger); + }); + }); +}); diff --git a/packages/twenty-server/src/modules/workflow/common/workspace-services/workflow-common.workspace-service.ts b/packages/twenty-server/src/modules/workflow/common/workspace-services/workflow-common.workspace-service.ts index f27e66f474..d89c67c39b 100644 --- a/packages/twenty-server/src/modules/workflow/common/workspace-services/workflow-common.workspace-service.ts +++ b/packages/twenty-server/src/modules/workflow/common/workspace-services/workflow-common.workspace-service.ts @@ -1,9 +1,11 @@ import { Injectable, Logger } from '@nestjs/common'; import { isDefined, isValidUuid } from 'twenty-shared/utils'; +import { FeatureFlagKey } from 'twenty-shared/types'; import { In } from 'typeorm'; import { type WorkspaceAuthContext } from 'src/engine/core-modules/auth/types/workspace-auth-context.type'; +import { FeatureFlagService } from 'src/engine/core-modules/feature-flag/services/feature-flag.service'; import { WorkflowVersionCoreSyncService } from 'src/engine/core-modules/workflow/services/workflow-version-core-sync.service'; import { CommandMenuItemService } from 'src/engine/metadata-modules/command-menu-item/command-menu-item.service'; import { WorkspaceManyOrAllFlatEntityMapsCacheService } from 'src/engine/metadata-modules/flat-entity/services/workspace-many-or-all-flat-entity-maps-cache.service'; @@ -58,6 +60,7 @@ export class WorkflowCommonWorkspaceService { private readonly workspaceManyOrAllFlatEntityMapsCacheService: WorkspaceManyOrAllFlatEntityMapsCacheService, private readonly commandMenuItemService: CommandMenuItemService, private readonly workflowVersionCoreSyncService: WorkflowVersionCoreSyncService, + private readonly featureFlagService: FeatureFlagService, ) {} async getWorkflowVersionOrFail({ @@ -91,12 +94,52 @@ export class WorkflowCommonWorkspaceService { }, }); - return this.getValidWorkflowVersionOrFail(workflowVersion); + const validWorkflowVersion = + await this.getValidWorkflowVersionOrFail(workflowVersion); + + return this.overlayCoreWorkflowVersionContent( + workspaceId, + validWorkflowVersion, + ); }, authContext, ); } + private async overlayCoreWorkflowVersionContent( + workspaceId: string, + workflowVersion: WorkflowVersionWorkspaceEntity, + ): Promise { + const isCoreReadEnabled = await this.featureFlagService.isFeatureEnabled( + FeatureFlagKey.IS_WORKFLOW_VERSION_IN_CORE_ENABLED, + workspaceId, + ); + + if ( + !isCoreReadEnabled || + !isDefined(workflowVersion.coreWorkflowVersionId) + ) { + return workflowVersion; + } + + const coreWorkflowVersion = + await this.workflowVersionCoreSyncService.findCoreVersionById( + workspaceId, + workflowVersion.coreWorkflowVersionId, + ); + + if (!isDefined(coreWorkflowVersion)) { + return workflowVersion; + } + + return { + ...workflowVersion, + trigger: coreWorkflowVersion.triggers?.[0] ?? null, + steps: coreWorkflowVersion.steps, + status: coreWorkflowVersion.status as unknown as WorkflowVersionStatus, + }; + } + async getValidWorkflowVersionOrFail( workflowVersion: WorkflowVersionWorkspaceEntity | null, ): Promise {