From 61c72942acd57ff96d247997eb012e052a5dcdc8 Mon Sep 17 00:00:00 2001 From: Thomas Trompette Date: Wed, 5 Aug 2026 10:58:19 +0200 Subject: [PATCH] feat(workflow): dispatch automated triggers from core behind a flag (#23775) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## Context Part of the workflow → core migration. Before we can stop writing workspace `trigger`/`steps`, automated-trigger dispatch must read from core. Dispatch currently reads the workspace `workflowAutomatedTrigger` table (populated from the workspace trigger), so it would go blank once those writes stop. This flips the dispatch reads behind a flag, mirroring the version-content read switch. ## What this does New flag `IS_WORKFLOW_DISPATCH_FROM_CORE_ENABLED` (per-workspace, default off). At each dispatch read site, flag-on reads the core-derived trigger map and flag-off keeps the current workspace query. - **DB-event listener** (`workflow-database-event-trigger.listener.ts`): extracted `getDatabaseEventListeners(workspaceId, eventName)`. Flag-on filters the core map (`getOrRecompute → byWorkflowId`, `type === DATABASE_EVENT && settings.eventName === name`); flag-off keeps the repo `find`. The evaluation type is broadened to the structural `{ workflowId, settings }` that both the entity and the map entry satisfy; the enqueue loop and `shouldTriggerJob` are unchanged. - **CRON job** (`workflow-cron-trigger-cron.job.ts`): extracted `getWorkspaceCronTriggers(workspaceId)`. Flag-on filters the core map for `type === CRON` → `{ workflowId, pattern }`; flag-off keeps the raw SQL. The redis cron cache, dedup and dispatch loop are unchanged; only the rebuild source swaps. ## Why it's safe - The core map is keyed by the workspace `workflowId`, and both sites enqueue `workflowId` only. Nothing consumes the map's core `workflowVersionId`, so `workflow-trigger.job.ts` still re-derives the version from workspace `lastPublishedVersionId` (no id translation). - Flag defaults off, per-workspace rollout. The drift cron's `checkAutomatedTriggerSync` already compares the core map against the workspace table, so it's the soak signal for flipping the flag. - The CRON source is only re-read on a cron-cache rebuild (cache miss), so a flag flip takes effect on the next rebuild: bounded by the cache TTL, or immediately on activation/deactivation, which invalidates the cache. Both sources emit identical `{ workflowId, pattern }` for a synced workspace, so the switch is a no-op in output. ## Prerequisite - The orphan-ACTIVE core-version cleanup (#23739) must land first: the core map is built from core ACTIVE versions, so a phantom orphan would become a live phantom trigger the moment this flag flips. ## Verification - Server unit specs cover both sites with the flag off (existing behavior) and on (reads the core map). - Live-verified on a dev instance: DB-event and CRON dispatch both fire from the core map with the flag on, and from the workspace entity with it off. Review in cubic --- .../src/metadata/generated/schema.graphql | 1 + .../src/metadata/generated/schema.ts | 5 +- .../src/generated-admin/graphql.ts | 1 + .../src/generated-metadata/graphql.ts | 1 + .../cached-workflow-automated-trigger.util.ts | 18 +++ .../workspace-entity-manager.spec.ts | 1 + .../automated-trigger.module.ts | 4 + .../workflow-cron-trigger-cron.job.spec.ts | 48 +++++++ .../jobs/workflow-cron-trigger-cron.job.ts | 73 +++++++--- ...ow-database-event-trigger.listener.spec.ts | 44 ++++++ ...orkflow-database-event-trigger.listener.ts | 128 ++++++++++++------ .../twenty-shared/src/types/FeatureFlagKey.ts | 1 + 12 files changed, 263 insertions(+), 62 deletions(-) create mode 100644 packages/twenty-server/src/engine/core-modules/workflow/utils/cached-workflow-automated-trigger.util.ts diff --git a/packages/twenty-client-sdk/src/metadata/generated/schema.graphql b/packages/twenty-client-sdk/src/metadata/generated/schema.graphql index e5044a483a..614e41286c 100644 --- a/packages/twenty-client-sdk/src/metadata/generated/schema.graphql +++ b/packages/twenty-client-sdk/src/metadata/generated/schema.graphql @@ -1793,6 +1793,7 @@ enum FeatureFlagKey { IS_LOGIC_FUNCTION_PREBUILT_MODE_ENABLED IS_SETTINGS_DISCOVERY_HERO_ENABLED IS_WORKFLOW_VERSION_IN_CORE_ENABLED + IS_WORKFLOW_DISPATCH_FROM_CORE_ENABLED } type WorkspaceUrls { diff --git a/packages/twenty-client-sdk/src/metadata/generated/schema.ts b/packages/twenty-client-sdk/src/metadata/generated/schema.ts index 46624ddfdc..db2fd54918 100644 --- a/packages/twenty-client-sdk/src/metadata/generated/schema.ts +++ b/packages/twenty-client-sdk/src/metadata/generated/schema.ts @@ -1434,7 +1434,7 @@ export interface FeatureFlag { __typename: 'FeatureFlag' } -export type FeatureFlagKey = 'IS_APP_CLAIMING_ENABLED' | 'IS_UNIQUE_INDEXES_ENABLED' | 'IS_JSON_FILTER_ENABLED' | 'IS_CALENDAR_WEEK_VIEW_ENABLED' | 'IS_EMAIL_GROUP_ENABLED' | 'IS_JUNCTION_RELATIONS_ENABLED' | 'IS_REST_METADATA_API_NEW_FORMAT_DIRECT' | 'IS_LOGIC_FUNCTION_PREBUILT_MODE_ENABLED' | 'IS_SETTINGS_DISCOVERY_HERO_ENABLED' | 'IS_WORKFLOW_VERSION_IN_CORE_ENABLED' +export type FeatureFlagKey = 'IS_APP_CLAIMING_ENABLED' | 'IS_UNIQUE_INDEXES_ENABLED' | 'IS_JSON_FILTER_ENABLED' | 'IS_CALENDAR_WEEK_VIEW_ENABLED' | 'IS_EMAIL_GROUP_ENABLED' | 'IS_JUNCTION_RELATIONS_ENABLED' | 'IS_REST_METADATA_API_NEW_FORMAT_DIRECT' | 'IS_LOGIC_FUNCTION_PREBUILT_MODE_ENABLED' | 'IS_SETTINGS_DISCOVERY_HERO_ENABLED' | 'IS_WORKFLOW_VERSION_IN_CORE_ENABLED' | 'IS_WORKFLOW_DISPATCH_FROM_CORE_ENABLED' export interface WorkspaceUrls { customUrl?: Scalars['String'] @@ -9621,7 +9621,8 @@ export const enumFeatureFlagKey = { IS_REST_METADATA_API_NEW_FORMAT_DIRECT: 'IS_REST_METADATA_API_NEW_FORMAT_DIRECT' as const, IS_LOGIC_FUNCTION_PREBUILT_MODE_ENABLED: 'IS_LOGIC_FUNCTION_PREBUILT_MODE_ENABLED' as const, IS_SETTINGS_DISCOVERY_HERO_ENABLED: 'IS_SETTINGS_DISCOVERY_HERO_ENABLED' as const, - IS_WORKFLOW_VERSION_IN_CORE_ENABLED: 'IS_WORKFLOW_VERSION_IN_CORE_ENABLED' as const + IS_WORKFLOW_VERSION_IN_CORE_ENABLED: 'IS_WORKFLOW_VERSION_IN_CORE_ENABLED' as const, + IS_WORKFLOW_DISPATCH_FROM_CORE_ENABLED: 'IS_WORKFLOW_DISPATCH_FROM_CORE_ENABLED' as const } export const enumIdentityProviderType = { diff --git a/packages/twenty-front/src/generated-admin/graphql.ts b/packages/twenty-front/src/generated-admin/graphql.ts index 9dee07bc19..b9e68a728e 100644 --- a/packages/twenty-front/src/generated-admin/graphql.ts +++ b/packages/twenty-front/src/generated-admin/graphql.ts @@ -329,6 +329,7 @@ export enum FeatureFlagKey { IS_REST_METADATA_API_NEW_FORMAT_DIRECT = 'IS_REST_METADATA_API_NEW_FORMAT_DIRECT', IS_SETTINGS_DISCOVERY_HERO_ENABLED = 'IS_SETTINGS_DISCOVERY_HERO_ENABLED', IS_UNIQUE_INDEXES_ENABLED = 'IS_UNIQUE_INDEXES_ENABLED', + IS_WORKFLOW_DISPATCH_FROM_CORE_ENABLED = 'IS_WORKFLOW_DISPATCH_FROM_CORE_ENABLED', IS_WORKFLOW_VERSION_IN_CORE_ENABLED = 'IS_WORKFLOW_VERSION_IN_CORE_ENABLED' } diff --git a/packages/twenty-front/src/generated-metadata/graphql.ts b/packages/twenty-front/src/generated-metadata/graphql.ts index 07441f2c23..cfa88c1aea 100644 --- a/packages/twenty-front/src/generated-metadata/graphql.ts +++ b/packages/twenty-front/src/generated-metadata/graphql.ts @@ -1808,6 +1808,7 @@ export enum FeatureFlagKey { IS_REST_METADATA_API_NEW_FORMAT_DIRECT = 'IS_REST_METADATA_API_NEW_FORMAT_DIRECT', IS_SETTINGS_DISCOVERY_HERO_ENABLED = 'IS_SETTINGS_DISCOVERY_HERO_ENABLED', IS_UNIQUE_INDEXES_ENABLED = 'IS_UNIQUE_INDEXES_ENABLED', + IS_WORKFLOW_DISPATCH_FROM_CORE_ENABLED = 'IS_WORKFLOW_DISPATCH_FROM_CORE_ENABLED', IS_WORKFLOW_VERSION_IN_CORE_ENABLED = 'IS_WORKFLOW_VERSION_IN_CORE_ENABLED' } diff --git a/packages/twenty-server/src/engine/core-modules/workflow/utils/cached-workflow-automated-trigger.util.ts b/packages/twenty-server/src/engine/core-modules/workflow/utils/cached-workflow-automated-trigger.util.ts new file mode 100644 index 0000000000..d0d8078c00 --- /dev/null +++ b/packages/twenty-server/src/engine/core-modules/workflow/utils/cached-workflow-automated-trigger.util.ts @@ -0,0 +1,18 @@ +import { type CachedWorkflowAutomatedTrigger } from 'src/engine/core-modules/workflow/types/workflow-automated-trigger-maps.type'; +import { AutomatedTriggerType } from 'src/modules/workflow/common/standard-objects/workflow-automated-trigger.workspace-entity'; +import { + type BaseDatabaseEventTriggerSettings, + type CronTriggerSettings, +} from 'src/modules/workflow/workflow-trigger/automated-trigger/constants/automated-trigger-settings'; + +export const isCachedCronTrigger = ( + trigger: CachedWorkflowAutomatedTrigger, +): trigger is CachedWorkflowAutomatedTrigger & { + settings: CronTriggerSettings; +} => trigger.type === AutomatedTriggerType.CRON; + +export const isCachedDatabaseEventTrigger = ( + trigger: CachedWorkflowAutomatedTrigger, +): trigger is CachedWorkflowAutomatedTrigger & { + settings: BaseDatabaseEventTriggerSettings; +} => trigger.type === AutomatedTriggerType.DATABASE_EVENT; diff --git a/packages/twenty-server/src/engine/twenty-orm/entity-manager/workspace-entity-manager.spec.ts b/packages/twenty-server/src/engine/twenty-orm/entity-manager/workspace-entity-manager.spec.ts index 22ddc4c471..935b924e94 100644 --- a/packages/twenty-server/src/engine/twenty-orm/entity-manager/workspace-entity-manager.spec.ts +++ b/packages/twenty-server/src/engine/twenty-orm/entity-manager/workspace-entity-manager.spec.ts @@ -249,6 +249,7 @@ describe('WorkspaceEntityManager', () => { IS_LOGIC_FUNCTION_PREBUILT_MODE_ENABLED: false, IS_SETTINGS_DISCOVERY_HERO_ENABLED: false, IS_WORKFLOW_VERSION_IN_CORE_ENABLED: false, + IS_WORKFLOW_DISPATCH_FROM_CORE_ENABLED: false, }, userWorkspaceRoleMap: {}, apiKeyRoleMap: {}, diff --git a/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/automated-trigger.module.ts b/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/automated-trigger.module.ts index d39f53dc0e..aa154db8c0 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/automated-trigger.module.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/automated-trigger.module.ts @@ -3,7 +3,9 @@ import { TypeOrmModule } from '@nestjs/typeorm'; import { CacheStorageModule } from 'src/engine/core-modules/cache-storage/cache-storage.module'; import { CronModule } from 'src/engine/core-modules/cron/cron.module'; +import { FeatureFlagModule } from 'src/engine/core-modules/feature-flag/feature-flag.module'; import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity'; +import { WorkspaceCacheModule } from 'src/engine/workspace-cache/workspace-cache.module'; import { WorkspaceDataSourceModule } from 'src/engine/workspace-datasource/workspace-datasource.module'; import { WorkflowCommonModule } from 'src/modules/workflow/common/workflow-common.module'; import { AutomatedTriggerWorkspaceService } from 'src/modules/workflow/workflow-trigger/automated-trigger/automated-trigger.workspace-service'; @@ -16,7 +18,9 @@ import { WorkflowDatabaseEventTriggerListener } from 'src/modules/workflow/workf TypeOrmModule.forFeature([WorkspaceEntity]), CacheStorageModule, CronModule, + FeatureFlagModule, WorkflowCommonModule, + WorkspaceCacheModule, WorkspaceDataSourceModule, ], providers: [ diff --git a/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/crons/jobs/__tests__/workflow-cron-trigger-cron.job.spec.ts b/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/crons/jobs/__tests__/workflow-cron-trigger-cron.job.spec.ts index 2f2b6099b2..1ee6697ac1 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/crons/jobs/__tests__/workflow-cron-trigger-cron.job.spec.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/crons/jobs/__tests__/workflow-cron-trigger-cron.job.spec.ts @@ -4,7 +4,9 @@ import { getDataSourceToken, getRepositoryToken } from '@nestjs/typeorm'; import { CacheStorageNamespace } from 'src/engine/core-modules/cache-storage/types/cache-storage-namespace.enum'; import { CronTriggerDeduplicationService } from 'src/engine/core-modules/cron/services/cron-trigger-deduplication.service'; import { ExceptionHandlerService } from 'src/engine/core-modules/exception-handler/exception-handler.service'; +import { FeatureFlagService } from 'src/engine/core-modules/feature-flag/services/feature-flag.service'; import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity'; +import { WorkspaceCacheService } from 'src/engine/workspace-cache/services/workspace-cache.service'; import { WORKFLOW_CRON_TRIGGER_CACHE_KEY } from 'src/modules/workflow/workflow-trigger/automated-trigger/crons/constants/workflow-cron-trigger-cache-key.constant'; import { WORKFLOW_CRON_TRIGGER_CACHE_TTL_MS } from 'src/modules/workflow/workflow-trigger/automated-trigger/crons/constants/workflow-cron-trigger-cache-ttl.constant'; import { WorkflowCronTriggerCronJob } from 'src/modules/workflow/workflow-trigger/automated-trigger/crons/jobs/workflow-cron-trigger-cron.job'; @@ -40,6 +42,14 @@ const mockCronTriggerDeduplicationService = { shouldDispatch: jest.fn(), }; +const mockFeatureFlagService = { + isFeatureEnabled: jest.fn(), +}; + +const mockWorkspaceCacheService = { + getOrRecompute: jest.fn(), +}; + describe('WorkflowCronTriggerCronJob', () => { let job: WorkflowCronTriggerCronJob; @@ -48,6 +58,8 @@ describe('WorkflowCronTriggerCronJob', () => { jest.useFakeTimers(); jest.setSystemTime(new Date('2026-04-02T15:00:30.000Z')); mockCronTriggerDeduplicationService.shouldDispatch.mockResolvedValue(true); + // Default flag off so the existing suite exercises the workspace-table path. + mockFeatureFlagService.isFeatureEnabled.mockResolvedValue(false); const module: TestingModule = await Test.createTestingModule({ providers: [ @@ -76,6 +88,14 @@ describe('WorkflowCronTriggerCronJob', () => { provide: CronTriggerDeduplicationService, useValue: mockCronTriggerDeduplicationService, }, + { + provide: FeatureFlagService, + useValue: mockFeatureFlagService, + }, + { + provide: WorkspaceCacheService, + useValue: mockWorkspaceCacheService, + }, ], }).compile(); @@ -262,6 +282,34 @@ describe('WorkflowCronTriggerCronJob', () => { expect(mockCacheStorageService.hashSet).not.toHaveBeenCalled(); expect(mockCacheStorageService.hashSetWithExpire).not.toHaveBeenCalled(); }); + + it('reads cron triggers from the core trigger map when dispatch-from-core is enabled', async () => { + mockFeatureFlagService.isFeatureEnabled.mockResolvedValue(true); + mockCacheStorageService.hashGetValues.mockResolvedValue([]); + mockWorkspaceRepository.find.mockResolvedValue([{ id: WORKSPACE_1 }]); + mockWorkspaceCacheService.getOrRecompute.mockResolvedValue({ + workflowAutomatedTriggerMaps: { + byWorkflowId: { + 'workflow-1': { + workflowId: 'workflow-1', + workflowVersionId: 'version-1', + type: 'CRON', + settings: { pattern: '* * * * *' }, + }, + }, + }, + } as any); + + await job.handle(); + + // Source is the core map, not the workspace table. + expect(mockCoreDataSource.query).not.toHaveBeenCalled(); + expect(mockMessageQueueService.add).toHaveBeenCalledWith( + WorkflowTriggerJob.name, + { workspaceId: WORKSPACE_1, workflowId: 'workflow-1', payload: {} }, + { retryLimit: 3 }, + ); + }); }); describe('error handling', () => { diff --git a/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/crons/jobs/workflow-cron-trigger-cron.job.ts b/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/crons/jobs/workflow-cron-trigger-cron.job.ts index 61b3c5a945..6b14917a1b 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/crons/jobs/workflow-cron-trigger-cron.job.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/crons/jobs/workflow-cron-trigger-cron.job.ts @@ -1,6 +1,7 @@ import { Logger } from '@nestjs/common'; import { InjectDataSource, InjectRepository } from '@nestjs/typeorm'; +import { FeatureFlagKey } from 'twenty-shared/types'; import { isDefined } from 'twenty-shared/utils'; import { WorkspaceActivationStatus } from 'twenty-shared/workspace'; import { DataSource, Repository } from 'typeorm'; @@ -11,12 +12,15 @@ import { CacheStorageNamespace } from 'src/engine/core-modules/cache-storage/typ import { SentryCronMonitor } from 'src/engine/core-modules/cron/sentry-cron-monitor.decorator'; import { CronTriggerDeduplicationService } from 'src/engine/core-modules/cron/services/cron-trigger-deduplication.service'; import { ExceptionHandlerService } from 'src/engine/core-modules/exception-handler/exception-handler.service'; +import { FeatureFlagService } from 'src/engine/core-modules/feature-flag/services/feature-flag.service'; import { InjectMessageQueue } from 'src/engine/core-modules/message-queue/decorators/message-queue.decorator'; import { Process } from 'src/engine/core-modules/message-queue/decorators/process.decorator'; import { Processor } from 'src/engine/core-modules/message-queue/decorators/processor.decorator'; import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants'; import { MessageQueueService } from 'src/engine/core-modules/message-queue/services/message-queue.service'; import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity'; +import { isCachedCronTrigger } from 'src/engine/core-modules/workflow/utils/cached-workflow-automated-trigger.util'; +import { WorkspaceCacheService } from 'src/engine/workspace-cache/services/workspace-cache.service'; import { getWorkspaceSchemaName } from 'src/engine/workspace-datasource/utils/get-workspace-schema-name.util'; import { AutomatedTriggerType } from 'src/modules/workflow/common/standard-objects/workflow-automated-trigger.workspace-entity'; import { type CronTriggerSettings } from 'src/modules/workflow/workflow-trigger/automated-trigger/constants/automated-trigger-settings'; @@ -45,6 +49,8 @@ export class WorkflowCronTriggerCronJob { @InjectCacheStorage(CacheStorageNamespace.ModuleWorkflow) private readonly cacheStorageService: CacheStorageService, private readonly cronTriggerDeduplicationService: CronTriggerDeduplicationService, + private readonly featureFlagService: FeatureFlagService, + private readonly workspaceCacheService: WorkspaceCacheService, ) {} @Process(WorkflowCronTriggerCronJob.name) @@ -152,57 +158,51 @@ export class WorkflowCronTriggerCronJob { now: Date, ): Promise { try { - const schemaName = getWorkspaceSchemaName(workspaceId); + const cronTriggers = await this.getWorkspaceCronTriggers(workspaceId); - const workflowAutomatedCronTriggers = await this.coreDataSource.query( - `SELECT * FROM ${schemaName}."workflowAutomatedTrigger" WHERE type = '${AutomatedTriggerType.CRON}'`, - ); - - if (workflowAutomatedCronTriggers.length === 0) { + if (cronTriggers.length === 0) { return []; } this.logger.log( - `Workspace ${workspaceId}: found ${workflowAutomatedCronTriggers.length} cron triggers`, + `Workspace ${workspaceId}: found ${cronTriggers.length} cron triggers`, ); const triggersToCache: CachedCronTrigger[] = []; - for (const trigger of workflowAutomatedCronTriggers) { - const settings = trigger.settings as CronTriggerSettings; - - if (!isDefined(settings.pattern)) { + for (const { workflowId, pattern } of cronTriggers) { + if (!isDefined(pattern)) { this.logger.warn( - `Trigger ${trigger.id}: skipping - pattern not defined`, + `Workflow ${workflowId}: skipping - cron pattern not defined`, ); continue; } const cachedTrigger: CachedCronTrigger = { workspaceId, - workflowId: trigger.workflowId, - pattern: settings.pattern, + workflowId, + pattern, }; triggersToCache.push(cachedTrigger); const shouldDispatch = await this.cronTriggerDeduplicationService.shouldDispatch( - `workflow-cron:${workspaceId}:${trigger.workflowId}`, - settings.pattern, + `workflow-cron:${workspaceId}:${workflowId}`, + pattern, now, ); if (shouldDispatch) { this.logger.log( - `Trigger ${trigger.id}: enqueuing WorkflowTriggerJob for workflow ${trigger.workflowId}`, + `Enqueuing WorkflowTriggerJob for workflow ${workflowId}`, ); await this.messageQueueService.add( WorkflowTriggerJob.name, { workspaceId, - workflowId: trigger.workflowId, + workflowId, payload: {}, }, { retryLimit: 3 }, @@ -220,4 +220,41 @@ export class WorkflowCronTriggerCronJob { return []; } } + + private async getWorkspaceCronTriggers( + workspaceId: string, + ): Promise> { + const isDispatchFromCoreEnabled = + await this.featureFlagService.isFeatureEnabled( + FeatureFlagKey.IS_WORKFLOW_DISPATCH_FROM_CORE_ENABLED, + workspaceId, + ); + + if (isDispatchFromCoreEnabled) { + const { workflowAutomatedTriggerMaps } = + await this.workspaceCacheService.getOrRecompute(workspaceId, [ + 'workflowAutomatedTriggerMaps', + ]); + + return Object.values(workflowAutomatedTriggerMaps.byWorkflowId) + .filter(isCachedCronTrigger) + .map((trigger) => ({ + workflowId: trigger.workflowId, + pattern: trigger.settings.pattern, + })); + } + + const schemaName = getWorkspaceSchemaName(workspaceId); + + const rows = await this.coreDataSource.query( + `SELECT "workflowId", settings FROM ${schemaName}."workflowAutomatedTrigger" WHERE type = '${AutomatedTriggerType.CRON}'`, + ); + + return rows.map( + (row: { workflowId: string; settings: CronTriggerSettings }) => ({ + workflowId: row.workflowId, + pattern: row.settings?.pattern, + }), + ); + } } diff --git a/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/listeners/__tests__/workflow-database-event-trigger.listener.spec.ts b/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/listeners/__tests__/workflow-database-event-trigger.listener.spec.ts index 8baccfc812..870e0e023f 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/listeners/__tests__/workflow-database-event-trigger.listener.spec.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/listeners/__tests__/workflow-database-event-trigger.listener.spec.ts @@ -1,8 +1,10 @@ import { Test, type TestingModule } from '@nestjs/testing'; +import { FeatureFlagService } from 'src/engine/core-modules/feature-flag/services/feature-flag.service'; import { MessageQueueService } from 'src/engine/core-modules/message-queue/services/message-queue.service'; import { type FlatObjectMetadata } from 'src/engine/metadata-modules/flat-object-metadata/types/flat-object-metadata.type'; import { GlobalWorkspaceOrmManager } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager'; +import { WorkspaceCacheService } from 'src/engine/workspace-cache/services/workspace-cache.service'; import { type WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type'; import { AutomatedTriggerType } from 'src/modules/workflow/common/standard-objects/workflow-automated-trigger.workspace-entity'; import { WorkflowCommonWorkspaceService } from 'src/modules/workflow/common/workspace-services/workflow-common.workspace-service'; @@ -13,6 +15,8 @@ describe('WorkflowDatabaseEventTriggerListener', () => { let listener: WorkflowDatabaseEventTriggerListener; let globalWorkspaceOrmManager: jest.Mocked; let messageQueueService: jest.Mocked; + let featureFlagService: jest.Mocked; + let workspaceCacheService: jest.Mocked; const mockRepository = { find: jest.fn(), @@ -58,6 +62,15 @@ describe('WorkflowDatabaseEventTriggerListener', () => { add: jest.fn(), } as any; + // Default flag off so the existing suite exercises the workspace-entity path. + featureFlagService = { + isFeatureEnabled: jest.fn().mockResolvedValue(false), + } as any; + + workspaceCacheService = { + getOrRecompute: jest.fn(), + } as any; + const module: TestingModule = await Test.createTestingModule({ providers: [ WorkflowDatabaseEventTriggerListener, @@ -69,6 +82,14 @@ describe('WorkflowDatabaseEventTriggerListener', () => { provide: MessageQueueService, useValue: messageQueueService, }, + { + provide: FeatureFlagService, + useValue: featureFlagService, + }, + { + provide: WorkspaceCacheService, + useValue: workspaceCacheService, + }, { provide: 'MESSAGE_QUEUE_workflow-queue', useValue: messageQueueService, @@ -140,6 +161,29 @@ describe('WorkflowDatabaseEventTriggerListener', () => { ); }); + it('reads listeners from the core trigger map when dispatch-from-core is enabled', async () => { + featureFlagService.isFeatureEnabled.mockResolvedValue(true); + workspaceCacheService.getOrRecompute.mockResolvedValue({ + workflowAutomatedTriggerMaps: { + byWorkflowId: { [workflowId]: mockEventListeners[0] }, + }, + } as any); + + await listener.handleObjectRecordUpdateEvent(mockPayload); + + // Dispatch is driven by the core map, not the workspace entity. + expect(mockRepository.find).not.toHaveBeenCalled(); + expect(messageQueueService.add).toHaveBeenCalledWith( + WorkflowTriggerJob.name, + { + workspaceId, + workflowId, + payload: mockPayload.events[0], + }, + { retryLimit: 3 }, + ); + }); + it('should trigger workflow when no fields are specified', async () => { mockRepository.find.mockResolvedValue([ { diff --git a/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/listeners/workflow-database-event-trigger.listener.ts b/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/listeners/workflow-database-event-trigger.listener.ts index c2b3c05a45..346b536471 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/listeners/workflow-database-event-trigger.listener.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/listeners/workflow-database-event-trigger.listener.ts @@ -8,13 +8,14 @@ import { type ObjectRecordUpdateEvent, type ObjectRecordUpsertEvent, } from 'twenty-shared/database-events'; -import { type ObjectRecord } from 'twenty-shared/types'; +import { FeatureFlagKey, type ObjectRecord } from 'twenty-shared/types'; import { isDefined, isNonEmptyArray } from 'twenty-shared/utils'; import { TRIGGER_STEP_ID } from 'twenty-shared/workflow'; import { In, Raw } from 'typeorm'; import { OnDatabaseBatchEvent } from 'src/engine/api/graphql/graphql-query-runner/decorators/on-database-batch-event.decorator'; import { DatabaseEventAction } from 'src/engine/api/graphql/graphql-query-runner/enums/database-event-action'; +import { FeatureFlagService } from 'src/engine/core-modules/feature-flag/services/feature-flag.service'; import { InjectMessageQueue } from 'src/engine/core-modules/message-queue/decorators/message-queue.decorator'; import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants'; import { MessageQueueService } from 'src/engine/core-modules/message-queue/services/message-queue.service'; @@ -26,6 +27,8 @@ import { buildFieldMapsFromFlatObjectMetadata } from 'src/engine/metadata-module import { type FlatObjectMetadata } from 'src/engine/metadata-modules/flat-object-metadata/types/flat-object-metadata.type'; import { GlobalWorkspaceOrmManager } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager'; import { buildSystemAuthContext } from 'src/engine/twenty-orm/utils/build-system-auth-context.util'; +import { isCachedDatabaseEventTrigger } from 'src/engine/core-modules/workflow/utils/cached-workflow-automated-trigger.util'; +import { WorkspaceCacheService } from 'src/engine/workspace-cache/services/workspace-cache.service'; import { type WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type'; import { AutomatedTriggerType, @@ -34,6 +37,7 @@ import { import { WorkflowCommonWorkspaceService } from 'src/modules/workflow/common/workspace-services/workflow-common.workspace-service'; import { evaluateStepFilters } from 'src/modules/workflow/workflow-executor/workflow-actions/filter/utils/evaluate-step-filters.util'; import { + type AutomatedTriggerSettings, type BaseDatabaseEventTriggerSettings, type UpdateEventTriggerSettings, } from 'src/modules/workflow/workflow-trigger/automated-trigger/constants/automated-trigger-settings'; @@ -42,9 +46,16 @@ import { type WorkflowTriggerJobData, } from 'src/modules/workflow/workflow-trigger/jobs/workflow-trigger.job'; +// Both the workspace workflowAutomatedTrigger entity and the core-derived +// trigger-map entry satisfy this shape, so dispatch can evaluate either source. +type DatabaseEventTriggerListener = { + workflowId: string; + settings: AutomatedTriggerSettings; +}; + type TriggerEvaluationArgs = { eventPayload: ObjectRecordEvent; - eventListener: WorkflowAutomatedTriggerWorkspaceEntity; + eventListener: DatabaseEventTriggerListener; action: DatabaseEventAction; }; @@ -59,6 +70,8 @@ export class WorkflowDatabaseEventTriggerListener { @InjectMessageQueue(MessageQueue.workflowQueue) private readonly messageQueueService: MessageQueueService, private readonly workflowCommonWorkspaceService: WorkflowCommonWorkspaceService, + private readonly featureFlagService: FeatureFlagService, + private readonly workspaceCacheService: WorkspaceCacheService, ) {} @OnDatabaseBatchEvent('*', DatabaseEventAction.CREATED) @@ -343,51 +356,82 @@ export class WorkflowDatabaseEventTriggerListener { }) { const workspaceId = payload.workspaceId; const databaseEventName = payload.name; - const automatedTriggerTableName = 'workflowAutomatedTrigger'; - const authContext = buildSystemAuthContext(workspaceId); + const eventListeners = await this.getDatabaseEventListeners( + workspaceId, + databaseEventName, + ); - await this.globalWorkspaceOrmManager.executeInWorkspaceContext(async () => { - const workflowAutomatedTriggerRepository = - await this.globalWorkspaceOrmManager.getRepository( - workspaceId, - automatedTriggerTableName, - { shouldBypassPermissionChecks: true }, - ); + for (const eventListener of eventListeners) { + for (const eventPayload of payload.events) { + const shouldTriggerJob = this.shouldTriggerJob({ + eventPayload, + eventListener, + action, + }); - const eventListeners = await workflowAutomatedTriggerRepository.find({ - where: { - type: AutomatedTriggerType.DATABASE_EVENT, - settings: Raw( - () => - `"${automatedTriggerTableName}"."settings"->>'eventName' = :eventName`, - { eventName: databaseEventName }, - ), - }, - }); - - for (const eventListener of eventListeners) { - for (const eventPayload of payload.events) { - const shouldTriggerJob = this.shouldTriggerJob({ - eventPayload, - eventListener, - action, - }); - - if (shouldTriggerJob) { - await this.messageQueueService.add( - WorkflowTriggerJob.name, - { - workspaceId, - workflowId: eventListener.workflowId, - payload: eventPayload, - }, - { retryLimit: 3 }, - ); - } + if (shouldTriggerJob) { + await this.messageQueueService.add( + WorkflowTriggerJob.name, + { + workspaceId, + workflowId: eventListener.workflowId, + payload: eventPayload, + }, + { retryLimit: 3 }, + ); } } - }, authContext); + } + } + + private async getDatabaseEventListeners( + workspaceId: string, + databaseEventName: string, + ): Promise { + const isDispatchFromCoreEnabled = + await this.featureFlagService.isFeatureEnabled( + FeatureFlagKey.IS_WORKFLOW_DISPATCH_FROM_CORE_ENABLED, + workspaceId, + ); + + if (isDispatchFromCoreEnabled) { + const { workflowAutomatedTriggerMaps } = + await this.workspaceCacheService.getOrRecompute(workspaceId, [ + 'workflowAutomatedTriggerMaps', + ]); + + return Object.values(workflowAutomatedTriggerMaps.byWorkflowId).filter( + (trigger) => + isCachedDatabaseEventTrigger(trigger) && + trigger.settings.eventName === databaseEventName, + ); + } + + const automatedTriggerTableName = 'workflowAutomatedTrigger'; + + return this.globalWorkspaceOrmManager.executeInWorkspaceContext( + async () => { + const workflowAutomatedTriggerRepository = + await this.globalWorkspaceOrmManager.getRepository( + workspaceId, + automatedTriggerTableName, + { shouldBypassPermissionChecks: true }, + ); + + return workflowAutomatedTriggerRepository.find({ + where: { + type: AutomatedTriggerType.DATABASE_EVENT, + settings: Raw( + () => + `"${automatedTriggerTableName}"."settings"->>'eventName' = :eventName`, + { eventName: databaseEventName }, + ), + }, + }); + }, + buildSystemAuthContext(workspaceId), + ); } private shouldTriggerJob({ diff --git a/packages/twenty-shared/src/types/FeatureFlagKey.ts b/packages/twenty-shared/src/types/FeatureFlagKey.ts index 2e441a870f..46d9899ca4 100644 --- a/packages/twenty-shared/src/types/FeatureFlagKey.ts +++ b/packages/twenty-shared/src/types/FeatureFlagKey.ts @@ -9,4 +9,5 @@ export enum FeatureFlagKey { IS_LOGIC_FUNCTION_PREBUILT_MODE_ENABLED = 'IS_LOGIC_FUNCTION_PREBUILT_MODE_ENABLED', IS_SETTINGS_DISCOVERY_HERO_ENABLED = 'IS_SETTINGS_DISCOVERY_HERO_ENABLED', IS_WORKFLOW_VERSION_IN_CORE_ENABLED = 'IS_WORKFLOW_VERSION_IN_CORE_ENABLED', + IS_WORKFLOW_DISPATCH_FROM_CORE_ENABLED = 'IS_WORKFLOW_DISPATCH_FROM_CORE_ENABLED', }