From 198e3969df8d8a3319f66e43edaee67475d63400 Mon Sep 17 00:00:00 2001 From: Thomas Trompette Date: Wed, 29 Jul 2026 10:07:08 +0200 Subject: [PATCH] feat(workflow): read workflow version content from core in the engine, behind a flag (#23403) ## Scope: engine only Behind the existing `IS_WORKFLOW_VERSION_IN_CORE_ENABLED` flag (**off by default**), `getWorkflowVersionOrFail` sources a version's `trigger` / `steps` / `status` from `core.workflowVersion` instead of the workspace row. That covers its 11 call sites: workflow build, validation, schema, run, and trigger dispatch. **The UI is not switched here.** `useWorkflowVersion` and the content fetch in `useWorkflowWithCurrentVersion` still go through generic object CRUD against the workspace columns. That work needs the client to stop treating the workspace record as the home for content, and is planned separately around `flowComponentState` (the jotai state the builder already reads from). It is the read that must land before the workspace `trigger` / `steps` columns can be dropped. ## Why an overlay, not a repository swap `core.workflowVersion` is a content-only projection: its own `id` (not the workspace version id), `workflowId`, `triggers[]`, `steps[]`, `status`. No `name`, no `position`. So core cannot fully back the entity. The flip is therefore an overlay: identity, name and position stay from the workspace row, and only content comes from core (`triggers[0] -> trigger`). ## Reversibility Falls back to workspace content when the flag is off (default), when the version has no `coreWorkflowVersionId` soft-ref, or when the core row is missing. Merging changes nothing until the flag is enabled per workspace, and it can be flipped back at any time. ## Verification - `nx typecheck twenty-server` green, `oxfmt` + `oxlint --type-aware` green - Unit test covering four branches: flag off, flag on with the core row present (overlays content **and** preserves `id` / `name` from the workspace row), flag on with the core row missing (falls back on `trigger`, `steps` and `status`), flag on with no soft-ref (skips the core read) ## Enablement gate Do not enable the flag in any workspace until the drift dashboard reports zero drift for `core.workflowVersion` and legacy drift is repaired. Enabling before that turns latent drift into live dispatch behaviour. --- .../workflow-version-core-sync.service.ts | 9 ++ .../workflow-common.workspace-service.spec.ts | 140 ++++++++++++++++++ .../workflow-common.workspace-service.ts | 45 +++++- 3 files changed, 193 insertions(+), 1 deletion(-) create mode 100644 packages/twenty-server/src/modules/workflow/common/workspace-services/__tests__/workflow-common.workspace-service.spec.ts 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 {