Files
twenty/packages/twenty-server/src/engine/api/common/common-query-runners/common-create-many-query-runner/common-create-many-query-runner.service.ts
T
Etienne 54aa52d11c feat(index): support composite unique indexes in create-many upsert conflict resolution (#22604)
## Context

The `createMany` upsert path resolved conflicts by scanning individual
field
metadata and only treating a field as a conflict target when it was
flagged
`isUnique` (plus the primary `id`). This ignored **composite unique
indexes**
(multi-column unique constraints), so upserting against a multi-field
unique
key never matched an existing row and could either insert a duplicate or
fail.

## What changed

- **Conflict groups are now derived from unique indexes**, not from
per-field
`isUnique` flags. `getConflictingFields` reads the object's index
metadata
(`flatIndexMaps`) and builds one `ConflictingFieldGroup` per unique
index
  (the primary `id` remains its own group).
- Each index group correctly expands its fields into DB columns,
handling:
- **Composite field types** — expands to the sub-columns included in the
unique constraint (or a specific sub-field when the index targets one).
- **`MANY_TO_ONE` relation fields** — resolves to the join column name.
  - **Scalar fields** — used directly.
- `ConflictingFieldGroup.baseField: string` → **`baseFields:
string[]`**, since
  a composite index spans multiple fields.
- **Clearer multi-match error message**: conflicting values are now
grouped per
index (`baseFields (fullPath: value, ...)`, groups joined by `;`) so
it's
  obvious which unique key caused the ambiguity when a payload matches
  different rows across different indexes.
- `CommonCreateManyQueryRunnerService` now fetches `flatIndexMaps` via
  `WorkspaceManyOrAllFlatEntityMapsCacheService` and passes them into
`getConflictingFields`; the cache module is wired into
`CoreCommonApiModule`.

## Tests

- New integration suite
`composite-unique-index-upsert.integration-spec.ts`:
- single composite unique index — insert, update-on-match, and
insert-when-key-differs
- **two independent composite unique indexes** — happy path (single row
matches
both) and failure path (payload matches different rows across the two
indexes → `Multiple records found with the same unique field values` /
    `BAD_USER_INPUT`).
- Updated unit specs for `get-conflicting-fields`,
`get-matching-record-id`,
  `build-where-conditions`, and `categorize-records` to reflect the
  index-driven grouping and the `baseFields[]` shape.

## Test plan

- [ ] `npx nx run twenty-server:test:integration:with-db-reset --
composite-unique-index-upsert`
- [ ] `npx nx test twenty-server -- get-conflicting-fields
get-matching-record-id build-where-conditions categorize-records`
- [ ] Manual: upsert against a composite unique index updates the
matching row instead of inserting a duplicate.

fixes
https://github.com/twentyhq/twenty/issues/22580#issuecomment-4894266699

<!-- This is an auto-generated description by cubic. -->
<a
href="https://cubic.dev/pr/twentyhq/twenty/pull/22604?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. -->
2026-07-07 13:54:33 +02:00

571 lines
18 KiB
TypeScript

import { Injectable } from '@nestjs/common';
import { msg } from '@lingui/core/macro';
import { QUERY_MAX_RECORDS } from 'twenty-shared/constants';
import { ObjectRecord } from 'twenty-shared/types';
import { isDefined } from 'twenty-shared/utils';
import { FindOptionsRelations, In, InsertResult, ObjectLiteral } from 'typeorm';
import { CommonBaseQueryRunnerService } from 'src/engine/api/common/common-query-runners/common-base-query-runner.service';
import { type ConflictingFieldGroup } from 'src/engine/api/common/common-query-runners/common-create-many-query-runner/types/conflicting-field-group.type';
import { PartialObjectRecordWithId } from 'src/engine/api/common/common-query-runners/common-create-many-query-runner/types/partial-object-record-with-id.type';
import { buildWhereConditions } from 'src/engine/api/common/common-query-runners/common-create-many-query-runner/utils/build-where-conditions.util';
import { categorizeRecords } from 'src/engine/api/common/common-query-runners/common-create-many-query-runner/utils/categorize-records.util';
import { getConflictingFields } from 'src/engine/api/common/common-query-runners/common-create-many-query-runner/utils/get-conflicting-fields.util';
import {
CommonQueryRunnerException,
CommonQueryRunnerExceptionCode,
} from 'src/engine/api/common/common-query-runners/errors/common-query-runner.exception';
import { STANDARD_ERROR_MESSAGE } from 'src/engine/api/common/common-query-runners/errors/standard-error-message.constant';
import { CommonBaseQueryRunnerContext } from 'src/engine/api/common/types/common-base-query-runner-context.type';
import { CommonExtendedQueryRunnerContext } from 'src/engine/api/common/types/common-extended-query-runner-context.type';
import {
CommonExtendedInput,
CommonInput,
CommonQueryNames,
CreateManyQueryArgs,
} from 'src/engine/api/common/types/common-query-args.type';
import { CommonSelectedFieldsResult } from 'src/engine/api/common/types/common-selected-fields-result.type';
import { buildColumnsToReturn } from 'src/engine/api/graphql/graphql-query-runner/utils/build-columns-to-return';
import { buildColumnsToSelect } from 'src/engine/api/graphql/graphql-query-runner/utils/build-columns-to-select';
import { assertIsValidUuid } from 'src/engine/api/graphql/workspace-query-runner/utils/assert-is-valid-uuid.util';
import { getAllSelectableColumnNames } from 'src/engine/api/utils/get-all-selectable-column-names.utils';
import { WorkspaceAuthContext } from 'src/engine/core-modules/auth/types/workspace-auth-context.type';
import { RecordPositionService } from 'src/engine/core-modules/record-position/services/record-position.service';
import { type 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 { type FlatFieldMetadata } from 'src/engine/metadata-modules/flat-field-metadata/types/flat-field-metadata.type';
import { buildFieldMapsFromFlatObjectMetadata } from 'src/engine/metadata-modules/flat-field-metadata/utils/build-field-maps-from-flat-object-metadata.util';
import { type FlatIndexMetadata } from 'src/engine/metadata-modules/flat-index-metadata/types/flat-index-metadata.type';
import { type 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 { GlobalWorkspaceDataSource } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource';
import { WorkspaceRepository } from 'src/engine/twenty-orm/repository/workspace.repository';
import { RolePermissionConfig } from 'src/engine/twenty-orm/types/role-permission-config';
@Injectable()
export class CommonCreateManyQueryRunnerService extends CommonBaseQueryRunnerService<
CreateManyQueryArgs,
ObjectRecord[]
> {
protected readonly operationName = CommonQueryNames.CREATE_MANY;
constructor(private readonly recordPositionService: RecordPositionService) {
super();
}
async run(
args: CommonExtendedInput<CreateManyQueryArgs>,
queryRunnerContext: CommonExtendedQueryRunnerContext,
): Promise<ObjectRecord[]> {
if (args.data.length > QUERY_MAX_RECORDS) {
throw new CommonQueryRunnerException(
`Maximum number of records to upsert is ${QUERY_MAX_RECORDS}.`,
CommonQueryRunnerExceptionCode.TOO_MANY_RECORDS_TO_UPDATE,
{
userFriendlyMessage: msg`Maximum number of records to upsert is ${QUERY_MAX_RECORDS}.`,
},
);
}
const {
repository,
authContext,
rolePermissionConfig,
flatObjectMetadata,
flatObjectMetadataMaps,
flatFieldMetadataMaps,
flatIndexMaps,
workspaceDataSource,
} = queryRunnerContext;
if (!isDefined(flatIndexMaps)) {
throw new CommonQueryRunnerException(
`Missing flatIndexMaps in queryRunnerContext`,
CommonQueryRunnerExceptionCode.MISSING_FLAT_INDEX_MAPS,
{ userFriendlyMessage: STANDARD_ERROR_MESSAGE },
);
}
const objectRecords = await this.insertOrUpsertRecords({
repository,
flatObjectMetadata,
flatObjectMetadataMaps,
flatFieldMetadataMaps,
flatIndexMaps,
args,
workspaceId: authContext.workspace.id,
});
const upsertedRecords = await this.fetchUpsertedRecords({
objectRecords,
flatObjectMetadata,
flatObjectMetadataMaps,
flatFieldMetadataMaps,
repository,
selectedFieldsResult: args.selectedFieldsResult,
});
await this.processNestedRelationsIfNeeded({
args,
records: upsertedRecords,
flatObjectMetadata,
flatObjectMetadataMaps,
flatFieldMetadataMaps,
authContext,
workspaceDataSource,
rolePermissionConfig,
});
return upsertedRecords;
}
private async processNestedRelationsIfNeeded({
args,
records,
flatObjectMetadata,
flatObjectMetadataMaps,
flatFieldMetadataMaps,
authContext,
workspaceDataSource,
rolePermissionConfig,
}: {
args: CommonExtendedInput<CreateManyQueryArgs>;
records: ObjectRecord[];
flatObjectMetadata: FlatObjectMetadata;
flatObjectMetadataMaps: FlatEntityMaps<FlatObjectMetadata>;
flatFieldMetadataMaps: FlatEntityMaps<FlatFieldMetadata>;
authContext: WorkspaceAuthContext;
workspaceDataSource: GlobalWorkspaceDataSource;
rolePermissionConfig?: RolePermissionConfig;
}): Promise<void> {
if (!args.selectedFieldsResult.relations) {
return;
}
await this.processNestedRelationsHelper.processNestedRelations({
flatObjectMetadataMaps,
flatFieldMetadataMaps,
parentObjectMetadataItem: flatObjectMetadata,
parentObjectRecords: records,
relations: args.selectedFieldsResult.relations as Record<
string,
FindOptionsRelations<ObjectLiteral>
>,
limit: QUERY_MAX_RECORDS,
authContext,
workspaceDataSource,
rolePermissionConfig,
selectedFields: args.selectedFieldsResult.select,
});
}
async computeArgs(
args: CommonInput<CreateManyQueryArgs>,
queryRunnerContext: CommonBaseQueryRunnerContext,
): Promise<CommonInput<CreateManyQueryArgs>> {
const {
authContext,
flatObjectMetadata,
flatFieldMetadataMaps,
flatObjectMetadataMaps,
} = queryRunnerContext;
return {
...args,
data: await this.dataArgProcessor.process({
partialRecordInputs: args.data,
authContext,
flatObjectMetadata,
flatFieldMetadataMaps,
flatObjectMetadataMaps,
shouldBackfillPositionIfUndefined: !args.upsert,
}),
};
}
async validate(
args: CommonInput<CreateManyQueryArgs>,
queryRunnerContext: CommonBaseQueryRunnerContext,
): Promise<void> {
const { flatObjectMetadata } = queryRunnerContext;
assertMutationNotOnRemoteObject(flatObjectMetadata);
args.data.forEach((record) => {
if (record?.id) {
assertIsValidUuid(record.id);
}
});
}
private async insertOrUpsertRecords({
repository,
flatObjectMetadata,
flatObjectMetadataMaps,
flatFieldMetadataMaps,
flatIndexMaps,
args,
workspaceId,
}: {
repository: WorkspaceRepository<ObjectLiteral>;
flatObjectMetadata: FlatObjectMetadata;
flatObjectMetadataMaps: FlatEntityMaps<FlatObjectMetadata>;
flatFieldMetadataMaps: FlatEntityMaps<FlatFieldMetadata>;
flatIndexMaps: FlatEntityMaps<FlatIndexMetadata>;
args: CommonExtendedInput<CreateManyQueryArgs>;
workspaceId: string;
}): Promise<InsertResult> {
const { selectedFieldsResult } = args;
if (!args.upsert) {
const selectedColumns = buildColumnsToReturn({
select: selectedFieldsResult.select,
relations: selectedFieldsResult.relations,
flatObjectMetadata,
flatObjectMetadataMaps,
flatFieldMetadataMaps,
});
return await repository.insert(args.data, undefined, selectedColumns);
}
return this.performUpsertOperation({
repository,
flatObjectMetadata,
flatObjectMetadataMaps,
flatFieldMetadataMaps,
flatIndexMaps,
args,
selectedFieldsResult,
workspaceId,
});
}
private async performUpsertOperation({
repository,
flatObjectMetadata,
flatObjectMetadataMaps,
flatFieldMetadataMaps,
flatIndexMaps,
args,
selectedFieldsResult,
workspaceId,
}: {
repository: WorkspaceRepository<ObjectLiteral>;
flatObjectMetadata: FlatObjectMetadata;
flatObjectMetadataMaps: FlatEntityMaps<FlatObjectMetadata>;
flatFieldMetadataMaps: FlatEntityMaps<FlatFieldMetadata>;
flatIndexMaps: FlatEntityMaps<FlatIndexMetadata>;
args: CreateManyQueryArgs;
selectedFieldsResult: CommonSelectedFieldsResult;
workspaceId: string;
}): Promise<InsertResult> {
const conflictingFieldGroups = getConflictingFields(
flatObjectMetadata,
flatFieldMetadataMaps,
flatIndexMaps,
);
const existingRecords = await this.findExistingRecords({
repository,
flatObjectMetadata,
flatFieldMetadataMaps,
args,
conflictingFieldGroups,
});
const { recordsToUpdate, recordsToInsert } = categorizeRecords(
args.data,
conflictingFieldGroups,
existingRecords,
);
const recordsToInsertWithPosition = await this.backfillPositionForInserts({
recordsToInsert,
flatObjectMetadata,
flatFieldMetadataMaps,
workspaceId,
});
const result: InsertResult = {
identifiers: [],
generatedMaps: [],
raw: [],
};
const columnsToReturn = buildColumnsToReturn({
select: selectedFieldsResult.select,
relations: selectedFieldsResult.relations,
flatObjectMetadata,
flatObjectMetadataMaps,
flatFieldMetadataMaps,
});
if (recordsToUpdate.length > 0) {
await this.processRecordsToUpdate({
partialRecordsToUpdate: recordsToUpdate,
repository,
flatObjectMetadata,
flatFieldMetadataMaps,
result,
columnsToReturn,
});
}
await this.processRecordsToInsert({
recordsToInsert: recordsToInsertWithPosition,
repository,
result,
columnsToReturn,
});
return result;
}
private async backfillPositionForInserts({
recordsToInsert,
flatObjectMetadata,
flatFieldMetadataMaps,
workspaceId,
}: {
recordsToInsert: Partial<ObjectRecord>[];
flatObjectMetadata: FlatObjectMetadata;
flatFieldMetadataMaps: FlatEntityMaps<FlatFieldMetadata>;
workspaceId: string;
}): Promise<Partial<ObjectRecord>[]> {
if (recordsToInsert.length === 0) {
return recordsToInsert;
}
const { fieldIdByName } = buildFieldMapsFromFlatObjectMetadata(
flatFieldMetadataMaps,
flatObjectMetadata,
);
return this.recordPositionService.overridePositionOnRecords({
partialRecordInputs: recordsToInsert,
workspaceId,
objectMetadata: {
isCustom: flatObjectMetadata.isCustom ?? false,
nameSingular: flatObjectMetadata.nameSingular,
fieldIdByName,
},
shouldBackfillPositionIfUndefined: true,
});
}
private async findExistingRecords({
repository,
flatObjectMetadata,
flatFieldMetadataMaps,
args,
conflictingFieldGroups,
}: {
repository: WorkspaceRepository<ObjectLiteral>;
flatObjectMetadata: FlatObjectMetadata;
flatFieldMetadataMaps: FlatEntityMaps<FlatFieldMetadata>;
args: CreateManyQueryArgs;
conflictingFieldGroups: ConflictingFieldGroup[];
}): Promise<PartialObjectRecordWithId[]> {
const queryBuilder = repository.createQueryBuilder(
flatObjectMetadata.nameSingular,
);
const whereConditions = buildWhereConditions(
args.data,
conflictingFieldGroups,
);
if (whereConditions.length === 0) {
return [];
}
whereConditions.forEach((condition) => {
queryBuilder.orWhere(condition);
});
const restrictedFields =
repository.objectRecordsPermissions?.[flatObjectMetadata.id]
?.restrictedFields;
const selectOptions = getAllSelectableColumnNames({
restrictedFields: restrictedFields ?? {},
objectMetadata: {
objectMetadataMapItem: flatObjectMetadata,
flatFieldMetadataMaps,
},
});
return (await queryBuilder
.withDeleted()
.setFindOptions({
select: selectOptions,
})
.getMany()) as PartialObjectRecordWithId[];
}
private async processRecordsToUpdate({
partialRecordsToUpdate,
repository,
flatObjectMetadata,
flatFieldMetadataMaps,
result,
columnsToReturn,
}: {
partialRecordsToUpdate: PartialObjectRecordWithId[];
repository: WorkspaceRepository<ObjectLiteral>;
flatObjectMetadata: FlatObjectMetadata;
flatFieldMetadataMaps: FlatEntityMaps<FlatFieldMetadata>;
result: InsertResult;
columnsToReturn: string[];
}): Promise<void> {
const partialRecordsToUpdateWithoutCreatedByUpdate =
partialRecordsToUpdate.map((record) =>
this.getRecordWithoutCreatedBy(
record,
flatObjectMetadata,
flatFieldMetadataMaps,
),
);
const savedRecords = await repository.updateMany(
partialRecordsToUpdateWithoutCreatedByUpdate.map((record) => ({
criteria: record.id,
partialEntity: { ...record, deletedAt: null },
})),
undefined,
columnsToReturn,
);
result.identifiers.push(
...savedRecords.generatedMaps.map((record) => ({ id: record.id })),
);
result.generatedMaps.push(
...savedRecords.generatedMaps.map((record) => ({ id: record.id })),
);
}
private async processRecordsToInsert({
recordsToInsert,
repository,
result,
columnsToReturn,
}: {
recordsToInsert: Partial<ObjectRecord>[];
repository: WorkspaceRepository<ObjectLiteral>;
result: InsertResult;
columnsToReturn: string[];
}): Promise<void> {
if (recordsToInsert.length > 0) {
const insertResult = await repository.insert(
recordsToInsert,
undefined,
columnsToReturn,
);
result.identifiers.push(...insertResult.identifiers);
result.generatedMaps.push(...insertResult.generatedMaps);
result.raw.push(...insertResult.raw);
}
}
private async fetchUpsertedRecords({
objectRecords,
flatObjectMetadata,
flatObjectMetadataMaps,
flatFieldMetadataMaps,
repository,
selectedFieldsResult,
}: {
objectRecords: InsertResult;
flatObjectMetadata: FlatObjectMetadata;
flatObjectMetadataMaps: FlatEntityMaps<FlatObjectMetadata>;
flatFieldMetadataMaps: FlatEntityMaps<FlatFieldMetadata>;
repository: WorkspaceRepository<ObjectLiteral>;
selectedFieldsResult: CommonSelectedFieldsResult;
}): Promise<ObjectRecord[]> {
const queryBuilder = repository.createQueryBuilder(
flatObjectMetadata.nameSingular,
);
const columnsToSelect = buildColumnsToSelect({
select: selectedFieldsResult.select,
relations: selectedFieldsResult.relations,
flatObjectMetadata,
flatObjectMetadataMaps,
flatFieldMetadataMaps,
});
const orderedIds = objectRecords.generatedMaps.map((record) => record.id);
const upsertedRecords = await queryBuilder
.setFindOptions({
select: columnsToSelect,
})
.where({
id: In(orderedIds),
})
.withDeleted()
.take(QUERY_MAX_RECORDS)
.getMany();
const orderIndex = new Map(orderedIds.map((id, index) => [id, index]));
upsertedRecords.sort(
(a, b) => (orderIndex.get(a.id) ?? 0) - (orderIndex.get(b.id) ?? 0),
);
return upsertedRecords as ObjectRecord[];
}
async processQueryResult(
queryResult: ObjectRecord[],
flatObjectMetadata: FlatObjectMetadata,
flatObjectMetadataMaps: FlatEntityMaps<FlatObjectMetadata>,
flatFieldMetadataMaps: FlatEntityMaps<FlatFieldMetadata>,
authContext: WorkspaceAuthContext,
): Promise<ObjectRecord[]> {
return await this.commonResultGettersService.processRecordArray(
queryResult,
flatObjectMetadata,
flatObjectMetadataMaps,
flatFieldMetadataMaps,
authContext.workspace.id,
);
}
private getRecordWithoutCreatedBy(
record: PartialObjectRecordWithId,
flatObjectMetadata: FlatObjectMetadata,
flatFieldMetadataMaps: FlatEntityMaps<FlatFieldMetadata>,
): Omit<PartialObjectRecordWithId, 'createdBy'> {
let recordWithoutCreatedByUpdate = record;
const { fieldIdByName } = buildFieldMapsFromFlatObjectMetadata(
flatFieldMetadataMaps,
flatObjectMetadata,
);
const createdByFieldMetadata = findFlatEntityByIdInFlatEntityMaps({
flatEntityId: fieldIdByName['createdBy'],
flatEntityMaps: flatFieldMetadataMaps,
});
if (!isDefined(createdByFieldMetadata)) {
throw new CommonQueryRunnerException(
`Missing createdBy field metadata for object ${flatObjectMetadata.nameSingular}`,
CommonQueryRunnerExceptionCode.MISSING_SYSTEM_FIELD,
{ userFriendlyMessage: STANDARD_ERROR_MESSAGE },
);
}
if ('createdBy' in record && createdByFieldMetadata.isSystem === true) {
const { createdBy: _createdBy, ...recordWithoutCreatedBy } = record;
recordWithoutCreatedByUpdate = recordWithoutCreatedBy;
}
return recordWithoutCreatedByUpdate;
}
}