From afaf47cdf514adc59bc721689ac9a1d0eaf90ddf Mon Sep 17 00:00:00 2001 From: Marie <51697796+ijreilly@users.noreply.github.com> Date: Fri, 19 Dec 2025 10:48:18 +0100 Subject: [PATCH] Read some endpoints on replica db (behind feature flag) (#16677) Next steps: - move some workers' activities to replica db --- .../src/generated-metadata/graphql.ts | 1 + .../twenty-front/src/generated/graphql.ts | 1 + .../common-base-query-runner.service.ts | 20 +++++-- .../common-find-many-query-runner.service.ts | 1 + .../common-group-by-query-runner.service.ts | 1 + .../enums/feature-flag-key.enum.ts | 1 + .../twenty-config/config-variables.ts | 36 ++++++++++++ .../workspace-entity-manager.spec.ts | 1 + .../global-workspace-datasource.module.ts | 2 - .../global-workspace-datasource.service.ts | 58 +++++++++++++++++-- .../global-workspace-orm.manager.ts | 14 +++++ 11 files changed, 125 insertions(+), 11 deletions(-) diff --git a/packages/twenty-front/src/generated-metadata/graphql.ts b/packages/twenty-front/src/generated-metadata/graphql.ts index 4080b14f47..575df342bf 100644 --- a/packages/twenty-front/src/generated-metadata/graphql.ts +++ b/packages/twenty-front/src/generated-metadata/graphql.ts @@ -1314,6 +1314,7 @@ export enum FeatureFlagKey { IS_PAGE_LAYOUT_ENABLED = 'IS_PAGE_LAYOUT_ENABLED', IS_POSTGRESQL_INTEGRATION_ENABLED = 'IS_POSTGRESQL_INTEGRATION_ENABLED', IS_PUBLIC_DOMAIN_ENABLED = 'IS_PUBLIC_DOMAIN_ENABLED', + IS_READ_ON_REPLICA_ENABLED = 'IS_READ_ON_REPLICA_ENABLED', IS_RECORD_PAGE_LAYOUT_ENABLED = 'IS_RECORD_PAGE_LAYOUT_ENABLED', IS_STRIPE_INTEGRATION_ENABLED = 'IS_STRIPE_INTEGRATION_ENABLED', IS_TIMELINE_ACTIVITY_MIGRATED = 'IS_TIMELINE_ACTIVITY_MIGRATED', diff --git a/packages/twenty-front/src/generated/graphql.ts b/packages/twenty-front/src/generated/graphql.ts index 78994f1ef5..73ad142850 100644 --- a/packages/twenty-front/src/generated/graphql.ts +++ b/packages/twenty-front/src/generated/graphql.ts @@ -1297,6 +1297,7 @@ export enum FeatureFlagKey { IS_PAGE_LAYOUT_ENABLED = 'IS_PAGE_LAYOUT_ENABLED', IS_POSTGRESQL_INTEGRATION_ENABLED = 'IS_POSTGRESQL_INTEGRATION_ENABLED', IS_PUBLIC_DOMAIN_ENABLED = 'IS_PUBLIC_DOMAIN_ENABLED', + IS_READ_ON_REPLICA_ENABLED = 'IS_READ_ON_REPLICA_ENABLED', IS_RECORD_PAGE_LAYOUT_ENABLED = 'IS_RECORD_PAGE_LAYOUT_ENABLED', IS_STRIPE_INTEGRATION_ENABLED = 'IS_STRIPE_INTEGRATION_ENABLED', IS_TIMELINE_ACTIVITY_MIGRATED = 'IS_TIMELINE_ACTIVITY_MIGRATED', diff --git a/packages/twenty-server/src/engine/api/common/common-query-runners/common-base-query-runner.service.ts b/packages/twenty-server/src/engine/api/common/common-query-runners/common-base-query-runner.service.ts index d3d18a0945..1112ed78a3 100644 --- a/packages/twenty-server/src/engine/api/common/common-query-runners/common-base-query-runner.service.ts +++ b/packages/twenty-server/src/engine/api/common/common-query-runners/common-base-query-runner.service.ts @@ -32,6 +32,7 @@ import { WorkspacePreQueryHookPayload } from 'src/engine/api/graphql/workspace-q import { WorkspaceQueryHookService } from 'src/engine/api/graphql/workspace-query-runner/workspace-query-hook/workspace-query-hook.service'; import { ApiKeyRoleService } from 'src/engine/core-modules/api-key/services/api-key-role.service'; import { AuthContext } from 'src/engine/core-modules/auth/types/auth-context.type'; +import { FeatureFlagKey } from 'src/engine/core-modules/feature-flag/enums/feature-flag-key.enum'; import { FeatureFlagService } from 'src/engine/core-modules/feature-flag/services/feature-flag.service'; import { MetricsService } from 'src/engine/core-modules/metrics/metrics.service'; import { MetricsKeys } from 'src/engine/core-modules/metrics/types/metrics-keys.type'; @@ -87,6 +88,8 @@ export abstract class CommonBaseQueryRunnerService< protected abstract readonly operationName: CommonQueryNames; + protected readonly isReadOnly: boolean = false; + public async execute( args: CommonInput, queryRunnerContext: CommonBaseQueryRunnerContext, @@ -334,15 +337,22 @@ export abstract class CommonBaseQueryRunnerService< intersectionOf: [roleId], }; - const repository = await this.globalWorkspaceOrmManager.getRepository( - workspaceId, + const isReadOnReplicaEnabled = + await this.featureFlagService.isFeatureEnabled( + FeatureFlagKey.IS_READ_ON_REPLICA_ENABLED, + workspaceId, + ); + + const globalWorkspaceDataSource = + this.isReadOnly && isReadOnReplicaEnabled + ? await this.globalWorkspaceOrmManager.getGlobalWorkspaceDataSourceReplica() + : await this.globalWorkspaceOrmManager.getGlobalWorkspaceDataSource(); + + const repository = globalWorkspaceDataSource.getRepository( queryRunnerContext.flatObjectMetadata.nameSingular, rolePermissionConfig, ); - const globalWorkspaceDataSource = - await this.globalWorkspaceOrmManager.getGlobalWorkspaceDataSource(); - return { ...queryRunnerContext, authContext, diff --git a/packages/twenty-server/src/engine/api/common/common-query-runners/common-find-many-query-runner.service.ts b/packages/twenty-server/src/engine/api/common/common-query-runners/common-find-many-query-runner.service.ts index cc1b152913..236ac58724 100644 --- a/packages/twenty-server/src/engine/api/common/common-query-runners/common-find-many-query-runner.service.ts +++ b/packages/twenty-server/src/engine/api/common/common-query-runners/common-find-many-query-runner.service.ts @@ -43,6 +43,7 @@ export class CommonFindManyQueryRunnerService extends CommonBaseQueryRunnerServi CommonFindManyOutput > { protected readonly operationName = CommonQueryNames.FIND_MANY; + protected readonly isReadOnly = true; async run( args: CommonExtendedInput, diff --git a/packages/twenty-server/src/engine/api/common/common-query-runners/common-group-by-query-runner.service.ts b/packages/twenty-server/src/engine/api/common/common-query-runners/common-group-by-query-runner.service.ts index 3b22c12198..011c2810db 100644 --- a/packages/twenty-server/src/engine/api/common/common-query-runners/common-group-by-query-runner.service.ts +++ b/packages/twenty-server/src/engine/api/common/common-query-runners/common-group-by-query-runner.service.ts @@ -71,6 +71,7 @@ export class CommonGroupByQueryRunnerService extends CommonBaseQueryRunnerServic } protected readonly operationName = CommonQueryNames.GROUP_BY; + protected readonly isReadOnly = true; async run( args: CommonExtendedInput, diff --git a/packages/twenty-server/src/engine/core-modules/feature-flag/enums/feature-flag-key.enum.ts b/packages/twenty-server/src/engine/core-modules/feature-flag/enums/feature-flag-key.enum.ts index 515de1309f..7f427d0c88 100644 --- a/packages/twenty-server/src/engine/core-modules/feature-flag/enums/feature-flag-key.enum.ts +++ b/packages/twenty-server/src/engine/core-modules/feature-flag/enums/feature-flag-key.enum.ts @@ -15,4 +15,5 @@ export enum FeatureFlagKey { IS_DASHBOARD_V2_ENABLED = 'IS_DASHBOARD_V2_ENABLED', IS_TIMELINE_ACTIVITY_MIGRATED = 'IS_TIMELINE_ACTIVITY_MIGRATED', IS_GLOBAL_WORKSPACE_DATASOURCE_ENABLED = 'IS_GLOBAL_WORKSPACE_DATASOURCE_ENABLED', + IS_READ_ON_REPLICA_ENABLED = 'IS_READ_ON_REPLICA_ENABLED', } diff --git a/packages/twenty-server/src/engine/core-modules/twenty-config/config-variables.ts b/packages/twenty-server/src/engine/core-modules/twenty-config/config-variables.ts index 75c4140a5c..0dfe87abca 100644 --- a/packages/twenty-server/src/engine/core-modules/twenty-config/config-variables.ts +++ b/packages/twenty-server/src/engine/core-modules/twenty-config/config-variables.ts @@ -834,6 +834,22 @@ export class ConfigVariables { }) PG_DATABASE_URL: string; + @ConfigVariablesMetadata({ + group: ConfigVariablesGroup.SERVER_CONFIG, + isSensitive: true, + description: 'Database connection URL', + type: ConfigVariableType.STRING, + isEnvOnly: true, + }) + @IsOptional() + @IsUrl({ + protocols: ['postgres', 'postgresql'], + require_tld: false, + allow_underscores: true, + require_host: false, + }) + PG_DATABASE_REPLICA_URL: string; + @ConfigVariablesMetadata({ group: ConfigVariablesGroup.SERVER_CONFIG, description: @@ -1449,6 +1465,26 @@ export class ConfigVariables { }) @IsOptional() AWS_SES_ACCOUNT_ID: string; + + @ConfigVariablesMetadata({ + group: ConfigVariablesGroup.SERVER_CONFIG, + description: 'Timeout in milliseconds for primary database queries', + type: ConfigVariableType.NUMBER, + isEnvOnly: true, + }) + @CastToPositiveNumber() + @IsOptional() + PG_DATABASE_PRIMARY_TIMEOUT_MS: number = 10000; + + @ConfigVariablesMetadata({ + group: ConfigVariablesGroup.SERVER_CONFIG, + description: 'Timeout in milliseconds for replica database queries', + type: ConfigVariableType.NUMBER, + isEnvOnly: true, + }) + @CastToPositiveNumber() + @IsOptional() + PG_DATABASE_REPLICA_TIMEOUT_MS: number = 10000; } export const validate = (config: Record): ConfigVariables => { diff --git a/packages/twenty-server/src/engine/twenty-orm/entity-manager/workspace-entity-manager.spec.ts b/packages/twenty-server/src/engine/twenty-orm/entity-manager/workspace-entity-manager.spec.ts index 5f2d534124..79be790e92 100644 --- a/packages/twenty-server/src/engine/twenty-orm/entity-manager/workspace-entity-manager.spec.ts +++ b/packages/twenty-server/src/engine/twenty-orm/entity-manager/workspace-entity-manager.spec.ts @@ -214,6 +214,7 @@ describe('WorkspaceEntityManager', () => { IS_DASHBOARD_V2_ENABLED: false, IS_TIMELINE_ACTIVITY_MIGRATED: false, IS_GLOBAL_WORKSPACE_DATASOURCE_ENABLED: false, + IS_READ_ON_REPLICA_ENABLED: false, }, eventEmitterService: { emitMutationEvent: jest.fn(), diff --git a/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource.module.ts b/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource.module.ts index b7caa35e0f..6b827a0683 100644 --- a/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource.module.ts +++ b/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource.module.ts @@ -1,7 +1,6 @@ import { Global, Module } from '@nestjs/common'; import { TypeOrmModule } from '@nestjs/typeorm'; -import { FeatureFlagModule } from 'src/engine/core-modules/feature-flag/feature-flag.module'; import { TwentyConfigModule } from 'src/engine/core-modules/twenty-config/twenty-config.module'; import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity'; import { DataSourceModule } from 'src/engine/metadata-modules/data-source/data-source.module'; @@ -32,7 +31,6 @@ import { WorkspaceEventEmitterModule } from 'src/engine/workspace-event-emitter/ WorkspaceCacheStorageModule, WorkspaceManyOrAllFlatEntityMapsCacheModule, WorkspaceFeatureFlagsMapCacheModule, - FeatureFlagModule, TwentyConfigModule, WorkspaceEventEmitterModule, WorkspaceCacheModule, diff --git a/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource.service.ts b/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource.service.ts index f011a1718a..fb9bfadc95 100644 --- a/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource.service.ts +++ b/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource.service.ts @@ -1,10 +1,11 @@ import { Injectable, - Logger, OnApplicationShutdown, OnModuleInit, } from '@nestjs/common'; +import { isDefined } from 'twenty-shared/utils'; + import { TwentyConfigService } from 'src/engine/core-modules/twenty-config/twenty-config.service'; import { GlobalWorkspaceDataSource } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource'; import { WorkspaceEventEmitter } from 'src/engine/workspace-event-emitter/workspace-event-emitter'; @@ -13,8 +14,9 @@ import { WorkspaceEventEmitter } from 'src/engine/workspace-event-emitter/worksp export class GlobalWorkspaceDataSourceService implements OnModuleInit, OnApplicationShutdown { - private readonly logger = new Logger(GlobalWorkspaceDataSourceService.name); private globalWorkspaceDataSource: GlobalWorkspaceDataSource | null = null; + private globalWorkspaceDataSourceReplica: GlobalWorkspaceDataSource | null = + null; constructor( private readonly twentyConfigService: TwentyConfigService, @@ -35,7 +37,9 @@ export class GlobalWorkspaceDataSourceService : undefined, poolSize: this.twentyConfigService.get('PG_POOL_MAX_CONNECTIONS'), extra: { - query_timeout: 10000, // 10 seconds, + query_timeout: this.twentyConfigService.get( + 'PG_DATABASE_PRIMARY_TIMEOUT_MS', + ), idleTimeoutMillis: this.twentyConfigService.get( 'PG_POOL_IDLE_TIMEOUT_MS', ), @@ -48,10 +52,44 @@ export class GlobalWorkspaceDataSourceService ); await this.globalWorkspaceDataSource.initialize(); + + const shouldInitializeReplicaDataSource = isDefined( + this.twentyConfigService.get('PG_DATABASE_REPLICA_URL'), + ); + + if (shouldInitializeReplicaDataSource) { + this.globalWorkspaceDataSourceReplica = new GlobalWorkspaceDataSource( + { + url: this.twentyConfigService.get('PG_DATABASE_REPLICA_URL'), + type: 'postgres', + logging: this.twentyConfigService.getLoggingConfig(), + entities: [], + ssl: this.twentyConfigService.get('PG_SSL_ALLOW_SELF_SIGNED') + ? { + rejectUnauthorized: false, + } + : undefined, + poolSize: this.twentyConfigService.get('PG_POOL_MAX_CONNECTIONS'), + extra: { + query_timeout: this.twentyConfigService.get( + 'PG_DATABASE_REPLICA_TIMEOUT_MS', + ), + idleTimeoutMillis: this.twentyConfigService.get( + 'PG_POOL_IDLE_TIMEOUT_MS', + ), + allowExitOnIdle: this.twentyConfigService.get( + 'PG_POOL_ALLOW_EXIT_ON_IDLE', + ), + }, + }, + this.workspaceEventEmitter, + ); + await this.globalWorkspaceDataSourceReplica.initialize(); + } } public getGlobalWorkspaceDataSource(): GlobalWorkspaceDataSource { - if (!this.globalWorkspaceDataSource) { + if (!isDefined(this.globalWorkspaceDataSource)) { throw new Error( 'GlobalWorkspaceDataSource has not been initialized. Make sure the module has been initialized.', ); @@ -60,10 +98,22 @@ export class GlobalWorkspaceDataSourceService return this.globalWorkspaceDataSource; } + public getGlobalWorkspaceDataSourceReplica(): GlobalWorkspaceDataSource { + if (!isDefined(this.globalWorkspaceDataSourceReplica)) { + return this.getGlobalWorkspaceDataSource(); + } + + return this.globalWorkspaceDataSourceReplica; + } + async onApplicationShutdown(): Promise { if (this.globalWorkspaceDataSource) { await this.globalWorkspaceDataSource.destroy(); this.globalWorkspaceDataSource = null; } + if (this.globalWorkspaceDataSourceReplica) { + await this.globalWorkspaceDataSourceReplica.destroy(); + this.globalWorkspaceDataSourceReplica = null; + } } } diff --git a/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager.ts b/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager.ts index 926116b342..7fff7d0f3a 100644 --- a/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager.ts +++ b/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager.ts @@ -4,6 +4,7 @@ import { type ObjectLiteral } from 'typeorm'; import { type WorkspaceAuthContext } from 'src/engine/api/common/interfaces/workspace-auth-context.interface'; +import { FeatureFlagKey } from 'src/engine/core-modules/feature-flag/enums/feature-flag-key.enum'; import { buildObjectIdByNameMaps } from 'src/engine/metadata-modules/flat-object-metadata/utils/build-object-id-by-name-maps.util'; import { GlobalWorkspaceDataSource } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource'; import { GlobalWorkspaceDataSourceService } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource.service'; @@ -62,6 +63,10 @@ export class GlobalWorkspaceOrmManager { return this.globalWorkspaceDataSourceService.getGlobalWorkspaceDataSource(); } + async getGlobalWorkspaceDataSourceReplica(): Promise { + return this.globalWorkspaceDataSourceService.getGlobalWorkspaceDataSourceReplica(); + } + async executeInWorkspaceContext( authContext: WorkspaceAuthContext, fn: () => T | Promise, @@ -109,4 +114,13 @@ export class GlobalWorkspaceOrmManager { userWorkspaceRoleMap, }; } + + private async isReadOnReplicaEnabled(workspaceId: string): Promise { + const { featureFlagsMap } = await this.workspaceCacheService.getOrRecompute( + workspaceId, + ['featureFlagsMap'], + ); + + return !!featureFlagsMap[FeatureFlagKey.IS_READ_ON_REPLICA_ENABLED]; + } }