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(); + } + }; +};