make mergeMany atomic and optimize relation/field-map handling (#21885)

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.


<!-- This is an auto-generated description by cubic. -->
<a
href="https://cubic.dev/pr/twentyhq/twenty/pull/21885?utm_source=github"
target="_blank" rel="noopener noreferrer"
data-no-image-dialog="true"><picture><source
media="(prefers-color-scheme: dark)"
srcset="https://www.cubic.dev/buttons/review-in-cubic-dark.svg"><source
media="(prefers-color-scheme: light)"
srcset="https://www.cubic.dev/buttons/review-in-cubic-light.svg"><img
alt="Review in cubic"
src="https://www.cubic.dev/buttons/review-in-cubic-dark.svg"></picture></a>
<!-- End of auto-generated description by cubic. -->

---------

Co-authored-by: Félix Malfait <felix@twenty.com>
Co-authored-by: Charles Bochet <charles@twenty.com>
This commit is contained in:
Abdul Rahman
2026-06-22 16:35:06 +05:30
committed by GitHub
parent eeca9cd42e
commit 02a966bb7f
2 changed files with 307 additions and 63 deletions
@@ -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<MergeManyQueryArgs>,
queryRunnerContext: CommonExtendedQueryRunnerContext,
): Promise<ObjectRecord> {
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<MergeManyQueryArgs>;
queryRunnerContext: CommonExtendedQueryRunnerContext;
idsToDelete: string[];
priorityRecordId: string;
mergedData: Partial<ObjectRecord>;
},
): Promise<ObjectRecord> {
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<ObjectRecord> {
const mergedResult: Partial<ObjectRecord> = {};
const { fieldIdByName } = buildFieldMapsFromFlatObjectMetadata(
flatFieldMetadataMaps,
flatObjectMetadata,
);
const allFieldNames = new Set<string>();
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<string, string>,
flatFieldMetadataMaps: FlatEntityMaps<FlatFieldMetadata>,
): 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<MergeManyQueryArgs>,
queryRunnerContext: CommonExtendedQueryRunnerContext,
repository: WorkspaceRepository<ObjectLiteral>,
priorityRecordId: string,
mergedData: Partial<ObjectRecord>,
): Promise<ObjectRecord> {
@@ -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();
}
}
@@ -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]);
});
});
});