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 e387c44831..1231811c76 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 @@ -29,12 +29,12 @@ export const SSEQuerySubscribeEffect = () => { const sseEventStreamReady = useAtomStateValue(sseEventStreamReadyState); const [addQueryToEventStream] = useMutation< - boolean, + { addQueryToEventStream: boolean }, { input: AddQuerySubscriptionInput } >(ADD_QUERY_TO_EVENT_STREAM_MUTATION); const [removeQueryFromEventStream] = useMutation< - void, + { removeQueryFromEventStream: boolean }, { input: RemoveQueryFromEventStreamInput } >(REMOVE_QUERY_FROM_EVENT_STREAM_MUTATION); @@ -66,7 +66,7 @@ export const SSEQuerySubscribeEffect = () => { try { for (const queryListenerToAdd of queryListenersToAdd) { - await addQueryToEventStream({ + const result = await addQueryToEventStream({ variables: { input: { eventStreamId: sseEventStreamId, @@ -75,10 +75,17 @@ export const SSEQuerySubscribeEffect = () => { }, }, }); + + if (result.data?.addQueryToEventStream === false) { + store.set(activeQueryListenersState.atom, []); + store.set(shouldDestroyEventStreamState.atom, true); + + return; + } } for (const queryListenerToRemove of queryListenersToRemove) { - await removeQueryFromEventStream({ + const result = await removeQueryFromEventStream({ variables: { input: { eventStreamId: sseEventStreamId, @@ -86,6 +93,13 @@ export const SSEQuerySubscribeEffect = () => { }, }, }); + + if (result.data?.removeQueryFromEventStream === false) { + store.set(activeQueryListenersState.atom, []); + store.set(shouldDestroyEventStreamState.atom, true); + + return; + } } } catch (error) { if (CombinedGraphQLErrors.is(error)) { @@ -99,6 +113,7 @@ export const SSEQuerySubscribeEffect = () => { ) { store.set(activeQueryListenersState.atom, []); store.set(shouldDestroyEventStreamState.atom, true); + return; } diff --git a/packages/twenty-front/src/modules/sse-db-event/utils/isGracefullyHandledEventStreamError.ts b/packages/twenty-front/src/modules/sse-db-event/utils/isGracefullyHandledEventStreamError.ts index a272170c11..981c4ad984 100644 --- a/packages/twenty-front/src/modules/sse-db-event/utils/isGracefullyHandledEventStreamError.ts +++ b/packages/twenty-front/src/modules/sse-db-event/utils/isGracefullyHandledEventStreamError.ts @@ -12,7 +12,6 @@ export const isGracefullyHandledEventStreamError = ({ } return ( - subCode === 'EVENT_STREAM_DOES_NOT_EXIST' || subCode === 'EVENT_STREAM_ALREADY_EXISTS' || subCode === 'NOT_AUTHORIZED' || code === 'UNAUTHENTICATED' || diff --git a/packages/twenty-server/src/engine/subscriptions/event-stream-exception.filter.ts b/packages/twenty-server/src/engine/subscriptions/event-stream-exception.filter.ts index 4474bbec54..0715f964f5 100644 --- a/packages/twenty-server/src/engine/subscriptions/event-stream-exception.filter.ts +++ b/packages/twenty-server/src/engine/subscriptions/event-stream-exception.filter.ts @@ -14,7 +14,6 @@ export class EventStreamExceptionFilter 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, 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 9d2dd1927c..b82f2c34e7 100644 --- a/packages/twenty-server/src/engine/subscriptions/event-stream.exception.ts +++ b/packages/twenty-server/src/engine/subscriptions/event-stream.exception.ts @@ -7,7 +7,6 @@ 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 = ( @@ -15,7 +14,6 @@ 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.resolver.ts b/packages/twenty-server/src/engine/subscriptions/event-stream.resolver.ts index 481bd1e953..e895ec8936 100644 --- a/packages/twenty-server/src/engine/subscriptions/event-stream.resolver.ts +++ b/packages/twenty-server/src/engine/subscriptions/event-stream.resolver.ts @@ -151,11 +151,9 @@ export class EventStreamResolver { ); if (!isDefined(streamData)) { - throw new EventStreamException( - 'Event stream does not exist', - EventStreamExceptionCode.EVENT_STREAM_DOES_NOT_EXIST, - ); + return false; } + const isAuthorized = await this.eventStreamService.isAuthorized({ streamData, authContext: { @@ -170,6 +168,7 @@ export class EventStreamResolver { EventStreamExceptionCode.NOT_AUTHORIZED, ); } + await this.eventStreamService.addQuery({ workspaceId: workspace.id, eventStreamChannelId, @@ -197,10 +196,7 @@ export class EventStreamResolver { ); if (!isDefined(streamData)) { - throw new EventStreamException( - 'Event stream does not exist', - EventStreamExceptionCode.EVENT_STREAM_DOES_NOT_EXIST, - ); + return false; } const isAuthorized = await this.eventStreamService.isAuthorized({