From abe4d7491c6c53858320b25e9ac14284199d1dfd Mon Sep 17 00:00:00 2001 From: Weiko Date: Fri, 24 Jul 2026 15:27:26 +0200 Subject: [PATCH] Cap nested relation query concurrency (#23252) ## Context Common API queries load selected relations after fetching the root records. Relation loading is batched: one query pipeline loads a relation for all parent records, so this is not an N+1 problem. However, every sibling relation currently starts concurrently through `Promise.all`. Nested relations repeat the same behavior recursively. A wide selection can therefore submit many independent relation query pipelines at once. Existing query complexity and record limits restrict what can be requested, but they do not limit how much database work starts concurrently. ## What this changes This PR adds a request-local FIFO concurrency limiter for nested relation loading. - At most four `findRelations` pipelines execute concurrently. - One limiter is created for the outer relation-loading call. - The same limiter is shared by every recursive level. - Queued work starts as permits become available. - Permits are released in `finally`, including when a query fails. Conceptually: ```text Before: all sibling relations -> database concurrently nested siblings -> more database work concurrently After: all sibling relations -> FIFO queue -> at most 4 database pipelines nested siblings -> same FIFO queue and same limit ``` Note: Also addressing https://github.com/twentyhq/twenty/pull/23251#discussion_r3644510597 --- .../database-pool-metrics.service.spec.ts | 19 ++- .../typeorm/database-pool-metrics.service.ts | 26 ++- .../process-nested-relations-v2.helper.ts | 151 +++++++++++------- .../create-concurrency-limiter.util.spec.ts | 97 +++++++++++ .../utils/create-concurrency-limiter.util.ts | 46 ++++++ 5 files changed, 272 insertions(+), 67 deletions(-) create mode 100644 packages/twenty-server/src/engine/api/common/common-nested-relations-processor/utils/__tests__/create-concurrency-limiter.util.spec.ts create mode 100644 packages/twenty-server/src/engine/api/common/common-nested-relations-processor/utils/create-concurrency-limiter.util.ts diff --git a/packages/twenty-server/src/database/typeorm/database-pool-metrics.service.spec.ts b/packages/twenty-server/src/database/typeorm/database-pool-metrics.service.spec.ts index fa3e09d436..414979fc60 100644 --- a/packages/twenty-server/src/database/typeorm/database-pool-metrics.service.spec.ts +++ b/packages/twenty-server/src/database/typeorm/database-pool-metrics.service.spec.ts @@ -59,16 +59,23 @@ describe('DatabasePoolMetricsService', () => { let service: DatabasePoolMetricsService; let gaugeCallbacks: Map; let histogramRecord: jest.Mock; + let counterAdd: jest.Mock; + let createCounter: jest.Mock; beforeEach(() => { gaugeCallbacks = new Map(); histogramRecord = jest.fn(); + counterAdd = jest.fn(); + createCounter = jest.fn().mockReturnValue({ + add: counterAdd, + }); const metricsService = { getMeter: jest.fn().mockReturnValue({ createHistogram: jest.fn().mockReturnValue({ record: histogramRecord, }), + createCounter, }), createMultiObservableGauge: jest .fn() @@ -191,6 +198,7 @@ describe('DatabasePoolMetricsService', () => { expect(histogramRecord).toHaveBeenCalledWith(0.25, { pool: DatabasePoolName.WorkspacePrimary, }); + expect(counterAdd).not.toHaveBeenCalled(); }); it('records failed connection acquisitions', async () => { @@ -204,12 +212,21 @@ describe('DatabasePoolMetricsService', () => { dataSource: dataSource.dataSource, }); - await expect(dataSource.driver.obtainMasterConnection()).rejects.toThrow( + await expect(dataSource.driver.obtainMasterConnection()).rejects.toBe( error, ); expect(histogramRecord).toHaveBeenCalledWith(0, { pool: DatabasePoolName.Core, }); + expect(counterAdd).toHaveBeenCalledWith(1, { + pool: DatabasePoolName.Core, + }); + expect(createCounter).toHaveBeenCalledWith( + 'twenty_database_pool_acquisition_failures', + { + description: 'Number of failed PostgreSQL pool connection acquisitions', + }, + ); }); it('does not instrument a data source more than once', async () => { diff --git a/packages/twenty-server/src/database/typeorm/database-pool-metrics.service.ts b/packages/twenty-server/src/database/typeorm/database-pool-metrics.service.ts index 09614ca3a6..2a24b96286 100644 --- a/packages/twenty-server/src/database/typeorm/database-pool-metrics.service.ts +++ b/packages/twenty-server/src/database/typeorm/database-pool-metrics.service.ts @@ -1,6 +1,6 @@ import { Injectable } from '@nestjs/common'; -import { type Histogram } from '@opentelemetry/api'; +import { type Counter, type Histogram } from '@opentelemetry/api'; import { type Pool } from 'pg'; import { type DataSource } from 'typeorm'; import { type PostgresDriver } from 'typeorm/driver/postgres/PostgresDriver'; @@ -46,18 +46,28 @@ export class DatabasePoolMetricsService { private readonly pools = new Map(); private readonly instrumentedDrivers = new WeakSet(); private readonly acquisitionDurationHistogram: Histogram; + private readonly acquisitionFailureCounter: Counter; constructor(private readonly metricsService: MetricsService) { - this.acquisitionDurationHistogram = this.metricsService - .getMeter() - .createHistogram('twenty_database_pool_acquisition_duration_seconds', { + const meter = this.metricsService.getMeter(); + + this.acquisitionDurationHistogram = meter.createHistogram( + 'twenty_database_pool_acquisition_duration_seconds', + { description: 'Time spent acquiring a connection from the PostgreSQL pool', unit: 's', advice: { explicitBucketBoundaries: ACQUISITION_DURATION_BUCKETS_SECONDS, }, - }); + }, + ); + this.acquisitionFailureCounter = meter.createCounter( + 'twenty_database_pool_acquisition_failures', + { + description: 'Number of failed PostgreSQL pool connection acquisitions', + }, + ); for (const gauge of POOL_GAUGES) { this.metricsService.createMultiObservableGauge({ @@ -99,6 +109,12 @@ export class DatabasePoolMetricsService { try { return await obtainMasterConnection(); + } catch (error) { + this.acquisitionFailureCounter.add(1, { + pool: poolName, + }); + + throw error; } finally { this.acquisitionDurationHistogram.record( (performance.now() - start) / 1000, diff --git a/packages/twenty-server/src/engine/api/common/common-nested-relations-processor/process-nested-relations-v2.helper.ts b/packages/twenty-server/src/engine/api/common/common-nested-relations-processor/process-nested-relations-v2.helper.ts index b0d939f59c..5e4f23e915 100644 --- a/packages/twenty-server/src/engine/api/common/common-nested-relations-processor/process-nested-relations-v2.helper.ts +++ b/packages/twenty-server/src/engine/api/common/common-nested-relations-processor/process-nested-relations-v2.helper.ts @@ -7,6 +7,10 @@ import { type FindOptionsRelations, type ObjectLiteral } from 'typeorm'; import { computeMorphOrRelationFieldJoinColumnName } from 'src/engine/metadata-modules/field-metadata/utils/compute-morph-or-relation-field-join-column-name.util'; import { RelationType } from 'src/engine/metadata-modules/field-metadata/interfaces/relation-type.interface'; +import { + type ConcurrencyLimiter, + createConcurrencyLimiter, +} from 'src/engine/api/common/common-nested-relations-processor/utils/create-concurrency-limiter.util'; import { STANDARD_ERROR_MESSAGE } from 'src/engine/api/common/common-query-runners/errors/standard-error-message.constant'; import { GraphqlQueryRunnerException, @@ -33,40 +37,57 @@ import { isFieldMetadataEntityOfType } from 'src/engine/utils/is-field-metadata- const EMPTY_RELATION_SENTINEL_RECORD_ID = '00000000-0000-0000-0000-000000000000'; +const NESTED_RELATION_QUERY_MAX_CONCURRENCY = 4; + +type ProcessNestedRelationsArgs = { + flatObjectMetadataMaps: FlatEntityMaps; + flatFieldMetadataMaps: FlatEntityMaps; + parentObjectMetadataItem: FlatObjectMetadata; + parentObjectRecords: T[]; + // oxlint-disable-next-line typescript/no-explicit-any + parentObjectRecordsAggregatedValues?: Record; + relations: Record>; + aggregate?: Record; + limit: number; + authContext: WorkspaceAuthContext; + workspaceDataSource: GlobalWorkspaceDataSource; + rolePermissionConfig?: RolePermissionConfig; + // oxlint-disable-next-line typescript/no-explicit-any + selectedFields: Record; +}; @Injectable() export class ProcessNestedRelationsV2Helper { constructor() {} - public async processNestedRelations({ - flatObjectMetadataMaps, - flatFieldMetadataMaps, - parentObjectMetadataItem, - parentObjectRecords, - parentObjectRecordsAggregatedValues = {}, - relations, - aggregate = {}, - limit, - authContext, - workspaceDataSource, - rolePermissionConfig, - selectedFields, - }: { - flatObjectMetadataMaps: FlatEntityMaps; - flatFieldMetadataMaps: FlatEntityMaps; - parentObjectMetadataItem: FlatObjectMetadata; - parentObjectRecords: T[]; - // oxlint-disable-next-line typescript/no-explicit-any - parentObjectRecordsAggregatedValues?: Record; - relations: Record>; - aggregate?: Record; - limit: number; - authContext: WorkspaceAuthContext; - workspaceDataSource: GlobalWorkspaceDataSource; - rolePermissionConfig?: RolePermissionConfig; - // oxlint-disable-next-line typescript/no-explicit-any - selectedFields: Record; - }): Promise { + public async processNestedRelations( + args: ProcessNestedRelationsArgs, + ): Promise { + await this.processNestedRelationsWithLimiter( + args, + createConcurrencyLimiter(NESTED_RELATION_QUERY_MAX_CONCURRENCY), + ); + } + + private async processNestedRelationsWithLimiter< + T extends ObjectRecord = ObjectRecord, + >( + { + flatObjectMetadataMaps, + flatFieldMetadataMaps, + parentObjectMetadataItem, + parentObjectRecords, + parentObjectRecordsAggregatedValues = {}, + relations, + aggregate = {}, + limit, + authContext, + workspaceDataSource, + rolePermissionConfig, + selectedFields, + }: ProcessNestedRelationsArgs, + relationQueryLimiter: ConcurrencyLimiter, + ): Promise { const processRelationTasks = Object.entries(relations).map( ([sourceFieldName, nestedRelations]) => this.processRelation({ @@ -82,6 +103,7 @@ export class ProcessNestedRelationsV2Helper { authContext, workspaceDataSource, rolePermissionConfig, + relationQueryLimiter, selectedFields: selectedFields[sourceFieldName] instanceof Object ? selectedFields[sourceFieldName] @@ -105,6 +127,7 @@ export class ProcessNestedRelationsV2Helper { authContext, workspaceDataSource, rolePermissionConfig, + relationQueryLimiter, selectedFields, }: { flatObjectMetadataMaps: FlatEntityMaps; @@ -120,6 +143,7 @@ export class ProcessNestedRelationsV2Helper { authContext: WorkspaceAuthContext; workspaceDataSource: GlobalWorkspaceDataSource; rolePermissionConfig?: RolePermissionConfig; + relationQueryLimiter: ConcurrencyLimiter; selectedFields: Record; }): Promise { const fieldMaps = buildFieldMapsFromFlatObjectMetadata( @@ -230,21 +254,23 @@ export class ProcessNestedRelationsV2Helper { }); const { relationResults, relationAggregatedFieldsResult } = - await this.findRelations({ - referenceQueryBuilder: targetObjectQueryBuilder, - targetObjectRepository, - column: - relationType === RelationType.ONE_TO_MANY - ? `"${fieldMetadataTargetRelationColumnName}"` - : 'id', - ids: relationIds, - relationType, - perParentLimit: limit, - parentRecordsCount: parentObjectRecords.length, - aggregate, - sourceFieldName, - targetObjectNameSingular, - }); + await relationQueryLimiter(() => + this.findRelations({ + referenceQueryBuilder: targetObjectQueryBuilder, + targetObjectRepository, + column: + relationType === RelationType.ONE_TO_MANY + ? `"${fieldMetadataTargetRelationColumnName}"` + : 'id', + ids: relationIds, + relationType, + perParentLimit: limit, + parentRecordsCount: parentObjectRecords.length, + aggregate, + sourceFieldName, + targetObjectNameSingular, + }), + ); this.assignRelationResults({ parentRecords: parentObjectRecords, @@ -262,23 +288,26 @@ export class ProcessNestedRelationsV2Helper { }); if (Object.keys(nestedRelations).length > 0) { - await this.processNestedRelations({ - flatObjectMetadataMaps, - flatFieldMetadataMaps, - parentObjectMetadataItem: targetObjectMetadata, - parentObjectRecords: relationResults as ObjectRecord[], - parentObjectRecordsAggregatedValues: relationAggregatedFieldsResult, - relations: nestedRelations as Record< - string, - FindOptionsRelations - >, - aggregate, - limit, - authContext, - workspaceDataSource, - rolePermissionConfig, - selectedFields, - }); + await this.processNestedRelationsWithLimiter( + { + flatObjectMetadataMaps, + flatFieldMetadataMaps, + parentObjectMetadataItem: targetObjectMetadata, + parentObjectRecords: relationResults as ObjectRecord[], + parentObjectRecordsAggregatedValues: relationAggregatedFieldsResult, + relations: nestedRelations as Record< + string, + FindOptionsRelations + >, + aggregate, + limit, + authContext, + workspaceDataSource, + rolePermissionConfig, + selectedFields, + }, + relationQueryLimiter, + ); } } diff --git a/packages/twenty-server/src/engine/api/common/common-nested-relations-processor/utils/__tests__/create-concurrency-limiter.util.spec.ts b/packages/twenty-server/src/engine/api/common/common-nested-relations-processor/utils/__tests__/create-concurrency-limiter.util.spec.ts new file mode 100644 index 0000000000..81034b5361 --- /dev/null +++ b/packages/twenty-server/src/engine/api/common/common-nested-relations-processor/utils/__tests__/create-concurrency-limiter.util.spec.ts @@ -0,0 +1,97 @@ +import { createConcurrencyLimiter } from 'src/engine/api/common/common-nested-relations-processor/utils/create-concurrency-limiter.util'; + +const createDeferred = () => { + let resolve: (value: T | PromiseLike) => void; + let reject: (reason?: unknown) => void; + const promise = new Promise((promiseResolve, promiseReject) => { + resolve = promiseResolve; + reject = promiseReject; + }); + + return { promise, resolve: resolve!, reject: reject! }; +}; + +describe('createConcurrencyLimiter', () => { + it('should limit concurrent tasks and start queued tasks in order', async () => { + const limitConcurrency = createConcurrencyLimiter(4); + const startedTaskIndexes: number[] = []; + const taskStartedDeferreds = Array.from({ length: 8 }, () => + createDeferred(), + ); + const taskFinishedDeferreds = Array.from({ length: 8 }, () => + createDeferred(), + ); + let activeTaskCount = 0; + let maximumActiveTaskCount = 0; + + const tasks = Array.from({ length: 8 }, (_, taskIndex) => + limitConcurrency(async () => { + startedTaskIndexes.push(taskIndex); + activeTaskCount++; + maximumActiveTaskCount = Math.max( + maximumActiveTaskCount, + activeTaskCount, + ); + taskStartedDeferreds[taskIndex].resolve(); + + await taskFinishedDeferreds[taskIndex].promise; + + activeTaskCount--; + + return taskIndex; + }), + ); + + await Promise.all( + taskStartedDeferreds.slice(0, 4).map(({ promise }) => promise), + ); + + expect(startedTaskIndexes).toEqual([0, 1, 2, 3]); + expect(activeTaskCount).toBe(4); + + for (let taskIndex = 0; taskIndex < 8; taskIndex++) { + await taskStartedDeferreds[taskIndex].promise; + taskFinishedDeferreds[taskIndex].resolve(); + } + + await expect(Promise.all(tasks)).resolves.toEqual([0, 1, 2, 3, 4, 5, 6, 7]); + expect(startedTaskIndexes).toEqual([0, 1, 2, 3, 4, 5, 6, 7]); + expect(maximumActiveTaskCount).toBe(4); + }); + + it('should release capacity after a task rejects', async () => { + const limitConcurrency = createConcurrencyLimiter(1); + const rejectedTask = limitConcurrency(async () => { + throw new Error('Task failed'); + }); + const nextTask = limitConcurrency(async () => 'completed'); + + await expect(rejectedTask).rejects.toThrow('Task failed'); + await expect(nextTask).resolves.toBe('completed'); + }); + + it('should not share capacity between limiter instances', async () => { + const firstLimiter = createConcurrencyLimiter(1); + const secondLimiter = createConcurrencyLimiter(1); + const firstTaskFinished = createDeferred(); + const firstTask = firstLimiter(async () => { + await firstTaskFinished.promise; + }); + + await expect(secondLimiter(async () => 'completed')).resolves.toBe( + 'completed', + ); + + firstTaskFinished.resolve(); + await firstTask; + }); + + it('should reject invalid concurrency values', () => { + expect(() => createConcurrencyLimiter(0)).toThrow( + 'Maximum concurrency must be a positive integer', + ); + expect(() => createConcurrencyLimiter(1.5)).toThrow( + 'Maximum concurrency must be a positive integer', + ); + }); +}); diff --git a/packages/twenty-server/src/engine/api/common/common-nested-relations-processor/utils/create-concurrency-limiter.util.ts b/packages/twenty-server/src/engine/api/common/common-nested-relations-processor/utils/create-concurrency-limiter.util.ts new file mode 100644 index 0000000000..8c9b893fa8 --- /dev/null +++ b/packages/twenty-server/src/engine/api/common/common-nested-relations-processor/utils/create-concurrency-limiter.util.ts @@ -0,0 +1,46 @@ +export type ConcurrencyLimiter = (task: () => Promise) => Promise; + +export const createConcurrencyLimiter = ( + maxConcurrency: number, +): ConcurrencyLimiter => { + if (!Number.isInteger(maxConcurrency) || maxConcurrency < 1) { + throw new Error('Maximum concurrency must be a positive integer'); + } + + let activeTaskCount = 0; + const waitingTaskResolvers: Array<() => void> = []; + + const acquire = (): Promise => { + if (activeTaskCount < maxConcurrency) { + activeTaskCount++; + + return Promise.resolve(); + } + + return new Promise((resolve) => { + waitingTaskResolvers.push(resolve); + }); + }; + + const release = () => { + const nextTaskResolver = waitingTaskResolvers.shift(); + + if (nextTaskResolver) { + nextTaskResolver(); + + return; + } + + activeTaskCount--; + }; + + return async (task: () => Promise): Promise => { + await acquire(); + + try { + return await task(); + } finally { + release(); + } + }; +};