Implemented SSE subscription mechanism on the frontend (#17017)

This PR is a follow-up of https://github.com/twentyhq/twenty/pull/16966
and implements a new mechanism to handle SSE events.

It creates only one event stream per browser tab, then use mutations to
tell the backend which query to listen to, without re-mounting the event
stream connexion.

Then each event that comes from this unique subscription is then
dispatched in a new JavaScript CustomEvent, per queryId, on which
specific hooks add an event listener.

This PR introduces the generic tooling as well as the handling of update
events on table.
This commit is contained in:
Lucas Bordeau
2026-01-12 18:58:46 +01:00
committed by GitHub
parent d7638a9075
commit 04a370e043
38 changed files with 747 additions and 300 deletions
@@ -3100,6 +3100,7 @@ export type ObjectPermissionInput = {
export type ObjectRecordEvent = {
__typename?: 'ObjectRecordEvent';
action: DatabaseEventAction;
objectNameSingular: Scalars['String'];
properties: ObjectRecordEventProperties;
recordId: Scalars['String'];
File diff suppressed because one or more lines are too long
@@ -14,7 +14,7 @@ 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 { SSEProvider } from '@/sse-db-event/components/SSEProvider';
import { SupportChatEffect } from '@/support/components/SupportChatEffect';
import { DialogManager } from '@/ui/feedback/dialog-manager/components/DialogManager';
import { DialogComponentInstanceContext } from '@/ui/feedback/dialog-manager/contexts/DialogComponentInstanceContext';
@@ -46,8 +46,8 @@ export const AppRouterProviders = () => {
<ChromeExtensionSidecarProvider>
<UserProvider>
<AuthProvider>
<SubscriptionProvider>
<ApolloCoreProvider>
<ApolloCoreProvider>
<SSEProvider>
<ObjectMetadataItemsLoadEffect />
<ObjectMetadataItemsProvider>
<PrefetchDataProvider>
@@ -73,8 +73,8 @@ export const AppRouterProviders = () => {
</PrefetchDataProvider>
<PageChangeEffect />
</ObjectMetadataItemsProvider>
</ApolloCoreProvider>
</SubscriptionProvider>
</SSEProvider>
</ApolloCoreProvider>
</AuthProvider>
</UserProvider>
</ChromeExtensionSidecarProvider>
@@ -9,6 +9,7 @@ import { RecordTableCellPortals } from '@/object-record/record-table/record-tabl
import { RecordTableAggregateFooter } from '@/object-record/record-table/record-table-footer/components/RecordTableAggregateFooter';
import { isRecordTableInitialLoadingComponentState } from '@/object-record/record-table/states/isRecordTableInitialLoadingComponentState';
import { RecordTableVirtualizedDataChangedEffect } from '@/object-record/record-table/virtualization/components/RecordTableVirtualizedDataChangedEffect';
import { RecordTableVirtualizedOnObjectRecordEventsEffect } from '@/object-record/record-table/virtualization/components/RecordTableVirtualizedOnObjectRecordEventsEffect';
import { RecordTableVirtualizedRowTreadmillEffect } from '@/object-record/record-table/virtualization/components/RecordTableVirtualizedRowTreadmillEffect';
import { useRecoilComponentValue } from '@/ui/utilities/state/component-state/hooks/useRecoilComponentValue';
@@ -37,6 +38,7 @@ export const RecordTableNoRecordGroupBody = () => {
)}
<RecordTableVirtualizedRowTreadmillEffect />
<RecordTableVirtualizedDataChangedEffect />
<RecordTableVirtualizedOnObjectRecordEventsEffect />
</RecordTableBodyNoRecordGroupDragDropContextProvider>
</RecordTableNoRecordGroupBodyContextProvider>
);
@@ -11,8 +11,8 @@ import { RecordTableRowArrowKeysEffect } from '@/object-record/record-table/reco
import { RecordTableRowHotkeyEffect } from '@/object-record/record-table/record-table-row/components/RecordTableRowHotkeyEffect';
import { isRecordTableRowFocusActiveComponentState } from '@/object-record/record-table/states/isRecordTableRowFocusActiveComponentState';
import { isRecordTableRowFocusedComponentFamilyState } from '@/object-record/record-table/states/isRecordTableRowFocusedComponentFamilyState';
import { ListenRecordUpdatesEffect } from '@/subscription/components/ListenRecordUpdatesEffect';
import { getDefaultRecordFieldsToListen } from '@/subscription/utils/getDefaultRecordFieldsToListen.util';
import { ListenRecordUpdatesEffect } from '@/sse-db-event/components/ListenRecordUpdatesEffect';
import { getDefaultRecordFieldsToListen } from '@/sse-db-event/utils/getDefaultRecordFieldsToListen';
import { useRecoilComponentFamilyValue } from '@/ui/utilities/state/component-state/hooks/useRecoilComponentFamilyValue';
import { useRecoilComponentValue } from '@/ui/utilities/state/component-state/hooks/useRecoilComponentValue';
@@ -68,46 +68,74 @@ export const RecordTableVirtualizedDataChangedEffect = () => {
return;
}
if (lastObjectOperation.data.type === 'update-one') {
const updateInput = lastObjectOperation.data.result.updateInput;
const updatedFieldNames = new Set<string>();
const updatedFieldNames = Object.keys(updateInput ?? {}) ?? [];
let thereIsAnUpdateOnAFilteredField = false;
let thereIsAnUpdateOnASortedField = false;
const updatedFieldMetadataItems = activeFieldMetadataItems.filter(
(fieldMetadataItemToFilter) =>
updatedFieldNames.includes(fieldMetadataItemToFilter.name) ||
(fieldMetadataItemToFilter.type === FieldMetadataType.RELATION &&
updatedFieldNames.includes(
`${fieldMetadataItemToFilter.name}Id`,
)),
);
if (
lastObjectOperation.data.type === 'update-one' ||
lastObjectOperation.data.type === 'update-many'
) {
const updateInputs =
lastObjectOperation.data.type === 'update-one'
? [lastObjectOperation.data.result.updateInput]
: lastObjectOperation.data.result.updateInputs;
const updatedFieldMetadataItemIds =
updatedFieldMetadataItems.map(mapById);
for (const updateInput of updateInputs) {
const fieldNamesForUpdateInput =
Object.keys(updateInput ?? {}) ?? [];
const thereIsAnUpdateOnAFilteredField = currentRecordFilters.some(
(recordFilter) =>
updatedFieldMetadataItemIds.includes(
recordFilter.fieldMetadataId,
),
);
for (const fieldName of fieldNamesForUpdateInput) {
updatedFieldNames.add(fieldName);
}
const thereIsAnUpdateOnASortedField = currentRecordSorts.some(
(recordSort) =>
const updatedFieldMetadataItems = activeFieldMetadataItems.filter(
(fieldMetadataItemToFilter) =>
fieldNamesForUpdateInput.includes(
fieldMetadataItemToFilter.name,
) ||
(fieldMetadataItemToFilter.type ===
FieldMetadataType.RELATION &&
fieldNamesForUpdateInput.includes(
`${fieldMetadataItemToFilter.name}Id`,
)),
);
const updatedFieldMetadataItemIds =
updatedFieldMetadataItems.map(mapById);
const updateOnAFilteredField = currentRecordFilters.some(
(recordFilter) =>
updatedFieldMetadataItemIds.includes(
recordFilter.fieldMetadataId,
),
);
const updateOnASortedField = currentRecordSorts.some((recordSort) =>
updatedFieldMetadataItemIds.includes(recordSort.fieldMetadataId),
);
);
if (updatedFieldNames.includes('position')) {
resetVirtualizationBecauseDataChanged();
} else if (
thereIsAnUpdateOnAFilteredField ||
thereIsAnUpdateOnASortedField
) {
resetVirtualizationBecauseDataChanged();
if (updateOnAFilteredField) {
thereIsAnUpdateOnAFilteredField = true;
}
if (updateOnASortedField) {
thereIsAnUpdateOnASortedField = true;
}
}
} else {
resetVirtualizationBecauseDataChanged();
}
if (updatedFieldNames.has('position')) {
resetVirtualizationBecauseDataChanged();
} else if (
thereIsAnUpdateOnAFilteredField ||
thereIsAnUpdateOnASortedField
) {
resetVirtualizationBecauseDataChanged();
}
}
}
}, [
@@ -0,0 +1,165 @@
import { triggerUpdateRecordOptimisticEffect } from '@/apollo/optimistic-effect/utils/triggerUpdateRecordOptimisticEffect';
import { useApolloCoreClient } from '@/object-metadata/hooks/useApolloCoreClient';
import { useObjectMetadataItems } from '@/object-metadata/hooks/useObjectMetadataItems';
import { useGetRecordFromCache } from '@/object-record/cache/hooks/useGetRecordFromCache';
import { getObjectTypename } from '@/object-record/cache/utils/getObjectTypename';
import { getRecordNodeFromRecord } from '@/object-record/cache/utils/getRecordNodeFromRecord';
import { useObjectPermissions } from '@/object-record/hooks/useObjectPermissions';
import { useRefetchAggregateQueries } from '@/object-record/hooks/useRefetchAggregateQueries';
import { useRegisterObjectOperation } from '@/object-record/hooks/useRegisterObjectOperation';
import { turnSortsIntoOrderBy } from '@/object-record/object-sort-dropdown/utils/turnSortsIntoOrderBy';
import { useRecordsFieldVisibleGqlFields } from '@/object-record/record-field/hooks/useRecordsFieldVisibleGqlFields';
import { currentRecordFilterGroupsComponentState } from '@/object-record/record-filter-group/states/currentRecordFilterGroupsComponentState';
import { useFilterValueDependencies } from '@/object-record/record-filter/hooks/useFilterValueDependencies';
import { currentRecordFiltersComponentState } from '@/object-record/record-filter/states/currentRecordFiltersComponentState';
import { useRecordIndexContextOrThrow } from '@/object-record/record-index/contexts/RecordIndexContext';
import { currentRecordSortsComponentState } from '@/object-record/record-sort/states/currentRecordSortsComponentState';
import { useUpsertRecordsInStore } from '@/object-record/record-store/hooks/useUpsertRecordsInStore';
import { useRecordTableContextOrThrow } from '@/object-record/record-table/contexts/RecordTableContext';
import { type ObjectRecord } from '@/object-record/types/ObjectRecord';
import { useListenToObjectRecordEventsForQuery } from '@/sse-db-event/hooks/useListenToObjectRecordEventsForQuery';
import { useRecoilComponentValue } from '@/ui/utilities/state/component-state/hooks/useRecoilComponentValue';
import {
computeRecordGqlOperationFilter,
isDefined,
isNonEmptyArray,
} from 'twenty-shared/utils';
import { useDebouncedCallback } from 'use-debounce';
import { DatabaseEventAction } from '~/generated/graphql';
export const RecordTableVirtualizedOnObjectRecordEventsEffect = () => {
const { objectMetadataItem } = useRecordIndexContextOrThrow();
const { objectNameSingular } = useRecordTableContextOrThrow();
const { registerObjectOperation } = useRegisterObjectOperation();
const { refetchAggregateQueries } = useRefetchAggregateQueries({
objectMetadataNamePlural: objectMetadataItem.namePlural,
});
const apolloCoreClient = useApolloCoreClient();
const getRecordFromCache = useGetRecordFromCache({
objectNameSingular,
});
const { upsertRecordsInStore } = useUpsertRecordsInStore();
const { objectMetadataItems } = useObjectMetadataItems();
const recordGqlFields = useRecordsFieldVisibleGqlFields({
objectMetadataItem,
});
const { objectPermissionsByObjectMetadataId } = useObjectPermissions();
const debouncedRefetchAggregateQueries = useDebouncedCallback(
refetchAggregateQueries,
200,
);
const currentRecordFilters = useRecoilComponentValue(
currentRecordFiltersComponentState,
);
const currentRecordSorts = useRecoilComponentValue(
currentRecordSortsComponentState,
);
const currentRecordFilterGroups = useRecoilComponentValue(
currentRecordFilterGroupsComponentState,
);
const { filterValueDependencies } = useFilterValueDependencies();
const queryId = `record-table-virtualized-${objectMetadataItem.nameSingular}`;
useListenToObjectRecordEventsForQuery({
queryId,
operationSignature: {
objectNameSingular: objectMetadataItem.nameSingular,
variables: {
filter: computeRecordGqlOperationFilter({
fields: objectMetadataItem.fields,
recordFilters: currentRecordFilters,
recordFilterGroups: currentRecordFilterGroups,
filterValueDependencies,
}),
orderBy: turnSortsIntoOrderBy(objectMetadataItem, currentRecordSorts),
},
},
onObjectRecordEvents: (objectRecordEvents) => {
const cache = apolloCoreClient.cache;
const updateEvents = objectRecordEvents.filter(
(objectRecordEvent) =>
objectRecordEvent.action === DatabaseEventAction.UPDATED,
);
const updatedRecordsWithUpdatedFieldsOnly: ObjectRecord[] = [];
for (const updateEvent of updateEvents) {
const updatedRecord = updateEvent.properties.after;
upsertRecordsInStore({ partialRecords: [updatedRecord] });
const cachedRecord = getRecordFromCache<ObjectRecord>(updatedRecord.id);
const cachedRecordWithConnection =
getRecordNodeFromRecord<ObjectRecord>({
record: cachedRecord,
objectMetadataItem,
objectMetadataItems,
recordGqlFields: recordGqlFields,
computeReferences: false,
});
const computedOptimisticRecord = {
...updatedRecord,
id: updatedRecord.id,
__typename: getObjectTypename(objectMetadataItem.nameSingular),
};
const computedOptimisticRecordWithConnection =
getRecordNodeFromRecord<ObjectRecord>({
record: computedOptimisticRecord,
objectMetadataItem,
objectMetadataItems,
recordGqlFields: recordGqlFields,
});
if (
!isDefined(cachedRecordWithConnection) ||
!isDefined(computedOptimisticRecordWithConnection)
) {
continue;
}
triggerUpdateRecordOptimisticEffect({
cache,
objectMetadataItem,
currentRecord: cachedRecordWithConnection,
updatedRecord: computedOptimisticRecordWithConnection,
objectMetadataItems,
objectPermissionsByObjectMetadataId,
upsertRecordsInStore,
});
const updatedFields = updateEvent.properties?.updatedFields ?? [];
updatedRecordsWithUpdatedFieldsOnly.push({
...Object.fromEntries(
updatedFields.map((fieldName) => {
return [fieldName, updatedRecord[fieldName]];
}),
),
id: updatedRecord.id,
__typename: getObjectTypename(objectMetadataItem.nameSingular),
});
}
registerObjectOperation(objectMetadataItem, {
type: 'update-many',
result: { updateInputs: updatedRecordsWithUpdatedFieldsOnly },
});
if (isNonEmptyArray(updateEvents)) {
debouncedRefetchAggregateQueries();
}
},
});
return null;
};
@@ -11,7 +11,7 @@ import { useObjectPermissions } from '@/object-record/hooks/useObjectPermissions
import { useUpsertRecordsInStore } from '@/object-record/record-store/hooks/useUpsertRecordsInStore';
import { recordStoreFamilyState } from '@/object-record/record-store/states/recordStoreFamilyState';
import { type ObjectRecord } from '@/object-record/types/ObjectRecord';
import { useOnDbEvent } from '@/subscription/hooks/useOnDbEvent';
import { useOnDbEvent } from '@/sse-db-event/hooks/useOnDbEvent';
import { useRecoilCallback } from 'recoil';
import { isDefined } from 'twenty-shared/utils';
import { DatabaseEventAction } from '~/generated/graphql';
@@ -0,0 +1,21 @@
import { SSEProviderEffect } from '@/sse-db-event/components/SSEProviderEffect';
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 { type ReactNode } from 'react';
type SSEProviderProps = {
children: ReactNode;
};
export const SSEProvider = ({ children }: SSEProviderProps) => {
const { sseClient } = useSseClient();
return (
<SseClientContext.Provider value={sseClient}>
<SSEProviderEffect />
<SSEQuerySubscribeEffect />
{children}
</SseClientContext.Provider>
);
};
@@ -0,0 +1,62 @@
import { SseClientContext } from '@/sse-db-event/contexts/SseClientContext';
import { ON_EVENT_SUBSCRIPTION } from '@/sse-db-event/graphql/subscriptions/OnEventSubscription';
import { sseEventStreamIdState } from '@/sse-db-event/states/sseEventStreamIdState';
import { dispatchObjectRecordEventsWithQueryIds } from '@/sse-db-event/utils/dispatchObjectRecordEventsWithQueryIds';
import { isNonEmptyString } from '@sniptt/guards';
import { type ExecutionResult, print } from 'graphql';
import { useContext, useEffect } from 'react';
import { useRecoilState } from 'recoil';
import { isDefined } from 'twenty-shared/utils';
import { v4 } from 'uuid';
import { type EventSubscription } from '~/generated/graphql';
export const SSEProviderEffect = () => {
const sseClient = useContext(SseClientContext);
const [sseEventStreamId, setSseEventStreamId] = useRecoilState(
sseEventStreamIdState,
);
useEffect(() => {
if (!isDefined(sseClient)) {
return;
}
if (!isNonEmptyString(sseEventStreamId)) {
setSseEventStreamId(v4());
return;
}
const unsubscribe = sseClient.subscribe(
{
query: print(ON_EVENT_SUBSCRIPTION),
variables: {
eventStreamId: sseEventStreamId,
},
},
{
next: (
value: ExecutionResult<{ onEventSubscription: EventSubscription }>,
) => {
const objectRecordEventsWithQueryIds =
value?.data?.onEventSubscription?.eventWithQueryIdsList ?? [];
dispatchObjectRecordEventsWithQueryIds(
objectRecordEventsWithQueryIds,
);
},
error: (error) => {
// eslint-disable-next-line no-console
console.error('Subscription error:', error);
},
complete: () => {},
},
);
return () => {
unsubscribe();
};
}, [sseClient, sseEventStreamId, setSseEventStreamId]);
return null;
};
@@ -0,0 +1,128 @@
import { useApolloCoreClient } from '@/object-metadata/hooks/useApolloCoreClient';
import { ADD_QUERY_TO_EVENT_STREAM_MUTATION } from '@/sse-db-event/graphql/mutations/AddQueryToEventStreamMutation';
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 { sseEventStreamIdState } from '@/sse-db-event/states/sseEventStreamIdState';
import { getSnapshotValue } from '@/ui/utilities/state/utils/getSnapshotValue';
import { useMutation } from '@apollo/client';
import { useEffect } from 'react';
import { useRecoilCallback, useRecoilValue } from 'recoil';
import {
compareArraysOfObjectsByProperty,
isDefined,
} from 'twenty-shared/utils';
import { useDebouncedCallback } from 'use-debounce';
import {
type AddQuerySubscriptionInput,
type RemoveQueryFromEventStreamInput,
} from '~/generated/graphql';
export const SSEQuerySubscribeEffect = () => {
const sseEventStreamId = useRecoilValue(sseEventStreamIdState);
const apolloCoreClient = useApolloCoreClient();
const [addQueryToEventStream] = useMutation<
boolean,
{ input: AddQuerySubscriptionInput }
>(ADD_QUERY_TO_EVENT_STREAM_MUTATION, { client: apolloCoreClient });
const [removeQueryFromEventStream] = useMutation<
void,
{ input: RemoveQueryFromEventStreamInput }
>(REMOVE_QUERY_FROM_EVENT_STREAM_MUTATION, { client: apolloCoreClient });
const requiredQueryListeners = useRecoilValue(requiredQueryListenersState);
const activeQueryListeners = useRecoilValue(activeQueryListenersState);
const updateQueryListeners = useRecoilCallback(
({ set, snapshot }) =>
async () => {
if (!isDefined(sseEventStreamId)) {
return;
}
const requiredQueryListeners = getSnapshotValue(
snapshot,
requiredQueryListenersState,
);
const activeQueryListeners = getSnapshotValue(
snapshot,
activeQueryListenersState,
);
const queryListenersToAdd = requiredQueryListeners.filter(
(listener) =>
!activeQueryListeners.some(
(activeListener) => activeListener.queryId === listener.queryId,
),
);
const queryListenersToRemove = activeQueryListeners.filter(
(listener) =>
!requiredQueryListeners.some(
(requiredListener) =>
requiredListener.queryId === listener.queryId,
),
);
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,
},
},
});
}
set(activeQueryListenersState, requiredQueryListeners);
},
[addQueryToEventStream, removeQueryFromEventStream, sseEventStreamId],
);
const debouncedUpdateQueryListeners = useDebouncedCallback(
updateQueryListeners,
1000,
{ leading: true },
);
useEffect(() => {
if (!sseEventStreamId) {
return;
}
const areRequiredQueryListenersDifferentFromActiveQueryListeners =
compareArraysOfObjectsByProperty(
requiredQueryListeners,
activeQueryListeners,
'queryId',
);
if (areRequiredQueryListenersDifferentFromActiveQueryListeners) {
debouncedUpdateQueryListeners();
}
}, [
sseEventStreamId,
requiredQueryListeners,
activeQueryListeners,
debouncedUpdateQueryListeners,
]);
return null;
};
@@ -0,0 +1,7 @@
import { gql } from '@apollo/client';
export const ADD_QUERY_TO_EVENT_STREAM_MUTATION = gql`
mutation AddQueryToEventStream($input: AddQuerySubscriptionInput!) {
addQueryToEventStream(input: $input)
}
`;
@@ -0,0 +1,9 @@
import { gql } from '@apollo/client';
export const REMOVE_QUERY_FROM_EVENT_STREAM_MUTATION = gql`
mutation RemoveQueryFromEventStream(
$input: RemoveQueryFromEventStreamInput!
) {
removeQueryFromEventStream(input: $input)
}
`;
@@ -0,0 +1,25 @@
import { gql } from '@apollo/client';
export const ON_EVENT_SUBSCRIPTION = gql`
subscription OnEventSubscription($eventStreamId: String!) {
onEventSubscription(eventStreamId: $eventStreamId) {
eventStreamId
eventWithQueryIdsList {
event {
action
objectNameSingular
recordId
userId
workspaceMemberId
properties {
updatedFields
before
after
diff
}
}
queryIds
}
}
}
`;
@@ -1,9 +1,9 @@
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';
import { useOnDbEvent } from '@/subscription/hooks/useOnDbEvent';
jest.mock('~/modules/apollo/utils/getTokenPair');
jest.mock('graphql-sse');
@@ -0,0 +1,76 @@
import { requiredQueryListenersState } from '@/sse-db-event/states/requiredQueryListenersState';
import { getObjectRecordEventsForQueryEventName } from '@/sse-db-event/utils/getObjectRecordEventsForQueryEventName';
import { getSnapshotValue } from '@/ui/utilities/state/utils/getSnapshotValue';
import { useEffect } from 'react';
import { useRecoilCallback } from 'recoil';
import { type RecordGqlOperationSignature } from 'twenty-shared/types';
import { type ObjectRecordEvent } from '~/generated/graphql';
export const useListenToObjectRecordEventsForQuery = ({
queryId,
operationSignature,
onObjectRecordEvents,
}: {
queryId: string;
operationSignature: RecordGqlOperationSignature;
onObjectRecordEvents: (objectRecordEvents: ObjectRecordEvent[]) => void;
}) => {
useEffect(() => {
const eventName = getObjectRecordEventsForQueryEventName(queryId);
const handleOnObjectRecordEventsForQuery = (event: Event) => {
const objectRecordEvents = (event as CustomEvent<ObjectRecordEvent[]>)
.detail;
onObjectRecordEvents(objectRecordEvents);
};
window.addEventListener(eventName, handleOnObjectRecordEventsForQuery);
return () => {
window.removeEventListener(eventName, handleOnObjectRecordEventsForQuery);
};
}, [onObjectRecordEvents, queryId]);
const changeQueryIdListenState = useRecoilCallback(
({ set, snapshot }) =>
(shouldListen: boolean, queryId: string) => {
const currentRequiredQueryListeners = getSnapshotValue(
snapshot,
requiredQueryListenersState,
);
const listeningForThisQueryIsActive =
currentRequiredQueryListeners.some(
(listener) => listener.queryId === queryId,
);
if (shouldListen === listeningForThisQueryIsActive) {
return;
}
if (shouldListen) {
set(requiredQueryListenersState, [
...currentRequiredQueryListeners,
{ queryId, operationSignature },
]);
} else {
set(
requiredQueryListenersState,
currentRequiredQueryListeners.filter(
(listener) => listener.queryId !== queryId,
),
);
}
},
[operationSignature],
);
useEffect(() => {
changeQueryIdListenState(true, queryId);
return () => {
changeQueryIdListenState(false, queryId);
};
}, [changeQueryIdListenState, queryId]);
};
@@ -1,5 +1,5 @@
import { ON_DB_EVENT } from '@/subscription/graphql/subscriptions/onDbEvent';
import { useSseClient } from '@/subscription/hooks/useSseClient.util';
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 {
type Subscription,
@@ -0,0 +1,9 @@
import { type RecordGqlOperationSignature } from 'twenty-shared/types';
import { createState } from 'twenty-ui/utilities';
export const activeQueryListenersState = createState<
{ queryId: string; operationSignature: RecordGqlOperationSignature }[]
>({
key: 'activeQueryListenersState',
defaultValue: [],
});
@@ -0,0 +1,9 @@
import { type RecordGqlOperationSignature } from 'twenty-shared/types';
import { createState } from 'twenty-ui/utilities';
export const requiredQueryListenersState = createState<
{ queryId: string; operationSignature: RecordGqlOperationSignature }[]
>({
key: 'requiredQueryListenersState',
defaultValue: [],
});
@@ -0,0 +1,6 @@
import { createState } from 'twenty-ui/utilities';
export const sseEventStreamIdState = createState<string | null>({
key: 'sseEventStreamIdState',
defaultValue: null,
});
@@ -0,0 +1,3 @@
import { type ObjectRecordEvent } from '~/generated/graphql';
export type ObjectRecordEventsByQueryId = Record<string, ObjectRecordEvent[]>;
@@ -0,0 +1,30 @@
import { type ObjectRecordEventsByQueryId } from '@/sse-db-event/types/ObjectRecordEventsByQueryId';
import { getObjectRecordEventsForQueryEventName } from '@/sse-db-event/utils/getObjectRecordEventsForQueryEventName';
import { isDefined } from 'twenty-shared/utils';
import { type EventWithQueryIds } from '~/generated/graphql';
export const dispatchObjectRecordEventsWithQueryIds = (
objectRecordEventsWithQueryIds: EventWithQueryIds[],
) => {
const objectRecordEventsByQueryId: ObjectRecordEventsByQueryId = {};
for (const objectRecordEventWithQueryIds of objectRecordEventsWithQueryIds) {
for (const queryId of objectRecordEventWithQueryIds.queryIds) {
if (!isDefined(objectRecordEventsByQueryId[queryId])) {
objectRecordEventsByQueryId[queryId] = [];
}
objectRecordEventsByQueryId[queryId].push(
objectRecordEventWithQueryIds.event,
);
}
}
for (const queryId in objectRecordEventsByQueryId) {
window.dispatchEvent(
new CustomEvent(getObjectRecordEventsForQueryEventName(queryId), {
detail: objectRecordEventsByQueryId[queryId],
}),
);
}
};
@@ -0,0 +1,3 @@
export const getObjectRecordEventsForQueryEventName = (queryId: string) => {
return `object-record-events-${queryId}`;
};
@@ -1,21 +0,0 @@
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 (
<SseClientContext.Provider value={sseClient}>
<SubscriptionProviderEffect />
{children}
</SseClientContext.Provider>
);
};
@@ -1,68 +0,0 @@
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';
import { type SubscriptionMatches } from '~/generated/graphql';
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: SubscriptionMatches;
} | null;
if (!data?.onSubscriptionMatch?.matches) {
return;
}
for (const match of data.onSubscriptionMatch.matches) {
for (const subscriptionId of match.subscriptionIds) {
const entry = subscriptionRegistry.get(subscriptionId);
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;
};
@@ -1,18 +0,0 @@
import { gql } from '@apollo/client';
export const ON_SUBSCRIPTION_MATCH = gql`
subscription OnSubscriptionMatch($subscriptions: [SubscriptionInput!]!) {
onSubscriptionMatch(subscriptions: $subscriptions) {
matches {
subscriptionIds
event {
action
objectNameSingular
eventDate
record
updatedFields
}
}
}
}
`;
@@ -1,50 +0,0 @@
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<string, unknown>;
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]);
};
@@ -1,14 +0,0 @@
import { createState } from 'twenty-ui/utilities';
export type SubscriptionEntry = {
id: string;
query: string;
onRefetch: () => void;
};
export const subscriptionRegistryState = createState<
Map<string, SubscriptionEntry>
>({
key: 'subscriptionRegistryState',
defaultValue: new Map(),
});
@@ -1,5 +1,5 @@
import { SKELETON_LOADER_HEIGHT_SIZES } from '@/activities/components/SkeletonLoader';
import { ListenRecordUpdatesEffect } from '@/subscription/components/ListenRecordUpdatesEffect';
import { ListenRecordUpdatesEffect } from '@/sse-db-event/components/ListenRecordUpdatesEffect';
import { useTargetRecord } from '@/ui/layout/contexts/useTargetRecord';
import { getWorkflowVisualizerComponentInstanceId } from '@/workflow/utils/getWorkflowVisualizerComponentInstanceId';
import { WorkflowRunVisualizer } from '@/workflow/workflow-diagram/components/WorkflowRunVisualizer';
@@ -1,9 +1,13 @@
import { Field, ObjectType } from '@nestjs/graphql';
import { DatabaseEventAction } from 'src/engine/api/graphql/graphql-query-runner/enums/database-event-action';
import { ObjectRecordEventPropertiesDTO } from 'src/engine/subscriptions/dtos/object-record-event-properties.dto';
@ObjectType('ObjectRecordEvent')
export class ObjectRecordEventDTO {
@Field(() => DatabaseEventAction)
action: DatabaseEventAction;
@Field(() => String)
objectNameSingular: string;
@@ -23,6 +23,7 @@ import { SubscriptionChannel } from 'src/engine/subscriptions/enums/subscription
import { EventStreamService } from 'src/engine/subscriptions/event-stream.service';
import { SubscriptionService } from 'src/engine/subscriptions/subscription.service';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import { parseEventNameOrThrow } from 'src/engine/workspace-event-emitter/utils/parse-event-name';
@Resolver()
@UseGuards(WorkspaceAuthGuard, UserAuthGuard, NoPermissionGuard)
@@ -144,9 +145,16 @@ export class WorkspaceEventEmitterResolver {
}[] = [];
for (const event of payload.workspaceEventBatch.events) {
const eventName = parseEventNameOrThrow(
payload.workspaceEventBatch.name,
);
const action = eventName.action;
const eventWithObjectName = {
objectNameSingular,
...event,
objectNameSingular,
action,
};
const matchedQueryIds =
@@ -0,0 +1,92 @@
import { compareArraysOfObjectsByProperty } from '@/utils/array/compareArraysOfObjectsByProperty';
type TestObject = {
id: string;
name: string;
};
describe('compareArraysOfObjectsByProperty', () => {
it('should return false when both arrays are empty', () => {
expect(compareArraysOfObjectsByProperty([], [], 'id')).toBe(false);
});
it('should return false when arrays have same objects by property', () => {
const arrayA: TestObject[] = [
{ id: '1', name: 'Test 1' },
{ id: '2', name: 'Test 2' },
];
const arrayB: TestObject[] = [
{ id: '1', name: 'Different Name' },
{ id: '2', name: 'Another Name' },
];
expect(compareArraysOfObjectsByProperty(arrayA, arrayB, 'id')).toBe(false);
});
it('should return true when arrays have different lengths', () => {
const arrayA: TestObject[] = [
{ id: '1', name: 'Test 1' },
{ id: '2', name: 'Test 2' },
];
const arrayB: TestObject[] = [{ id: '1', name: 'Test 1' }];
expect(compareArraysOfObjectsByProperty(arrayA, arrayB, 'id')).toBe(true);
});
it('should return true when arrayA has items not in arrayB', () => {
const arrayA: TestObject[] = [
{ id: '1', name: 'Test 1' },
{ id: '2', name: 'Test 2' },
];
const arrayB: TestObject[] = [
{ id: '1', name: 'Test 1' },
{ id: '3', name: 'Test 3' },
];
expect(compareArraysOfObjectsByProperty(arrayA, arrayB, 'id')).toBe(true);
});
it('should return true when arrayB has items not in arrayA', () => {
const arrayA: TestObject[] = [
{ id: '1', name: 'Test 1' },
{ id: '3', name: 'Test 3' },
];
const arrayB: TestObject[] = [
{ id: '1', name: 'Test 1' },
{ id: '2', name: 'Test 2' },
];
expect(compareArraysOfObjectsByProperty(arrayA, arrayB, 'id')).toBe(true);
});
it('should return false when arrays have same items in different order', () => {
const arrayA: TestObject[] = [
{ id: '1', name: 'Test 1' },
{ id: '2', name: 'Test 2' },
{ id: '3', name: 'Test 3' },
];
const arrayB: TestObject[] = [
{ id: '3', name: 'Test 3' },
{ id: '1', name: 'Test 1' },
{ id: '2', name: 'Test 2' },
];
expect(compareArraysOfObjectsByProperty(arrayA, arrayB, 'id')).toBe(false);
});
it('should compare by the specified property', () => {
const arrayA: TestObject[] = [
{ id: '1', name: 'Alpha' },
{ id: '2', name: 'Beta' },
];
const arrayB: TestObject[] = [
{ id: '3', name: 'Alpha' },
{ id: '4', name: 'Beta' },
];
expect(compareArraysOfObjectsByProperty(arrayA, arrayB, 'name')).toBe(
false,
);
expect(compareArraysOfObjectsByProperty(arrayA, arrayB, 'id')).toBe(true);
});
});
@@ -0,0 +1,15 @@
export const compareArraysOfObjectsByProperty = <T, K extends keyof T>(
arrayA: T[],
arrayB: T[],
property: K,
) => {
return (
arrayA.length !== arrayB.length ||
arrayA.some(
(item) => !arrayB.some((itemB) => itemB[property] === item[property]),
) ||
arrayB.some(
(item) => !arrayA.some((itemA) => itemA[property] === item[property]),
)
);
};
@@ -8,6 +8,7 @@
*/
export { applyDiff } from './applyDiff';
export { compareArraysOfObjectsByProperty } from './array/compareArraysOfObjectsByProperty';
export { filterOutByProperty } from './array/filterOutByProperty';
export { findById } from './array/findById';
export { findByProperty } from './array/findByProperty';