diff --git a/packages/twenty-front/src/generated-metadata/graphql.ts b/packages/twenty-front/src/generated-metadata/graphql.ts index 32d87ca0b0..f1dd7785ab 100644 --- a/packages/twenty-front/src/generated-metadata/graphql.ts +++ b/packages/twenty-front/src/generated-metadata/graphql.ts @@ -4292,12 +4292,13 @@ export enum SubscriptionInterval { export type SubscriptionMatch = { __typename?: 'SubscriptionMatch'; - id: Scalars['String']; + event: OnDbEvent; + subscriptionIds: Array; }; export type SubscriptionMatches = { __typename?: 'SubscriptionMatches'; - subscriptions: Array; + matches: Array; }; export enum SubscriptionStatus { diff --git a/packages/twenty-front/src/generated/graphql.ts b/packages/twenty-front/src/generated/graphql.ts index c7311e252e..f8bf0ba957 100644 --- a/packages/twenty-front/src/generated/graphql.ts +++ b/packages/twenty-front/src/generated/graphql.ts @@ -4084,12 +4084,13 @@ export enum SubscriptionInterval { export type SubscriptionMatch = { __typename?: 'SubscriptionMatch'; - id: Scalars['String']; + event: OnDbEvent; + subscriptionIds: Array; }; export type SubscriptionMatches = { __typename?: 'SubscriptionMatches'; - subscriptions: Array; + matches: Array; }; export enum SubscriptionStatus { @@ -5112,7 +5113,7 @@ export type OnSubscriptionMatchSubscriptionVariables = Exact<{ }>; -export type OnSubscriptionMatchSubscription = { __typename?: 'Subscription', onSubscriptionMatch?: { __typename?: 'SubscriptionMatches', subscriptions: Array<{ __typename?: 'SubscriptionMatch', id: string }> } | null }; +export type OnSubscriptionMatchSubscription = { __typename?: 'Subscription', onSubscriptionMatch?: { __typename?: 'SubscriptionMatches', matches: Array<{ __typename?: 'SubscriptionMatch', subscriptionIds: Array, event: { __typename?: 'OnDbEvent', action: DatabaseEventAction, objectNameSingular: string, eventDate: string, record: any, updatedFields?: Array | null } }> } | null }; export type ViewFieldFragmentFragment = { __typename?: 'CoreViewField', id: any, fieldMetadataId: any, viewId: any, isVisible: boolean, position: number, size: number, aggregateOperation?: AggregateOperations | null, createdAt: string, updatedAt: string, deletedAt?: string | null }; @@ -5783,8 +5784,15 @@ export type OnDbEventSubscriptionResult = Apollo.SubscriptionResult { const subscriptionRegistry = useRecoilValue(subscriptionRegistryState); @@ -31,21 +32,23 @@ export const SubscriptionProviderEffect = () => { { next: (value) => { const data = value.data as { - onSubscriptionMatch: { subscriptions: { id: string }[] }; + onSubscriptionMatch: SubscriptionMatches; } | null; - if (!data?.onSubscriptionMatch?.subscriptions) { + if (!data?.onSubscriptionMatch?.matches) { return; } - for (const subscription of data.onSubscriptionMatch.subscriptions) { - const entry = subscriptionRegistry.get(subscription.id); + for (const match of data.onSubscriptionMatch.matches) { + for (const subscriptionId of match.subscriptionIds) { + const entry = subscriptionRegistry.get(subscriptionId); - if (!isDefined(entry)) { - continue; + if (!isDefined(entry)) { + continue; + } + + entry.onRefetch(); } - - entry.onRefetch(); } }, error: (error) => { diff --git a/packages/twenty-front/src/modules/subscription/graphql/subscriptions/onSubscriptionMatch.ts b/packages/twenty-front/src/modules/subscription/graphql/subscriptions/onSubscriptionMatch.ts index 7903db228b..98a968f92a 100644 --- a/packages/twenty-front/src/modules/subscription/graphql/subscriptions/onSubscriptionMatch.ts +++ b/packages/twenty-front/src/modules/subscription/graphql/subscriptions/onSubscriptionMatch.ts @@ -3,8 +3,15 @@ import { gql } from '@apollo/client'; export const ON_SUBSCRIPTION_MATCH = gql` subscription OnSubscriptionMatch($subscriptions: [SubscriptionInput!]!) { onSubscriptionMatch(subscriptions: $subscriptions) { - subscriptions { - id + matches { + subscriptionIds + event { + action + objectNameSingular + eventDate + record + updatedFields + } } } } diff --git a/packages/twenty-server/src/engine/subscriptions/dtos/subscription-matches.dto.ts b/packages/twenty-server/src/engine/subscriptions/dtos/subscription-matches.dto.ts index c352449b7e..b21a68e657 100644 --- a/packages/twenty-server/src/engine/subscriptions/dtos/subscription-matches.dto.ts +++ b/packages/twenty-server/src/engine/subscriptions/dtos/subscription-matches.dto.ts @@ -1,13 +1,18 @@ import { Field, ObjectType } from '@nestjs/graphql'; +import { OnDbEventDTO } from './on-db-event.dto'; + @ObjectType('SubscriptionMatch') export class SubscriptionMatchDTO { - @Field(() => String) - id: string; + @Field(() => [String]) + subscriptionIds: string[]; + + @Field(() => OnDbEventDTO) + event: OnDbEventDTO; } @ObjectType('SubscriptionMatches') export class SubscriptionMatchesDTO { @Field(() => [SubscriptionMatchDTO]) - subscriptions: SubscriptionMatchDTO[]; + matches: SubscriptionMatchDTO[]; } diff --git a/packages/twenty-server/src/engine/subscriptions/enums/subscription-channel.enum.ts b/packages/twenty-server/src/engine/subscriptions/enums/subscription-channel.enum.ts index 26ff6acc16..cc3047530c 100644 --- a/packages/twenty-server/src/engine/subscriptions/enums/subscription-channel.enum.ts +++ b/packages/twenty-server/src/engine/subscriptions/enums/subscription-channel.enum.ts @@ -1,4 +1,5 @@ export enum SubscriptionChannel { DATABASE_EVENT_CHANNEL = 'DATABASE_EVENT_CHANNEL', + DATABASE_BATCH_EVENTS_CHANNEL = 'DATABASE_BATCH_EVENTS_CHANNEL', SERVERLESS_FUNCTION_LOGS_CHANNEL = 'SERVERLESS_FUNCTION_LOGS_CHANNEL', } diff --git a/packages/twenty-server/src/engine/workspace-event-emitter/workspace-event-emitter.resolver.ts b/packages/twenty-server/src/engine/workspace-event-emitter/workspace-event-emitter.resolver.ts index c986970e14..199c9b8443 100644 --- a/packages/twenty-server/src/engine/workspace-event-emitter/workspace-event-emitter.resolver.ts +++ b/packages/twenty-server/src/engine/workspace-event-emitter/workspace-event-emitter.resolver.ts @@ -61,36 +61,41 @@ export class WorkspaceEventEmitterResolver { nullable: true, resolve: async function ( this: WorkspaceEventEmitterResolver, - payload: { onDbEvent: OnDbEventDTO }, + payload: { onDbEvents: OnDbEventDTO[] }, args: { subscriptions: SubscriptionInput[] }, context: { req: { workspace: { id: string } } }, ): Promise { const workspaceId = context.req.workspace.id; - const matchedSubscriptionIds = await Promise.all( - args.subscriptions.map(async (subscription) => { - const matches = - await this.subscriptionService.isSubscriptionMatchingEvent( - subscription, - payload.onDbEvent, - workspaceId, - ); + const matches: { subscriptionIds: string[]; event: OnDbEventDTO }[] = []; - return matches ? subscription.id : null; - }), - ); + for (const event of payload.onDbEvents) { + const matchedSubscriptionIds = await Promise.all( + args.subscriptions.map(async (subscription) => { + const isMatch = + await this.subscriptionService.isSubscriptionMatchingEvent( + subscription, + event, + workspaceId, + ); - const filteredIds = matchedSubscriptionIds.filter( - (id): id is string => id !== null, - ); + return isMatch ? subscription.id : null; + }), + ); - if (filteredIds.length === 0) { - return { subscriptions: [] }; + const filteredIds = matchedSubscriptionIds.filter( + (id): id is string => id !== null, + ); + + if (filteredIds.length > 0) { + matches.push({ + subscriptionIds: filteredIds, + event, + }); + } } - return { - subscriptions: filteredIds.map((id) => ({ id })), - }; + return { matches }; }, }) onSubscriptionMatch( @@ -99,7 +104,7 @@ export class WorkspaceEventEmitterResolver { @AuthWorkspace() workspace: WorkspaceEntity, ) { return this.subscriptionService.subscribe({ - channel: SubscriptionChannel.DATABASE_EVENT_CHANNEL, + channel: SubscriptionChannel.DATABASE_BATCH_EVENTS_CHANNEL, workspaceId: workspace.id, }); } diff --git a/packages/twenty-server/src/engine/workspace-event-emitter/workspace-event-emitter.service.ts b/packages/twenty-server/src/engine/workspace-event-emitter/workspace-event-emitter.service.ts index 4083bf0f93..1a570c150e 100644 --- a/packages/twenty-server/src/engine/workspace-event-emitter/workspace-event-emitter.service.ts +++ b/packages/twenty-server/src/engine/workspace-event-emitter/workspace-event-emitter.service.ts @@ -2,10 +2,10 @@ import { Injectable } from '@nestjs/common'; import { type ObjectRecordEvent } from 'twenty-shared/database-events'; -import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type'; import { transformEventToWebhookEvent } from 'src/engine/core-modules/webhook/utils/transform-event-to-webhook-event'; -import { SubscriptionService } from 'src/engine/subscriptions/subscription.service'; import { SubscriptionChannel } from 'src/engine/subscriptions/enums/subscription-channel.enum'; +import { SubscriptionService } from 'src/engine/subscriptions/subscription.service'; +import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type'; @Injectable() export class WorkspaceEventEmitterService { @@ -16,25 +16,36 @@ export class WorkspaceEventEmitterService { ): Promise { const [nameSingular, operation] = workspaceEventBatch.name.split('.'); + const batchEvents = []; + for (const eventData of workspaceEventBatch.events) { const { record, updatedFields } = transformEventToWebhookEvent({ eventName: workspaceEventBatch.name, event: eventData, }); + const event = { + action: operation, + objectNameSingular: nameSingular, + eventDate: new Date(), + record, + ...(updatedFields && { updatedFields }), + }; + + batchEvents.push(event); + + // Publish individual events to legacy channel (onDbEvent) await this.subscriptionService.publish({ channel: SubscriptionChannel.DATABASE_EVENT_CHANNEL, workspaceId: workspaceEventBatch.workspaceId, - payload: { - onDbEvent: { - action: operation, - objectNameSingular: nameSingular, - eventDate: new Date(), - record, - ...(updatedFields && { updatedFields }), - }, - }, + payload: { onDbEvent: event }, }); } + + await this.subscriptionService.publish({ + channel: SubscriptionChannel.DATABASE_BATCH_EVENTS_CHANNEL, + workspaceId: workspaceEventBatch.workspaceId, + payload: { onDbEvents: batchEvents }, + }); } }