diff --git a/packages/twenty-front/src/modules/sse-db-event/components/SSEKeepAliveEffect.tsx b/packages/twenty-front/src/modules/sse-db-event/components/SSEKeepAliveEffect.tsx new file mode 100644 index 0000000000..fcba4785a3 --- /dev/null +++ b/packages/twenty-front/src/modules/sse-db-event/components/SSEKeepAliveEffect.tsx @@ -0,0 +1,43 @@ +import { SSE_LIVENESS_CHECK_INTERVAL_IN_MS } from '@/sse-db-event/constants/SseLivenessCheckIntervalInMs'; +import { SSE_LIVENESS_TIMEOUT_IN_MS } from '@/sse-db-event/constants/SseLivenessTimeoutInMs'; +import { activeQueryListenersState } from '@/sse-db-event/states/activeQueryListenersState'; +import { lastSseEventReceivedTimestampState } from '@/sse-db-event/states/lastSseEventReceivedTimestampState'; +import { shouldDestroyEventStreamState } from '@/sse-db-event/states/shouldDestroyEventStreamState'; +import { sseEventStreamIdState } from '@/sse-db-event/states/sseEventStreamIdState'; +import { sseEventStreamReadyState } from '@/sse-db-event/states/sseEventStreamReadyState'; +import { useAtomStateValue } from '@/ui/utilities/state/jotai/hooks/useAtomStateValue'; +import { isNonEmptyString } from '@sniptt/guards'; +import { useStore } from 'jotai'; +import { useEffect } from 'react'; +import { isDefined } from 'twenty-shared/utils'; + +export const SSEKeepAliveEffect = () => { + const store = useStore(); + const sseEventStreamReady = useAtomStateValue(sseEventStreamReadyState); + const sseEventStreamId = useAtomStateValue(sseEventStreamIdState); + + useEffect(() => { + if (!sseEventStreamReady || !isNonEmptyString(sseEventStreamId)) { + return; + } + + const interval = setInterval(() => { + const lastTimestamp = store.get(lastSseEventReceivedTimestampState.atom); + + if (!isDefined(lastTimestamp)) { + return; + } + + const timeSinceLastEvent = Date.now() - lastTimestamp; + + if (timeSinceLastEvent > SSE_LIVENESS_TIMEOUT_IN_MS) { + store.set(activeQueryListenersState.atom, []); + store.set(shouldDestroyEventStreamState.atom, true); + } + }, SSE_LIVENESS_CHECK_INTERVAL_IN_MS); + + return () => clearInterval(interval); + }, [sseEventStreamReady, sseEventStreamId, store]); + + 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 2f352ff729..83df8bd327 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,6 +1,7 @@ import { MetadataStoreSSEEffect } from '@/metadata-store/effect-components/MetadataStoreSSEEffect'; import { SSEClientEffect } from '@/sse-db-event/components/SSEClientEffect'; import { SSEEventStreamEffect } from '@/sse-db-event/components/SSEEventStreamEffect'; +import { SSEKeepAliveEffect } from '@/sse-db-event/components/SSEKeepAliveEffect'; import { SSEQuerySubscribeEffect } from '@/sse-db-event/components/SSEQuerySubscribeEffect'; import { type ReactNode } from 'react'; @@ -14,6 +15,7 @@ export const SSEProvider = ({ children }: SSEProviderProps) => { + {children} diff --git a/packages/twenty-front/src/modules/sse-db-event/constants/SseLivenessCheckIntervalInMs.ts b/packages/twenty-front/src/modules/sse-db-event/constants/SseLivenessCheckIntervalInMs.ts new file mode 100644 index 0000000000..c2e3a71a39 --- /dev/null +++ b/packages/twenty-front/src/modules/sse-db-event/constants/SseLivenessCheckIntervalInMs.ts @@ -0,0 +1 @@ +export const SSE_LIVENESS_CHECK_INTERVAL_IN_MS = 30_000; diff --git a/packages/twenty-front/src/modules/sse-db-event/constants/SseLivenessTimeoutInMs.ts b/packages/twenty-front/src/modules/sse-db-event/constants/SseLivenessTimeoutInMs.ts new file mode 100644 index 0000000000..0fcfcfc1c8 --- /dev/null +++ b/packages/twenty-front/src/modules/sse-db-event/constants/SseLivenessTimeoutInMs.ts @@ -0,0 +1 @@ +export const SSE_LIVENESS_TIMEOUT_IN_MS = 90_000; diff --git a/packages/twenty-front/src/modules/sse-db-event/hooks/useTriggerEventStreamCreation.ts b/packages/twenty-front/src/modules/sse-db-event/hooks/useTriggerEventStreamCreation.ts index 71d86d3b77..f3a324ccbf 100644 --- a/packages/twenty-front/src/modules/sse-db-event/hooks/useTriggerEventStreamCreation.ts +++ b/packages/twenty-front/src/modules/sse-db-event/hooks/useTriggerEventStreamCreation.ts @@ -5,6 +5,7 @@ import { useTriggerOptimisticEffectFromSseEvents } from '@/sse-db-event/hooks/us 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 { lastSseEventReceivedTimestampState } from '@/sse-db-event/states/lastSseEventReceivedTimestampState'; import { shouldDestroyEventStreamState } from '@/sse-db-event/states/shouldDestroyEventStreamState'; import { sseClientState } from '@/sse-db-event/states/sseClientState'; import { sseEventStreamIdState } from '@/sse-db-event/states/sseEventStreamIdState'; @@ -81,6 +82,8 @@ export const useTriggerEventStreamCreation = () => { onEventSubscription: EventSubscription; }>, ) => { + store.set(lastSseEventReceivedTimestampState.atom, Date.now()); + if (isDefined(value?.errors) && Array.isArray(value.errors)) { const extensions = getGraphqlErrorExtensionsFromError( value.errors[0], @@ -132,8 +135,11 @@ export const useTriggerEventStreamCreation = () => { }, error: (error) => { captureException(error); + store.set(shouldDestroyEventStreamState.atom, true); + }, + complete: () => { + store.set(shouldDestroyEventStreamState.atom, true); }, - complete: () => {}, }, { message: ({ data, event }) => { @@ -143,6 +149,8 @@ export const useTriggerEventStreamCreation = () => { try { if (event === 'next') { + store.set(lastSseEventReceivedTimestampState.atom, Date.now()); + if (isDefined(result?.errors)) { const extensions = getGraphqlErrorExtensionsFromError( result.errors[0], diff --git a/packages/twenty-front/src/modules/sse-db-event/states/lastSseEventReceivedTimestampState.ts b/packages/twenty-front/src/modules/sse-db-event/states/lastSseEventReceivedTimestampState.ts new file mode 100644 index 0000000000..0812d3311d --- /dev/null +++ b/packages/twenty-front/src/modules/sse-db-event/states/lastSseEventReceivedTimestampState.ts @@ -0,0 +1,8 @@ +import { createAtomState } from '@/ui/utilities/state/jotai/utils/createAtomState'; + +export const lastSseEventReceivedTimestampState = createAtomState< + number | null +>({ + key: 'lastSseEventReceivedTimestampState', + defaultValue: null, +}); diff --git a/packages/twenty-server/src/engine/subscriptions/constants/application-keepalive-interval-ms.constant.ts b/packages/twenty-server/src/engine/subscriptions/constants/application-keepalive-interval-ms.constant.ts new file mode 100644 index 0000000000..730697e7be --- /dev/null +++ b/packages/twenty-server/src/engine/subscriptions/constants/application-keepalive-interval-ms.constant.ts @@ -0,0 +1 @@ +export const APPLICATION_KEEPALIVE_INTERVAL_MS = 30 * 1_000; // 30 seconds 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 e895ec8936..ad5e9ee544 100644 --- a/packages/twenty-server/src/engine/subscriptions/event-stream.resolver.ts +++ b/packages/twenty-server/src/engine/subscriptions/event-stream.resolver.ts @@ -16,6 +16,7 @@ import { AuthWorkspace } from 'src/engine/decorators/auth/auth-workspace.decorat 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 { APPLICATION_KEEPALIVE_INTERVAL_MS } from 'src/engine/subscriptions/constants/application-keepalive-interval-ms.constant'; import { EVENT_STREAM_TTL_MS } from 'src/engine/subscriptions/constants/event-stream-ttl.constant'; import { AddQuerySubscriptionInput } from 'src/engine/subscriptions/dtos/add-query-subscription.input'; import { EventSubscriptionDTO } from 'src/engine/subscriptions/dtos/event-subscription.dto'; @@ -116,17 +117,36 @@ export class EventStreamResolver { throw error; } + let lastTtlRefreshAt = 0; + return wrapAsyncIteratorWithLifecycle(iterator, { initialValue: { objectRecordEventsWithQueryIds: [], metadataEvents: [], }, - onHeartbeat: () => - this.eventStreamService.refreshEventStreamTTL({ + onHeartbeat: async () => { + const now = Date.now(); + + if (now - lastTtlRefreshAt > EVENT_STREAM_TTL_MS / 5) { + lastTtlRefreshAt = now; + await this.eventStreamService.refreshEventStreamTTL({ + workspaceId: workspace.id, + eventStreamChannelId, + }); + } + + await this.subscriptionService.publishToEventStream({ workspaceId: workspace.id, eventStreamChannelId, - }), - heartbeatIntervalMs: EVENT_STREAM_TTL_MS / 5, + payload: { + objectRecordEventsWithQueryIds: [], + metadataEvents: [], + }, + }); + + return true; + }, + heartbeatIntervalMs: APPLICATION_KEEPALIVE_INTERVAL_MS, onCleanup: () => this.eventStreamService.destroyEventStream({ workspaceId: workspace.id,