diff --git a/packages/twenty-front/src/modules/auth/hooks/useAuth.ts b/packages/twenty-front/src/modules/auth/hooks/useAuth.ts index 9f557a8ca9..d754f082ef 100644 --- a/packages/twenty-front/src/modules/auth/hooks/useAuth.ts +++ b/packages/twenty-front/src/modules/auth/hooks/useAuth.ts @@ -58,6 +58,8 @@ import { useRedirectToWorkspaceDomain } from '@/domain-manager/hooks/useRedirect import { domainConfigurationState } from '@/domain-manager/states/domainConfigurationState'; import { useLoadMockedObjectMetadataItems } from '@/object-metadata/hooks/useLoadMockedObjectMetadataItems'; import { useRefreshObjectMetadataItems } from '@/object-metadata/hooks/useRefreshObjectMetadataItems'; +import { sseClientState } from '@/sse-db-event/states/sseClientState'; +import { getSnapshotValue } from '@/ui/utilities/state/utils/getSnapshotValue'; import { useLoadCurrentUser } from '@/users/hooks/useLoadCurrentUser'; import { workspaceAuthProvidersState } from '@/workspace/states/workspaceAuthProvidersState'; import { i18n } from '@lingui/core'; @@ -124,6 +126,10 @@ export const useAuth = () => { const clearSession = useRecoilCallback( ({ snapshot, set }) => async () => { + const sseClient = getSnapshotValue(snapshot, sseClientState); + + sseClient?.dispose(); + const emptySnapshot = snapshot_UNSTABLE(); const iconsValue = snapshot.getLoadable(iconsState).getValue(); diff --git a/packages/twenty-front/src/modules/sse-db-event/components/SSEClientEffect.tsx b/packages/twenty-front/src/modules/sse-db-event/components/SSEClientEffect.tsx new file mode 100644 index 0000000000..a308c5adab --- /dev/null +++ b/packages/twenty-front/src/modules/sse-db-event/components/SSEClientEffect.tsx @@ -0,0 +1,71 @@ +import { getTokenPair } from '@/apollo/utils/getTokenPair'; +import { useIsLogged } from '@/auth/hooks/useIsLogged'; +import { SSE_CONNECTION_RETRY_MAX_WAIT_TIME_IN_MS } from '@/sse-db-event/constants/SseConnectionRetryMaxWaitTimeInMs'; +import { SSE_CONNECTION_RETRY_WAIT_TIME_IN_MS_FOR_DEV_MODE } from '@/sse-db-event/constants/SseConnectionRetryWaitTimeInMsForDevMode'; +import { activeQueryListenersState } from '@/sse-db-event/states/activeQueryListenersState'; +import { sseClientState } from '@/sse-db-event/states/sseClientState'; +import { getSnapshotValue } from '@/ui/utilities/state/utils/getSnapshotValue'; +import { isNonEmptyArray } from '@sniptt/guards'; +import { createClient } from 'graphql-sse'; +import { useEffect } from 'react'; +import { useRecoilCallback, useRecoilState } from 'recoil'; +import { isDefined } from 'twenty-shared/utils'; +import { REACT_APP_SERVER_BASE_URL } from '~/config'; +import { getIsDevelopmentEnvironment } from '~/utils/getIsDevelopmentEnvironment'; + +import { sleep } from '~/utils/sleep'; + +export const SSEClientEffect = () => { + const isLoggedIn = useIsLogged(); + const [sseClient, setSseClient] = useRecoilState(sseClientState); + + const handleSSEClientConnected = useRecoilCallback( + ({ snapshot, set }) => + () => { + const currentActiveQueryListeners = getSnapshotValue( + snapshot, + activeQueryListenersState, + ); + + if (isNonEmptyArray(currentActiveQueryListeners)) { + set(activeQueryListenersState, []); + } + }, + [], + ); + + useEffect(() => { + if (isLoggedIn && !isDefined(sseClient)) { + const tokenPair = getTokenPair(); + const token = tokenPair?.accessOrWorkspaceAgnosticToken?.token; + + const newSseClient = createClient({ + url: `${REACT_APP_SERVER_BASE_URL}/graphql`, + headers: { + Authorization: token ? `Bearer ${token}` : '', + }, + on: { + connected: handleSSEClientConnected, + }, + retryAttempts: Infinity, + retry: async () => { + const randomWaitTimeInMsToSpaceAllClientsReconnection = Math.round( + Math.random() * SSE_CONNECTION_RETRY_MAX_WAIT_TIME_IN_MS, + ); + + const isDevelopmentEnvironment = getIsDevelopmentEnvironment(); + + const waitTimeInMs = isDevelopmentEnvironment + ? SSE_CONNECTION_RETRY_WAIT_TIME_IN_MS_FOR_DEV_MODE + : randomWaitTimeInMsToSpaceAllClientsReconnection; + + await sleep(waitTimeInMs); + }, + }); + + setSseClient(newSseClient); + } + }, [handleSSEClientConnected, isLoggedIn, setSseClient, sseClient]); + + return null; +}; diff --git a/packages/twenty-front/src/modules/sse-db-event/components/SSEEventStreamEffect.tsx b/packages/twenty-front/src/modules/sse-db-event/components/SSEEventStreamEffect.tsx index 2b7c7441ae..4f878ed91b 100644 --- a/packages/twenty-front/src/modules/sse-db-event/components/SSEEventStreamEffect.tsx +++ b/packages/twenty-front/src/modules/sse-db-event/components/SSEEventStreamEffect.tsx @@ -1,22 +1,82 @@ -import { useObjectMetadataItems } from '@/object-metadata/hooks/useObjectMetadataItems'; -import { SseClientContext } from '@/sse-db-event/contexts/SseClientContext'; -import { useSubscribeToSseEventStream } from '@/sse-db-event/hooks/useSubscribeToSseEventStream'; -import { useContext, useEffect } from 'react'; +import { useIsLogged } from '@/auth/hooks/useIsLogged'; +import { currentUserState } from '@/auth/states/currentUserState'; +import { objectMetadataItemsState } from '@/object-metadata/states/objectMetadataItemsState'; +import { useTriggerEventStreamCreation } from '@/sse-db-event/hooks/useTriggerEventStreamCreation'; +import { useTriggerEventStreamDestroy } from '@/sse-db-event/hooks/useTriggerEventStreamDestroy'; +import { isCreatingSseEventStreamState } from '@/sse-db-event/states/isCreatingSseEventStreamState'; +import { isDestroyingEventStreamState } from '@/sse-db-event/states/isDestroyingEventStreamState'; +import { shouldDestroyEventStreamState } from '@/sse-db-event/states/shouldDestroyEventStreamState'; +import { sseClientState } from '@/sse-db-event/states/sseClientState'; +import { sseEventStreamIdState } from '@/sse-db-event/states/sseEventStreamIdState'; +import { useIsFeatureEnabled } from '@/workspace/hooks/useIsFeatureEnabled'; +import { isNonEmptyArray } from '@apollo/client/utilities'; +import { isNonEmptyString } from '@sniptt/guards'; +import { useEffect } from 'react'; +import { useRecoilValue } from 'recoil'; import { isDefined } from 'twenty-shared/utils'; +import { FeatureFlagKey, OnboardingStatus } from '~/generated/graphql'; export const SSEEventStreamEffect = () => { - const sseClient = useContext(SseClientContext); - const { objectMetadataItems } = useObjectMetadataItems(); + const sseClient = useRecoilValue(sseClientState); + const objectMetadataItems = useRecoilValue(objectMetadataItemsState); - const { subscribeToSseEventStream } = useSubscribeToSseEventStream(); + const sseEventStreamId = useRecoilValue(sseEventStreamIdState); + const isCreatingSseEventStream = useRecoilValue( + isCreatingSseEventStreamState, + ); + const shouldDestroyEventStream = useRecoilValue( + shouldDestroyEventStreamState, + ); + const isDestroyingEventStream = useRecoilValue(isDestroyingEventStreamState); + + const isLoggedIn = useIsLogged(); + const isSseDbEventsEnabled = useIsFeatureEnabled( + FeatureFlagKey.IS_SSE_DB_EVENTS_ENABLED, + ); + const currentUser = useRecoilValue(currentUserState); + + const { triggerEventStreamCreation } = useTriggerEventStreamCreation(); + const { triggerEventStreamDestroy } = useTriggerEventStreamDestroy(); useEffect(() => { - if (!isDefined(sseClient) || objectMetadataItems.length === 0) { - return; - } + const isSseClientAvailabble = + isDefined(sseClient) && + !isCreatingSseEventStream && + !isDestroyingEventStream; - subscribeToSseEventStream(sseClient); - }, [sseClient, subscribeToSseEventStream, objectMetadataItems]); + const willCreateEventStream = + isSseClientAvailabble && + isLoggedIn && + isSseDbEventsEnabled && + isDefined(currentUser) && + currentUser.onboardingStatus === OnboardingStatus.COMPLETED && + !shouldDestroyEventStream && + !isNonEmptyString(sseEventStreamId) && + isNonEmptyArray(objectMetadataItems); + + const willDestroyEventStream = + isSseClientAvailabble && + isNonEmptyString(sseEventStreamId) && + shouldDestroyEventStream; + + if (willDestroyEventStream) { + triggerEventStreamDestroy(); + } else if (willCreateEventStream) { + triggerEventStreamCreation(sseClient); + } + }, [ + isCreatingSseEventStream, + triggerEventStreamCreation, + isLoggedIn, + currentUser, + isSseDbEventsEnabled, + isDestroyingEventStream, + triggerEventStreamDestroy, + sseClient, + shouldDestroyEventStream, + sseEventStreamId, + objectMetadataItems, + ]); return null; }; diff --git a/packages/twenty-front/src/modules/sse-db-event/components/SSEProvider.tsx b/packages/twenty-front/src/modules/sse-db-event/components/SSEProvider.tsx index 74e8738af0..478c044db8 100644 --- a/packages/twenty-front/src/modules/sse-db-event/components/SSEProvider.tsx +++ b/packages/twenty-front/src/modules/sse-db-event/components/SSEProvider.tsx @@ -1,9 +1,12 @@ +import { getTokenPair } from '@/apollo/utils/getTokenPair'; +import { SSEClientEffect } from '@/sse-db-event/components/SSEClientEffect'; import { SSEEventStreamEffect } from '@/sse-db-event/components/SSEEventStreamEffect'; import { SSEQuerySubscribeEffect } from '@/sse-db-event/components/SSEQuerySubscribeEffect'; import { SseClientContext } from '@/sse-db-event/contexts/SseClientContext'; -import { useSseClient } from '@/sse-db-event/hooks/useSseClient.util'; import { useIsFeatureEnabled } from '@/workspace/hooks/useIsFeatureEnabled'; +import { createClient } from 'graphql-sse'; import { type ReactNode } from 'react'; +import { REACT_APP_SERVER_BASE_URL } from '~/config'; import { FeatureFlagKey } from '~/generated/graphql'; type SSEProviderProps = { @@ -14,21 +17,31 @@ export const SSEProvider = ({ children }: SSEProviderProps) => { const isSseDbEventsEnabled = useIsFeatureEnabled( FeatureFlagKey.IS_SSE_DB_EVENTS_ENABLED, ); - const { sseClient } = useSseClient(); if (!isSseDbEventsEnabled) { + const tokenPair = getTokenPair(); + const token = tokenPair?.accessOrWorkspaceAgnosticToken?.token; + + const sseClient = createClient({ + url: `${REACT_APP_SERVER_BASE_URL}/graphql`, + headers: { + Authorization: token ? `Bearer ${token}` : '', + }, + }); + return ( - + {children} ); } return ( - + <> + {children} - + ); }; diff --git a/packages/twenty-front/src/modules/sse-db-event/components/SSEQuerySubscribeEffect.tsx b/packages/twenty-front/src/modules/sse-db-event/components/SSEQuerySubscribeEffect.tsx index 1b8bdc70f8..12cced5ece 100644 --- a/packages/twenty-front/src/modules/sse-db-event/components/SSEQuerySubscribeEffect.tsx +++ b/packages/twenty-front/src/modules/sse-db-event/components/SSEQuerySubscribeEffect.tsx @@ -3,9 +3,10 @@ import { ADD_QUERY_TO_EVENT_STREAM_MUTATION } from '@/sse-db-event/graphql/mutat import { REMOVE_QUERY_FROM_EVENT_STREAM_MUTATION } from '@/sse-db-event/graphql/mutations/RemoveQueryFromEventStreamMutation'; import { activeQueryListenersState } from '@/sse-db-event/states/activeQueryListenersState'; import { requiredQueryListenersState } from '@/sse-db-event/states/requiredQueryListenersState'; +import { shouldDestroyEventStreamState } from '@/sse-db-event/states/shouldDestroyEventStreamState'; import { sseEventStreamIdState } from '@/sse-db-event/states/sseEventStreamIdState'; import { getSnapshotValue } from '@/ui/utilities/state/utils/getSnapshotValue'; -import { useMutation } from '@apollo/client'; +import { ApolloError, useMutation } from '@apollo/client'; import { isNonEmptyString } from '@sniptt/guards'; import { useEffect } from 'react'; import { useRecoilCallback, useRecoilValue } from 'recoil'; @@ -69,27 +70,47 @@ export const SSEQuerySubscribeEffect = () => { ), ); - for (const queryListenerToAdd of queryListenersToAdd) { - await addQueryToEventStream({ - variables: { - input: { - eventStreamId: sseEventStreamId, - queryId: queryListenerToAdd.queryId, - operationSignature: queryListenerToAdd.operationSignature, + try { + for (const queryListenerToAdd of queryListenersToAdd) { + await addQueryToEventStream({ + variables: { + input: { + eventStreamId: sseEventStreamId, + queryId: queryListenerToAdd.queryId, + operationSignature: queryListenerToAdd.operationSignature, + }, }, - }, - }); - } + }); + } - for (const queryListenerToRemove of queryListenersToRemove) { - await removeQueryFromEventStream({ - variables: { - input: { - eventStreamId: sseEventStreamId, - queryId: queryListenerToRemove.queryId, + for (const queryListenerToRemove of queryListenersToRemove) { + await removeQueryFromEventStream({ + variables: { + input: { + eventStreamId: sseEventStreamId, + queryId: queryListenerToRemove.queryId, + }, }, - }, - }); + }); + } + } catch (error) { + if (error instanceof ApolloError) { + const subCode = error.graphQLErrors[0]?.extensions?.subCode; + + switch (subCode) { + case 'EVENT_STREAM_DOES_NOT_EXIST': + case 'EVENT_STREAM_ALREADY_EXISTS': { + set(activeQueryListenersState, []); + set(shouldDestroyEventStreamState, true); + return; + } + default: { + throw new Error( + `Unhandled error for event stream: ${error.message}`, + ); + } + } + } } set(activeQueryListenersState, requiredQueryListeners); diff --git a/packages/twenty-front/src/modules/sse-db-event/constants/SseConnectionRetryMaxWaitTimeInMs.ts b/packages/twenty-front/src/modules/sse-db-event/constants/SseConnectionRetryMaxWaitTimeInMs.ts new file mode 100644 index 0000000000..7429b645ea --- /dev/null +++ b/packages/twenty-front/src/modules/sse-db-event/constants/SseConnectionRetryMaxWaitTimeInMs.ts @@ -0,0 +1 @@ +export const SSE_CONNECTION_RETRY_MAX_WAIT_TIME_IN_MS = 2 * 60 * 1_000; diff --git a/packages/twenty-front/src/modules/sse-db-event/constants/SseConnectionRetryWaitTimeInMsForDevMode.ts b/packages/twenty-front/src/modules/sse-db-event/constants/SseConnectionRetryWaitTimeInMsForDevMode.ts new file mode 100644 index 0000000000..3c634a23e0 --- /dev/null +++ b/packages/twenty-front/src/modules/sse-db-event/constants/SseConnectionRetryWaitTimeInMsForDevMode.ts @@ -0,0 +1 @@ +export const SSE_CONNECTION_RETRY_WAIT_TIME_IN_MS_FOR_DEV_MODE = 1_000; diff --git a/packages/twenty-front/src/modules/sse-db-event/hooks/__tests__/useOnDbEvent.test.tsx b/packages/twenty-front/src/modules/sse-db-event/hooks/__tests__/useOnDbEvent.test.tsx deleted file mode 100644 index 147fbcfc35..0000000000 --- a/packages/twenty-front/src/modules/sse-db-event/hooks/__tests__/useOnDbEvent.test.tsx +++ /dev/null @@ -1,298 +0,0 @@ -import { renderHook } from '@testing-library/react'; -import { createClient } from 'graphql-sse'; - -import { useOnDbEvent } from '@/sse-db-event/hooks/useOnDbEvent'; -import { DatabaseEventAction } from '~/generated/graphql'; -import { getTokenPair } from '~/modules/apollo/utils/getTokenPair'; - -jest.mock('~/modules/apollo/utils/getTokenPair'); -jest.mock('graphql-sse'); - -const mockGetTokenPair = getTokenPair as jest.MockedFunction< - typeof getTokenPair ->; -const mockCreateClient = createClient as jest.MockedFunction< - typeof createClient ->; - -// Mock environment variable -const mockServerBaseUrl = 'http://localhost:3000'; -jest.mock('~/config', () => ({ - REACT_APP_SERVER_BASE_URL: 'http://localhost:3000', -})); - -describe('useOnDbEvent', () => { - const mockUnsubscribe = jest.fn(); - const mockClient = { - subscribe: jest.fn(() => mockUnsubscribe), - dispose: jest.fn(), - }; - - beforeEach(() => { - jest.clearAllMocks(); - mockCreateClient.mockReturnValue(mockClient as any); - }); - - describe('token safety checks', () => { - it('should handle undefined tokenPair gracefully', () => { - mockGetTokenPair.mockReturnValue(undefined); - - renderHook(() => - useOnDbEvent({ - input: { - objectNameSingular: 'test', - action: DatabaseEventAction.CREATED, - }, - onData: jest.fn(), - }), - ); - - expect(mockCreateClient).toHaveBeenCalledWith({ - url: `${mockServerBaseUrl}/graphql`, - headers: { - Authorization: '', - }, - }); - }); - - it('should handle tokenPair with undefined accessOrWorkspaceAgnosticToken gracefully', () => { - mockGetTokenPair.mockReturnValue({ - accessOrWorkspaceAgnosticToken: undefined, - refreshToken: { - token: 'refresh-token', - expiresAt: '2024-01-02T00:00:00Z', - }, - } as any); - - renderHook(() => - useOnDbEvent({ - input: { - objectNameSingular: 'test', - action: DatabaseEventAction.CREATED, - }, - onData: jest.fn(), - }), - ); - - expect(mockCreateClient).toHaveBeenCalledWith({ - url: `${mockServerBaseUrl}/graphql`, - headers: { - Authorization: '', - }, - }); - }); - - it('should handle tokenPair with accessOrWorkspaceAgnosticToken but undefined token gracefully', () => { - mockGetTokenPair.mockReturnValue({ - accessOrWorkspaceAgnosticToken: { - token: undefined, - expiresAt: '2024-01-01T00:00:00Z', - }, - refreshToken: { - token: 'refresh-token', - expiresAt: '2024-01-02T00:00:00Z', - }, - } as any); - - renderHook(() => - useOnDbEvent({ - input: { - objectNameSingular: 'test', - action: DatabaseEventAction.CREATED, - }, - onData: jest.fn(), - }), - ); - - expect(mockCreateClient).toHaveBeenCalledWith({ - url: `${mockServerBaseUrl}/graphql`, - headers: { - Authorization: '', - }, - }); - }); - - it('should handle tokenPair with null token gracefully', () => { - mockGetTokenPair.mockReturnValue({ - accessOrWorkspaceAgnosticToken: { - token: null, - expiresAt: '2024-01-01T00:00:00Z', - }, - refreshToken: { - token: 'refresh-token', - expiresAt: '2024-01-02T00:00:00Z', - }, - } as any); - - renderHook(() => - useOnDbEvent({ - input: { - objectNameSingular: 'test', - action: DatabaseEventAction.CREATED, - }, - onData: jest.fn(), - }), - ); - - expect(mockCreateClient).toHaveBeenCalledWith({ - url: `${mockServerBaseUrl}/graphql`, - headers: { - Authorization: '', - }, - }); - }); - - it('should properly set authorization header when token is valid', () => { - const validToken = 'valid-access-token'; - mockGetTokenPair.mockReturnValue({ - accessOrWorkspaceAgnosticToken: { - token: validToken, - expiresAt: '2024-01-01T00:00:00Z', - }, - refreshToken: { - token: 'refresh-token', - expiresAt: '2024-01-02T00:00:00Z', - }, - }); - - renderHook(() => - useOnDbEvent({ - input: { - objectNameSingular: 'test', - action: DatabaseEventAction.CREATED, - }, - onData: jest.fn(), - }), - ); - - expect(mockCreateClient).toHaveBeenCalledWith({ - url: `${mockServerBaseUrl}/graphql`, - headers: { - Authorization: `Bearer ${validToken}`, - }, - }); - }); - - it('should handle empty string token gracefully', () => { - mockGetTokenPair.mockReturnValue({ - accessOrWorkspaceAgnosticToken: { - token: '', - expiresAt: '2024-01-01T00:00:00Z', - }, - refreshToken: { - token: 'refresh-token', - expiresAt: '2024-01-02T00:00:00Z', - }, - }); - - renderHook(() => - useOnDbEvent({ - input: { - objectNameSingular: 'test', - action: DatabaseEventAction.CREATED, - }, - onData: jest.fn(), - }), - ); - - expect(mockCreateClient).toHaveBeenCalledWith({ - url: `${mockServerBaseUrl}/graphql`, - headers: { - Authorization: '', - }, - }); - }); - }); - - describe('basic functionality', () => { - const validToken = 'test-token'; - - beforeEach(() => { - mockGetTokenPair.mockReturnValue({ - accessOrWorkspaceAgnosticToken: { - token: validToken, - expiresAt: '2024-01-01T00:00:00Z', - }, - refreshToken: { - token: 'refresh-token', - expiresAt: '2024-01-02T00:00:00Z', - }, - }); - }); - - it('should create SSE client with correct URL and headers', () => { - renderHook(() => - useOnDbEvent({ - input: { - objectNameSingular: 'test', - action: DatabaseEventAction.CREATED, - }, - onData: jest.fn(), - }), - ); - - expect(mockCreateClient).toHaveBeenCalledWith({ - url: `${mockServerBaseUrl}/graphql`, - headers: { - Authorization: `Bearer ${validToken}`, - }, - }); - }); - - it('should not subscribe when skip is true', () => { - renderHook(() => - useOnDbEvent({ - input: { - objectNameSingular: 'test', - action: DatabaseEventAction.CREATED, - }, - onData: jest.fn(), - skip: true, - }), - ); - - expect(mockClient.subscribe).not.toHaveBeenCalled(); - }); - - it('should subscribe when skip is false or undefined', () => { - const mockOnData = jest.fn(); - - renderHook(() => - useOnDbEvent({ - input: { - objectNameSingular: 'test', - action: DatabaseEventAction.CREATED, - }, - onData: mockOnData, - skip: false, - }), - ); - - expect(mockClient.subscribe).toHaveBeenCalled(); - }); - - it('should pass correct parameters to subscription', () => { - const mockOnData = jest.fn(); - - renderHook(() => - useOnDbEvent({ - input: { - objectNameSingular: 'person', - action: DatabaseEventAction.CREATED, - }, - onData: mockOnData, - }), - ); - - expect(mockClient.subscribe).toHaveBeenCalledWith( - expect.objectContaining({ - query: expect.stringContaining('subscription'), - }), - expect.objectContaining({ - next: expect.any(Function), - error: expect.any(Function), - }), - ); - }); - }); -}); diff --git a/packages/twenty-front/src/modules/sse-db-event/hooks/useOnDbEvent.ts b/packages/twenty-front/src/modules/sse-db-event/hooks/useOnDbEvent.ts index 22acb82e0e..34e08cbe3d 100644 --- a/packages/twenty-front/src/modules/sse-db-event/hooks/useOnDbEvent.ts +++ b/packages/twenty-front/src/modules/sse-db-event/hooks/useOnDbEvent.ts @@ -1,6 +1,6 @@ +import { SseClientContext } from '@/sse-db-event/contexts/SseClientContext'; import { ON_DB_EVENT } from '@/sse-db-event/graphql/subscriptions/onDbEvent'; -import { useSseClient } from '@/sse-db-event/hooks/useSseClient.util'; -import { useEffect } from 'react'; +import { useContext, useEffect } from 'react'; import { type Subscription, type SubscriptionOnDbEventArgs, @@ -20,7 +20,7 @@ export const useOnDbEvent = ({ input, skip = false, }: OnDbEventArgs) => { - const { sseClient } = useSseClient(); + const sseClient = useContext(SseClientContext); useEffect(() => { if (skip === true) { @@ -29,7 +29,8 @@ export const useOnDbEvent = ({ const next = (value: { data: Subscription }) => onData?.(value.data); const error = (err: unknown) => onError?.(err); const complete = () => onComplete?.(); - const unsubscribe = sseClient.subscribe( + + const unsubscribe = sseClient?.subscribe( { query: ON_DB_EVENT.loc?.source.body || '', variables: { input }, @@ -42,7 +43,7 @@ export const useOnDbEvent = ({ ); return () => { - unsubscribe(); + unsubscribe?.(); }; }, [input, onComplete, onData, onError, skip, sseClient]); }; diff --git a/packages/twenty-front/src/modules/sse-db-event/hooks/useSubscribeToSseEventStream.ts b/packages/twenty-front/src/modules/sse-db-event/hooks/useTriggerEventStreamCreation.ts similarity index 55% rename from packages/twenty-front/src/modules/sse-db-event/hooks/useSubscribeToSseEventStream.ts rename to packages/twenty-front/src/modules/sse-db-event/hooks/useTriggerEventStreamCreation.ts index bb37966bc2..177ceee3e4 100644 --- a/packages/twenty-front/src/modules/sse-db-event/hooks/useSubscribeToSseEventStream.ts +++ b/packages/twenty-front/src/modules/sse-db-event/hooks/useTriggerEventStreamCreation.ts @@ -1,40 +1,65 @@ import { ON_EVENT_SUBSCRIPTION } from '@/sse-db-event/graphql/subscriptions/OnEventSubscription'; import { useDispatchObjectRecordEventsFromSseToBrowserEvents } from '@/sse-db-event/hooks/useDispatchObjectRecordEventsFromSseToBrowserEvents'; import { useTriggerOptimisticEffectFromSseEvents } from '@/sse-db-event/hooks/useTriggerOptimisticEffectFromSseEvents'; +import { disposeFunctionForEventStreamState } from '@/sse-db-event/states/disposeFunctionByEventStreamMapState'; +import { isCreatingSseEventStreamState } from '@/sse-db-event/states/isCreatingSseEventStreamState'; +import { isDestroyingEventStreamState } from '@/sse-db-event/states/isDestroyingEventStreamState'; +import { shouldDestroyEventStreamState } from '@/sse-db-event/states/shouldDestroyEventStreamState'; import { sseEventStreamIdState } from '@/sse-db-event/states/sseEventStreamIdState'; import { getSnapshotValue } from '@/ui/utilities/state/utils/getSnapshotValue'; import { captureException } from '@sentry/react'; import { isNonEmptyString } from '@sniptt/guards'; import { print, type ExecutionResult } from 'graphql'; import { type Client } from 'graphql-sse'; -import { useRecoilCallback } from 'recoil'; + +import { useRecoilCallback, useSetRecoilState } from 'recoil'; +import { isDefined } from 'twenty-shared/utils'; import { v4 } from 'uuid'; import { type EventSubscription } from '~/generated/graphql'; -export const useSubscribeToSseEventStream = () => { +export const useTriggerEventStreamCreation = () => { + const setIsCreatingSseEventStream = useSetRecoilState( + isCreatingSseEventStreamState, + ); + const { dispatchObjectRecordEventsFromSseToBrowserEvents } = useDispatchObjectRecordEventsFromSseToBrowserEvents(); const { triggerOptimisticEffectFromSseEvents } = useTriggerOptimisticEffectFromSseEvents(); - const subscribeToSseEventStream = useRecoilCallback( - ({ set, snapshot }) => - (sseClientConnected: Client) => { + const triggerEventStreamCreation = useRecoilCallback( + ({ snapshot, set }) => + (sseClient: Client) => { + const isCreatingSseEventStream = snapshot + .getLoadable(isCreatingSseEventStreamState) + .getValue(); + + const isDestroyingEventStream = snapshot + .getLoadable(isDestroyingEventStreamState) + .getValue(); + const currentSseEventStreamId = getSnapshotValue( snapshot, sseEventStreamIdState, ); - if (isNonEmptyString(currentSseEventStreamId)) { + if ( + isCreatingSseEventStream || + isDestroyingEventStream || + !isDefined(sseClient) || + isNonEmptyString(currentSseEventStreamId) + ) { return; } + setIsCreatingSseEventStream(true); + const newSseEventStreamId = v4(); set(sseEventStreamIdState, newSseEventStreamId); - sseClientConnected.subscribe( + const dispose = sseClient.subscribe( { query: print(ON_EVENT_SUBSCRIPTION), variables: { @@ -69,15 +94,37 @@ export const useSubscribeToSseEventStream = () => { }, complete: () => {}, }, + { + message: ({ data, event }) => { + if (event === 'next') { + if (isDefined(data?.errors)) { + const subCode = data.errors[0]?.extensions?.subCode; + + switch (subCode) { + case 'EVENT_STREAM_ALREADY_EXISTS': { + set(shouldDestroyEventStreamState, true); + break; + } + } + } + } + }, + }, ); + + set(disposeFunctionForEventStreamState, { dispose }); + + setIsCreatingSseEventStream(false); }, [ - triggerOptimisticEffectFromSseEvents, dispatchObjectRecordEventsFromSseToBrowserEvents, + setIsCreatingSseEventStream, + + triggerOptimisticEffectFromSseEvents, ], ); return { - subscribeToSseEventStream, + triggerEventStreamCreation, }; }; diff --git a/packages/twenty-front/src/modules/sse-db-event/hooks/useTriggerEventStreamDestroy.ts b/packages/twenty-front/src/modules/sse-db-event/hooks/useTriggerEventStreamDestroy.ts new file mode 100644 index 0000000000..b12432d5af --- /dev/null +++ b/packages/twenty-front/src/modules/sse-db-event/hooks/useTriggerEventStreamDestroy.ts @@ -0,0 +1,55 @@ +import { disposeFunctionForEventStreamState } from '@/sse-db-event/states/disposeFunctionByEventStreamMapState'; +import { isCreatingSseEventStreamState } from '@/sse-db-event/states/isCreatingSseEventStreamState'; +import { isDestroyingEventStreamState } from '@/sse-db-event/states/isDestroyingEventStreamState'; +import { shouldDestroyEventStreamState } from '@/sse-db-event/states/shouldDestroyEventStreamState'; +import { sseEventStreamIdState } from '@/sse-db-event/states/sseEventStreamIdState'; +import { isNonEmptyString } from '@sniptt/guards'; +import { useRecoilCallback, useSetRecoilState } from 'recoil'; + +export const useTriggerEventStreamDestroy = () => { + const setIsDestroyingEventStream = useSetRecoilState( + isDestroyingEventStreamState, + ); + + const triggerEventStreamDestroy = useRecoilCallback( + ({ snapshot, set }) => + () => { + const isDestroyingEventStream = snapshot + .getLoadable(isDestroyingEventStreamState) + .getValue(); + + const isCreatingSseEventStream = snapshot + .getLoadable(isCreatingSseEventStreamState) + .getValue(); + + if (isDestroyingEventStream || isCreatingSseEventStream) { + return; + } + + setIsDestroyingEventStream(true); + + const eventStreamId = snapshot + .getLoadable(sseEventStreamIdState) + .getValue(); + + const disposeFunctionForEventStream = snapshot + .getLoadable(disposeFunctionForEventStreamState) + .getValue(); + + if (isNonEmptyString(eventStreamId)) { + disposeFunctionForEventStream?.dispose(); + + set(sseEventStreamIdState, null); + set(disposeFunctionForEventStreamState, null); + set(shouldDestroyEventStreamState, false); + } + + setIsDestroyingEventStream(false); + }, + [setIsDestroyingEventStream], + ); + + return { + triggerEventStreamDestroy, + }; +}; diff --git a/packages/twenty-front/src/modules/sse-db-event/states/disposeFunctionByEventStreamMapState.ts b/packages/twenty-front/src/modules/sse-db-event/states/disposeFunctionByEventStreamMapState.ts new file mode 100644 index 0000000000..5c0d5dfde5 --- /dev/null +++ b/packages/twenty-front/src/modules/sse-db-event/states/disposeFunctionByEventStreamMapState.ts @@ -0,0 +1,8 @@ +import { createState } from 'twenty-ui/utilities'; + +export const disposeFunctionForEventStreamState = createState<{ + dispose: () => void; +} | null>({ + key: 'disposeFunctionForEventStreamState', + defaultValue: null, +}); diff --git a/packages/twenty-front/src/modules/sse-db-event/states/isCreatingSseEventStreamState.ts b/packages/twenty-front/src/modules/sse-db-event/states/isCreatingSseEventStreamState.ts new file mode 100644 index 0000000000..6209ef4c8a --- /dev/null +++ b/packages/twenty-front/src/modules/sse-db-event/states/isCreatingSseEventStreamState.ts @@ -0,0 +1,6 @@ +import { createState } from 'twenty-ui/utilities'; + +export const isCreatingSseEventStreamState = createState({ + key: 'isCreatingSseEventStreamState', + defaultValue: false, +}); diff --git a/packages/twenty-front/src/modules/sse-db-event/states/isDestroyingEventStreamState.ts b/packages/twenty-front/src/modules/sse-db-event/states/isDestroyingEventStreamState.ts new file mode 100644 index 0000000000..939d059b0a --- /dev/null +++ b/packages/twenty-front/src/modules/sse-db-event/states/isDestroyingEventStreamState.ts @@ -0,0 +1,6 @@ +import { createState } from 'twenty-ui/utilities'; + +export const isDestroyingEventStreamState = createState({ + key: 'isDestroyingEventStreamState', + defaultValue: false, +}); diff --git a/packages/twenty-front/src/modules/sse-db-event/states/shouldDestroyEventStreamState.ts b/packages/twenty-front/src/modules/sse-db-event/states/shouldDestroyEventStreamState.ts new file mode 100644 index 0000000000..265d24981d --- /dev/null +++ b/packages/twenty-front/src/modules/sse-db-event/states/shouldDestroyEventStreamState.ts @@ -0,0 +1,6 @@ +import { createState } from 'twenty-ui/utilities'; + +export const shouldDestroyEventStreamState = createState({ + key: 'shouldDestroyEventStreamState', + defaultValue: false, +}); diff --git a/packages/twenty-front/src/modules/sse-db-event/states/sseClientState.ts b/packages/twenty-front/src/modules/sse-db-event/states/sseClientState.ts new file mode 100644 index 0000000000..2b40354bba --- /dev/null +++ b/packages/twenty-front/src/modules/sse-db-event/states/sseClientState.ts @@ -0,0 +1,7 @@ +import { type Client } from 'graphql-sse'; +import { createState } from 'twenty-ui/utilities'; + +export const sseClientState = createState({ + key: 'sseClientState', + defaultValue: null, +}); diff --git a/packages/twenty-front/src/utils/getIsDevelopmentEnvironment.ts b/packages/twenty-front/src/utils/getIsDevelopmentEnvironment.ts new file mode 100644 index 0000000000..8900ae8f5b --- /dev/null +++ b/packages/twenty-front/src/utils/getIsDevelopmentEnvironment.ts @@ -0,0 +1,3 @@ +export const getIsDevelopmentEnvironment = () => { + return process.env.IS_DEV_ENV === 'true'; +}; diff --git a/packages/twenty-front/vite.config.ts b/packages/twenty-front/vite.config.ts index 89acf61b64..38ab9e49d0 100644 --- a/packages/twenty-front/vite.config.ts +++ b/packages/twenty-front/vite.config.ts @@ -264,6 +264,7 @@ export default defineConfig(({ command, mode }) => { 'process.env': { REACT_APP_SERVER_BASE_URL, IS_DEBUG_MODE, + IS_DEV_ENV: mode === 'development' ? 'true' : 'false', }, }, css: { diff --git a/packages/twenty-server/src/engine/core-modules/graphql/hooks/use-graphql-error-handler.hook.ts b/packages/twenty-server/src/engine/core-modules/graphql/hooks/use-graphql-error-handler.hook.ts index e05a44cb69..294b8cfbcf 100644 --- a/packages/twenty-server/src/engine/core-modules/graphql/hooks/use-graphql-error-handler.hook.ts +++ b/packages/twenty-server/src/engine/core-modules/graphql/hooks/use-graphql-error-handler.hook.ts @@ -79,6 +79,7 @@ export const useGraphQLErrorHandlerHook = < } return { + // TODO: define onSubscribe here to handle subscription errors too async onExecute({ args }) { const exceptionHandlerService = options.exceptionHandlerService; const rootOperation = args.document.definitions.find( diff --git a/packages/twenty-server/src/engine/subscriptions/event-stream.exception.ts b/packages/twenty-server/src/engine/subscriptions/event-stream.exception.ts index b82f2c34e7..9d2dd1927c 100644 --- a/packages/twenty-server/src/engine/subscriptions/event-stream.exception.ts +++ b/packages/twenty-server/src/engine/subscriptions/event-stream.exception.ts @@ -7,6 +7,7 @@ import { CustomException } from 'src/utils/custom-exception'; export enum EventStreamExceptionCode { EVENT_STREAM_ALREADY_EXISTS = 'EVENT_STREAM_ALREADY_EXISTS', NOT_AUTHORIZED = 'NOT_AUTHORIZED', + EVENT_STREAM_DOES_NOT_EXIST = 'EVENT_STREAM_DOES_NOT_EXIST', } const getEventStreamExceptionUserFriendlyMessage = ( @@ -14,6 +15,7 @@ const getEventStreamExceptionUserFriendlyMessage = ( ) => { switch (code) { case EventStreamExceptionCode.EVENT_STREAM_ALREADY_EXISTS: + case EventStreamExceptionCode.EVENT_STREAM_DOES_NOT_EXIST: return msg`Failed to receive real time updates.`; case EventStreamExceptionCode.NOT_AUTHORIZED: return msg`You are not authorized to perform this action.`; diff --git a/packages/twenty-server/src/engine/subscriptions/event-stream.service.ts b/packages/twenty-server/src/engine/subscriptions/event-stream.service.ts index 8e0a5c1b53..57a6641843 100644 --- a/packages/twenty-server/src/engine/subscriptions/event-stream.service.ts +++ b/packages/twenty-server/src/engine/subscriptions/event-stream.service.ts @@ -157,23 +157,12 @@ export class EventStreamService implements OnModuleInit { } async isAuthorized({ - workspaceId, - eventStreamChannelId, authContext, + streamData, }: { - workspaceId: string; - eventStreamChannelId: string; authContext: SerializableAuthContext; + streamData: EventStreamData; }): Promise { - const streamData = await this.getStreamData( - workspaceId, - eventStreamChannelId, - ); - - if (!isDefined(streamData)) { - return false; - } - if (isDefined(authContext.userWorkspaceId)) { return ( streamData.authContext.userWorkspaceId === authContext.userWorkspaceId @@ -262,7 +251,7 @@ export class EventStreamService implements OnModuleInit { return `workspace:${workspaceId}:activeStreams`; } - private async getStreamData( + async getStreamData( workspaceId: string, eventStreamChannelId: string, ): Promise { diff --git a/packages/twenty-server/src/engine/workspace-event-emitter/workspace-event-emitter-exception.filter.ts b/packages/twenty-server/src/engine/workspace-event-emitter/workspace-event-emitter-exception.filter.ts new file mode 100644 index 0000000000..50eeff50cc --- /dev/null +++ b/packages/twenty-server/src/engine/workspace-event-emitter/workspace-event-emitter-exception.filter.ts @@ -0,0 +1,30 @@ +import { Catch } from '@nestjs/common'; +import { GqlExceptionFilter } from '@nestjs/graphql'; + +import { assertUnreachable } from 'twenty-shared/utils'; + +import { InternalServerError } from 'src/engine/core-modules/graphql/utils/graphql-errors.util'; +import { + EventStreamException, + EventStreamExceptionCode, +} from 'src/engine/subscriptions/event-stream.exception'; + +@Catch(EventStreamException) +export class WorkspaceEventEmitterExceptionFilter + implements GqlExceptionFilter +{ + catch(exception: EventStreamException) { + switch (exception.code) { + case EventStreamExceptionCode.EVENT_STREAM_ALREADY_EXISTS: + case EventStreamExceptionCode.EVENT_STREAM_DOES_NOT_EXIST: + case EventStreamExceptionCode.NOT_AUTHORIZED: + throw new InternalServerError(exception.message, { + subCode: exception.code, + userFriendlyMessage: exception.userFriendlyMessage, + }); + default: { + throw assertUnreachable(exception.code); + } + } + } +} 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 a231a4eb33..1a7b3d26dd 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 @@ -32,13 +32,17 @@ import { import { EventStreamService } from 'src/engine/subscriptions/event-stream.service'; import { SubscriptionService } from 'src/engine/subscriptions/subscription.service'; import { wrapAsyncIteratorWithLifecycle } from 'src/engine/workspace-event-emitter/utils/wrap-async-iterator-with-lifecycle'; +import { WorkspaceEventEmitterExceptionFilter } from 'src/engine/workspace-event-emitter/workspace-event-emitter-exception.filter'; import { eventStreamIdToChannelId } from './utils/get-channel-id-from-event-stream-id'; @Resolver() @UseGuards(WorkspaceAuthGuard, UserAuthGuard, NoPermissionGuard) @UsePipes(ResolverValidationPipe) -@UseFilters(PreventNestToAutoLogGraphqlErrorsFilter) +@UseFilters( + WorkspaceEventEmitterExceptionFilter, + PreventNestToAutoLogGraphqlErrorsFilter, +) export class WorkspaceEventEmitterResolver { constructor( private readonly subscriptionService: SubscriptionService, @@ -99,6 +103,18 @@ export class WorkspaceEventEmitterResolver { ) { const eventStreamChannelId = eventStreamIdToChannelId(eventStreamId); + const streamData = await this.eventStreamService.getStreamData( + workspace.id, + eventStreamChannelId, + ); + + if (isDefined(streamData)) { + throw new EventStreamException( + 'Event stream already exists', + EventStreamExceptionCode.EVENT_STREAM_ALREADY_EXISTS, + ); + } + await this.eventStreamService.createEventStream({ workspaceId: workspace.id, eventStreamChannelId, @@ -143,14 +159,24 @@ export class WorkspaceEventEmitterResolver { async addQueryToEventStream( @Args('input') input: AddQuerySubscriptionInput, @AuthWorkspace() workspace: WorkspaceEntity, + @AuthUser({ allowUndefined: true }) user: UserEntity | undefined, @AuthUserWorkspaceId() userWorkspaceId: string | undefined, @AuthApiKey() apiKey: ApiKeyEntity | undefined, ): Promise { const eventStreamChannelId = eventStreamIdToChannelId(input.eventStreamId); - - const isAuthorized = await this.eventStreamService.isAuthorized({ - workspaceId: workspace.id, + const streamData = await this.eventStreamService.getStreamData( + workspace.id, eventStreamChannelId, + ); + + if (!isDefined(streamData)) { + throw new EventStreamException( + 'Event stream does not exist', + EventStreamExceptionCode.EVENT_STREAM_DOES_NOT_EXIST, + ); + } + const isAuthorized = await this.eventStreamService.isAuthorized({ + streamData, authContext: { userWorkspaceId, apiKeyId: apiKey?.id, @@ -163,7 +189,6 @@ export class WorkspaceEventEmitterResolver { EventStreamExceptionCode.NOT_AUTHORIZED, ); } - await this.eventStreamService.addQuery({ workspaceId: workspace.id, eventStreamChannelId, @@ -178,14 +203,26 @@ export class WorkspaceEventEmitterResolver { async removeQueryFromEventStream( @Args('input') input: RemoveQueryFromEventStreamInput, @AuthWorkspace() workspace: WorkspaceEntity, + @AuthUser({ allowUndefined: true }) user: UserEntity | undefined, @AuthUserWorkspaceId() userWorkspaceId: string | undefined, @AuthApiKey() apiKey: ApiKeyEntity | undefined, ): Promise { const eventStreamChannelId = eventStreamIdToChannelId(input.eventStreamId); - const isAuthorized = await this.eventStreamService.isAuthorized({ - workspaceId: workspace.id, + const streamData = await this.eventStreamService.getStreamData( + workspace.id, eventStreamChannelId, + ); + + if (!isDefined(streamData)) { + throw new EventStreamException( + 'Event stream does not exist', + EventStreamExceptionCode.EVENT_STREAM_DOES_NOT_EXIST, + ); + } + + const isAuthorized = await this.eventStreamService.isAuthorized({ + streamData, authContext: { userWorkspaceId, apiKeyId: apiKey?.id,