From 63555573568fae2f27284a1e80c730bfabd62294 Mon Sep 17 00:00:00 2001 From: Thomas Trompette Date: Fri, 19 Dec 2025 18:17:02 +0100 Subject: [PATCH] [POC] Real-Time (#16633) https://github.com/user-attachments/assets/3ad7a4f3-d08b-4a1c-b3c2-bb42ef9f3575 Real time POC. Next steps: - find out the correct API for subscription. Should be triggered directly in query hooks (useFindManyRecords, useLazy...) - improve query matching --- .../src/generated-metadata/graphql.ts | 22 +++++ .../twenty-front/src/generated/graphql.ts | 61 +++++++++++++ .../app/components/AppRouterProviders.tsx | 59 ++++++------ .../components/SubscriptionProvider.tsx | 21 +++++ .../components/SubscriptionProviderEffect.tsx | 65 +++++++++++++ .../subscription/contexts/SseClientContext.ts | 4 + .../subscriptions/onSubscriptionMatch.ts | 11 +++ .../subscription/hooks/useOnDbEvent.ts | 19 +--- .../subscription/hooks/useSseClient.util.ts | 24 +++++ .../hooks/useSubscribeToRefetch.ts | 50 ++++++++++ .../states/subscriptionRegistryState.ts | 14 +++ .../listeners/entity-events-to-db.listener.ts | 8 +- .../dtos/subscription-matches.dto.ts | 13 +++ .../subscriptions/dtos/subscription.input.ts | 15 +++ .../subscriptions/subscription.service.ts | 91 ++++++++++++++++++- .../subscriptions/subscriptions.module.ts | 5 +- .../workspace-event-emitter.resolver.ts | 55 ++++++++++- 17 files changed, 482 insertions(+), 55 deletions(-) create mode 100644 packages/twenty-front/src/modules/subscription/components/SubscriptionProvider.tsx create mode 100644 packages/twenty-front/src/modules/subscription/components/SubscriptionProviderEffect.tsx create mode 100644 packages/twenty-front/src/modules/subscription/contexts/SseClientContext.ts create mode 100644 packages/twenty-front/src/modules/subscription/graphql/subscriptions/onSubscriptionMatch.ts create mode 100644 packages/twenty-front/src/modules/subscription/hooks/useSseClient.util.ts create mode 100644 packages/twenty-front/src/modules/subscription/hooks/useSubscribeToRefetch.ts create mode 100644 packages/twenty-front/src/modules/subscription/states/subscriptionRegistryState.ts create mode 100644 packages/twenty-server/src/engine/subscriptions/dtos/subscription-matches.dto.ts create mode 100644 packages/twenty-server/src/engine/subscriptions/dtos/subscription.input.ts diff --git a/packages/twenty-front/src/generated-metadata/graphql.ts b/packages/twenty-front/src/generated-metadata/graphql.ts index 21df1267b9..478cf3ece7 100644 --- a/packages/twenty-front/src/generated-metadata/graphql.ts +++ b/packages/twenty-front/src/generated-metadata/graphql.ts @@ -4066,6 +4066,7 @@ export type SubmitFormStepInput = { export type Subscription = { __typename?: 'Subscription'; onDbEvent: OnDbEvent; + onSubscriptionMatch?: Maybe; serverlessFunctionLogs: ServerlessFunctionLogs; }; @@ -4075,15 +4076,36 @@ export type SubscriptionOnDbEventArgs = { }; +export type SubscriptionOnSubscriptionMatchArgs = { + subscriptions: Array; +}; + + export type SubscriptionServerlessFunctionLogsArgs = { input: ServerlessFunctionLogsInput; }; +export type SubscriptionInput = { + id: Scalars['String']; + query: Scalars['String']; + selectedEventActions?: InputMaybe>; +}; + export enum SubscriptionInterval { Month = 'Month', Year = 'Year' } +export type SubscriptionMatch = { + __typename?: 'SubscriptionMatch'; + id: Scalars['String']; +}; + +export type SubscriptionMatches = { + __typename?: 'SubscriptionMatches'; + subscriptions: Array; +}; + export enum SubscriptionStatus { Active = 'Active', Canceled = 'Canceled', diff --git a/packages/twenty-front/src/generated/graphql.ts b/packages/twenty-front/src/generated/graphql.ts index 54d317958c..d2ce952a8d 100644 --- a/packages/twenty-front/src/generated/graphql.ts +++ b/packages/twenty-front/src/generated/graphql.ts @@ -3903,6 +3903,7 @@ export type SubmitFormStepInput = { export type Subscription = { __typename?: 'Subscription'; onDbEvent: OnDbEvent; + onSubscriptionMatch?: Maybe; serverlessFunctionLogs: ServerlessFunctionLogs; }; @@ -3912,15 +3913,36 @@ export type SubscriptionOnDbEventArgs = { }; +export type SubscriptionOnSubscriptionMatchArgs = { + subscriptions: Array; +}; + + export type SubscriptionServerlessFunctionLogsArgs = { input: ServerlessFunctionLogsInput; }; +export type SubscriptionInput = { + id: Scalars['String']; + query: Scalars['String']; + selectedEventActions?: InputMaybe>; +}; + export enum SubscriptionInterval { Month = 'Month', Year = 'Year' } +export type SubscriptionMatch = { + __typename?: 'SubscriptionMatch'; + id: Scalars['String']; +}; + +export type SubscriptionMatches = { + __typename?: 'SubscriptionMatches'; + subscriptions: Array; +}; + export enum SubscriptionStatus { Active = 'Active', Canceled = 'Canceled', @@ -4888,6 +4910,13 @@ export type OnDbEventSubscriptionVariables = Exact<{ export type OnDbEventSubscription = { __typename?: 'Subscription', onDbEvent: { __typename?: 'OnDbEvent', eventDate: string, action: DatabaseEventAction, objectNameSingular: string, updatedFields?: Array | null, record: any } }; +export type OnSubscriptionMatchSubscriptionVariables = Exact<{ + subscriptions: Array | SubscriptionInput; +}>; + + +export type OnSubscriptionMatchSubscription = { __typename?: 'Subscription', onSubscriptionMatch?: { __typename?: 'SubscriptionMatches', subscriptions: Array<{ __typename?: 'SubscriptionMatch', id: string }> } | 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 }; export type ViewFilterFragmentFragment = { __typename?: 'CoreViewFilter', id: any, fieldMetadataId: any, operand: ViewFilterOperand, value: any, viewFilterGroupId?: any | null, positionInViewFilterGroup?: number | null, subFieldName?: string | null, viewId: any, createdAt: string, updatedAt: string, deletedAt?: string | null }; @@ -5546,6 +5575,38 @@ export function useOnDbEventSubscription(baseOptions: Apollo.SubscriptionHookOpt } export type OnDbEventSubscriptionHookResult = ReturnType; export type OnDbEventSubscriptionResult = Apollo.SubscriptionResult; +export const OnSubscriptionMatchDocument = gql` + subscription OnSubscriptionMatch($subscriptions: [SubscriptionInput!]!) { + onSubscriptionMatch(subscriptions: $subscriptions) { + subscriptions { + id + } + } +} + `; + +/** + * __useOnSubscriptionMatchSubscription__ + * + * To run a query within a React component, call `useOnSubscriptionMatchSubscription` and pass it any options that fit your needs. + * When your component renders, `useOnSubscriptionMatchSubscription` returns an object from Apollo Client that contains loading, error, and data properties + * you can use to render your UI. + * + * @param baseOptions options that will be passed into the subscription, supported options are listed on: https://www.apollographql.com/docs/react/api/react-hooks/#options; + * + * @example + * const { data, loading, error } = useOnSubscriptionMatchSubscription({ + * variables: { + * subscriptions: // value for 'subscriptions' + * }, + * }); + */ +export function useOnSubscriptionMatchSubscription(baseOptions: Apollo.SubscriptionHookOptions) { + const options = {...defaultOptions, ...baseOptions} + return Apollo.useSubscription(OnSubscriptionMatchDocument, options); + } +export type OnSubscriptionMatchSubscriptionHookResult = ReturnType; +export type OnSubscriptionMatchSubscriptionResult = Apollo.SubscriptionResult; export const CreateCoreViewDocument = gql` mutation CreateCoreView($input: CreateViewInput!) { createCoreView(input: $input) { diff --git a/packages/twenty-front/src/modules/app/components/AppRouterProviders.tsx b/packages/twenty-front/src/modules/app/components/AppRouterProviders.tsx index 98e6340582..d8a3ab5c38 100644 --- a/packages/twenty-front/src/modules/app/components/AppRouterProviders.tsx +++ b/packages/twenty-front/src/modules/app/components/AppRouterProviders.tsx @@ -14,6 +14,8 @@ import { ApolloCoreProvider } from '@/object-metadata/components/ApolloCoreProvi import { ObjectMetadataItemsLoadEffect } from '@/object-metadata/components/ObjectMetadataItemsLoadEffect'; import { ObjectMetadataItemsProvider } from '@/object-metadata/components/ObjectMetadataItemsProvider'; import { PrefetchDataProvider } from '@/prefetch/components/PrefetchDataProvider'; +import { SubscriptionProvider } from '@/subscription/components/SubscriptionProvider'; +import { SupportChatEffect } from '@/support/components/SupportChatEffect'; import { DialogManager } from '@/ui/feedback/dialog-manager/components/DialogManager'; import { DialogComponentInstanceContext } from '@/ui/feedback/dialog-manager/contexts/DialogComponentInstanceContext'; import { SnackBarProvider } from '@/ui/feedback/snack-bar-manager/components/SnackBarProvider'; @@ -21,7 +23,6 @@ import { BaseThemeProvider } from '@/ui/theme/components/BaseThemeProvider'; import { UserThemeProviderEffect } from '@/ui/theme/components/UserThemeProviderEffect'; import { PageFavicon } from '@/ui/utilities/page-favicon/components/PageFavicon'; import { PageTitle } from '@/ui/utilities/page-title/components/PageTitle'; -import { SupportChatEffect } from '@/support/components/SupportChatEffect'; import { UserAndViewsProviderEffect } from '@/users/components/UserAndViewsProviderEffect'; import { UserProvider } from '@/users/components/UserProvider'; import { WorkspaceProviderEffect } from '@/workspace/components/WorkspaceProviderEffect'; @@ -45,33 +46,35 @@ export const AppRouterProviders = () => { - - - - - - - - - - - - - - - - - - - - - - - - - + + + + + + + + + + + + + + + + + + + + + + + + + + + diff --git a/packages/twenty-front/src/modules/subscription/components/SubscriptionProvider.tsx b/packages/twenty-front/src/modules/subscription/components/SubscriptionProvider.tsx new file mode 100644 index 0000000000..78171da6ca --- /dev/null +++ b/packages/twenty-front/src/modules/subscription/components/SubscriptionProvider.tsx @@ -0,0 +1,21 @@ +import { SubscriptionProviderEffect } from '@/subscription/components/SubscriptionProviderEffect'; +import { SseClientContext } from '@/subscription/contexts/SseClientContext'; +import { useSseClient } from '@/subscription/hooks/useSseClient.util'; +import { type ReactNode } from 'react'; + +type SubscriptionProviderProps = { + children: ReactNode; +}; + +export const SubscriptionProvider = ({ + children, +}: SubscriptionProviderProps) => { + const { sseClient } = useSseClient(); + + return ( + + + {children} + + ); +}; diff --git a/packages/twenty-front/src/modules/subscription/components/SubscriptionProviderEffect.tsx b/packages/twenty-front/src/modules/subscription/components/SubscriptionProviderEffect.tsx new file mode 100644 index 0000000000..867cf82ef9 --- /dev/null +++ b/packages/twenty-front/src/modules/subscription/components/SubscriptionProviderEffect.tsx @@ -0,0 +1,65 @@ +import { ON_SUBSCRIPTION_MATCH } from '@/subscription/graphql/subscriptions/onSubscriptionMatch'; +import { useSseClient } from '@/subscription/hooks/useSseClient.util'; +import { subscriptionRegistryState } from '@/subscription/states/subscriptionRegistryState'; +import { print } from 'graphql'; +import { useEffect, useMemo } from 'react'; +import { useRecoilValue } from 'recoil'; +import { isDefined } from 'twenty-shared/utils'; + +export const SubscriptionProviderEffect = () => { + const subscriptionRegistry = useRecoilValue(subscriptionRegistryState); + + const { sseClient } = useSseClient(); + + const subscriptions = useMemo(() => { + return Array.from(subscriptionRegistry.values()).map((entry) => ({ + id: entry.id, + query: entry.query, + })); + }, [subscriptionRegistry]); + + useEffect(() => { + if (!sseClient || subscriptions.length === 0) { + return; + } + + const unsubscribe = sseClient.subscribe( + { + query: print(ON_SUBSCRIPTION_MATCH), + variables: { subscriptions }, + }, + { + next: (value) => { + const data = value.data as { + onSubscriptionMatch: { subscriptions: { id: string }[] }; + } | null; + + if (!data?.onSubscriptionMatch?.subscriptions) { + return; + } + + for (const subscription of data.onSubscriptionMatch.subscriptions) { + const entry = subscriptionRegistry.get(subscription.id); + + if (!isDefined(entry)) { + continue; + } + + entry.onRefetch(); + } + }, + error: (error) => { + // eslint-disable-next-line no-console + console.error('Subscription error:', error); + }, + complete: () => {}, + }, + ); + + return () => { + unsubscribe(); + }; + }, [sseClient, subscriptions, subscriptionRegistry]); + + return null; +}; diff --git a/packages/twenty-front/src/modules/subscription/contexts/SseClientContext.ts b/packages/twenty-front/src/modules/subscription/contexts/SseClientContext.ts new file mode 100644 index 0000000000..4063965604 --- /dev/null +++ b/packages/twenty-front/src/modules/subscription/contexts/SseClientContext.ts @@ -0,0 +1,4 @@ +import { type Client } from 'graphql-sse'; +import { createContext } from 'react'; + +export const SseClientContext = createContext(null); diff --git a/packages/twenty-front/src/modules/subscription/graphql/subscriptions/onSubscriptionMatch.ts b/packages/twenty-front/src/modules/subscription/graphql/subscriptions/onSubscriptionMatch.ts new file mode 100644 index 0000000000..7903db228b --- /dev/null +++ b/packages/twenty-front/src/modules/subscription/graphql/subscriptions/onSubscriptionMatch.ts @@ -0,0 +1,11 @@ +import { gql } from '@apollo/client'; + +export const ON_SUBSCRIPTION_MATCH = gql` + subscription OnSubscriptionMatch($subscriptions: [SubscriptionInput!]!) { + onSubscriptionMatch(subscriptions: $subscriptions) { + subscriptions { + id + } + } + } +`; diff --git a/packages/twenty-front/src/modules/subscription/hooks/useOnDbEvent.ts b/packages/twenty-front/src/modules/subscription/hooks/useOnDbEvent.ts index 02efde9f43..aecd5e62cf 100644 --- a/packages/twenty-front/src/modules/subscription/hooks/useOnDbEvent.ts +++ b/packages/twenty-front/src/modules/subscription/hooks/useOnDbEvent.ts @@ -1,8 +1,6 @@ -import { getTokenPair } from '@/apollo/utils/getTokenPair'; import { ON_DB_EVENT } from '@/subscription/graphql/subscriptions/onDbEvent'; -import { createClient } from 'graphql-sse'; -import { useEffect, useMemo } from 'react'; -import { REACT_APP_SERVER_BASE_URL } from '~/config'; +import { useSseClient } from '@/subscription/hooks/useSseClient.util'; +import { useEffect } from 'react'; import { type Subscription, type SubscriptionOnDbEventArgs, @@ -22,18 +20,7 @@ export const useOnDbEvent = ({ input, skip = false, }: OnDbEventArgs) => { - const tokenPair = getTokenPair(); - - const sseClient = useMemo(() => { - const token = tokenPair?.accessOrWorkspaceAgnosticToken?.token; - - return createClient({ - url: `${REACT_APP_SERVER_BASE_URL}/graphql`, - headers: { - Authorization: token ? `Bearer ${token}` : '', - }, - }); - }, [tokenPair?.accessOrWorkspaceAgnosticToken?.token]); + const { sseClient } = useSseClient(); useEffect(() => { if (skip === true) { diff --git a/packages/twenty-front/src/modules/subscription/hooks/useSseClient.util.ts b/packages/twenty-front/src/modules/subscription/hooks/useSseClient.util.ts new file mode 100644 index 0000000000..8c33152631 --- /dev/null +++ b/packages/twenty-front/src/modules/subscription/hooks/useSseClient.util.ts @@ -0,0 +1,24 @@ +import { getTokenPair } from '@/apollo/utils/getTokenPair'; +import { createClient } from 'graphql-sse'; +import { useMemo } from 'react'; +import { REACT_APP_SERVER_BASE_URL } from '~/config'; + +export const useSseClient = () => { + const tokenPair = getTokenPair(); + const token = tokenPair?.accessOrWorkspaceAgnosticToken?.token; + + const sseClient = useMemo( + () => + createClient({ + url: `${REACT_APP_SERVER_BASE_URL}/graphql`, + headers: { + Authorization: token ? `Bearer ${token}` : '', + }, + }), + [token], + ); + + return { + sseClient, + }; +}; diff --git a/packages/twenty-front/src/modules/subscription/hooks/useSubscribeToRefetch.ts b/packages/twenty-front/src/modules/subscription/hooks/useSubscribeToRefetch.ts new file mode 100644 index 0000000000..beb52d6952 --- /dev/null +++ b/packages/twenty-front/src/modules/subscription/hooks/useSubscribeToRefetch.ts @@ -0,0 +1,50 @@ +import { subscriptionRegistryState } from '@/subscription/states/subscriptionRegistryState'; +import { type DocumentNode, print } from 'graphql'; +import { useEffect, useId } from 'react'; +import { useSetRecoilState } from 'recoil'; + +type UseSubscribeToRefetchParams = { + query: DocumentNode; + variables?: Record; + refetch: () => void; + skip?: boolean; +}; + +export const useSubscribeToRefetch = ({ + query, + variables, + refetch, + skip = false, +}: UseSubscribeToRefetchParams) => { + const setRegistry = useSetRecoilState(subscriptionRegistryState); + const subscriptionId = useId(); + + const queryString = JSON.stringify({ + query: print(query), + variables, + }); + + useEffect(() => { + if (skip) { + return; + } + + setRegistry((prev) => { + const next = new Map(prev); + next.set(subscriptionId, { + id: subscriptionId, + query: queryString, + onRefetch: refetch, + }); + return next; + }); + + return () => { + setRegistry((prev) => { + const next = new Map(prev); + next.delete(subscriptionId); + return next; + }); + }; + }, [subscriptionId, queryString, refetch, skip, setRegistry]); +}; diff --git a/packages/twenty-front/src/modules/subscription/states/subscriptionRegistryState.ts b/packages/twenty-front/src/modules/subscription/states/subscriptionRegistryState.ts new file mode 100644 index 0000000000..cc5e2d1533 --- /dev/null +++ b/packages/twenty-front/src/modules/subscription/states/subscriptionRegistryState.ts @@ -0,0 +1,14 @@ +import { createState } from 'twenty-ui/utilities'; + +export type SubscriptionEntry = { + id: string; + query: string; + onRefetch: () => void; +}; + +export const subscriptionRegistryState = createState< + Map +>({ + key: 'subscriptionRegistryState', + defaultValue: new Map(), +}); 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 3b9974ecef..8c828aae23 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 @@ -4,10 +4,10 @@ import { type ObjectRecordCreateEvent, type ObjectRecordDeleteEvent, type ObjectRecordDestroyEvent, - type ObjectRecordRestoreEvent, - type ObjectRecordUpdateEvent, type ObjectRecordEvent, type ObjectRecordNonDestructiveEvent, + type ObjectRecordRestoreEvent, + type ObjectRecordUpdateEvent, } from 'twenty-shared/database-events'; import { OnDatabaseBatchEvent } from 'src/engine/api/graphql/graphql-query-runner/decorators/on-database-batch-event.decorator'; @@ -17,11 +17,11 @@ import { InjectMessageQueue } from 'src/engine/core-modules/message-queue/decora 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 { CallWebhookJobsJob } from 'src/engine/core-modules/webhook/jobs/call-webhook-jobs.job'; +import { WorkspaceEventBatchForWebhook } from 'src/engine/core-modules/webhook/types/workspace-event-batch-for-webhook.type'; import { CallDatabaseEventTriggerJobsJob } from 'src/engine/metadata-modules/database-event-trigger/jobs/call-database-event-trigger-jobs.job'; import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type'; -import { UpsertTimelineActivityFromInternalEvent } from 'src/modules/timeline/jobs/upsert-timeline-activity-from-internal-event.job'; -import { WorkspaceEventBatchForWebhook } from 'src/engine/core-modules/webhook/types/workspace-event-batch-for-webhook.type'; import { WorkspaceEventEmitterService } from 'src/engine/workspace-event-emitter/workspace-event-emitter.service'; +import { UpsertTimelineActivityFromInternalEvent } from 'src/modules/timeline/jobs/upsert-timeline-activity-from-internal-event.job'; @Injectable() export class EntityEventsToDbListener { 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 new file mode 100644 index 0000000000..c352449b7e --- /dev/null +++ b/packages/twenty-server/src/engine/subscriptions/dtos/subscription-matches.dto.ts @@ -0,0 +1,13 @@ +import { Field, ObjectType } from '@nestjs/graphql'; + +@ObjectType('SubscriptionMatch') +export class SubscriptionMatchDTO { + @Field(() => String) + id: string; +} + +@ObjectType('SubscriptionMatches') +export class SubscriptionMatchesDTO { + @Field(() => [SubscriptionMatchDTO]) + subscriptions: SubscriptionMatchDTO[]; +} diff --git a/packages/twenty-server/src/engine/subscriptions/dtos/subscription.input.ts b/packages/twenty-server/src/engine/subscriptions/dtos/subscription.input.ts new file mode 100644 index 0000000000..f0acbf77a8 --- /dev/null +++ b/packages/twenty-server/src/engine/subscriptions/dtos/subscription.input.ts @@ -0,0 +1,15 @@ +import { Field, InputType } from '@nestjs/graphql'; + +import { DatabaseEventAction } from 'src/engine/api/graphql/graphql-query-runner/enums/database-event-action'; + +@InputType() +export class SubscriptionInput { + @Field() + id: string; + + @Field() + query: string; + + @Field(() => [DatabaseEventAction], { nullable: true }) + selectedEventActions?: DatabaseEventAction[]; +} diff --git a/packages/twenty-server/src/engine/subscriptions/subscription.service.ts b/packages/twenty-server/src/engine/subscriptions/subscription.service.ts index f9f5828b2e..e0ed7dae62 100644 --- a/packages/twenty-server/src/engine/subscriptions/subscription.service.ts +++ b/packages/twenty-server/src/engine/subscriptions/subscription.service.ts @@ -1,11 +1,20 @@ import { Injectable } from '@nestjs/common'; -import { SubscriptionChannel } from 'src/engine/subscriptions/enums/subscription-channel.enum'; +import { FieldNode, OperationDefinitionNode, parse } from 'graphql'; + import { RedisClientService } from 'src/engine/core-modules/redis-client/redis-client.service'; +import { WorkspaceManyOrAllFlatEntityMapsCacheService } from 'src/engine/metadata-modules/flat-entity/services/workspace-many-or-all-flat-entity-maps-cache.service'; +import { FlatObjectMetadata } from 'src/engine/metadata-modules/flat-object-metadata/types/flat-object-metadata.type'; +import { OnDbEventDTO } from 'src/engine/subscriptions/dtos/on-db-event.dto'; +import { SubscriptionInput } from 'src/engine/subscriptions/dtos/subscription.input'; +import { SubscriptionChannel } from 'src/engine/subscriptions/enums/subscription-channel.enum'; @Injectable() export class SubscriptionService { - constructor(private readonly redisClient: RedisClientService) {} + constructor( + private readonly redisClient: RedisClientService, + private readonly workspaceManyOrAllFlatEntityMapsCacheService: WorkspaceManyOrAllFlatEntityMapsCacheService, + ) {} private getSubscriptionChannel({ channel, @@ -47,4 +56,82 @@ export class SubscriptionService { payload, ); } + + public async isSubscriptionMatchingEvent( + subscription: SubscriptionInput, + event: OnDbEventDTO, + workspaceId: string, + ): Promise { + const objectName = this.parseQueryObjectName(subscription.query); + + if (!objectName) { + return false; + } + + if ( + subscription.selectedEventActions && + !subscription.selectedEventActions.includes(event.action) + ) { + return false; + } + + const { flatObjectMetadataMaps } = + await this.workspaceManyOrAllFlatEntityMapsCacheService.getOrRecomputeManyOrAllFlatEntityMaps( + { + workspaceId, + flatMapsKeys: ['flatObjectMetadataMaps'], + }, + ); + + const eventObjectMetadata = Object.values(flatObjectMetadataMaps.byId).find( + (metadata: FlatObjectMetadata) => + metadata.nameSingular === event.objectNameSingular, + ); + + if (!eventObjectMetadata) { + return false; + } + + const queryObjectNameLower = objectName.toLowerCase(); + const eventNameSingularLower = + eventObjectMetadata.nameSingular.toLowerCase(); + const eventNamePluralLower = eventObjectMetadata.namePlural.toLowerCase(); + + return ( + queryObjectNameLower === eventNameSingularLower || + queryObjectNameLower === eventNamePluralLower + ); + } + + private parseQueryObjectName(queryString: string): string | null { + try { + const { query } = JSON.parse(queryString) as { + query: string; + variables?: Record; + }; + + const ast = parse(query); + + const firstOperation = ast.definitions.find( + (def): def is OperationDefinitionNode => + def.kind === 'OperationDefinition', + ); + + if (!firstOperation) { + return null; + } + + const rootSelection = firstOperation.selectionSet.selections[0]; + + if (rootSelection.kind !== 'Field') { + return null; + } + + const rootField = rootSelection as FieldNode; + + return rootField.name.value; + } catch { + return null; + } + } } diff --git a/packages/twenty-server/src/engine/subscriptions/subscriptions.module.ts b/packages/twenty-server/src/engine/subscriptions/subscriptions.module.ts index 72d334b13a..8c57e092f5 100644 --- a/packages/twenty-server/src/engine/subscriptions/subscriptions.module.ts +++ b/packages/twenty-server/src/engine/subscriptions/subscriptions.module.ts @@ -1,10 +1,11 @@ import { Module } from '@nestjs/common'; -import { SubscriptionService } from 'src/engine/subscriptions/subscription.service'; import { RedisClientModule } from 'src/engine/core-modules/redis-client/redis-client.module'; +import { WorkspaceManyOrAllFlatEntityMapsCacheModule } from 'src/engine/metadata-modules/flat-entity/services/workspace-many-or-all-flat-entity-maps-cache.module'; +import { SubscriptionService } from 'src/engine/subscriptions/subscription.service'; @Module({ - imports: [RedisClientModule], + imports: [RedisClientModule, WorkspaceManyOrAllFlatEntityMapsCacheModule], providers: [SubscriptionService], exports: [SubscriptionService], }) 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 4e94571da6..c986970e14 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 @@ -5,15 +5,17 @@ import { isDefined } from 'twenty-shared/utils'; import { PreventNestToAutoLogGraphqlErrorsFilter } from 'src/engine/core-modules/graphql/filters/prevent-nest-to-auto-log-graphql-errors.filter'; import { ResolverValidationPipe } from 'src/engine/core-modules/graphql/pipes/resolver-validation.pipe'; +import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity'; +import { AuthWorkspace } from 'src/engine/decorators/auth/auth-workspace.decorator'; import { NoPermissionGuard } from 'src/engine/guards/no-permission.guard'; import { UserAuthGuard } from 'src/engine/guards/user-auth.guard'; import { WorkspaceAuthGuard } from 'src/engine/guards/workspace-auth.guard'; import { OnDbEventDTO } from 'src/engine/subscriptions/dtos/on-db-event.dto'; import { OnDbEventInput } from 'src/engine/subscriptions/dtos/on-db-event.input'; -import { SubscriptionService } from 'src/engine/subscriptions/subscription.service'; -import { AuthWorkspace } from 'src/engine/decorators/auth/auth-workspace.decorator'; -import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity'; +import { SubscriptionMatchesDTO } from 'src/engine/subscriptions/dtos/subscription-matches.dto'; +import { SubscriptionInput } from 'src/engine/subscriptions/dtos/subscription.input'; import { SubscriptionChannel } from 'src/engine/subscriptions/enums/subscription-channel.enum'; +import { SubscriptionService } from 'src/engine/subscriptions/subscription.service'; @Resolver() @UseGuards(WorkspaceAuthGuard, UserAuthGuard, NoPermissionGuard) @@ -54,4 +56,51 @@ export class WorkspaceEventEmitterResolver { workspaceId: workspace.id, }); } + + @Subscription(() => SubscriptionMatchesDTO, { + nullable: true, + resolve: async function ( + this: WorkspaceEventEmitterResolver, + payload: { onDbEvent: 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, + ); + + return matches ? subscription.id : null; + }), + ); + + const filteredIds = matchedSubscriptionIds.filter( + (id): id is string => id !== null, + ); + + if (filteredIds.length === 0) { + return { subscriptions: [] }; + } + + return { + subscriptions: filteredIds.map((id) => ({ id })), + }; + }, + }) + onSubscriptionMatch( + @Args('subscriptions', { type: () => [SubscriptionInput] }) + _: SubscriptionInput[], + @AuthWorkspace() workspace: WorkspaceEntity, + ) { + return this.subscriptionService.subscribe({ + channel: SubscriptionChannel.DATABASE_EVENT_CHANNEL, + workspaceId: workspace.id, + }); + } }