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.
This commit is contained in:
Thomas Trompette
2026-07-29 10:07:08 +02:00
committed by GitHub
parent b602294f1d
commit 198e3969df
3 changed files with 193 additions and 1 deletions
@@ -100,6 +100,15 @@ export class WorkflowVersionCoreSyncService {
await this.invalidateAutomatedTriggerMaps(workspaceId);
}
async findCoreVersionById(
workspaceId: string,
coreWorkflowVersionId: string,
): Promise<WorkflowVersionEntity | null> {
return this.coreWorkflowVersionRepository.findOne(workspaceId, {
where: { id: coreWorkflowVersionId },
});
}
async mirrorWorkflowVersionWrite({
workspaceId,
entityManager,
@@ -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);
});
});
});
@@ -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<WorkflowVersionWorkspaceEntity> {
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<WorkflowVersionWorkspaceEntity> {