From 02a966bb7f71c5318b5cc73ed7bf799c721d6d08 Mon Sep 17 00:00:00 2001 From: Abdul Rahman <81605929+abdulrahmancodes@users.noreply.github.com> Date: Mon, 22 Jun 2026 16:35:06 +0530 Subject: [PATCH] make mergeMany atomic and optimize relation/field-map handling (#21885) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Closes [core-team-issue#2333](https://github.com/twentyhq/core-team-issues/issues/2333) ## Summary Hardens and optimizes `CommonMergeManyQueryRunnerService`: - **Atomicity**: wrap relation migration + duplicate deletion + survivor update in a single transaction so a mid-merge failure rolls back fully (previously failures were swallowed and could leave orphaned/half-merged data). - **Perf**: drop the redundant `find`-before-`update` in relation migration (2N → N queries, no row hydration) and hoist `buildFieldMapsFromFlatObjectMetadata` out of the per-field loops. ### Why a transaction (not parallelization) The relation migrations could be parallelized with `Promise.all`, but merge is a destructive operation: a partial failure leaves orphaned or half-merged records. We prioritize correctness, so the steps run inside one transaction. Review in cubic --------- Co-authored-by: Félix Malfait Co-authored-by: Charles Bochet --- .../common-merge-many-query-runner.service.ts | 148 +++++++----- .../companies-merge-many.integration-spec.ts | 222 ++++++++++++++++++ 2 files changed, 307 insertions(+), 63 deletions(-) diff --git a/packages/twenty-server/src/engine/api/common/common-query-runners/common-merge-many-query-runner.service.ts b/packages/twenty-server/src/engine/api/common/common-query-runners/common-merge-many-query-runner.service.ts index a4dd260f45..0800b8123c 100644 --- a/packages/twenty-server/src/engine/api/common/common-query-runners/common-merge-many-query-runner.service.ts +++ b/packages/twenty-server/src/engine/api/common/common-query-runners/common-merge-many-query-runner.service.ts @@ -1,4 +1,4 @@ -import { Injectable, Logger } from '@nestjs/common'; +import { Injectable } from '@nestjs/common'; import { msg } from '@lingui/core/macro'; import { @@ -15,7 +15,6 @@ import { isDefined } from 'twenty-shared/utils'; import { FindOptionsRelations, In, ObjectLiteral } from 'typeorm'; import { v4 as uuidv4 } from 'uuid'; -import { computeMorphOrRelationFieldJoinColumnName } from 'src/engine/metadata-modules/field-metadata/utils/compute-morph-or-relation-field-join-column-name.util'; import { CommonBaseQueryRunnerService } from 'src/engine/api/common/common-query-runners/common-base-query-runner.service'; import { CommonQueryRunnerException, @@ -35,6 +34,7 @@ import { buildColumnsToSelect } from 'src/engine/api/graphql/graphql-query-runne import { hasRecordFieldValue } from 'src/engine/api/graphql/graphql-query-runner/utils/has-record-field-value.util'; import { mergeFieldValues } from 'src/engine/api/graphql/graphql-query-runner/utils/merge-field-values.util'; import { WorkspaceAuthContext } from 'src/engine/core-modules/auth/types/workspace-auth-context.type'; +import { computeMorphOrRelationFieldJoinColumnName } from 'src/engine/metadata-modules/field-metadata/utils/compute-morph-or-relation-field-join-column-name.util'; import { FlatEntityMaps } from 'src/engine/metadata-modules/flat-entity/types/flat-entity-maps.type'; import { findFlatEntityByIdInFlatEntityMaps } from 'src/engine/metadata-modules/flat-entity/utils/find-flat-entity-by-id-in-flat-entity-maps.util'; import { FlatFieldMetadata } from 'src/engine/metadata-modules/flat-field-metadata/types/flat-field-metadata.type'; @@ -42,6 +42,8 @@ import { buildFieldMapsFromFlatObjectMetadata } from 'src/engine/metadata-module import { isFlatFieldMetadataOfType } from 'src/engine/metadata-modules/flat-field-metadata/utils/is-flat-field-metadata-of-type.util'; import { FlatObjectMetadata } from 'src/engine/metadata-modules/flat-object-metadata/types/flat-object-metadata.type'; import { assertMutationNotOnRemoteObject } from 'src/engine/metadata-modules/object-metadata/utils/assert-mutation-not-on-remote-object.util'; +import { WorkspaceEntityManager } from 'src/engine/twenty-orm/entity-manager/workspace-entity-manager'; +import { WorkspaceRepository } from 'src/engine/twenty-orm/repository/workspace.repository'; @Injectable() export class CommonMergeManyQueryRunnerService extends CommonBaseQueryRunnerService< @@ -50,16 +52,11 @@ export class CommonMergeManyQueryRunnerService extends CommonBaseQueryRunnerServ > { protected readonly operationName = CommonQueryNames.MERGE_MANY; - private readonly logger = new Logger(CommonMergeManyQueryRunnerService.name); async run( args: CommonExtendedInput, queryRunnerContext: CommonExtendedQueryRunnerContext, ): Promise { - const { - flatObjectMetadataMaps, - flatFieldMetadataMaps, - flatObjectMetadata, - } = queryRunnerContext; + const { flatFieldMetadataMaps, flatObjectMetadata } = queryRunnerContext; const recordsToMerge = await this.fetchRecordsToMerge( queryRunnerContext, @@ -86,15 +83,48 @@ export class CommonMergeManyQueryRunnerService extends CommonBaseQueryRunnerServ const idsToDelete = args.ids.filter((id) => id !== priorityRecord.id); - await this.migrateRelatedRecords( + const updatedRecord = + await queryRunnerContext.workspaceDataSource.transaction( + (transactionManager: WorkspaceEntityManager) => + this.executeMergeWithinTransaction(transactionManager, { + args, + queryRunnerContext, + idsToDelete, + priorityRecordId: priorityRecord.id, + mergedData, + }), + ); + + await this.processNestedRelations({ + args, + queryRunnerContext, + updatedRecords: [updatedRecord], + }); + + return updatedRecord; + } + + private async executeMergeWithinTransaction( + transactionManager: WorkspaceEntityManager, + { + args, queryRunnerContext, idsToDelete, - priorityRecord.id, - ); - - const queryBuilder = queryRunnerContext.repository.createQueryBuilder( - flatObjectMetadata.nameSingular, - ); + priorityRecordId, + mergedData, + }: { + args: CommonExtendedInput; + queryRunnerContext: CommonExtendedQueryRunnerContext; + idsToDelete: string[]; + priorityRecordId: string; + mergedData: Partial; + }, + ): Promise { + const { + flatObjectMetadata, + flatObjectMetadataMaps, + flatFieldMetadataMaps, + } = queryRunnerContext; const columnsToReturn = buildColumnsToReturn({ select: args.selectedFieldsResult.select, @@ -104,26 +134,33 @@ export class CommonMergeManyQueryRunnerService extends CommonBaseQueryRunnerServ flatFieldMetadataMaps, }); - await queryBuilder + const transactionRepository = transactionManager.getRepository( + flatObjectMetadata.nameSingular, + queryRunnerContext.rolePermissionConfig, + queryRunnerContext.authContext, + ); + + await this.migrateRelatedRecords( + transactionManager, + queryRunnerContext, + idsToDelete, + priorityRecordId, + ); + + await transactionRepository + .createQueryBuilder(flatObjectMetadata.nameSingular) .delete() .whereInIds(idsToDelete) .returning(columnsToReturn) .execute(); - const updatedRecord = await this.updatePriorityRecord( + return this.updatePriorityRecord( args, queryRunnerContext, - priorityRecord.id, + transactionRepository, + priorityRecordId, mergedData, ); - - await this.processNestedRelations({ - args, - queryRunnerContext, - updatedRecords: [updatedRecord], - }); - - return updatedRecord; } private async fetchRecordsToMerge( @@ -212,6 +249,11 @@ export class CommonMergeManyQueryRunnerService extends CommonBaseQueryRunnerServ ): Partial { const mergedResult: Partial = {}; + const { fieldIdByName } = buildFieldMapsFromFlatObjectMetadata( + flatFieldMetadataMaps, + flatObjectMetadata, + ); + const allFieldNames = new Set(); recordsToMerge.forEach((record) => { @@ -219,7 +261,7 @@ export class CommonMergeManyQueryRunnerService extends CommonBaseQueryRunnerServ if ( !this.shouldExcludeFieldFromMerge( fieldName, - flatObjectMetadata, + fieldIdByName, flatFieldMetadataMaps, ) ) { @@ -244,10 +286,6 @@ export class CommonMergeManyQueryRunnerService extends CommonBaseQueryRunnerServ } else if (recordsWithValues.length === 1) { mergedResult[fieldName] = recordsWithValues[0].value; } else { - const { fieldIdByName } = buildFieldMapsFromFlatObjectMetadata( - flatFieldMetadataMaps, - flatObjectMetadata, - ); const fieldMetadata = findFlatEntityByIdInFlatEntityMaps({ flatEntityId: fieldIdByName[fieldName], flatEntityMaps: flatFieldMetadataMaps, @@ -279,13 +317,9 @@ export class CommonMergeManyQueryRunnerService extends CommonBaseQueryRunnerServ private shouldExcludeFieldFromMerge( fieldName: string, - flatObjectMetadata: FlatObjectMetadata, + fieldIdByName: Record, flatFieldMetadataMaps: FlatEntityMaps, ): boolean { - const { fieldIdByName } = buildFieldMapsFromFlatObjectMetadata( - flatFieldMetadataMaps, - flatObjectMetadata, - ); const fieldMetadata = findFlatEntityByIdInFlatEntityMaps({ flatEntityId: fieldIdByName[fieldName], flatEntityMaps: flatFieldMetadataMaps, @@ -311,6 +345,7 @@ export class CommonMergeManyQueryRunnerService extends CommonBaseQueryRunnerServ private async updatePriorityRecord( args: CommonExtendedInput, queryRunnerContext: CommonExtendedQueryRunnerContext, + repository: WorkspaceRepository, priorityRecordId: string, mergedData: Partial, ): Promise { @@ -318,7 +353,6 @@ export class CommonMergeManyQueryRunnerService extends CommonBaseQueryRunnerServ flatObjectMetadata, flatObjectMetadataMaps, flatFieldMetadataMaps, - repository, } = queryRunnerContext; const queryBuilder = repository.createQueryBuilder( @@ -354,6 +388,7 @@ export class CommonMergeManyQueryRunnerService extends CommonBaseQueryRunnerServ } private async migrateRelatedRecords( + transactionManager: WorkspaceEntityManager, context: CommonExtendedQueryRunnerContext, fromIds: string[], toId: string, @@ -366,8 +401,6 @@ export class CommonMergeManyQueryRunnerService extends CommonBaseQueryRunnerServ const relationFieldsPointingToCurrentObject: Array<{ objectMetadata: FlatObjectMetadata; - fieldName: string; - fieldId: string; joinColumnName: string | undefined; }> = []; @@ -401,8 +434,6 @@ export class CommonMergeManyQueryRunnerService extends CommonBaseQueryRunnerServ relationFieldsPointingToCurrentObject.push({ objectMetadata: objMetadata, - fieldName: field.name, - fieldId: field.id, joinColumnName: computeMorphOrRelationFieldJoinColumnName({ name: field.name, }), @@ -414,29 +445,20 @@ export class CommonMergeManyQueryRunnerService extends CommonBaseQueryRunnerServ continue; } - try { - const repository = context.workspaceDataSource.getRepository( - relationField.objectMetadata.nameSingular, - context.rolePermissionConfig, - ); + const repository = transactionManager.getRepository( + relationField.objectMetadata.nameSingular, + context.rolePermissionConfig, + context.authContext, + ); - const whereCondition = { [relationField.joinColumnName]: In(fromIds) }; - - const existingRecords = await repository.find({ - where: whereCondition, - }); - - if (existingRecords.length > 0) { - await repository.update(whereCondition, { - [relationField.joinColumnName]: toId, - }); - } - } catch (error) { - this.logger.warn( - `Failed to migrate relation field "${relationField.fieldName}" (${relationField.joinColumnName}) in object "${relationField.objectMetadata.nameSingular}":`, - error.message, - ); - } + // repository.update() runs outside the transaction; build from the transaction-scoped repository so the migration rolls back with the merge. + await repository + .createQueryBuilder(relationField.objectMetadata.nameSingular) + .update() + .set({ [relationField.joinColumnName]: toId }) + .where({ [relationField.joinColumnName]: In(fromIds) }) + .returning('*') + .execute(); } } diff --git a/packages/twenty-server/test/integration/graphql/suites/object-generated/companies-merge-many.integration-spec.ts b/packages/twenty-server/test/integration/graphql/suites/object-generated/companies-merge-many.integration-spec.ts index ec16912b07..89e00a1d5b 100644 --- a/packages/twenty-server/test/integration/graphql/suites/object-generated/companies-merge-many.integration-spec.ts +++ b/packages/twenty-server/test/integration/graphql/suites/object-generated/companies-merge-many.integration-spec.ts @@ -1,5 +1,7 @@ import { COMPANY_GQL_FIELDS } from 'test/integration/constants/company-gql-fields.constants'; import { createManyOperationFactory } from 'test/integration/graphql/utils/create-many-operation-factory.util'; +import { findManyOperationFactory } from 'test/integration/graphql/utils/find-many-operation-factory.util'; +import { findOneOperationFactory } from 'test/integration/graphql/utils/find-one-operation-factory.util'; import { makeGraphqlAPIRequest } from 'test/integration/graphql/utils/make-graphql-api-request.util'; import { mergeManyOperationFactory } from 'test/integration/graphql/utils/merge-many-operation-factory.util'; import { deleteRecordsByIds } from 'test/integration/utils/delete-records-by-ids'; @@ -256,4 +258,224 @@ describe('companies merge resolvers (integration)', () => { ); }); }); + + describe('migrating related records', () => { + let createdPersonIds: string[] = []; + + afterEach(async () => { + if (createdPersonIds.length > 0) { + await deleteRecordsByIds('person', createdPersonIds); + createdPersonIds = []; + } + }); + + it('should re-parent related records onto the survivor and delete the duplicate', async () => { + const createCompaniesOperation = createManyOperationFactory({ + objectMetadataSingularName: 'company', + objectMetadataPluralName: 'companies', + gqlFields: COMPANY_GQL_FIELDS, + data: [{ name: 'Survivor Inc' }, { name: 'Duplicate Inc' }], + }); + + const createCompaniesResponse = await makeGraphqlAPIRequest( + createCompaniesOperation, + ); + + const survivorCompanyId = + createCompaniesResponse.body.data.createCompanies[0].id; + const duplicateCompanyId = + createCompaniesResponse.body.data.createCompanies[1].id; + + createdCompanyIds.push(survivorCompanyId, duplicateCompanyId); + + const createPeopleOperation = createManyOperationFactory({ + objectMetadataSingularName: 'person', + objectMetadataPluralName: 'people', + gqlFields: ` + id + company { + id + } + `, + data: [ + { + name: { firstName: 'Related', lastName: 'Contact' }, + companyId: duplicateCompanyId, + }, + ], + }); + + const createPeopleResponse = await makeGraphqlAPIRequest( + createPeopleOperation, + ); + + const relatedPerson = createPeopleResponse.body.data.createPeople[0]; + + createdPersonIds.push(relatedPerson.id); + + expect(relatedPerson.company.id).toBe(duplicateCompanyId); + + const mergeOperation = mergeManyOperationFactory({ + objectMetadataPluralName: 'companies', + gqlFields: COMPANY_GQL_FIELDS, + ids: [survivorCompanyId, duplicateCompanyId], + conflictPriorityIndex: 0, + }); + + const mergeResponse = await makeGraphqlAPIRequest(mergeOperation); + + expect(mergeResponse.body.errors).toBeUndefined(); + expect(mergeResponse.body.data.mergeCompanies.id).toBe(survivorCompanyId); + + const findPersonOperation = findOneOperationFactory({ + objectMetadataSingularName: 'person', + gqlFields: ` + id + company { + id + } + `, + filter: { id: { eq: relatedPerson.id } }, + }); + + const findPersonResponse = + await makeGraphqlAPIRequest(findPersonOperation); + + expect(findPersonResponse.body.data.person).not.toBeNull(); + expect(findPersonResponse.body.data.person.company.id).toBe( + survivorCompanyId, + ); + + const findDuplicateOperation = findOneOperationFactory({ + objectMetadataSingularName: 'company', + gqlFields: ` + id + `, + filter: { id: { eq: duplicateCompanyId } }, + }); + + const findDuplicateResponse = await makeGraphqlAPIRequest( + findDuplicateOperation, + ); + + expect(findDuplicateResponse.body.data.company).toBeNull(); + }); + + it('should re-parent children from every duplicate onto the survivor', async () => { + const createCompaniesOperation = createManyOperationFactory({ + objectMetadataSingularName: 'company', + objectMetadataPluralName: 'companies', + gqlFields: COMPANY_GQL_FIELDS, + data: [ + { name: 'Survivor Multi' }, + { name: 'Duplicate One' }, + { name: 'Duplicate Two' }, + ], + }); + + const createCompaniesResponse = await makeGraphqlAPIRequest( + createCompaniesOperation, + ); + + const survivorCompanyId = + createCompaniesResponse.body.data.createCompanies[0].id; + const duplicateOneId = + createCompaniesResponse.body.data.createCompanies[1].id; + const duplicateTwoId = + createCompaniesResponse.body.data.createCompanies[2].id; + + createdCompanyIds.push(survivorCompanyId, duplicateOneId, duplicateTwoId); + + const createPeopleOperation = createManyOperationFactory({ + objectMetadataSingularName: 'person', + objectMetadataPluralName: 'people', + gqlFields: ` + id + company { + id + } + `, + data: [ + { + name: { firstName: 'Already', lastName: 'OnSurvivor' }, + companyId: survivorCompanyId, + }, + { + name: { firstName: 'Child', lastName: 'OfDuplicateOne' }, + companyId: duplicateOneId, + }, + { + name: { firstName: 'Child', lastName: 'OfDuplicateTwo' }, + companyId: duplicateTwoId, + }, + ], + }); + + const createPeopleResponse = await makeGraphqlAPIRequest( + createPeopleOperation, + ); + + const createdPeopleIds = createPeopleResponse.body.data.createPeople.map( + (person: { id: string }) => person.id, + ); + + createdPersonIds.push(...createdPeopleIds); + + const mergeOperation = mergeManyOperationFactory({ + objectMetadataPluralName: 'companies', + gqlFields: COMPANY_GQL_FIELDS, + ids: [survivorCompanyId, duplicateOneId, duplicateTwoId], + conflictPriorityIndex: 0, + }); + + const mergeResponse = await makeGraphqlAPIRequest(mergeOperation); + + expect(mergeResponse.body.errors).toBeUndefined(); + expect(mergeResponse.body.data.mergeCompanies.id).toBe(survivorCompanyId); + + const findPeopleOperation = findManyOperationFactory({ + objectMetadataSingularName: 'person', + objectMetadataPluralName: 'people', + gqlFields: ` + id + company { + id + } + `, + filter: { id: { in: createdPeopleIds } }, + }); + + const findPeopleResponse = + await makeGraphqlAPIRequest(findPeopleOperation); + + const peopleAfterMerge = findPeopleResponse.body.data.people.edges; + + expect(peopleAfterMerge).toHaveLength(3); + peopleAfterMerge.forEach( + (edge: { node: { company: { id: string } } }) => { + expect(edge.node.company.id).toBe(survivorCompanyId); + }, + ); + + const findCompaniesOperation = findManyOperationFactory({ + objectMetadataSingularName: 'company', + objectMetadataPluralName: 'companies', + gqlFields: `id`, + filter: { + id: { in: [survivorCompanyId, duplicateOneId, duplicateTwoId] }, + }, + }); + + const findCompaniesResponse = await makeGraphqlAPIRequest( + findCompaniesOperation, + ); + + const remainingCompanyIds = + findCompaniesResponse.body.data.companies.edges.map( + (edge: { node: { id: string } }) => edge.node.id, + ); + + expect(remainingCompanyIds).toEqual([survivorCompanyId]); + }); + }); });