From d51c988a9e660f661e58e879890be378d4fe9fc3 Mon Sep 17 00:00:00 2001 From: Lucas Bordeau Date: Sun, 18 Jan 2026 17:14:36 +0100 Subject: [PATCH] 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 --- ...chAggregateQueriesForObjectMetadataItem.ts | 26 ++++ .../RecordBoardDataChangedEffect.tsx | 46 ++++++ .../components/RecordBoardEffects.tsx | 4 + .../RecordBoardSSESubscribeEffect.tsx | 28 ++++ ...uldInitializeRecordBoardForUpdateInputs.ts | 100 +++++++++++++ .../RecordCalendarSSESubscribeEffect.tsx | 47 ++++++ .../RecordIndexCalendarContainer.tsx | 2 + .../RecordShowPageSSESubscribeEffect.tsx | 26 ++++ .../RecordTableNoRecordGroupBody.tsx | 2 + ...cordTableVirtualizedSSESubscribeEffect.tsx | 46 ++++++ .../components/SSEEventStreamEffect.tsx | 22 +++ .../sse-db-event/components/SSEProvider.tsx | 4 +- .../components/SSEProviderEffect.tsx | 62 -------- .../components/SSEQuerySubscribeEffect.tsx | 3 +- ...bjectRecordEventsFromSseToBrowserEvents.ts | 56 ++++++++ .../useListenToObjectRecordEventsForQuery.ts | 21 --- .../hooks/useSubscribeToSseEventStream.ts | 83 +++++++++++ ...useTriggerOptimisticEffectFromSseEvents.ts | 69 +++++++++ ...ggerOptimisticEffectFromSseUpdateEvents.ts | 122 ++++++++++++++++ .../getObjectRecordOperationUpdateInputs.ts | 26 ++++ .../groupObjectRecordSseEventsByEventType.ts | 27 ++++ ...eEventsByObjectMetadataItemNameSingular.ts | 26 ++++ ...ventToObjectRecordOperationBrowserEvent.ts | 135 ++++++++++++++++++ .../pages/object-record/RecordShowPage.tsx | 5 + .../object-record-subscription-event.type.ts | 3 + .../twenty-orm/utils/format-result.util.ts | 24 +++- .../workspace-event-emitter.service.ts | 7 +- 27 files changed, 933 insertions(+), 89 deletions(-) create mode 100644 packages/twenty-front/src/modules/object-record/hooks/useRefetchAggregateQueriesForObjectMetadataItem.ts create mode 100644 packages/twenty-front/src/modules/object-record/record-board/components/RecordBoardDataChangedEffect.tsx create mode 100644 packages/twenty-front/src/modules/object-record/record-board/components/RecordBoardSSESubscribeEffect.tsx create mode 100644 packages/twenty-front/src/modules/object-record/record-board/hooks/useGetShouldInitializeRecordBoardForUpdateInputs.ts create mode 100644 packages/twenty-front/src/modules/object-record/record-calendar/components/RecordCalendarSSESubscribeEffect.tsx create mode 100644 packages/twenty-front/src/modules/object-record/record-show/components/RecordShowPageSSESubscribeEffect.tsx create mode 100644 packages/twenty-front/src/modules/object-record/record-table/virtualization/components/RecordTableVirtualizedSSESubscribeEffect.tsx create mode 100644 packages/twenty-front/src/modules/sse-db-event/components/SSEEventStreamEffect.tsx delete mode 100644 packages/twenty-front/src/modules/sse-db-event/components/SSEProviderEffect.tsx create mode 100644 packages/twenty-front/src/modules/sse-db-event/hooks/useDispatchObjectRecordEventsFromSseToBrowserEvents.ts create mode 100644 packages/twenty-front/src/modules/sse-db-event/hooks/useSubscribeToSseEventStream.ts create mode 100644 packages/twenty-front/src/modules/sse-db-event/hooks/useTriggerOptimisticEffectFromSseEvents.ts create mode 100644 packages/twenty-front/src/modules/sse-db-event/hooks/useTriggerOptimisticEffectFromSseUpdateEvents.ts create mode 100644 packages/twenty-front/src/modules/sse-db-event/utils/getObjectRecordOperationUpdateInputs.ts create mode 100644 packages/twenty-front/src/modules/sse-db-event/utils/groupObjectRecordSseEventsByEventType.ts create mode 100644 packages/twenty-front/src/modules/sse-db-event/utils/groupObjectRecordSseEventsByObjectMetadataItemNameSingular.ts create mode 100644 packages/twenty-front/src/modules/sse-db-event/utils/turnSseObjectRecordEventToObjectRecordOperationBrowserEvent.ts diff --git a/packages/twenty-front/src/modules/object-record/hooks/useRefetchAggregateQueriesForObjectMetadataItem.ts b/packages/twenty-front/src/modules/object-record/hooks/useRefetchAggregateQueriesForObjectMetadataItem.ts new file mode 100644 index 0000000000..c132d841ca --- /dev/null +++ b/packages/twenty-front/src/modules/object-record/hooks/useRefetchAggregateQueriesForObjectMetadataItem.ts @@ -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, + }; +}; diff --git a/packages/twenty-front/src/modules/object-record/record-board/components/RecordBoardDataChangedEffect.tsx b/packages/twenty-front/src/modules/object-record/record-board/components/RecordBoardDataChangedEffect.tsx new file mode 100644 index 0000000000..95f133260b --- /dev/null +++ b/packages/twenty-front/src/modules/object-record/record-board/components/RecordBoardDataChangedEffect.tsx @@ -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; +}; diff --git a/packages/twenty-front/src/modules/object-record/record-board/components/RecordBoardEffects.tsx b/packages/twenty-front/src/modules/object-record/record-board/components/RecordBoardEffects.tsx index 87c82fa808..16e3a0ffcc 100644 --- a/packages/twenty-front/src/modules/object-record/record-board/components/RecordBoardEffects.tsx +++ b/packages/twenty-front/src/modules/object-record/record-board/components/RecordBoardEffects.tsx @@ -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 = () => { + + diff --git a/packages/twenty-front/src/modules/object-record/record-board/components/RecordBoardSSESubscribeEffect.tsx b/packages/twenty-front/src/modules/object-record/record-board/components/RecordBoardSSESubscribeEffect.tsx new file mode 100644 index 0000000000..f9022387db --- /dev/null +++ b/packages/twenty-front/src/modules/object-record/record-board/components/RecordBoardSSESubscribeEffect.tsx @@ -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; +}; diff --git a/packages/twenty-front/src/modules/object-record/record-board/hooks/useGetShouldInitializeRecordBoardForUpdateInputs.ts b/packages/twenty-front/src/modules/object-record/record-board/hooks/useGetShouldInitializeRecordBoardForUpdateInputs.ts new file mode 100644 index 0000000000..00f80566c8 --- /dev/null +++ b/packages/twenty-front/src/modules/object-record/record-board/hooks/useGetShouldInitializeRecordBoardForUpdateInputs.ts @@ -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(); + 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, + }; +}; diff --git a/packages/twenty-front/src/modules/object-record/record-calendar/components/RecordCalendarSSESubscribeEffect.tsx b/packages/twenty-front/src/modules/object-record/record-calendar/components/RecordCalendarSSESubscribeEffect.tsx new file mode 100644 index 0000000000..d2e4c78ea1 --- /dev/null +++ b/packages/twenty-front/src/modules/object-record/record-calendar/components/RecordCalendarSSESubscribeEffect.tsx @@ -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; +}; diff --git a/packages/twenty-front/src/modules/object-record/record-index/components/RecordIndexCalendarContainer.tsx b/packages/twenty-front/src/modules/object-record/record-index/components/RecordIndexCalendarContainer.tsx index ba77adc3b1..b9ed6e063d 100644 --- a/packages/twenty-front/src/modules/object-record/record-index/components/RecordIndexCalendarContainer.tsx +++ b/packages/twenty-front/src/modules/object-record/record-index/components/RecordIndexCalendarContainer.tsx @@ -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 = ({ }} > + diff --git a/packages/twenty-front/src/modules/object-record/record-show/components/RecordShowPageSSESubscribeEffect.tsx b/packages/twenty-front/src/modules/object-record/record-show/components/RecordShowPageSSESubscribeEffect.tsx new file mode 100644 index 0000000000..0542c00e51 --- /dev/null +++ b/packages/twenty-front/src/modules/object-record/record-show/components/RecordShowPageSSESubscribeEffect.tsx @@ -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; +}; diff --git a/packages/twenty-front/src/modules/object-record/record-table/record-table-body/components/RecordTableNoRecordGroupBody.tsx b/packages/twenty-front/src/modules/object-record/record-table/record-table-body/components/RecordTableNoRecordGroupBody.tsx index 7a67457936..517ee62867 100644 --- a/packages/twenty-front/src/modules/object-record/record-table/record-table-body/components/RecordTableNoRecordGroupBody.tsx +++ b/packages/twenty-front/src/modules/object-record/record-table/record-table-body/components/RecordTableNoRecordGroupBody.tsx @@ -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 = () => { )} + ); diff --git a/packages/twenty-front/src/modules/object-record/record-table/virtualization/components/RecordTableVirtualizedSSESubscribeEffect.tsx b/packages/twenty-front/src/modules/object-record/record-table/virtualization/components/RecordTableVirtualizedSSESubscribeEffect.tsx new file mode 100644 index 0000000000..771f6f14df --- /dev/null +++ b/packages/twenty-front/src/modules/object-record/record-table/virtualization/components/RecordTableVirtualizedSSESubscribeEffect.tsx @@ -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; +}; diff --git a/packages/twenty-front/src/modules/sse-db-event/components/SSEEventStreamEffect.tsx b/packages/twenty-front/src/modules/sse-db-event/components/SSEEventStreamEffect.tsx new file mode 100644 index 0000000000..2b7c7441ae --- /dev/null +++ b/packages/twenty-front/src/modules/sse-db-event/components/SSEEventStreamEffect.tsx @@ -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; +}; 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 0e0dbab789..74e8738af0 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,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 ( - + {children} diff --git a/packages/twenty-front/src/modules/sse-db-event/components/SSEProviderEffect.tsx b/packages/twenty-front/src/modules/sse-db-event/components/SSEProviderEffect.tsx deleted file mode 100644 index b9806d08f1..0000000000 --- a/packages/twenty-front/src/modules/sse-db-event/components/SSEProviderEffect.tsx +++ /dev/null @@ -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; -}; 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 a75cc51f27..1b8bdc70f8 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 @@ -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; } diff --git a/packages/twenty-front/src/modules/sse-db-event/hooks/useDispatchObjectRecordEventsFromSseToBrowserEvents.ts b/packages/twenty-front/src/modules/sse-db-event/hooks/useDispatchObjectRecordEventsFromSseToBrowserEvents.ts new file mode 100644 index 0000000000..32e7c863f7 --- /dev/null +++ b/packages/twenty-front/src/modules/sse-db-event/hooks/useDispatchObjectRecordEventsFromSseToBrowserEvents.ts @@ -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 }; +}; diff --git a/packages/twenty-front/src/modules/sse-db-event/hooks/useListenToObjectRecordEventsForQuery.ts b/packages/twenty-front/src/modules/sse-db-event/hooks/useListenToObjectRecordEventsForQuery.ts index 52223b53e2..0d121223d1 100644 --- a/packages/twenty-front/src/modules/sse-db-event/hooks/useListenToObjectRecordEventsForQuery.ts +++ b/packages/twenty-front/src/modules/sse-db-event/hooks/useListenToObjectRecordEventsForQuery.ts @@ -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) - .detail; - - onObjectRecordEvents(objectRecordEvents); - }; - - window.addEventListener(eventName, handleOnObjectRecordEventsForQuery); - - return () => { - window.removeEventListener(eventName, handleOnObjectRecordEventsForQuery); - }; - }, [onObjectRecordEvents, queryId]); - const changeQueryIdListenState = useRecoilCallback( ({ set, snapshot }) => (shouldListen: boolean, queryId: string) => { diff --git a/packages/twenty-front/src/modules/sse-db-event/hooks/useSubscribeToSseEventStream.ts b/packages/twenty-front/src/modules/sse-db-event/hooks/useSubscribeToSseEventStream.ts new file mode 100644 index 0000000000..bb37966bc2 --- /dev/null +++ b/packages/twenty-front/src/modules/sse-db-event/hooks/useSubscribeToSseEventStream.ts @@ -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, + }; +}; diff --git a/packages/twenty-front/src/modules/sse-db-event/hooks/useTriggerOptimisticEffectFromSseEvents.ts b/packages/twenty-front/src/modules/sse-db-event/hooks/useTriggerOptimisticEffectFromSseEvents.ts new file mode 100644 index 0000000000..6a62b899d0 --- /dev/null +++ b/packages/twenty-front/src/modules/sse-db-event/hooks/useTriggerOptimisticEffectFromSseEvents.ts @@ -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 }; +}; diff --git a/packages/twenty-front/src/modules/sse-db-event/hooks/useTriggerOptimisticEffectFromSseUpdateEvents.ts b/packages/twenty-front/src/modules/sse-db-event/hooks/useTriggerOptimisticEffectFromSseUpdateEvents.ts new file mode 100644 index 0000000000..d332429ae2 --- /dev/null +++ b/packages/twenty-front/src/modules/sse-db-event/hooks/useTriggerOptimisticEffectFromSseUpdateEvents.ts @@ -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, + }; +}; diff --git a/packages/twenty-front/src/modules/sse-db-event/utils/getObjectRecordOperationUpdateInputs.ts b/packages/twenty-front/src/modules/sse-db-event/utils/getObjectRecordOperationUpdateInputs.ts new file mode 100644 index 0000000000..efb725d290 --- /dev/null +++ b/packages/twenty-front/src/modules/sse-db-event/utils/getObjectRecordOperationUpdateInputs.ts @@ -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); +}; diff --git a/packages/twenty-front/src/modules/sse-db-event/utils/groupObjectRecordSseEventsByEventType.ts b/packages/twenty-front/src/modules/sse-db-event/utils/groupObjectRecordSseEventsByEventType.ts new file mode 100644 index 0000000000..83f658bcd0 --- /dev/null +++ b/packages/twenty-front/src/modules/sse-db-event/utils/groupObjectRecordSseEventsByEventType.ts @@ -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 }; +}; diff --git a/packages/twenty-front/src/modules/sse-db-event/utils/groupObjectRecordSseEventsByObjectMetadataItemNameSingular.ts b/packages/twenty-front/src/modules/sse-db-event/utils/groupObjectRecordSseEventsByObjectMetadataItemNameSingular.ts new file mode 100644 index 0000000000..5af512957e --- /dev/null +++ b/packages/twenty-front/src/modules/sse-db-event/utils/groupObjectRecordSseEventsByObjectMetadataItemNameSingular.ts @@ -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; +}; diff --git a/packages/twenty-front/src/modules/sse-db-event/utils/turnSseObjectRecordEventToObjectRecordOperationBrowserEvent.ts b/packages/twenty-front/src/modules/sse-db-event/utils/turnSseObjectRecordEventToObjectRecordOperationBrowserEvent.ts new file mode 100644 index 0000000000..2c5cd0befa --- /dev/null +++ b/packages/twenty-front/src/modules/sse-db-event/utils/turnSseObjectRecordEventToObjectRecordOperationBrowserEvent.ts @@ -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); +}; diff --git a/packages/twenty-front/src/pages/object-record/RecordShowPage.tsx b/packages/twenty-front/src/pages/object-record/RecordShowPage.tsx index ed4fd05d9b..0f08d76883 100644 --- a/packages/twenty-front/src/pages/object-record/RecordShowPage.tsx +++ b/packages/twenty-front/src/pages/object-record/RecordShowPage.tsx @@ -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, }} /> + diff --git a/packages/twenty-server/src/engine/subscriptions/types/object-record-subscription-event.type.ts b/packages/twenty-server/src/engine/subscriptions/types/object-record-subscription-event.type.ts index 359c7fc594..e72be402d4 100644 --- a/packages/twenty-server/src/engine/subscriptions/types/object-record-subscription-event.type.ts +++ b/packages/twenty-server/src/engine/subscriptions/types/object-record-subscription-event.type.ts @@ -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; }; diff --git a/packages/twenty-server/src/engine/twenty-orm/utils/format-result.util.ts b/packages/twenty-server/src/engine/twenty-orm/utils/format-result.util.ts index 390118e996..673c2d6d0a 100644 --- a/packages/twenty-server/src/engine/twenty-orm/utils/format-result.util.ts +++ b/packages/twenty-server/src/engine/twenty-orm/utils/format-result.util.ts @@ -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( 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 diff --git a/packages/twenty-server/src/engine/workspace-event-emitter/workspace-event-emitter.service.ts b/packages/twenty-server/src/engine/workspace-event-emitter/workspace-event-emitter.service.ts index 78607021b8..879b9da49d 100644 --- a/packages/twenty-server/src/engine/workspace-event-emitter/workspace-event-emitter.service.ts +++ b/packages/twenty-server/src/engine/workspace-event-emitter/workspace-event-emitter.service.ts @@ -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, };