From 9ecab8fb8266e9d5e85bacad63d3e2746bfdbadd Mon Sep 17 00:00:00 2001 From: Lucas Bordeau Date: Fri, 23 Jan 2026 18:55:07 +0100 Subject: [PATCH] Fix event logic for soft-delete and restore (#17393) This PR changes the shape and logic of SSE events `DELETE` and `RESTORE`, because they behave like `UPDATE` events in practice, they should share the same logic. Before this PR, it was impossible for the frontend to obtain the `deletedAt` value, and the logic to handle soft-delete and restore would have been flawed. Because there is a typing confusion in the parameters of `formatTwentyOrmEventToDatabaseBatchEvent`, due to TypeORM, we also update this util to only accept an array of records, instead of `T | T[]`. We should improve our TypeORM layer in the future. Also the naming was not clear, so we clearly use `recordsAfter` and `recordsBefore` as much as possible, because that is what we have at the end in events. Events are sent from their respective query builders, so these last ones have been updated also. Because TypeORM `soft-remove` operation only returns record ids, we add `.getMany()` to fetch all fields for soft-removed records, so that our event can have before and after. --- .../listeners/entity-events-to-db.listener.ts | 5 +- .../workspace-entity-manager.ts | 81 ++++++- .../workspace-delete-query-builder.ts | 11 +- .../workspace-insert-query-builder.ts | 4 +- .../workspace-soft-delete-query-builder.ts | 13 +- .../workspace-update-query-builder.ts | 16 +- ...event-to-database-batch-event.util.spec.ts | 37 +--- ...-orm-event-to-database-batch-event.util.ts | 207 ++++++++++++------ ...orkflow-database-event-trigger.listener.ts | 8 +- .../object-record-delete.event.ts | 4 + .../object-record-restore.event.ts | 4 + .../object-record-update.event.ts | 6 +- 12 files changed, 262 insertions(+), 134 deletions(-) diff --git a/packages/twenty-server/src/engine/api/graphql/workspace-query-runner/listeners/entity-events-to-db.listener.ts b/packages/twenty-server/src/engine/api/graphql/workspace-query-runner/listeners/entity-events-to-db.listener.ts index 8c828aae23..e8ecc7a5b4 100644 --- a/packages/twenty-server/src/engine/api/graphql/workspace-query-runner/listeners/entity-events-to-db.listener.ts +++ b/packages/twenty-server/src/engine/api/graphql/workspace-query-runner/listeners/entity-events-to-db.listener.ts @@ -109,7 +109,10 @@ export class EntityEventsToDbListener { promises.push( this.entityEventsToDbQueueService.add< WorkspaceEventBatch - >(UpsertTimelineActivityFromInternalEvent.name, batchEvent), + >( + UpsertTimelineActivityFromInternalEvent.name, + batchEvent as WorkspaceEventBatch, + ), ); } } diff --git a/packages/twenty-server/src/engine/twenty-orm/entity-manager/workspace-entity-manager.ts b/packages/twenty-server/src/engine/twenty-orm/entity-manager/workspace-entity-manager.ts index deca11125a..e75c4969e4 100644 --- a/packages/twenty-server/src/engine/twenty-orm/entity-manager/workspace-entity-manager.ts +++ b/packages/twenty-server/src/engine/twenty-orm/entity-manager/workspace-entity-manager.ts @@ -45,6 +45,7 @@ import { PermissionsException, PermissionsExceptionCode, } from 'src/engine/metadata-modules/permissions/permissions.exception'; +import { type BaseWorkspaceEntity } from 'src/engine/twenty-orm/base.workspace-entity'; import { type DeepPartialWithNestedRelationFields } from 'src/engine/twenty-orm/entity-manager/types/deep-partial-entity-with-nested-relation-fields.type'; import { type QueryDeepPartialEntityWithNestedRelationFields } from 'src/engine/twenty-orm/entity-manager/types/query-deep-partial-entity-with-nested-relation-fields.type'; import { getEntityTarget } from 'src/engine/twenty-orm/entity-manager/utils/get-entity-target'; @@ -1272,8 +1273,8 @@ export class WorkspaceEntityManager extends EntityManager { objectMetadataItem, flatFieldMetadataMaps: this.internalContext.flatFieldMetadataMaps, workspaceId: this.internalContext.workspaceId, - entities: updatedEntities, - beforeEntities: updatedEntities.map( + recordsAfter: updatedEntities, + recordsBefore: updatedEntities.map( (entity) => beforeUpdateMapById[entity.id], ), }), @@ -1285,7 +1286,7 @@ export class WorkspaceEntityManager extends EntityManager { objectMetadataItem, flatFieldMetadataMaps: this.internalContext.flatFieldMetadataMaps, workspaceId: this.internalContext.workspaceId, - entities: createdEntities, + recordsAfter: createdEntities, }), ); @@ -1471,13 +1472,17 @@ export class WorkspaceEntityManager extends EntityManager { this.internalContext.flatFieldMetadataMaps, ); + const recordsBefore = Array.isArray(formattedResult) + ? formattedResult + : [formattedResult]; + this.internalContext.eventEmitterService.emitDatabaseBatchEvent( formatTwentyOrmEventToDatabaseBatchEvent({ action: DatabaseEventAction.DESTROYED, objectMetadataItem, flatFieldMetadataMaps: this.internalContext.flatFieldMetadataMaps, workspaceId: this.internalContext.workspaceId, - entities: formattedResult, + recordsBefore, }), ); @@ -1562,6 +1567,29 @@ export class WorkspaceEntityManager extends EntityManager { const entityTarget = target ?? (isEntityArray ? entity[0]?.constructor : entity.constructor); + const entityArray = isEntityArray ? entity : [entity]; + + const entityIds = entityArray + .map((entity) => (entity as { id: string }).id) + .filter(isDefined); + + const recordsBeforeFindResult = await this.find( + entityTarget, + { + where: { id: In(entityIds) }, + }, + { shouldBypassPermissionChecks: true }, // Bypass as this is for event emission + ); + + const beforeUpdateMapById = recordsBeforeFindResult.reduce( + (acc, e: BaseWorkspaceEntity) => { + acc[e.id] = e; + + return acc; + }, + {} as Record, + ); + const objectMetadataItem = getObjectMetadataFromEntityTarget( entityTarget, this.internalContext, @@ -1592,13 +1620,22 @@ export class WorkspaceEntityManager extends EntityManager { this.internalContext.flatFieldMetadataMaps, ); + const recordsAfter = Array.isArray(formattedResult) + ? formattedResult + : [formattedResult]; + + const recordsBefore = recordsAfter.map( + (record) => beforeUpdateMapById[record.id] as unknown as Entity, + ); + this.internalContext.eventEmitterService.emitDatabaseBatchEvent( formatTwentyOrmEventToDatabaseBatchEvent({ action: DatabaseEventAction.DELETED, objectMetadataItem, flatFieldMetadataMaps: this.internalContext.flatFieldMetadataMaps, workspaceId: this.internalContext.workspaceId, - entities: formattedResult, + recordsAfter, + recordsBefore, }), ); @@ -1679,6 +1716,29 @@ export class WorkspaceEntityManager extends EntityManager { const entityTarget = target ?? (isEntityArray ? entity[0]?.constructor : entity.constructor); + const entityArray = isEntityArray ? entity : [entity]; + + const entityIds = entityArray + .map((entity) => (entity as { id: string }).id) + .filter(isDefined); + + const recordsBeforeFindResult = await this.find( + entityTarget, + { + where: { id: In(entityIds) }, + }, + { shouldBypassPermissionChecks: true }, // Bypass as this is for event emission + ); + + const beforeUpdateMapById = recordsBeforeFindResult.reduce( + (acc, e: BaseWorkspaceEntity) => { + acc[e.id] = e; + + return acc; + }, + {} as Record, + ); + const objectMetadataItem = getObjectMetadataFromEntityTarget( entityTarget, this.internalContext, @@ -1709,13 +1769,22 @@ export class WorkspaceEntityManager extends EntityManager { this.internalContext.flatFieldMetadataMaps, ); + const recordsAfter = Array.isArray(formattedResult) + ? formattedResult + : [formattedResult]; + + const recordsBefore = recordsAfter.map( + (record) => beforeUpdateMapById[record.id] as unknown as Entity, + ); + this.internalContext.eventEmitterService.emitDatabaseBatchEvent( formatTwentyOrmEventToDatabaseBatchEvent({ action: DatabaseEventAction.RESTORED, objectMetadataItem, flatFieldMetadataMaps: this.internalContext.flatFieldMetadataMaps, workspaceId: this.internalContext.workspaceId, - entities: formattedResult, + recordsAfter, + recordsBefore, }), ); diff --git a/packages/twenty-server/src/engine/twenty-orm/repository/workspace-delete-query-builder.ts b/packages/twenty-server/src/engine/twenty-orm/repository/workspace-delete-query-builder.ts index 7b00df30e4..7495dd94b1 100644 --- a/packages/twenty-server/src/engine/twenty-orm/repository/workspace-delete-query-builder.ts +++ b/packages/twenty-server/src/engine/twenty-orm/repository/workspace-delete-query-builder.ts @@ -1,4 +1,5 @@ import { type ObjectsPermissions } from 'twenty-shared/types'; +import { isDefined } from 'twenty-shared/utils'; import { DeleteQueryBuilder, type DeleteResult, @@ -120,20 +121,26 @@ export class WorkspaceDeleteQueryBuilder< this.internalContext.flatFieldMetadataMaps, ); - const formattedBefore = formatResult( + const formattedBefore = formatResult( before, objectMetadata, this.internalContext.flatObjectMetadataMaps, this.internalContext.flatFieldMetadataMaps, ); + const recordsBefore = isDefined(formattedBefore) + ? Array.isArray(formattedBefore) + ? formattedBefore + : [formattedBefore] + : []; + this.internalContext.eventEmitterService.emitDatabaseBatchEvent( formatTwentyOrmEventToDatabaseBatchEvent({ action: DatabaseEventAction.DESTROYED, objectMetadataItem: objectMetadata, flatFieldMetadataMaps: this.internalContext.flatFieldMetadataMaps, workspaceId: this.internalContext.workspaceId, - entities: formattedBefore, + recordsBefore, authContext: this.authContext, }), ); diff --git a/packages/twenty-server/src/engine/twenty-orm/repository/workspace-insert-query-builder.ts b/packages/twenty-server/src/engine/twenty-orm/repository/workspace-insert-query-builder.ts index 64eefca8a2..d18ac2d940 100644 --- a/packages/twenty-server/src/engine/twenty-orm/repository/workspace-insert-query-builder.ts +++ b/packages/twenty-server/src/engine/twenty-orm/repository/workspace-insert-query-builder.ts @@ -213,7 +213,7 @@ export class WorkspaceInsertQueryBuilder< objectMetadataItem: objectMetadata, flatFieldMetadataMaps: this.internalContext.flatFieldMetadataMaps, workspaceId: this.internalContext.workspaceId, - entities: formattedResultForEvent, + recordsAfter: formattedResultForEvent, authContext: this.authContext, }), ); @@ -224,7 +224,7 @@ export class WorkspaceInsertQueryBuilder< objectMetadataItem: objectMetadata, flatFieldMetadataMaps: this.internalContext.flatFieldMetadataMaps, workspaceId: this.internalContext.workspaceId, - entities: formattedResultForEvent, + recordsAfter: formattedResultForEvent, authContext: this.authContext, }), ); diff --git a/packages/twenty-server/src/engine/twenty-orm/repository/workspace-soft-delete-query-builder.ts b/packages/twenty-server/src/engine/twenty-orm/repository/workspace-soft-delete-query-builder.ts index b362df1e02..9cda9f8b1d 100644 --- a/packages/twenty-server/src/engine/twenty-orm/repository/workspace-soft-delete-query-builder.ts +++ b/packages/twenty-server/src/engine/twenty-orm/repository/workspace-soft-delete-query-builder.ts @@ -111,10 +111,12 @@ export class WorkspaceSoftDeleteQueryBuilder< aliasName: objectMetadata.nameSingular, }) as WhereClause[]; - const after = await super.execute(); + const typeORMSoftRemoveResultWithOnlyIdColumn = await super.execute(); + + const afterWithAllFields = await beforeEventSelectQueryBuilder.getMany(); const formattedAfter = formatResult( - after.raw, + afterWithAllFields, objectMetadata, this.internalContext.flatObjectMetadataMaps, this.internalContext.flatFieldMetadataMaps, @@ -133,15 +135,16 @@ export class WorkspaceSoftDeleteQueryBuilder< objectMetadataItem: objectMetadata, flatFieldMetadataMaps: this.internalContext.flatFieldMetadataMaps, workspaceId: this.internalContext.workspaceId, - entities: formattedBefore, + recordsBefore: formattedBefore, + recordsAfter: formattedAfter, authContext: this.authContext, }), ); return { - raw: after.raw, + raw: typeORMSoftRemoveResultWithOnlyIdColumn.raw, generatedMaps: formattedAfter, - affected: after.affected, + affected: typeORMSoftRemoveResultWithOnlyIdColumn.affected, }; } catch (error) { throw await computeTwentyORMException(error); diff --git a/packages/twenty-server/src/engine/twenty-orm/repository/workspace-update-query-builder.ts b/packages/twenty-server/src/engine/twenty-orm/repository/workspace-update-query-builder.ts index 445c47f6ca..54424f7caf 100644 --- a/packages/twenty-server/src/engine/twenty-orm/repository/workspace-update-query-builder.ts +++ b/packages/twenty-server/src/engine/twenty-orm/repository/workspace-update-query-builder.ts @@ -209,8 +209,8 @@ export class WorkspaceUpdateQueryBuilder< objectMetadataItem: objectMetadata, flatFieldMetadataMaps: this.internalContext.flatFieldMetadataMaps, workspaceId: this.internalContext.workspaceId, - entities: formattedAfter, - beforeEntities: formattedBefore, + recordsAfter: formattedAfter, + recordsBefore: formattedBefore, authContext: this.authContext, }), ); @@ -221,8 +221,8 @@ export class WorkspaceUpdateQueryBuilder< objectMetadataItem: objectMetadata, flatFieldMetadataMaps: this.internalContext.flatFieldMetadataMaps, workspaceId: this.internalContext.workspaceId, - entities: formattedAfter, - beforeEntities: formattedBefore, + recordsAfter: formattedAfter, + recordsBefore: formattedBefore, authContext: this.authContext, }), ); @@ -388,8 +388,8 @@ export class WorkspaceUpdateQueryBuilder< objectMetadataItem: objectMetadata, flatFieldMetadataMaps: this.internalContext.flatFieldMetadataMaps, workspaceId: this.internalContext.workspaceId, - entities: formattedAfter, - beforeEntities: formattedBefore, + recordsAfter: formattedAfter, + recordsBefore: formattedBefore, authContext: this.authContext, }), ); @@ -400,8 +400,8 @@ export class WorkspaceUpdateQueryBuilder< objectMetadataItem: objectMetadata, flatFieldMetadataMaps: this.internalContext.flatFieldMetadataMaps, workspaceId: this.internalContext.workspaceId, - entities: formattedAfter, - beforeEntities: formattedBefore, + recordsAfter: formattedAfter, + recordsBefore: formattedBefore, authContext: this.authContext, }), ); diff --git a/packages/twenty-server/src/engine/twenty-orm/utils/__tests__/format-twenty-orm-event-to-database-batch-event.util.spec.ts b/packages/twenty-server/src/engine/twenty-orm/utils/__tests__/format-twenty-orm-event-to-database-batch-event.util.spec.ts index 6a30ed6990..f040bb5685 100644 --- a/packages/twenty-server/src/engine/twenty-orm/utils/__tests__/format-twenty-orm-event-to-database-batch-event.util.spec.ts +++ b/packages/twenty-server/src/engine/twenty-orm/utils/__tests__/format-twenty-orm-event-to-database-batch-event.util.spec.ts @@ -1,5 +1,5 @@ -import { FieldMetadataType } from 'twenty-shared/types'; import { type ObjectRecordUpdateEvent } from 'twenty-shared/database-events'; +import { FieldMetadataType } from 'twenty-shared/types'; import { DatabaseEventAction } from 'src/engine/api/graphql/graphql-query-runner/enums/database-event-action'; import { type FlatEntityMaps } from 'src/engine/metadata-modules/flat-entity/types/flat-entity-maps.type'; @@ -114,8 +114,8 @@ describe('formatTwentyOrmEventToDatabaseBatchEvent', () => { flatFieldMetadataMaps, workspaceId: mockWorkspaceId, authContext: mockAuthContext, - entities: afterEntities, - beforeEntities: beforeEntities, + recordsAfter: afterEntities, + recordsBefore: beforeEntities, }); } catch (error) { expect(error).toBeInstanceOf(TwentyORMException); @@ -157,8 +157,8 @@ describe('formatTwentyOrmEventToDatabaseBatchEvent', () => { flatFieldMetadataMaps, workspaceId: mockWorkspaceId, authContext: mockAuthContext, - entities: afterEntities, - beforeEntities: beforeEntities, + recordsAfter: afterEntities, + recordsBefore: beforeEntities, }); expect(result).toBeDefined(); @@ -180,32 +180,5 @@ describe('formatTwentyOrmEventToDatabaseBatchEvent', () => { expect(updateEvent2.properties?.before?.name).toBe('Jane Doe'); expect(updateEvent2.properties?.after?.name).toBe('Jane Doe Updated'); }); - - it('should handle single entity (non-array) for both before and after', () => { - const afterEntity = { - id: 'record-1', - name: 'John Doe Updated', - }; - - const beforeEntity = { - id: 'record-1', - name: 'John Doe', - }; - - const result = formatTwentyOrmEventToDatabaseBatchEvent({ - action: DatabaseEventAction.UPDATED, - objectMetadataItem: flatObjectMetadata, - flatFieldMetadataMaps, - workspaceId: mockWorkspaceId, - authContext: mockAuthContext, - entities: afterEntity, - beforeEntities: beforeEntity, - }); - - expect(result).toBeDefined(); - expect(result?.action).toBe(DatabaseEventAction.UPDATED); - expect(result?.events).toHaveLength(1); - expect(result?.events[0].recordId).toBe('record-1'); - }); }); }); diff --git a/packages/twenty-server/src/engine/twenty-orm/utils/format-twenty-orm-event-to-database-batch-event.util.ts b/packages/twenty-server/src/engine/twenty-orm/utils/format-twenty-orm-event-to-database-batch-event.util.ts index 31a47aba79..d474ef3d7d 100644 --- a/packages/twenty-server/src/engine/twenty-orm/utils/format-twenty-orm-event-to-database-batch-event.util.ts +++ b/packages/twenty-server/src/engine/twenty-orm/utils/format-twenty-orm-event-to-database-batch-event.util.ts @@ -1,13 +1,18 @@ -import { STANDARD_OBJECT_IDS } from 'twenty-shared/metadata'; -import { isDefined } from 'twenty-shared/utils'; import { ObjectRecordCreateEvent, ObjectRecordDeleteEvent, ObjectRecordDestroyEvent, + ObjectRecordRestoreEvent, ObjectRecordUpdateEvent, ObjectRecordUpsertEvent, type ObjectRecordDiff, } from 'twenty-shared/database-events'; +import { STANDARD_OBJECT_IDS } from 'twenty-shared/metadata'; +import { + assertUnreachable, + isDefined, + isNonEmptyArray, +} from 'twenty-shared/utils'; import type { ObjectLiteral } from 'typeorm'; @@ -31,68 +36,98 @@ export const formatTwentyOrmEventToDatabaseBatchEvent = < flatFieldMetadataMaps, workspaceId, authContext, - entities, - beforeEntities, + recordsAfter, + recordsBefore, }: { action: DatabaseEventAction; objectMetadataItem: FlatObjectMetadata; flatFieldMetadataMaps: FlatEntityMaps; workspaceId: string; authContext?: AuthContext; - entities: T | T[]; - beforeEntities?: T | T[]; + recordsAfter?: T[]; + recordsBefore?: T[]; }): DatabaseBatchEventInput | undefined => { if (objectMetadataItem.standardId === STANDARD_OBJECT_IDS.timelineActivity) { return; } const objectMetadataNameSingular = objectMetadataItem.nameSingular; - const entityArray = isDefined(entities) - ? Array.isArray(entities) - ? entities - : [entities] - : []; + let events: ( - | ObjectRecordCreateEvent - | ObjectRecordUpdateEvent | ObjectRecordDeleteEvent + | ObjectRecordRestoreEvent + | ObjectRecordUpdateEvent + | ObjectRecordCreateEvent + | ObjectRecordDestroyEvent | ObjectRecordUpsertEvent )[] = []; switch (action) { - case DatabaseEventAction.CREATED: - events = entityArray.map((after) => { - const event = new ObjectRecordCreateEvent(); + case DatabaseEventAction.CREATED: { + if (!isDefined(recordsAfter)) { + throw new Error( + `recordsAfter is required for ${action.toUpperCase()} action`, + ); + } - event.userId = authContext?.user?.id; - event.workspaceMemberId = authContext?.workspaceMemberId; - event.recordId = after.id; - event.properties = { after }; + if (!isNonEmptyArray(recordsAfter)) { + break; + } - return event; - }); + events = + recordsAfter?.map((recordAfter) => { + const event = new ObjectRecordCreateEvent(); + + event.userId = authContext?.user?.id; + event.workspaceMemberId = authContext?.workspaceMemberId; + event.recordId = recordAfter.id; + event.properties = { after: recordAfter }; + + return event; + }) ?? []; break; + } case DatabaseEventAction.UPDATED: - events = entityArray - .map((after) => { - if (!beforeEntities) { - throw new Error('beforeEntities is required for UPDATED action'); + case DatabaseEventAction.DELETED: + case DatabaseEventAction.RESTORED: { + if (!isDefined(recordsAfter)) { + throw new Error( + `recordsAfter is required for ${action.toUpperCase()} action`, + ); + } + + if (!isDefined(recordsBefore)) { + throw new Error( + `recordsBefore is required for ${action.toUpperCase()} action`, + ); + } + + if (!isNonEmptyArray(recordsAfter)) { + break; + } + + events = recordsAfter + .map((recordAfter) => { + if (!isNonEmptyArray(recordsBefore)) { + throw new Error( + `recordsBefore is required for ${action.toUpperCase()} action`, + ); } - const before = Array.isArray(beforeEntities) - ? beforeEntities.find((before) => before.id === after.id) - : beforeEntities; + const correspondingRecordBefore = recordsBefore.find( + (recordBeforeToFind) => recordBeforeToFind.id === recordAfter.id, + ); - if (!isDefined(before)) { + if (!isDefined(correspondingRecordBefore)) { throw new TwentyORMException( - 'Record mismatch detected while computing event data for UPDATED action', + `Record mismatch detected while computing event data for ${action.toUpperCase()} action`, TwentyORMExceptionCode.ORM_EVENT_DATA_CORRUPTED, ); } const diff = objectRecordChangedValues( - before, - after, + correspondingRecordBefore, + recordAfter, objectMetadataItem, flatFieldMetadataMaps, ) as Partial>; @@ -103,74 +138,102 @@ export const formatTwentyOrmEventToDatabaseBatchEvent = < return; } - const event = new ObjectRecordUpdateEvent(); + const eventPayload = { + userId: authContext?.user?.id, + workspaceMemberId: authContext?.workspaceMemberId, + recordId: recordAfter.id, + properties: { + before: correspondingRecordBefore, + after: recordAfter, + updatedFields, + diff, + }, + } satisfies + | ObjectRecordUpdateEvent + | ObjectRecordDeleteEvent + | ObjectRecordRestoreEvent; - event.userId = authContext?.user?.id; - event.workspaceMemberId = authContext?.workspaceMemberId; - event.recordId = after.id; - event.properties = { - before, - after, - updatedFields, - diff, - }; - - return event; + switch (action) { + case DatabaseEventAction.DELETED: + return Object.assign( + new ObjectRecordDeleteEvent(), + eventPayload, + ); + case DatabaseEventAction.UPDATED: + return Object.assign( + new ObjectRecordUpdateEvent(), + eventPayload, + ); + case DatabaseEventAction.RESTORED: + return Object.assign( + new ObjectRecordRestoreEvent(), + eventPayload, + ); + default: + return assertUnreachable(action); + } }) .filter(isDefined); break; - case DatabaseEventAction.DELETED: - events = entityArray.map((before) => { - const event = new ObjectRecordDeleteEvent(); + } + case DatabaseEventAction.DESTROYED: { + if (!isDefined(recordsBefore)) { + throw new Error(`recordsBefore is required for "${action}" action`); + } - event.userId = authContext?.user?.id; - event.workspaceMemberId = authContext?.workspaceMemberId; - event.recordId = before.id; - event.properties = { before }; + if (!isNonEmptyArray(recordsBefore)) { + break; + } - return event; - }); - break; - case DatabaseEventAction.DESTROYED: - events = entityArray.map((before) => { + events = recordsBefore.map((recordBefore) => { const event = new ObjectRecordDestroyEvent(); event.userId = authContext?.user?.id; event.workspaceMemberId = authContext?.workspaceMemberId; - event.recordId = before.id; - event.properties = { before }; + event.recordId = recordBefore.id; + event.properties = { before: recordBefore }; return event; }); break; - case DatabaseEventAction.UPSERTED: - events = entityArray.map((after) => { + } + case DatabaseEventAction.UPSERTED: { + if (!isDefined(recordsAfter)) { + throw new Error(`recordsAfter is required for "${action}" action`); + } + + if (!isNonEmptyArray(recordsAfter)) { + break; + } + + events = recordsAfter.map((recordAfter) => { const event = new ObjectRecordUpsertEvent(); event.userId = authContext?.user?.id; event.workspaceMemberId = authContext?.workspaceMemberId; - event.recordId = after.id; + event.recordId = recordAfter.id; - const before = beforeEntities - ? Array.isArray(beforeEntities) - ? beforeEntities.find((before) => before.id === after.id) - : beforeEntities - : undefined; + const correspondingRecordBefore = recordsBefore?.find( + (recordBeforeToFind) => recordBeforeToFind.id === recordAfter.id, + ); let updatedFields; let diff; diff = objectRecordChangedValues( - before ?? {}, - after, + correspondingRecordBefore ?? {}, + recordAfter, objectMetadataItem, flatFieldMetadataMaps, ) as Partial>; + updatedFields = Object.keys(diff); event.properties = { - after, - ...(before && { before }), + after: recordAfter, + ...(correspondingRecordBefore && { + before: correspondingRecordBefore, + }), ...(diff && { diff }), ...(updatedFields && { updatedFields }), }; @@ -178,6 +241,8 @@ export const formatTwentyOrmEventToDatabaseBatchEvent = < return event; }); break; + } + default: return; } diff --git a/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/listeners/workflow-database-event-trigger.listener.ts b/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/listeners/workflow-database-event-trigger.listener.ts index 2560eae3ce..3a236d137a 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/listeners/workflow-database-event-trigger.listener.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/listeners/workflow-database-event-trigger.listener.ts @@ -1,10 +1,10 @@ import { Injectable, Logger } from '@nestjs/common'; import { + ObjectRecordEvent, type ObjectRecordCreateEvent, type ObjectRecordDeleteEvent, type ObjectRecordDestroyEvent, - type ObjectRecordNonDestructiveEvent, type ObjectRecordUpdateEvent, type ObjectRecordUpsertEvent, } from 'twenty-shared/database-events'; @@ -308,7 +308,7 @@ export class WorkflowDatabaseEventTriggerListener { } private async shouldIgnoreEvent( - payload: WorkspaceEventBatch, + payload: WorkspaceEventBatch, ) { const workspaceId = payload.workspaceId; const databaseEventName = payload.name; @@ -330,7 +330,7 @@ export class WorkflowDatabaseEventTriggerListener { payload, action, }: { - payload: WorkspaceEventBatch; + payload: WorkspaceEventBatch; action: DatabaseEventAction; }) { const workspaceId = payload.workspaceId; @@ -390,7 +390,7 @@ export class WorkflowDatabaseEventTriggerListener { eventListener, action, }: { - eventPayload: ObjectRecordNonDestructiveEvent; + eventPayload: ObjectRecordEvent; eventListener: WorkflowAutomatedTriggerWorkspaceEntity; action: DatabaseEventAction; }) { diff --git a/packages/twenty-shared/src/database-events/object-record-delete.event.ts b/packages/twenty-shared/src/database-events/object-record-delete.event.ts index 288e2e72a2..7c98923410 100644 --- a/packages/twenty-shared/src/database-events/object-record-delete.event.ts +++ b/packages/twenty-shared/src/database-events/object-record-delete.event.ts @@ -1,3 +1,4 @@ +import { type ObjectRecordDiff } from '@/database-events/object-record-diff'; import { ObjectRecordBaseEvent } from '@/database-events/object-record.base.event'; export class ObjectRecordDeleteEvent< @@ -5,5 +6,8 @@ export class ObjectRecordDeleteEvent< > extends ObjectRecordBaseEvent { declare properties: { before: T; + after: T; + updatedFields: string[]; + diff: Partial>; }; } diff --git a/packages/twenty-shared/src/database-events/object-record-restore.event.ts b/packages/twenty-shared/src/database-events/object-record-restore.event.ts index fb38b3c107..5dffc3ef3a 100644 --- a/packages/twenty-shared/src/database-events/object-record-restore.event.ts +++ b/packages/twenty-shared/src/database-events/object-record-restore.event.ts @@ -1,9 +1,13 @@ import { ObjectRecordCreateEvent } from '@/database-events/object-record-create.event'; +import { type ObjectRecordDiff } from '@/database-events/object-record-diff'; export class ObjectRecordRestoreEvent< T = object, > extends ObjectRecordCreateEvent { declare properties: { + before: T; after: T; + updatedFields: string[]; + diff: Partial>; }; } diff --git a/packages/twenty-shared/src/database-events/object-record-update.event.ts b/packages/twenty-shared/src/database-events/object-record-update.event.ts index 49be950d99..158cbfb0d8 100644 --- a/packages/twenty-shared/src/database-events/object-record-update.event.ts +++ b/packages/twenty-shared/src/database-events/object-record-update.event.ts @@ -1,12 +1,12 @@ -import { ObjectRecordBaseEvent } from '@/database-events/object-record.base.event'; import { type ObjectRecordDiff } from '@/database-events/object-record-diff'; +import { ObjectRecordBaseEvent } from '@/database-events/object-record.base.event'; export class ObjectRecordUpdateEvent< T = object, > extends ObjectRecordBaseEvent { declare properties: { - updatedFields?: string[]; - diff?: Partial>; + updatedFields: string[]; + diff: Partial>; before: T; after: T; };