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 63eb769ba5..505fb3d42c 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 @@ -4,6 +4,7 @@ import { activeQueryListenersState } from '@/sse-db-event/states/activeQueryList 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 { sseEventStreamReadyState } from '@/sse-db-event/states/sseEventStreamReadyState'; import { getSnapshotValue } from '@/ui/utilities/state/utils/getSnapshotValue'; import { ApolloError, useMutation } from '@apollo/client'; import { isNonEmptyString } from '@sniptt/guards'; @@ -21,6 +22,7 @@ import { export const SSEQuerySubscribeEffect = () => { const sseEventStreamId = useRecoilValue(sseEventStreamIdState); + const sseEventStreamReady = useRecoilValue(sseEventStreamReadyState); const [addQueryToEventStream] = useMutation< boolean, @@ -122,7 +124,7 @@ export const SSEQuerySubscribeEffect = () => { ); useEffect(() => { - if (!isNonEmptyString(sseEventStreamId)) { + if (!isNonEmptyString(sseEventStreamId) || !sseEventStreamReady) { return; } @@ -138,6 +140,7 @@ export const SSEQuerySubscribeEffect = () => { } }, [ sseEventStreamId, + sseEventStreamReady, requiredQueryListeners, activeQueryListeners, debouncedUpdateQueryListeners, 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 897b9a91da..547d14904a 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 @@ -7,6 +7,7 @@ import { isDestroyingEventStreamState } from '@/sse-db-event/states/isDestroying 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 { sseEventStreamReadyState } from '@/sse-db-event/states/sseEventStreamReadyState'; import { getSnapshotValue } from '@/ui/utilities/state/utils/getSnapshotValue'; import { captureException } from '@sentry/react'; import { isNonEmptyString } from '@sniptt/guards'; @@ -60,6 +61,9 @@ export const useTriggerEventStreamCreation = () => { const newSseEventStreamId = v4(); set(sseEventStreamIdState, newSseEventStreamId); + set(sseEventStreamReadyState, false); + + let hasReceivedFirstEvent = false; const dispose = sseClient.subscribe( { @@ -74,6 +78,11 @@ export const useTriggerEventStreamCreation = () => { onEventSubscription: EventSubscription; }>, ) => { + if (!hasReceivedFirstEvent) { + hasReceivedFirstEvent = true; + set(sseEventStreamReadyState, true); + } + const objectRecordEventsWithQueryIds = value?.data?.onEventSubscription?.eventWithQueryIdsList ?? []; 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 index b12432d5af..d2552d3cce 100644 --- a/packages/twenty-front/src/modules/sse-db-event/hooks/useTriggerEventStreamDestroy.ts +++ b/packages/twenty-front/src/modules/sse-db-event/hooks/useTriggerEventStreamDestroy.ts @@ -3,6 +3,7 @@ import { isCreatingSseEventStreamState } from '@/sse-db-event/states/isCreatingS 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 { sseEventStreamReadyState } from '@/sse-db-event/states/sseEventStreamReadyState'; import { isNonEmptyString } from '@sniptt/guards'; import { useRecoilCallback, useSetRecoilState } from 'recoil'; @@ -40,6 +41,7 @@ export const useTriggerEventStreamDestroy = () => { disposeFunctionForEventStream?.dispose(); set(sseEventStreamIdState, null); + set(sseEventStreamReadyState, false); set(disposeFunctionForEventStreamState, null); set(shouldDestroyEventStreamState, false); } diff --git a/packages/twenty-front/src/modules/sse-db-event/states/sseEventStreamReadyState.ts b/packages/twenty-front/src/modules/sse-db-event/states/sseEventStreamReadyState.ts new file mode 100644 index 0000000000..6bce93928f --- /dev/null +++ b/packages/twenty-front/src/modules/sse-db-event/states/sseEventStreamReadyState.ts @@ -0,0 +1,6 @@ +import { createState } from 'twenty-ui/utilities'; + +export const sseEventStreamReadyState = createState({ + key: 'sseEventStreamReadyState', + defaultValue: false, +}); diff --git a/packages/twenty-server/src/engine/workspace-event-emitter/utils/wrap-async-iterator-with-lifecycle.ts b/packages/twenty-server/src/engine/workspace-event-emitter/utils/wrap-async-iterator-with-lifecycle.ts index ad304cd0d4..47583feb8f 100644 --- a/packages/twenty-server/src/engine/workspace-event-emitter/utils/wrap-async-iterator-with-lifecycle.ts +++ b/packages/twenty-server/src/engine/workspace-event-emitter/utils/wrap-async-iterator-with-lifecycle.ts @@ -1,6 +1,7 @@ import { isDefined } from 'twenty-shared/utils'; -type AsyncIteratorLifecycleOptions = { +type AsyncIteratorLifecycleOptions = { + initialValue?: T; onHeartbeat?: () => Promise; heartbeatIntervalMs?: number; onCleanup?: () => Promise; @@ -8,10 +9,11 @@ type AsyncIteratorLifecycleOptions = { export function wrapAsyncIteratorWithLifecycle( iterator: AsyncIterableIterator, - options: AsyncIteratorLifecycleOptions, + options: AsyncIteratorLifecycleOptions, ): AsyncIterableIterator { - const { onHeartbeat, heartbeatIntervalMs, onCleanup } = options; + const { initialValue, onHeartbeat, heartbeatIntervalMs, onCleanup } = options; let heartbeatInterval: NodeJS.Timeout | null = null; + let hasYieldedInitialValue = false; const startHeartbeat = () => { if (onHeartbeat && heartbeatIntervalMs) { @@ -41,6 +43,12 @@ export function wrapAsyncIteratorWithLifecycle( startHeartbeat(); } + if (isDefined(initialValue) && !hasYieldedInitialValue) { + hasYieldedInitialValue = true; + + return { done: false, value: initialValue }; + } + let result: IteratorResult; try { 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 dda6a5e314..526fcb0a9a 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 @@ -142,6 +142,7 @@ export class WorkspaceEventEmitterResolver { } return wrapAsyncIteratorWithLifecycle(iterator, { + initialValue: [], onHeartbeat: () => this.eventStreamService.refreshEventStreamTTL({ workspaceId: workspace.id,