Implemented SSE update across main front components (#17205)

This PR implements SSE update events across the main components of the
application : tables, boards, calendars, show pages.

There is still work to do on other event type and on making sure
everything works fine, but this first implementation should be robust
enough to start with.

Some problems encountered along the way : 
- Events are returning raw Postgres output, because they are not
normalized by the GraphQL layer, so we ended up with amountMicros as
string values, which the frontend does not like, so I implemented a
small util in our `formatResult` generic pipeline to turn amountMicros
to a number if it's a string value. We could implement other formatters
for composite fields if we see problems with events.
-`action` property was missing in SSE events, which is required by the
frontend to know what kind of event it is.

# QA


https://github.com/user-attachments/assets/393b45ac-59d2-48b0-855b-2ce9c4b8ae57


https://github.com/user-attachments/assets/cb214e7a-1595-4b85-bdb6-87630993d2e2
This commit is contained in:
Lucas Bordeau
2026-01-18 17:14:36 +01:00
committed by GitHub
parent a91bac9dd9
commit d51c988a9e
27 changed files with 933 additions and 89 deletions
@@ -0,0 +1,26 @@
import { useApolloCoreClient } from '@/object-metadata/hooks/useApolloCoreClient';
import { type ObjectMetadataItem } from '@/object-metadata/types/ObjectMetadataItem';
import { getGroupByAggregateQueryName } from '@/object-record/record-aggregate/utils/getGroupByAggregateQueryName';
import { getAggregateQueryName } from '@/object-record/utils/getAggregateQueryName';
export const useRefetchAggregateQueriesForObjectMetadataItem = () => {
const apolloCoreClient = useApolloCoreClient();
const refetchAggregateQueriesForObjectMetadataItem = async ({
objectMetadataItem,
}: {
objectMetadataItem: ObjectMetadataItem;
}) => {
const queryName = getAggregateQueryName(objectMetadataItem.namePlural);
const groupByAggregateQueryName = getGroupByAggregateQueryName({
objectMetadataNamePlural: objectMetadataItem.namePlural,
});
await apolloCoreClient.refetchQueries({
include: [queryName, groupByAggregateQueryName],
});
};
return {
refetchAggregateQueriesForObjectMetadataItem,
};
};
@@ -0,0 +1,46 @@
import { useListenToObjectRecordOperationBrowserEvent } from '@/object-record/hooks/useListenToObjectRecordOperationBrowserEvent';
import { useGetShouldInitializeRecordBoardForUpdateInputs } from '@/object-record/record-board/hooks/useGetShouldInitializeRecordBoardForUpdateInputs';
import { useTriggerRecordBoardInitialQuery } from '@/object-record/record-board/hooks/useTriggerRecordBoardInitialQuery';
import { useRecordIndexContextOrThrow } from '@/object-record/record-index/contexts/RecordIndexContext';
import { type ObjectRecordOperationBrowserEventDetail } from '@/object-record/types/ObjectRecordOperationBrowserEventDetail';
export const RecordBoardDataChangedEffect = () => {
const { objectMetadataItem } = useRecordIndexContextOrThrow();
const { triggerRecordBoardInitialQuery } =
useTriggerRecordBoardInitialQuery();
const { getShouldInitializeRecordBoardForUpdateInputs } =
useGetShouldInitializeRecordBoardForUpdateInputs();
const handleObjectRecordOperation = (
objectRecordOperationEventDetail: ObjectRecordOperationBrowserEventDetail,
) => {
const objectRecordOperation = objectRecordOperationEventDetail.operation;
const isUpdateOperation =
objectRecordOperation.type === 'update-one' ||
objectRecordOperation.type === 'update-many';
if (isUpdateOperation) {
const updateInputs =
objectRecordOperation.type === 'update-one'
? [objectRecordOperation.result.updateInput]
: objectRecordOperation.result.updateInputs;
const shouldInitializeForUpdateOperation =
getShouldInitializeRecordBoardForUpdateInputs(updateInputs);
if (shouldInitializeForUpdateOperation) {
triggerRecordBoardInitialQuery();
}
} else {
triggerRecordBoardInitialQuery();
}
};
useListenToObjectRecordOperationBrowserEvent({
onObjectRecordOperationBrowserEvent: handleObjectRecordOperation,
objectMetadataItemId: objectMetadataItem.id,
});
return null;
};
@@ -1,7 +1,9 @@
import { RecordBoardClickOutsideEffect } from '@/object-record/record-board/components/RecordBoardClickOutsideEffect';
import { RecordBoardDataChangedEffect } from '@/object-record/record-board/components/RecordBoardDataChangedEffect';
import { RecordBoardQueryEffect } from '@/object-record/record-board/components/RecordBoardQueryEffect';
import { RecordBoardScrollToFocusedCardEffect } from '@/object-record/record-board/components/RecordBoardScrollToFocusedCardEffect';
import { RecordBoardSelectRecordsEffect } from '@/object-record/record-board/components/RecordBoardSelectRecordsEffect';
import { RecordBoardSSESubscribeEffect } from '@/object-record/record-board/components/RecordBoardSSESubscribeEffect';
import { RecordBoardStickyHeaderEffect } from '@/object-record/record-board/components/RecordBoardStickyHeaderEffect';
import { RecordBoardDeactivateBoardCardEffect } from '@/object-record/record-board/record-board-card/components/RecordBoardDeactivateBoardCardEffect';
@@ -11,6 +13,8 @@ export const RecordBoardEffects = () => {
<RecordBoardStickyHeaderEffect />
<RecordBoardScrollToFocusedCardEffect />
<RecordBoardDeactivateBoardCardEffect />
<RecordBoardSSESubscribeEffect />
<RecordBoardDataChangedEffect />
<RecordBoardQueryEffect />
<RecordBoardSelectRecordsEffect />
<RecordBoardClickOutsideEffect />
@@ -0,0 +1,28 @@
import { useContext } from 'react';
import { RecordBoardContext } from '@/object-record/record-board/contexts/RecordBoardContext';
import { useRecordIndexContextOrThrow } from '@/object-record/record-index/contexts/RecordIndexContext';
import { useRecordIndexGroupCommonQueryVariables } from '@/object-record/record-index/hooks/useRecordIndexGroupCommonQueryVariables';
import { useListenToObjectRecordEventsForQuery } from '@/sse-db-event/hooks/useListenToObjectRecordEventsForQuery';
export const RecordBoardSSESubscribeEffect = () => {
const { recordBoardId } = useContext(RecordBoardContext);
const { objectMetadataItem } = useRecordIndexContextOrThrow();
const { combinedFilters, orderBy } =
useRecordIndexGroupCommonQueryVariables();
const queryId = `record-board-${recordBoardId}`;
useListenToObjectRecordEventsForQuery({
queryId,
operationSignature: {
objectNameSingular: objectMetadataItem.nameSingular,
variables: {
filter: combinedFilters,
orderBy,
},
},
});
return null;
};
@@ -0,0 +1,100 @@
import { useActiveFieldMetadataItems } from '@/object-metadata/hooks/useActiveFieldMetadataItems';
import { currentRecordFiltersComponentState } from '@/object-record/record-filter/states/currentRecordFiltersComponentState';
import { useRecordIndexContextOrThrow } from '@/object-record/record-index/contexts/RecordIndexContext';
import { recordIndexGroupFieldMetadataItemComponentState } from '@/object-record/record-index/states/recordIndexGroupFieldMetadataComponentState';
import { currentRecordSortsComponentState } from '@/object-record/record-sort/states/currentRecordSortsComponentState';
import { type ObjectRecordOperationUpdateInput } from '@/object-record/types/ObjectRecordOperationUpdateInput';
import { useRecoilComponentValue } from '@/ui/utilities/state/component-state/hooks/useRecoilComponentValue';
import { FieldMetadataType } from 'twenty-shared/types';
import { isDefined, mapById } from 'twenty-shared/utils';
export const useGetShouldInitializeRecordBoardForUpdateInputs = () => {
const { objectMetadataItem } = useRecordIndexContextOrThrow();
const { activeFieldMetadataItems } = useActiveFieldMetadataItems({
objectMetadataItem,
});
const currentRecordSorts = useRecoilComponentValue(
currentRecordSortsComponentState,
);
const currentRecordFilters = useRecoilComponentValue(
currentRecordFiltersComponentState,
);
const recordIndexGroupFieldMetadataItem = useRecoilComponentValue(
recordIndexGroupFieldMetadataItemComponentState,
);
const getShouldInitializeRecordBoardForUpdateInputs = (
updateInputs: ObjectRecordOperationUpdateInput[],
) => {
const updatedFieldNames = new Set<string>();
let thereIsAnUpdateOnAFilteredField = false;
let thereIsAnUpdateOnASortedField = false;
let thereIsAnUpdateOnAGroupField = false;
for (const updateInput of updateInputs) {
const fieldNamesForUpdateInput = updateInput.updatedFields.flatMap(
(updatedField) => Object.keys(updatedField ?? {}),
);
for (const fieldName of fieldNamesForUpdateInput) {
updatedFieldNames.add(fieldName);
}
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 (isDefined(recordIndexGroupFieldMetadataItem)) {
if (
updatedFieldMetadataItemIds.includes(
recordIndexGroupFieldMetadataItem.id,
)
) {
thereIsAnUpdateOnAGroupField = true;
}
}
if (updateOnAFilteredField) {
thereIsAnUpdateOnAFilteredField = true;
}
if (updateOnASortedField) {
thereIsAnUpdateOnASortedField = true;
}
}
if (updatedFieldNames.has('position')) {
return true;
}
return (
thereIsAnUpdateOnAFilteredField ||
thereIsAnUpdateOnASortedField ||
thereIsAnUpdateOnAGroupField
);
};
return {
getShouldInitializeRecordBoardForUpdateInputs,
};
};
@@ -0,0 +1,47 @@
import { hasObjectMetadataItemPositionField } from '@/object-metadata/utils/hasObjectMetadataItemPositionField';
import { useRecordCalendarContextOrThrow } from '@/object-record/record-calendar/contexts/RecordCalendarContext';
import { useRecordCalendarQueryDateRangeFilter } from '@/object-record/record-calendar/month/hooks/useRecordCalendarQueryDateRangeFilter';
import { RecordCalendarComponentInstanceContext } from '@/object-record/record-calendar/states/contexts/RecordCalendarComponentInstanceContext';
import { recordCalendarSelectedDateComponentState } from '@/object-record/record-calendar/states/recordCalendarSelectedDateComponentState';
import { useListenToObjectRecordEventsForQuery } from '@/sse-db-event/hooks/useListenToObjectRecordEventsForQuery';
import { useAvailableComponentInstanceIdOrThrow } from '@/ui/utilities/state/component-state/hooks/useAvailableComponentInstanceIdOrThrow';
import { useRecoilComponentValue } from '@/ui/utilities/state/component-state/hooks/useRecoilComponentValue';
import { type RecordGqlOperationOrderBy } from 'twenty-shared/types';
export const RecordCalendarSSESubscribeEffect = () => {
const recordCalendarId = useAvailableComponentInstanceIdOrThrow(
RecordCalendarComponentInstanceContext,
);
const { objectMetadataItem } = useRecordCalendarContextOrThrow();
const recordCalendarSelectedDate = useRecoilComponentValue(
recordCalendarSelectedDateComponentState,
);
const { dateRangeFilter } = useRecordCalendarQueryDateRangeFilter(
recordCalendarSelectedDate,
);
const orderBy: RecordGqlOperationOrderBy =
!objectMetadataItem.isRemote &&
hasObjectMetadataItemPositionField(objectMetadataItem)
? [
{
position: 'AscNullsFirst',
},
]
: [];
const queryId = `record-calendar-${recordCalendarId}`;
useListenToObjectRecordEventsForQuery({
queryId,
operationSignature: {
objectNameSingular: objectMetadataItem.nameSingular,
variables: {
filter: dateRangeFilter,
orderBy,
},
},
});
return null;
};
@@ -2,6 +2,7 @@ import { useObjectMetadataItem } from '@/object-metadata/hooks/useObjectMetadata
import { RecordComponentInstanceContextsWrapper } from '@/object-record/components/RecordComponentInstanceContextsWrapper';
import { useObjectPermissionsForObject } from '@/object-record/hooks/useObjectPermissionsForObject';
import { RecordCalendar } from '@/object-record/record-calendar/components/RecordCalendar';
import { RecordCalendarSSESubscribeEffect } from '@/object-record/record-calendar/components/RecordCalendarSSESubscribeEffect';
import { RecordIndexCalendarDataLoaderEffect } from '@/object-record/record-calendar/components/RecordIndexCalendarDataLoaderEffect';
import { RecordIndexCalendarSelectedDateInitEffect } from '@/object-record/record-calendar/components/RecordIndexCalendarSelectedDateInitEffect';
import { RecordCalendarContextProvider } from '@/object-record/record-calendar/contexts/RecordCalendarContext';
@@ -52,6 +53,7 @@ export const RecordIndexCalendarContainer = ({
}}
>
<RecordCalendar />
<RecordCalendarSSESubscribeEffect />
<RecordIndexCalendarDataLoaderEffect />
<RecordIndexCalendarSelectedDateInitEffect />
</RecordCalendarContextProvider>
@@ -0,0 +1,26 @@
import { useListenToObjectRecordEventsForQuery } from '@/sse-db-event/hooks/useListenToObjectRecordEventsForQuery';
type RecordShowPageSSESubscribeEffectProps = {
objectNameSingular: string;
recordId: string;
};
export const RecordShowPageSSESubscribeEffect = ({
objectNameSingular,
recordId,
}: RecordShowPageSSESubscribeEffectProps) => {
const queryId = `record-show-${objectNameSingular}-${recordId}`;
useListenToObjectRecordEventsForQuery({
queryId,
operationSignature: {
objectNameSingular,
variables: {
filter: { id: { eq: recordId } },
limit: 1,
},
},
});
return null;
};
@@ -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 { RecordTableVirtualizedSSESubscribeEffect } from '@/object-record/record-table/virtualization/components/RecordTableVirtualizedSSESubscribeEffect';
import { RecordTableVirtualizedRowTreadmillEffect } from '@/object-record/record-table/virtualization/components/RecordTableVirtualizedRowTreadmillEffect';
import { useRecoilComponentValue } from '@/ui/utilities/state/component-state/hooks/useRecoilComponentValue';
@@ -38,6 +39,7 @@ export const RecordTableNoRecordGroupBody = () => {
)}
<RecordTableVirtualizedRowTreadmillEffect />
<RecordTableVirtualizedDataChangedEffect />
<RecordTableVirtualizedSSESubscribeEffect />
</RecordTableBodyNoRecordGroupDragDropContextProvider>
</RecordTableNoRecordGroupBodyContextProvider>
);
@@ -0,0 +1,46 @@
import { turnSortsIntoOrderBy } from '@/object-record/object-sort-dropdown/utils/turnSortsIntoOrderBy';
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 { useListenToObjectRecordEventsForQuery } from '@/sse-db-event/hooks/useListenToObjectRecordEventsForQuery';
import { useRecoilComponentValue } from '@/ui/utilities/state/component-state/hooks/useRecoilComponentValue';
import { computeRecordGqlOperationFilter } from 'twenty-shared/utils';
export const RecordTableVirtualizedSSESubscribeEffect = () => {
const { objectMetadataItem } = useRecordIndexContextOrThrow();
const { filterValueDependencies } = useFilterValueDependencies();
const currentRecordFilters = useRecoilComponentValue(
currentRecordFiltersComponentState,
);
const currentRecordSorts = useRecoilComponentValue(
currentRecordSortsComponentState,
);
const currentRecordFilterGroups = useRecoilComponentValue(
currentRecordFilterGroupsComponentState,
);
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),
},
},
});
return null;
};
@@ -0,0 +1,22 @@
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 { isDefined } from 'twenty-shared/utils';
export const SSEEventStreamEffect = () => {
const sseClient = useContext(SseClientContext);
const { objectMetadataItems } = useObjectMetadataItems();
const { subscribeToSseEventStream } = useSubscribeToSseEventStream();
useEffect(() => {
if (!isDefined(sseClient) || objectMetadataItems.length === 0) {
return;
}
subscribeToSseEventStream(sseClient);
}, [sseClient, subscribeToSseEventStream, objectMetadataItems]);
return null;
};
@@ -1,4 +1,4 @@
import { SSEProviderEffect } from '@/sse-db-event/components/SSEProviderEffect';
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';
@@ -26,7 +26,7 @@ export const SSEProvider = ({ children }: SSEProviderProps) => {
return (
<SseClientContext.Provider value={sseClient}>
<SSEProviderEffect />
<SSEEventStreamEffect />
<SSEQuerySubscribeEffect />
{children}
</SseClientContext.Provider>
@@ -1,62 +0,0 @@
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;
};
@@ -6,6 +6,7 @@ import { requiredQueryListenersState } from '@/sse-db-event/states/requiredQuery
import { sseEventStreamIdState } from '@/sse-db-event/states/sseEventStreamIdState';
import { getSnapshotValue } from '@/ui/utilities/state/utils/getSnapshotValue';
import { useMutation } from '@apollo/client';
import { isNonEmptyString } from '@sniptt/guards';
import { useEffect } from 'react';
import { useRecoilCallback, useRecoilValue } from 'recoil';
import {
@@ -103,7 +104,7 @@ export const SSEQuerySubscribeEffect = () => {
);
useEffect(() => {
if (!sseEventStreamId) {
if (!isNonEmptyString(sseEventStreamId)) {
return;
}
@@ -0,0 +1,56 @@
import { useObjectMetadataItems } from '@/object-metadata/hooks/useObjectMetadataItems';
import { dispatchObjectRecordOperationBrowserEvent } from '@/object-record/utils/dispatchObjectRecordOperationBrowserEvent';
import { groupObjectRecordSseEventsByObjectMetadataItemNameSingular } from '@/sse-db-event/utils/groupObjectRecordSseEventsByObjectMetadataItemNameSingular';
import { turnSseObjectRecordEventsToObjectRecordOperationBrowserEvents } from '@/sse-db-event/utils/turnSseObjectRecordEventToObjectRecordOperationBrowserEvent';
import { useCallback } from 'react';
import { isDefined } from 'twenty-shared/utils';
import { type EventWithQueryIds } from '~/generated/graphql';
export const useDispatchObjectRecordEventsFromSseToBrowserEvents = () => {
const { objectMetadataItems } = useObjectMetadataItems();
const dispatchObjectRecordEventsFromSseToBrowserEvents = useCallback(
(eventsWithQueryIds: EventWithQueryIds[]) => {
const objectRecordEvents = eventsWithQueryIds.map((eventWithQueryIds) => {
return eventWithQueryIds.event;
});
const objectRecordEventsByObjectMetadataItemNameSingular =
groupObjectRecordSseEventsByObjectMetadataItemNameSingular({
objectRecordEvents,
});
const objectMetadataItemNamesSingular = Array.from(
objectRecordEventsByObjectMetadataItemNameSingular.keys(),
);
for (const objectMetadataItemNameSingular of objectMetadataItemNamesSingular) {
const objectRecordEventsForThisObjectMetadataItem =
objectRecordEventsByObjectMetadataItemNameSingular.get(
objectMetadataItemNameSingular,
) ?? [];
const objectMetadataItem = objectMetadataItems.find((metadataItem) => {
return metadataItem.nameSingular === objectMetadataItemNameSingular;
});
if (!isDefined(objectMetadataItem)) {
continue;
}
const objectRecordOperationBrowserEvents =
turnSseObjectRecordEventsToObjectRecordOperationBrowserEvents({
objectMetadataItem,
objectRecordEvents: objectRecordEventsForThisObjectMetadataItem,
});
for (const browserEvent of objectRecordOperationBrowserEvents) {
dispatchObjectRecordOperationBrowserEvent(browserEvent);
}
}
},
[objectMetadataItems],
);
return { dispatchObjectRecordEventsFromSseToBrowserEvents };
};
@@ -1,37 +1,16 @@
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) => {
@@ -0,0 +1,83 @@
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 { 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 { v4 } from 'uuid';
import { type EventSubscription } from '~/generated/graphql';
export const useSubscribeToSseEventStream = () => {
const { dispatchObjectRecordEventsFromSseToBrowserEvents } =
useDispatchObjectRecordEventsFromSseToBrowserEvents();
const { triggerOptimisticEffectFromSseEvents } =
useTriggerOptimisticEffectFromSseEvents();
const subscribeToSseEventStream = useRecoilCallback(
({ set, snapshot }) =>
(sseClientConnected: Client) => {
const currentSseEventStreamId = getSnapshotValue(
snapshot,
sseEventStreamIdState,
);
if (isNonEmptyString(currentSseEventStreamId)) {
return;
}
const newSseEventStreamId = v4();
set(sseEventStreamIdState, newSseEventStreamId);
sseClientConnected.subscribe(
{
query: print(ON_EVENT_SUBSCRIPTION),
variables: {
eventStreamId: newSseEventStreamId,
},
},
{
next: (
value: ExecutionResult<{
onEventSubscription: EventSubscription;
}>,
) => {
const objectRecordEventsWithQueryIds =
value?.data?.onEventSubscription?.eventWithQueryIdsList ?? [];
const objectRecordEvents = objectRecordEventsWithQueryIds.map(
(eventWithQueryIds) => {
return eventWithQueryIds.event;
},
);
triggerOptimisticEffectFromSseEvents({
objectRecordEvents,
});
dispatchObjectRecordEventsFromSseToBrowserEvents(
objectRecordEventsWithQueryIds,
);
},
error: (error) => {
captureException(error);
},
complete: () => {},
},
);
},
[
triggerOptimisticEffectFromSseEvents,
dispatchObjectRecordEventsFromSseToBrowserEvents,
],
);
return {
subscribeToSseEventStream,
};
};
@@ -0,0 +1,69 @@
import { useObjectMetadataItems } from '@/object-metadata/hooks/useObjectMetadataItems';
import { useTriggerOptimisticEffectFromSseUpdateEvents } from '@/sse-db-event/hooks/useTriggerOptimisticEffectFromSseUpdateEvents';
import { groupObjectRecordSseEventsByEventType } from '@/sse-db-event/utils/groupObjectRecordSseEventsByEventType';
import { groupObjectRecordSseEventsByObjectMetadataItemNameSingular } from '@/sse-db-event/utils/groupObjectRecordSseEventsByObjectMetadataItemNameSingular';
import { useCallback } from 'react';
import { isDefined } from 'twenty-shared/utils';
import {
DatabaseEventAction,
type ObjectRecordEvent,
} from '~/generated/graphql';
export const useTriggerOptimisticEffectFromSseEvents = () => {
const { objectMetadataItems } = useObjectMetadataItems();
const { triggerOptimisticEffectFromSseUpdateEvents } =
useTriggerOptimisticEffectFromSseUpdateEvents();
const triggerOptimisticEffectFromSseEvents = useCallback(
({ objectRecordEvents }: { objectRecordEvents: ObjectRecordEvent[] }) => {
const objectRecordEventsByObjectMetadataItemNameSingular =
groupObjectRecordSseEventsByObjectMetadataItemNameSingular({
objectRecordEvents,
});
const objectMetadataItemNamesSingular = Array.from(
objectRecordEventsByObjectMetadataItemNameSingular.keys(),
);
for (const objectMetadataItemNameSingular of objectMetadataItemNamesSingular) {
const objectRecordEventsForThisObjectMetadataItem =
objectRecordEventsByObjectMetadataItemNameSingular.get(
objectMetadataItemNameSingular,
) ?? [];
const objectMetadataItem = objectMetadataItems.find((metadataItem) => {
return metadataItem.nameSingular === objectMetadataItemNameSingular;
});
if (!isDefined(objectMetadataItem)) {
continue;
}
const { objectRecordEventsByEventType } =
groupObjectRecordSseEventsByEventType({
objectRecordEvents: objectRecordEventsForThisObjectMetadataItem,
});
const sseEventTypes = Array.from(objectRecordEventsByEventType.keys());
for (const sseEventType of sseEventTypes) {
const objectRecordEventsForThisEventType =
objectRecordEventsByEventType.get(sseEventType) ?? [];
switch (sseEventType) {
case DatabaseEventAction.UPDATED:
triggerOptimisticEffectFromSseUpdateEvents({
objectRecordEvents: objectRecordEventsForThisEventType,
objectMetadataItem,
});
break;
}
}
}
},
[objectMetadataItems, triggerOptimisticEffectFromSseUpdateEvents],
);
return { triggerOptimisticEffectFromSseEvents };
};
@@ -0,0 +1,122 @@
import { triggerUpdateRecordOptimisticEffect } from '@/apollo/optimistic-effect/utils/triggerUpdateRecordOptimisticEffect';
import { useApolloCoreClient } from '@/object-metadata/hooks/useApolloCoreClient';
import { useObjectMetadataItems } from '@/object-metadata/hooks/useObjectMetadataItems';
import { type ObjectMetadataItem } from '@/object-metadata/types/ObjectMetadataItem';
import { getObjectTypename } from '@/object-record/cache/utils/getObjectTypename';
import { getRecordFromCache } from '@/object-record/cache/utils/getRecordFromCache';
import { getRecordNodeFromRecord } from '@/object-record/cache/utils/getRecordNodeFromRecord';
import { generateDepthRecordGqlFieldsFromObject } from '@/object-record/graphql/record-gql-fields/utils/generateDepthRecordGqlFieldsFromObject';
import { useObjectPermissions } from '@/object-record/hooks/useObjectPermissions';
import { useRefetchAggregateQueriesForObjectMetadataItem } from '@/object-record/hooks/useRefetchAggregateQueriesForObjectMetadataItem';
import { useUpsertRecordsInStore } from '@/object-record/record-store/hooks/useUpsertRecordsInStore';
import { useCallback } from 'react';
import { isDefined, isNonEmptyArray } from 'twenty-shared/utils';
import {
DatabaseEventAction,
type ObjectRecordEvent,
} from '~/generated/graphql';
export const useTriggerOptimisticEffectFromSseUpdateEvents = () => {
const apolloCoreClient = useApolloCoreClient();
const { objectMetadataItems } = useObjectMetadataItems();
const { objectPermissionsByObjectMetadataId } = useObjectPermissions();
const { refetchAggregateQueriesForObjectMetadataItem } =
useRefetchAggregateQueriesForObjectMetadataItem();
const { upsertRecordsInStore } = useUpsertRecordsInStore();
const triggerOptimisticEffectFromSseUpdateEvents = useCallback(
({
objectRecordEvents,
objectMetadataItem,
}: {
objectRecordEvents: ObjectRecordEvent[];
objectMetadataItem: ObjectMetadataItem;
}) => {
const recordGqlFields = generateDepthRecordGqlFieldsFromObject({
objectMetadataItem,
objectMetadataItems,
depth: 1,
});
const updateEvents = objectRecordEvents.filter((objectRecordEvent) => {
return objectRecordEvent.action === DatabaseEventAction.UPDATED;
});
for (const updateEvent of updateEvents) {
const updatedRecord = updateEvent.properties.after;
if (!isDefined(updatedRecord)) {
continue;
}
upsertRecordsInStore({ partialRecords: [updatedRecord] });
const cachedRecord = getRecordFromCache({
cache: apolloCoreClient.cache,
objectMetadataItem,
objectMetadataItems,
recordId: updatedRecord.id,
recordGqlFields,
objectPermissionsByObjectMetadataId,
});
const cachedRecordWithConnection = getRecordNodeFromRecord({
record: cachedRecord,
objectMetadataItem,
objectMetadataItems,
recordGqlFields,
computeReferences: false,
});
const computedOptimisticRecord = {
...updatedRecord,
id: updatedRecord.id,
__typename: getObjectTypename(objectMetadataItem.nameSingular),
};
const computedOptimisticRecordWithConnection = getRecordNodeFromRecord({
record: computedOptimisticRecord,
objectMetadataItem,
objectMetadataItems,
recordGqlFields,
});
if (
!isDefined(cachedRecordWithConnection) ||
!isDefined(computedOptimisticRecordWithConnection)
) {
continue;
}
triggerUpdateRecordOptimisticEffect({
cache: apolloCoreClient.cache,
objectMetadataItem,
currentRecord: cachedRecordWithConnection,
updatedRecord: computedOptimisticRecordWithConnection,
objectMetadataItems,
objectPermissionsByObjectMetadataId,
upsertRecordsInStore,
});
}
if (isNonEmptyArray(updateEvents)) {
refetchAggregateQueriesForObjectMetadataItem({
objectMetadataItem,
});
}
return isNonEmptyArray(updateEvents);
},
[
apolloCoreClient.cache,
objectMetadataItems,
objectPermissionsByObjectMetadataId,
refetchAggregateQueriesForObjectMetadataItem,
upsertRecordsInStore,
],
);
return {
triggerOptimisticEffectFromSseUpdateEvents,
};
};
@@ -0,0 +1,26 @@
import { type ObjectRecordOperationUpdateInput } from '@/object-record/types/ObjectRecordOperationUpdateInput';
import { isDefined } from 'twenty-shared/utils';
import { type ObjectRecordEvent } from '~/generated/graphql';
export const getObjectRecordOperationUpdateInputs = (
events: ObjectRecordEvent[],
): ObjectRecordOperationUpdateInput[] => {
return events
.map((event) => {
const updatedFieldNames = event.properties?.updatedFields ?? [];
const updatedRecord = event.properties?.after;
if (!isDefined(updatedRecord)) {
return null;
}
return {
recordId: event.recordId,
updatedFields: updatedFieldNames.map((fieldName) => ({
[fieldName]: updatedRecord[fieldName],
})),
};
})
.filter(isDefined)
.filter((updateInput) => updateInput.updatedFields.length > 0);
};
@@ -0,0 +1,27 @@
import {
type DatabaseEventAction,
type ObjectRecordEvent,
} from '~/generated/graphql';
export const groupObjectRecordSseEventsByEventType = ({
objectRecordEvents,
}: {
objectRecordEvents: ObjectRecordEvent[];
}) => {
const objectRecordEventsByEventType = new Map<
DatabaseEventAction,
ObjectRecordEvent[]
>();
for (const objectRecordEvent of objectRecordEvents) {
const existingObjectRecordEvents =
objectRecordEventsByEventType.get(objectRecordEvent.action) ?? [];
objectRecordEventsByEventType.set(objectRecordEvent.action, [
...existingObjectRecordEvents,
objectRecordEvent,
]);
}
return { objectRecordEventsByEventType };
};
@@ -0,0 +1,26 @@
import { type ObjectRecordEvent } from '~/generated/graphql';
export const groupObjectRecordSseEventsByObjectMetadataItemNameSingular = ({
objectRecordEvents,
}: {
objectRecordEvents: ObjectRecordEvent[];
}) => {
const objectRecordEventsByObjectMetadataItemNameSingular = new Map<
string,
ObjectRecordEvent[]
>();
for (const objectRecordEvent of objectRecordEvents) {
const existingObjectRecordEvents =
objectRecordEventsByObjectMetadataItemNameSingular.get(
objectRecordEvent.objectNameSingular,
) ?? [];
objectRecordEventsByObjectMetadataItemNameSingular.set(
objectRecordEvent.objectNameSingular,
[...existingObjectRecordEvents, objectRecordEvent],
);
}
return objectRecordEventsByObjectMetadataItemNameSingular;
};
@@ -0,0 +1,135 @@
import { type ObjectMetadataItem } from '@/object-metadata/types/ObjectMetadataItem';
import { type ObjectRecordOperationBrowserEventDetail } from '@/object-record/types/ObjectRecordOperationBrowserEventDetail';
import { getObjectRecordOperationUpdateInputs } from '@/sse-db-event/utils/getObjectRecordOperationUpdateInputs';
import { groupObjectRecordSseEventsByEventType } from '@/sse-db-event/utils/groupObjectRecordSseEventsByEventType';
import { assertUnreachable, isDefined } from 'twenty-shared/utils';
import {
DatabaseEventAction,
type ObjectRecordEvent,
} from '~/generated/graphql';
export const turnSseObjectRecordEventsToObjectRecordOperationBrowserEvents = ({
objectMetadataItem,
objectRecordEvents,
}: {
objectMetadataItem: ObjectMetadataItem;
objectRecordEvents: ObjectRecordEvent[];
}): ObjectRecordOperationBrowserEventDetail[] => {
const { objectRecordEventsByEventType } =
groupObjectRecordSseEventsByEventType({
objectRecordEvents,
});
const eventTypes = Array.from(objectRecordEventsByEventType.keys());
const objectRecordOperationBrowserEvents: ObjectRecordOperationBrowserEventDetail[] =
[];
for (const eventType of eventTypes) {
const objectRecordEventsForThisEventType =
objectRecordEventsByEventType.get(eventType) ?? [];
const hasSingleEvent = objectRecordEventsForThisEventType.length === 1;
switch (eventType) {
case DatabaseEventAction.UPDATED: {
const updateInputs = getObjectRecordOperationUpdateInputs(
objectRecordEventsForThisEventType,
);
if (hasSingleEvent) {
objectRecordOperationBrowserEvents.push({
objectMetadataItem,
operation: {
type: 'update-one',
result: { updateInput: updateInputs[0] },
},
});
} else {
objectRecordOperationBrowserEvents.push({
objectMetadataItem,
operation: {
type: 'update-many',
result: { updateInputs },
},
});
}
break;
}
case DatabaseEventAction.DESTROYED:
if (hasSingleEvent) {
objectRecordOperationBrowserEvents.push({
objectMetadataItem,
operation: {
type: 'destroy-one',
},
});
} else {
objectRecordOperationBrowserEvents.push({
objectMetadataItem,
operation: {
type: 'destroy-many',
},
});
}
break;
case DatabaseEventAction.RESTORED:
if (hasSingleEvent) {
objectRecordOperationBrowserEvents.push({
objectMetadataItem,
operation: { type: 'restore-one' },
});
} else {
objectRecordOperationBrowserEvents.push({
objectMetadataItem,
operation: { type: 'restore-many' },
});
}
break;
case DatabaseEventAction.UPSERTED:
if (hasSingleEvent) {
objectRecordOperationBrowserEvents.push({
objectMetadataItem,
operation: { type: 'create-one' },
});
} else {
objectRecordOperationBrowserEvents.push({
objectMetadataItem,
operation: { type: 'create-many' },
});
}
break;
case DatabaseEventAction.CREATED:
if (hasSingleEvent) {
objectRecordOperationBrowserEvents.push({
objectMetadataItem,
operation: { type: 'create-one' },
});
} else {
objectRecordOperationBrowserEvents.push({
objectMetadataItem,
operation: { type: 'create-many' },
});
}
break;
case DatabaseEventAction.DELETED:
if (hasSingleEvent) {
objectRecordOperationBrowserEvents.push({
objectMetadataItem,
operation: { type: 'delete-one' },
});
} else {
objectRecordOperationBrowserEvents.push({
objectMetadataItem,
operation: { type: 'delete-many' },
});
}
break;
default: {
assertUnreachable(eventType);
}
}
}
return objectRecordOperationBrowserEvents.filter(isDefined);
};
@@ -8,6 +8,7 @@ import { ContextStoreComponentInstanceContext } from '@/context-store/states/con
import { MainContainerLayoutWithCommandMenu } from '@/object-record/components/MainContainerLayoutWithCommandMenu';
import { RecordComponentInstanceContextsWrapper } from '@/object-record/components/RecordComponentInstanceContextsWrapper';
import { PageLayoutDispatcher } from '@/object-record/record-show/components/PageLayoutDispatcher';
import { RecordShowPageSSESubscribeEffect } from '@/object-record/record-show/components/RecordShowPageSSESubscribeEffect';
import { useRecordShowPage } from '@/object-record/record-show/hooks/useRecordShowPage';
import { computeRecordShowComponentInstanceId } from '@/object-record/record-show/utils/computeRecordShowComponentInstanceId';
import { PageHeaderToggleCommandMenuButton } from '@/ui/layout/page-header/components/PageHeaderToggleCommandMenuButton';
@@ -63,6 +64,10 @@ export const RecordShowPage = () => {
targetObjectNameSingular: objectNameSingular,
}}
/>
<RecordShowPageSSESubscribeEffect
objectNameSingular={objectNameSingular}
recordId={objectRecordId}
/>
</TimelineActivityContext.Provider>
</MainContainerLayoutWithCommandMenu>
</PageContainer>
@@ -1,5 +1,8 @@
import { type ObjectRecordEvent } from 'twenty-shared/database-events';
import { type DatabaseEventAction } from 'src/engine/api/graphql/graphql-query-runner/enums/database-event-action';
export type ObjectRecordSubscriptionEvent = ObjectRecordEvent & {
action: DatabaseEventAction;
objectNameSingular: string;
};
@@ -1,6 +1,6 @@
import { isPlainObject } from '@nestjs/common/utils/shared.utils';
import { isNull } from '@sniptt/guards';
import { isNonEmptyString, isNull } from '@sniptt/guards';
import {
FieldActorSource,
FieldMetadataType,
@@ -155,7 +155,7 @@ export function formatResult<T>(
compositeProperty.name,
fieldMetadata,
)
: value;
: formatCompositeFieldValue(value, compositeProperty.name, fieldMetadata);
}
// After assembling composite fields, handle those with missing required subfields
@@ -259,6 +259,26 @@ function transformCompositeFieldNullValue(
);
}
function formatCompositeFieldValue(
value: unknown,
compositePropertyName: string,
fieldMetadata: FlatFieldMetadata,
) {
switch (fieldMetadata.type) {
case FieldMetadataType.CURRENCY: {
if (compositePropertyName === 'amountMicros') {
if (isNonEmptyString(value)) {
return parseInt(value);
}
return value;
}
}
}
return value;
}
/**
* Handles composite fields with missing required subfields.
* - For nullable fields: sets to null if all required subfields are null
@@ -8,7 +8,9 @@ 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 { type EventStreamData } from 'src/engine/subscriptions/types/event-stream-data.type';
import { ObjectRecordSubscriptionEvent } from 'src/engine/subscriptions/types/object-record-subscription-event.type';
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';
@Injectable()
export class WorkspaceEventEmitterService {
@@ -98,7 +100,10 @@ export class WorkspaceEventEmitterService {
const objectNameSingular = workspaceEventBatch.objectMetadata.nameSingular;
for (const event of workspaceEventBatch.events) {
const eventWithObjectName = {
const { action } = parseEventNameOrThrow(workspaceEventBatch.name);
const eventWithObjectName: ObjectRecordSubscriptionEvent = {
action,
objectNameSingular,
...event,
};