From 1801b5086af6af83855620bbc353ed88ee4b976e Mon Sep 17 00:00:00 2001 From: martmull Date: Tue, 2 Sep 2025 14:40:30 +0200 Subject: [PATCH] Remove subscriptions job (#14250) Remove subscriptions publish job --------- Co-authored-by: Charles Bochet --- .../src/generated-metadata/graphql.ts | 1 + .../twenty-front/src/generated/graphql.ts | 1 + .../listeners/entity-events-to-db.listener.ts | 85 +++++++++++-------- .../workspace-query-runner.module.ts | 2 + .../enums/feature-flag-key.enum.ts | 1 + .../message-queue/message-queue.constants.ts | 1 - .../subscriptions/subscriptions.module.ts | 6 +- ...ptions.job.ts => subscriptions.service.ts} | 12 +-- .../workspace-entity-manager.spec.ts | 1 + .../core/utils/seed-feature-flags.util.ts | 5 ++ 10 files changed, 67 insertions(+), 48 deletions(-) rename packages/twenty-server/src/engine/subscriptions/{subscriptions.job.ts => subscriptions.service.ts} (78%) diff --git a/packages/twenty-front/src/generated-metadata/graphql.ts b/packages/twenty-front/src/generated-metadata/graphql.ts index aad4f42d90d..3d9731eab5a 100644 --- a/packages/twenty-front/src/generated-metadata/graphql.ts +++ b/packages/twenty-front/src/generated-metadata/graphql.ts @@ -970,6 +970,7 @@ export enum FeatureFlagKey { IS_API_KEY_ROLES_ENABLED = 'IS_API_KEY_ROLES_ENABLED', IS_CORE_VIEW_ENABLED = 'IS_CORE_VIEW_ENABLED', IS_CORE_VIEW_SYNCING_ENABLED = 'IS_CORE_VIEW_SYNCING_ENABLED', + IS_DATABASE_EVENT_TRIGGER_ENABLED = 'IS_DATABASE_EVENT_TRIGGER_ENABLED', IS_IMAP_SMTP_CALDAV_ENABLED = 'IS_IMAP_SMTP_CALDAV_ENABLED', IS_JSON_FILTER_ENABLED = 'IS_JSON_FILTER_ENABLED', IS_MESSAGE_FOLDER_CONTROL_ENABLED = 'IS_MESSAGE_FOLDER_CONTROL_ENABLED', diff --git a/packages/twenty-front/src/generated/graphql.ts b/packages/twenty-front/src/generated/graphql.ts index ab86491a2b1..41a007dc250 100644 --- a/packages/twenty-front/src/generated/graphql.ts +++ b/packages/twenty-front/src/generated/graphql.ts @@ -934,6 +934,7 @@ export enum FeatureFlagKey { IS_API_KEY_ROLES_ENABLED = 'IS_API_KEY_ROLES_ENABLED', IS_CORE_VIEW_ENABLED = 'IS_CORE_VIEW_ENABLED', IS_CORE_VIEW_SYNCING_ENABLED = 'IS_CORE_VIEW_SYNCING_ENABLED', + IS_DATABASE_EVENT_TRIGGER_ENABLED = 'IS_DATABASE_EVENT_TRIGGER_ENABLED', IS_IMAP_SMTP_CALDAV_ENABLED = 'IS_IMAP_SMTP_CALDAV_ENABLED', IS_JSON_FILTER_ENABLED = 'IS_JSON_FILTER_ENABLED', IS_MESSAGE_FOLDER_CONTROL_ENABLED = 'IS_MESSAGE_FOLDER_CONTROL_ENABLED', diff --git a/packages/twenty-server/src/engine/api/graphql/workspace-query-runner/listeners/entity-events-to-db.listener.ts b/packages/twenty-server/src/engine/api/graphql/workspace-query-runner/listeners/entity-events-to-db.listener.ts index 8cab667a070..8417e507423 100644 --- a/packages/twenty-server/src/engine/api/graphql/workspace-query-runner/listeners/entity-events-to-db.listener.ts +++ b/packages/twenty-server/src/engine/api/graphql/workspace-query-runner/listeners/entity-events-to-db.listener.ts @@ -13,12 +13,14 @@ import { type ObjectRecordUpdateEvent } from 'src/engine/core-modules/event-emit 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'; -import { SubscriptionsJob } from 'src/engine/subscriptions/subscriptions.job'; import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type'; import { UpsertTimelineActivityFromInternalEvent } from 'src/modules/timeline/jobs/upsert-timeline-activity-from-internal-event.job'; import { CallWebhookJobsJob } from 'src/engine/core-modules/webhook/jobs/call-webhook-jobs.job'; import { type ObjectRecordEventForWebhook } from 'src/engine/core-modules/webhook/types/object-record-event-for-webhook.type'; import { CallDatabaseEventTriggerJobsJob } from 'src/engine/metadata-modules/trigger/jobs/call-database-event-trigger-jobs.job'; +import { SubscriptionsService } from 'src/engine/subscriptions/subscriptions.service'; +import { FeatureFlagService } from 'src/engine/core-modules/feature-flag/services/feature-flag.service'; +import { FeatureFlagKey } from 'src/engine/core-modules/feature-flag/enums/feature-flag-key.enum'; @Injectable() export class EntityEventsToDbListener { @@ -27,10 +29,10 @@ export class EntityEventsToDbListener { private readonly entityEventsToDbQueueService: MessageQueueService, @InjectMessageQueue(MessageQueue.webhookQueue) private readonly webhookQueueService: MessageQueueService, - @InjectMessageQueue(MessageQueue.subscriptionsQueue) - private readonly subscriptionsQueueService: MessageQueueService, @InjectMessageQueue(MessageQueue.triggerQueue) private readonly triggerQueueService: MessageQueueService, + private readonly subscriptionsService: SubscriptionsService, + private readonly featureFlagService: FeatureFlagService, ) {} @OnDatabaseBatchEvent('*', DatabaseEventAction.CREATED) @@ -79,17 +81,14 @@ export class EntityEventsToDbListener { }, })); - await Promise.all([ - this.subscriptionsQueueService.add>( - SubscriptionsJob.name, - batchEvent, - { retryLimit: 3 }, - ), - this.triggerQueueService.add>( - CallDatabaseEventTriggerJobsJob.name, - batchEvent, - { retryLimit: 3 }, - ), + const isDatabaseEventTriggerEnabled = + await this.featureFlagService.isFeatureEnabled( + FeatureFlagKey.IS_DATABASE_EVENT_TRIGGER_ENABLED, + batchEvent.workspaceId, + ); + + const promises = [ + this.subscriptionsService.publish(batchEvent), this.webhookQueueService.add< WorkspaceEventBatch >( @@ -99,27 +98,41 @@ export class EntityEventsToDbListener { retryLimit: 3, }, ), - ...(auditLogsEvents.length > 0 - ? [ - this.entityEventsToDbQueueService.add>( - CreateAuditLogFromInternalEvent.name, - { - ...batchEvent, - events: auditLogsEvents, - }, - ), - ] - : []), - ...(action !== DatabaseEventAction.DESTROYED && auditLogsEvents.length > 0 - ? [ - this.entityEventsToDbQueueService.add< - WorkspaceEventBatch - >(UpsertTimelineActivityFromInternalEvent.name, { - ...batchEvent, - events: auditLogsEvents, - }), - ] - : []), - ]); + ]; + + if (isDatabaseEventTriggerEnabled) { + promises.push( + this.triggerQueueService.add>( + CallDatabaseEventTriggerJobsJob.name, + batchEvent, + { retryLimit: 3 }, + ), + ); + } + + if (auditLogsEvents.length > 0) { + promises.push( + this.entityEventsToDbQueueService.add>( + CreateAuditLogFromInternalEvent.name, + { + ...batchEvent, + events: auditLogsEvents, + }, + ), + ); + + if (action !== DatabaseEventAction.DESTROYED) { + promises.push( + this.entityEventsToDbQueueService.add< + WorkspaceEventBatch + >(UpsertTimelineActivityFromInternalEvent.name, { + ...batchEvent, + events: auditLogsEvents, + }), + ); + } + } + + await Promise.all(promises); } } diff --git a/packages/twenty-server/src/engine/api/graphql/workspace-query-runner/workspace-query-runner.module.ts b/packages/twenty-server/src/engine/api/graphql/workspace-query-runner/workspace-query-runner.module.ts index 08505ee6fd7..c503bce302c 100644 --- a/packages/twenty-server/src/engine/api/graphql/workspace-query-runner/workspace-query-runner.module.ts +++ b/packages/twenty-server/src/engine/api/graphql/workspace-query-runner/workspace-query-runner.module.ts @@ -14,6 +14,7 @@ import { RecordPositionModule } from 'src/engine/core-modules/record-position/re import { RecordTransformerModule } from 'src/engine/core-modules/record-transformer/record-transformer.module'; import { TelemetryModule } from 'src/engine/core-modules/telemetry/telemetry.module'; import { WorkspaceDataSourceModule } from 'src/engine/workspace-datasource/workspace-datasource.module'; +import { SubscriptionsModule } from 'src/engine/subscriptions/subscriptions.module'; import { EntityEventsToDbListener } from './listeners/entity-events-to-db.listener'; @@ -30,6 +31,7 @@ import { EntityEventsToDbListener } from './listeners/entity-events-to-db.listen FeatureFlagModule, RecordTransformerModule, RecordPositionModule, + SubscriptionsModule, ], providers: [ ...workspaceQueryRunnerFactories, diff --git a/packages/twenty-server/src/engine/core-modules/feature-flag/enums/feature-flag-key.enum.ts b/packages/twenty-server/src/engine/core-modules/feature-flag/enums/feature-flag-key.enum.ts index 932eed79330..d2dd188fb70 100644 --- a/packages/twenty-server/src/engine/core-modules/feature-flag/enums/feature-flag-key.enum.ts +++ b/packages/twenty-server/src/engine/core-modules/feature-flag/enums/feature-flag-key.enum.ts @@ -16,4 +16,5 @@ export enum FeatureFlagKey { IS_PAGE_LAYOUT_ENABLED = 'IS_PAGE_LAYOUT_ENABLED', IS_MESSAGE_FOLDER_CONTROL_ENABLED = 'IS_MESSAGE_FOLDER_CONTROL_ENABLED', IS_WORKFLOW_ITERATOR_ENABLED = 'IS_WORKFLOW_ITERATOR_ENABLED', + IS_DATABASE_EVENT_TRIGGER_ENABLED = 'IS_DATABASE_EVENT_TRIGGER_ENABLED', } diff --git a/packages/twenty-server/src/engine/core-modules/message-queue/message-queue.constants.ts b/packages/twenty-server/src/engine/core-modules/message-queue/message-queue.constants.ts index d122e666a3d..ff9dfcb1c22 100644 --- a/packages/twenty-server/src/engine/core-modules/message-queue/message-queue.constants.ts +++ b/packages/twenty-server/src/engine/core-modules/message-queue/message-queue.constants.ts @@ -15,7 +15,6 @@ export enum MessageQueue { entityEventsToDbQueue = 'entity-events-to-db-queue', workflowQueue = 'workflow-queue', deleteCascadeQueue = 'delete-cascade-queue', - subscriptionsQueue = 'subscriptions-queue', serverlessFunctionQueue = 'serverless-function-queue', triggerQueue = 'trigger-queue', } diff --git a/packages/twenty-server/src/engine/subscriptions/subscriptions.module.ts b/packages/twenty-server/src/engine/subscriptions/subscriptions.module.ts index 815e44ef433..cf4afae38e4 100644 --- a/packages/twenty-server/src/engine/subscriptions/subscriptions.module.ts +++ b/packages/twenty-server/src/engine/subscriptions/subscriptions.module.ts @@ -4,10 +4,9 @@ import { RedisPubSub } from 'graphql-redis-subscriptions'; import { RedisClientService } from 'src/engine/core-modules/redis-client/redis-client.service'; import { SubscriptionsResolver } from 'src/engine/subscriptions/subscriptions.resolver'; -import { SubscriptionsJob } from 'src/engine/subscriptions/subscriptions.job'; +import { SubscriptionsService } from 'src/engine/subscriptions/subscriptions.service'; @Module({ - exports: ['PUB_SUB'], providers: [ { provide: 'PUB_SUB', @@ -20,8 +19,9 @@ import { SubscriptionsJob } from 'src/engine/subscriptions/subscriptions.job'; }), }, SubscriptionsResolver, - SubscriptionsJob, + SubscriptionsService, ], + exports: ['PUB_SUB', SubscriptionsService], }) export class SubscriptionsModule implements OnModuleDestroy { constructor(@Inject('PUB_SUB') private readonly pubSub: RedisPubSub) {} diff --git a/packages/twenty-server/src/engine/subscriptions/subscriptions.job.ts b/packages/twenty-server/src/engine/subscriptions/subscriptions.service.ts similarity index 78% rename from packages/twenty-server/src/engine/subscriptions/subscriptions.job.ts rename to packages/twenty-server/src/engine/subscriptions/subscriptions.service.ts index 67445bb4f18..14b16fad143 100644 --- a/packages/twenty-server/src/engine/subscriptions/subscriptions.job.ts +++ b/packages/twenty-server/src/engine/subscriptions/subscriptions.service.ts @@ -1,21 +1,17 @@ -import { Inject } from '@nestjs/common'; +import { Inject, Injectable } from '@nestjs/common'; import { RedisPubSub } from 'graphql-redis-subscriptions'; import { isDefined } from 'twenty-shared/utils'; import { type ObjectRecordEvent } from 'src/engine/core-modules/event-emitter/types/object-record-event.event'; -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 { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event.type'; import { removeSecretFromWebhookRecord } from 'src/utils/remove-secret-from-webhook-record'; -@Processor(MessageQueue.subscriptionsQueue) -export class SubscriptionsJob { +@Injectable() +export class SubscriptionsService { constructor(@Inject('PUB_SUB') private readonly pubSub: RedisPubSub) {} - @Process(SubscriptionsJob.name) - async handle( + async publish( workspaceEventBatch: WorkspaceEventBatch, ): Promise { for (const eventData of workspaceEventBatch.events) { 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 047c6e4ebf8..18701da25cb 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 @@ -139,6 +139,7 @@ describe('WorkspaceEntityManager', () => { IS_PAGE_LAYOUT_ENABLED: false, IS_MESSAGE_FOLDER_CONTROL_ENABLED: false, IS_WORKFLOW_ITERATOR_ENABLED: false, + IS_DATABASE_EVENT_TRIGGER_ENABLED: false, }, eventEmitterService: { emitMutationEvent: jest.fn(), diff --git a/packages/twenty-server/src/engine/workspace-manager/dev-seeder/core/utils/seed-feature-flags.util.ts b/packages/twenty-server/src/engine/workspace-manager/dev-seeder/core/utils/seed-feature-flags.util.ts index caa96f75a76..a7436e08daa 100644 --- a/packages/twenty-server/src/engine/workspace-manager/dev-seeder/core/utils/seed-feature-flags.util.ts +++ b/packages/twenty-server/src/engine/workspace-manager/dev-seeder/core/utils/seed-feature-flags.util.ts @@ -85,6 +85,11 @@ export const seedFeatureFlags = async ( workspaceId: workspaceId, value: false, }, + { + key: FeatureFlagKey.IS_DATABASE_EVENT_TRIGGER_ENABLED, + workspaceId: workspaceId, + value: false, + }, ]) .execute(); };