feat(workflow): dispatch automated triggers from core behind a flag (#23775)

## 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.

<!-- This is an auto-generated description by cubic. -->
<a
href="https://cubic.dev/pr/twentyhq/twenty/pull/23775?utm_source=github"
target="_blank" rel="noopener noreferrer"
data-no-image-dialog="true"><picture><source
media="(prefers-color-scheme: dark)"
srcset="https://www.cubic.dev/buttons/review-in-cubic-dark.svg"><source
media="(prefers-color-scheme: light)"
srcset="https://www.cubic.dev/buttons/review-in-cubic-light.svg"><img
alt="Review in cubic"
src="https://www.cubic.dev/buttons/review-in-cubic-dark.svg"></picture></a>
<!-- End of auto-generated description by cubic. -->
This commit is contained in:
Thomas Trompette
2026-08-05 10:58:19 +02:00
committed by GitHub
parent 42aa566e32
commit 61c72942ac
12 changed files with 263 additions and 62 deletions
@@ -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 {
@@ -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 = {
@@ -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'
}
@@ -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'
}
@@ -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;
@@ -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: {},
@@ -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: [
@@ -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', () => {
@@ -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<CachedCronTrigger[]> {
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<WorkflowTriggerJobData>(
WorkflowTriggerJob.name,
{
workspaceId,
workflowId: trigger.workflowId,
workflowId,
payload: {},
},
{ retryLimit: 3 },
@@ -220,4 +220,41 @@ export class WorkflowCronTriggerCronJob {
return [];
}
}
private async getWorkspaceCronTriggers(
workspaceId: string,
): Promise<Array<{ workflowId: string; pattern?: string }>> {
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,
}),
);
}
}
@@ -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<GlobalWorkspaceOrmManager>;
let messageQueueService: jest.Mocked<MessageQueueService>;
let featureFlagService: jest.Mocked<FeatureFlagService>;
let workspaceCacheService: jest.Mocked<WorkspaceCacheService>;
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([
{
@@ -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<WorkflowAutomatedTriggerWorkspaceEntity>(
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<WorkflowTriggerJobData>(
WorkflowTriggerJob.name,
{
workspaceId,
workflowId: eventListener.workflowId,
payload: eventPayload,
},
{ retryLimit: 3 },
);
}
if (shouldTriggerJob) {
await this.messageQueueService.add<WorkflowTriggerJobData>(
WorkflowTriggerJob.name,
{
workspaceId,
workflowId: eventListener.workflowId,
payload: eventPayload,
},
{ retryLimit: 3 },
);
}
}
}, authContext);
}
}
private async getDatabaseEventListeners(
workspaceId: string,
databaseEventName: string,
): Promise<DatabaseEventTriggerListener[]> {
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<WorkflowAutomatedTriggerWorkspaceEntity>(
workspaceId,
automatedTriggerTableName,
{ shouldBypassPermissionChecks: true },
);
return workflowAutomatedTriggerRepository.find({
where: {
type: AutomatedTriggerType.DATABASE_EVENT,
settings: Raw(
() =>
`"${automatedTriggerTableName}"."settings"->>'eventName' = :eventName`,
{ eventName: databaseEventName },
),
},
});
},
buildSystemAuthContext(workspaceId),
);
}
private shouldTriggerJob({
@@ -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',
}