Read some endpoints on replica db (behind feature flag) (#16677)

Next steps: 
- move some workers' activities to replica db
This commit is contained in:
Marie
2025-12-19 10:48:18 +01:00
committed by GitHub
parent 1d2aba5b22
commit afaf47cdf5
11 changed files with 125 additions and 11 deletions
@@ -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',
@@ -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',
@@ -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<Args>,
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,
@@ -43,6 +43,7 @@ export class CommonFindManyQueryRunnerService extends CommonBaseQueryRunnerServi
CommonFindManyOutput
> {
protected readonly operationName = CommonQueryNames.FIND_MANY;
protected readonly isReadOnly = true;
async run(
args: CommonExtendedInput<FindManyQueryArgs>,
@@ -71,6 +71,7 @@ export class CommonGroupByQueryRunnerService extends CommonBaseQueryRunnerServic
}
protected readonly operationName = CommonQueryNames.GROUP_BY;
protected readonly isReadOnly = true;
async run(
args: CommonExtendedInput<GroupByQueryArgs>,
@@ -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',
}
@@ -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<string, unknown>): ConfigVariables => {
@@ -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(),
@@ -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,
@@ -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<void> {
if (this.globalWorkspaceDataSource) {
await this.globalWorkspaceDataSource.destroy();
this.globalWorkspaceDataSource = null;
}
if (this.globalWorkspaceDataSourceReplica) {
await this.globalWorkspaceDataSourceReplica.destroy();
this.globalWorkspaceDataSourceReplica = null;
}
}
}
@@ -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<GlobalWorkspaceDataSource> {
return this.globalWorkspaceDataSourceService.getGlobalWorkspaceDataSourceReplica();
}
async executeInWorkspaceContext<T>(
authContext: WorkspaceAuthContext,
fn: () => T | Promise<T>,
@@ -109,4 +114,13 @@ export class GlobalWorkspaceOrmManager {
userWorkspaceRoleMap,
};
}
private async isReadOnReplicaEnabled(workspaceId: string): Promise<boolean> {
const { featureFlagsMap } = await this.workspaceCacheService.getOrRecompute(
workspaceId,
['featureFlagsMap'],
);
return !!featureFlagsMap[FeatureFlagKey.IS_READ_ON_REPLICA_ENABLED];
}
}