From 02fda92a93148dec403b743811e361f7f8aaee1c Mon Sep 17 00:00:00 2001 From: Weiko Date: Mon, 10 Nov 2025 17:32:49 +0100 Subject: [PATCH] Global workspace datasource poc (#15744) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## Context This PR introduces a Global Workspace DataSource that consolidates workspace-specific database access through a single TypeORM DataSource instance with AsyncLocalStorage-based context management instead of N workspace datasources. - Created GlobalWorkspaceDataSource extending TypeORM's DataSource to manage multiple workspaces with entity metadata caching (1-hour TTL) - Implemented AsyncLocalStorage for workspace context propagation (WorkspaceContextForStorage) containing workspace ID, metadata, permissions, and feature flags - Modified query execution flow to wrap operations in workspace context via GlobalWorkspaceOrmManager.executeInWorkspaceContext() - Added schema name to entity schemas for proper multi-tenant database separation Next: - use the new global workspace datasource everywhere and deprecate workspace datasource factory - improve metadata caching using a short TTL to avoid multiple calls to redis - Leverage the new WorkspaceContextALS and put it higher in the request hierarchy to have access to permission, metadata and featureflag everywhere --- build it manually for commands --- find a way to propagate it in jobs? - Remove PG_POOL patch once we have a unique datasource and increase global datasource pool size ## Implementation Why ALS: 1. Automatic Per-Request Isolation With schema-based multi-tenancy, each workspace has its own PostgreSQL schema (e.g. workspace_20202020-1c25-4d02-bf25-6aeccf7ea419). The critical challenge is ensuring that concurrent requests from different tenants don't interfere with each other. ```typescript // Request A (Workspace 1) and Request B (Workspace 2) executing concurrently // Without ALS: Race condition — they'd share the same global state! // With ALS: Each request has isolated context ✓ ``` ALS automatically isolates context per async execution chain, so: - Request from Tenant A → ALS stores workspaceId: "tenant-a" → Queries hit workspace_tenant_a schema - Request from Tenant B → ALS stores workspaceId: "tenant-b" → Queries hit workspace_tenant_b schema ✅ No interference, even when executing simultaneously on the same Node.js event loop. 2. No Manual Context Passing Before ALS, you'd need to pass workspaceId through every function call: ```typescript // ❌ Without ALS - Context threading nightmare getRepository(workspaceId, entity) → createEntityManager(workspaceId) → getMetadata(workspaceId, target) → findInCache(workspaceId, cacheKey) ``` With ALS: ```typescript // ✅ With ALS - Clean, implicit context getRepository(entity) // Reads workspaceId from ALS → createEntityManager() // Reads workspaceId from ALS → getMetadata(target) // Reads workspaceId from ALS → findInCache(cacheKey) // Reads workspaceId from ALS ``` example ```typescript override findMetadata(target: EntityTarget): EntityMetadata | undefined { const context = getWorkspaceContext(); // 👈 Automatically gets the right workspace! const { workspaceId, metadataVersion } = context; const cacheKey = `${workspaceId}-${metadataVersion}`; // ... returns metadata for THIS workspace's schema } ``` 3. Async Chain Propagation Node.js operations are heavily async. ALS automatically propagates context through: - async/await chains - Promise chains - Callbacks ```typescript executeInWorkspaceContext(workspaceId, async () => { await prepareContext(); // Has context ✓ const results = await run(); // Has context ✓ await enrichResults(); // Has context ✓ // Even nested async operations maintain context! await Promise.all([ saveToCache(), // Has context ✓ emitEvent(), // Has context ✓ logMetrics(), // Has context ✓ ]); }); ``` 5. Schema-Specific Metadata Caching The implementation caches entity metadata per workspace + version: ```typescript // Cache key format: "workspaceId-metadataVersion" const cacheKey = `${workspaceId}-${metadataVersion}`; ``` Why this matters with schemas: - Each workspace has different table structures (custom fields, objects) - EntitySchema includes schema: "workspace_xxx" property - Each cached metadata points to the correct schema ALS ensures getWorkspaceContext() returns the right workspaceId, so you always get the correct schema's metadata from cache. 6. Single DataSource for All Tenants The key change here: ```typescript // ❌ Old approach: One DataSource per tenant const dataSourceTenantA = new DataSource({ schema: 'workspace_a' }); const dataSourceTenantB = new DataSource({ schema: 'workspace_b' }); // Problem: Hundreds of DB connection pools! ``` ```typescript // ✅ New approach: One shared DataSource + ALS context const globalDataSource = new GlobalWorkspaceDataSource(); // ALS determines which schema to use at runtime ``` When you call: ```typescript globalDataSource.getRepository('person'); ``` It internally does: ```typescript const context = getWorkspaceContext(); // Gets current tenant from ALS const metadata = this.findMetadata('person'); // Finds metadata for THIS tenant's schema // EntityMetadata includes: schema: "workspace_20202020-1c25..." // TypeORM automatically queries: SELECT * FROM "workspace_20202020-1c25...".person ``` 7. Request Lifecycle Example ```typescript // 1. GraphQL request arrives: "query people { ... }" // 2. Middleware extracts authContext.workspace.id = "tenant-a" // 3. Query runner wraps execution in ALS: executeInWorkspaceContext("tenant-a", async () => { // 4. Everything inside has access to workspace context: const repo = getRepository('person'); // ALS → tenant-a const metadata = getMetadata('person'); // ALS → tenant-a → cache["tenant-a-v5"] // 5. TypeORM builds query with correct schema: // SELECT * FROM "workspace_tenant_a"."person" WHERE ... // 6. Even nested calls work: await saveAuditLog(); // ALS → tenant-a → correct audit schema await emitWebhook(); // ALS → tenant-a → correct tenant webhook }); // 7. Request completes, ALS context automatically cleaned up ``` 8. Safety & Error Prevention ```typescript // If you forget to set context: const context = getWorkspaceContext(); // ❌ Throws: "Workspace context not set..." // Fails fast rather than querying wrong schema! // Can't accidentally query wrong tenant: // Context is immutable within execution scope ``` --- .../src/generated-metadata/graphql.ts | 1 + .../twenty-front/src/generated/graphql.ts | 1 + .../common-base-query-runner.service.ts | 109 +++++- .../api/common/core-common-api.module.ts | 4 + .../enums/feature-flag-key.enum.ts | 1 + .../workspace-entity-manager.spec.ts | 2 + .../workspace-entity-manager.ts | 3 +- .../factories/entity-schema.factory.ts | 5 +- .../factories/workspace-datasource.factory.ts | 1 - .../global-workspace-datasource.module.ts | 41 ++ .../global-workspace-datasource.service.ts | 70 ++++ .../global-workspace-datasource.ts | 364 ++++++++++++++++++ .../global-workspace-orm.manager.ts | 118 ++++++ .../storage/workspace-context.storage.ts | 42 ++ 14 files changed, 745 insertions(+), 17 deletions(-) create mode 100644 packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource.module.ts create mode 100644 packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource.service.ts create mode 100644 packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource.ts create mode 100644 packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager.ts create mode 100644 packages/twenty-server/src/engine/twenty-orm/storage/workspace-context.storage.ts diff --git a/packages/twenty-front/src/generated-metadata/graphql.ts b/packages/twenty-front/src/generated-metadata/graphql.ts index 6582e47c30..6bdd3d6721 100644 --- a/packages/twenty-front/src/generated-metadata/graphql.ts +++ b/packages/twenty-front/src/generated-metadata/graphql.ts @@ -1263,6 +1263,7 @@ export enum FeatureFlagKey { IS_APPLICATION_ENABLED = 'IS_APPLICATION_ENABLED', IS_DASHBOARD_V2_ENABLED = 'IS_DASHBOARD_V2_ENABLED', IS_EMAILING_DOMAIN_ENABLED = 'IS_EMAILING_DOMAIN_ENABLED', + IS_GLOBAL_WORKSPACE_DATASOURCE_ENABLED = 'IS_GLOBAL_WORKSPACE_DATASOURCE_ENABLED', IS_IMAP_SMTP_CALDAV_ENABLED = 'IS_IMAP_SMTP_CALDAV_ENABLED', IS_JSON_FILTER_ENABLED = 'IS_JSON_FILTER_ENABLED', IS_MESSAGE_FOLDER_CONTROL_ENABLED = 'IS_MESSAGE_FOLDER_CONTROL_ENABLED', diff --git a/packages/twenty-front/src/generated/graphql.ts b/packages/twenty-front/src/generated/graphql.ts index b412884320..bf490f1848 100644 --- a/packages/twenty-front/src/generated/graphql.ts +++ b/packages/twenty-front/src/generated/graphql.ts @@ -1227,6 +1227,7 @@ export enum FeatureFlagKey { IS_APPLICATION_ENABLED = 'IS_APPLICATION_ENABLED', IS_DASHBOARD_V2_ENABLED = 'IS_DASHBOARD_V2_ENABLED', IS_EMAILING_DOMAIN_ENABLED = 'IS_EMAILING_DOMAIN_ENABLED', + IS_GLOBAL_WORKSPACE_DATASOURCE_ENABLED = 'IS_GLOBAL_WORKSPACE_DATASOURCE_ENABLED', IS_IMAP_SMTP_CALDAV_ENABLED = 'IS_IMAP_SMTP_CALDAV_ENABLED', IS_JSON_FILTER_ENABLED = 'IS_JSON_FILTER_ENABLED', IS_MESSAGE_FOLDER_CONTROL_ENABLED = 'IS_MESSAGE_FOLDER_CONTROL_ENABLED', 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 69b5036483..be327175aa 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 @@ -31,6 +31,8 @@ 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/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'; import { ThrottlerService } from 'src/engine/core-modules/throttler/throttler.service'; @@ -46,6 +48,8 @@ import { ObjectMetadataItemWithFieldMaps } from 'src/engine/metadata-modules/typ import { ObjectMetadataMaps } from 'src/engine/metadata-modules/types/object-metadata-maps'; import { UserRoleService } from 'src/engine/metadata-modules/user-role/user-role.service'; import { WorkspacePermissionsCacheService } from 'src/engine/metadata-modules/workspace-permissions-cache/workspace-permissions-cache.service'; +import { WorkspaceDataSource } from 'src/engine/twenty-orm/datasource/workspace.datasource'; +import { GlobalWorkspaceOrmManager } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager'; import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager'; @Injectable() @@ -62,6 +66,8 @@ export abstract class CommonBaseQueryRunnerService< @Inject() protected readonly twentyORMGlobalManager: TwentyORMGlobalManager; @Inject() + protected readonly globalWorkspaceOrmManager: GlobalWorkspaceOrmManager; + @Inject() protected readonly processNestedRelationsHelper: ProcessNestedRelationsHelper; @Inject() protected readonly permissionsService: PermissionsService; @@ -81,6 +87,8 @@ export abstract class CommonBaseQueryRunnerService< protected readonly twentyConfigService: TwentyConfigService; @Inject() protected readonly metricsService: MetricsService; + @Inject() + protected readonly featureFlagService: FeatureFlagService; protected abstract readonly operationName: CommonQueryNames; @@ -121,24 +129,33 @@ export abstract class CommonBaseQueryRunnerService< commonQueryParser, ); - const extendedQueryRunnerContext = - await this.prepareExtendedQueryRunnerContext( - authContext, - queryRunnerContext, + const isGlobalDatasourceEnabled = + await this.featureFlagService.isFeatureEnabled( + FeatureFlagKey.IS_GLOBAL_WORKSPACE_DATASOURCE_ENABLED, + authContext.workspace.id, ); - const results = await this.run(processedArgs, { - ...extendedQueryRunnerContext, - commonQueryParser, - }); + if (isGlobalDatasourceEnabled) { + return this.globalWorkspaceOrmManager.executeInWorkspaceContext( + authContext.workspace.id, + async () => + this.executeQueryAndEnrichResults( + processedArgs, + authContext, + queryRunnerContext, + commonQueryParser, + isGlobalDatasourceEnabled, + ), + ); + } - return this.enrichResultsWithGettersAndHooks({ - results, - operationName: this.operationName, + return this.executeQueryAndEnrichResults( + processedArgs, authContext, - objectMetadataItemWithFieldMaps, - objectMetadataMaps, - }); + queryRunnerContext, + commonQueryParser, + isGlobalDatasourceEnabled, + ); } protected abstract run( @@ -190,6 +207,38 @@ export abstract class CommonBaseQueryRunnerService< }; } + private async executeQueryAndEnrichResults( + processedArgs: CommonExtendedInput, + authContext: WorkspaceAuthContext, + queryRunnerContext: CommonBaseQueryRunnerContext, + commonQueryParser: GraphqlQueryParser, + isGlobalDatasourceEnabled: boolean, + ): Promise { + const extendedQueryRunnerContext = isGlobalDatasourceEnabled + ? await this.prepareExtendedQueryRunnerContextWithGlobalDatasource( + authContext, + queryRunnerContext, + ) + : await this.prepareExtendedQueryRunnerContext( + authContext, + queryRunnerContext, + ); + + const results = await this.run(processedArgs, { + ...extendedQueryRunnerContext, + commonQueryParser, + }); + + return this.enrichResultsWithGettersAndHooks({ + results, + operationName: this.operationName, + authContext, + objectMetadataItemWithFieldMaps: + queryRunnerContext.objectMetadataItemWithFieldMaps, + objectMetadataMaps: queryRunnerContext.objectMetadataMaps, + }); + } + private async enrichResultsWithGettersAndHooks({ results, operationName, @@ -335,6 +384,38 @@ export abstract class CommonBaseQueryRunnerService< }; } + private async prepareExtendedQueryRunnerContextWithGlobalDatasource( + authContext: WorkspaceAuthContext, + queryRunnerContext: CommonBaseQueryRunnerContext, + ): Promise> { + const workspaceId = authContext.workspace.id; + + const { roleId } = await this.getRoleIdAndObjectsPermissions( + authContext, + workspaceId, + ); + + const rolePermissionConfig = { unionOf: [roleId] }; + + const repository = await this.globalWorkspaceOrmManager.getRepository( + workspaceId, + queryRunnerContext.objectMetadataItemWithFieldMaps.nameSingular, + rolePermissionConfig, + ); + + const globalWorkspaceDataSource = + await this.globalWorkspaceOrmManager.getGlobalWorkspaceDataSource(); + + return { + ...queryRunnerContext, + authContext, + workspaceDataSource: + globalWorkspaceDataSource as unknown as WorkspaceDataSource, + rolePermissionConfig, + repository, + }; + } + private async throttleQueryExecution(authContext: WorkspaceAuthContext) { try { if (!isDefined(authContext.apiKey)) return; diff --git a/packages/twenty-server/src/engine/api/common/core-common-api.module.ts b/packages/twenty-server/src/engine/api/common/core-common-api.module.ts index 76a89f25c2..be29c40913 100644 --- a/packages/twenty-server/src/engine/api/common/core-common-api.module.ts +++ b/packages/twenty-server/src/engine/api/common/core-common-api.module.ts @@ -11,6 +11,7 @@ import { ProcessNestedRelationsHelper } from 'src/engine/api/graphql/graphql-que import { WorkspaceQueryHookModule } from 'src/engine/api/graphql/workspace-query-runner/workspace-query-hook/workspace-query-hook.module'; import { WorkspaceQueryRunnerModule } from 'src/engine/api/graphql/workspace-query-runner/workspace-query-runner.module'; import { ApiKeyModule } from 'src/engine/core-modules/api-key/api-key.module'; +import { FeatureFlagModule } from 'src/engine/core-modules/feature-flag/feature-flag.module'; import { FileModule } from 'src/engine/core-modules/file/file.module'; import { MetricsModule } from 'src/engine/core-modules/metrics/metrics.module'; import { ThrottlerModule } from 'src/engine/core-modules/throttler/throttler.module'; @@ -21,6 +22,7 @@ import { ViewFilterGroupModule } from 'src/engine/metadata-modules/view-filter-g import { ViewFilterModule } from 'src/engine/metadata-modules/view-filter/view-filter.module'; import { ViewModule } from 'src/engine/metadata-modules/view/view.module'; import { WorkspacePermissionsCacheModule } from 'src/engine/metadata-modules/workspace-permissions-cache/workspace-permissions-cache.module'; +import { GlobalWorkspaceDataSourceModule } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource.module'; @Module({ imports: [ @@ -37,6 +39,8 @@ import { WorkspacePermissionsCacheModule } from 'src/engine/metadata-modules/wor ViewFilterGroupModule, ThrottlerModule, MetricsModule, + GlobalWorkspaceDataSourceModule, + FeatureFlagModule, ], providers: [ ProcessNestedRelationsHelper, 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 60eb50ae95..c40cb97e00 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 @@ -16,4 +16,5 @@ export enum FeatureFlagKey { IS_EMAILING_DOMAIN_ENABLED = 'IS_EMAILING_DOMAIN_ENABLED', IS_WORKFLOW_RUN_STOPPAGE_ENABLED = 'IS_WORKFLOW_RUN_STOPPAGE_ENABLED', IS_DASHBOARD_V2_ENABLED = 'IS_DASHBOARD_V2_ENABLED', + IS_GLOBAL_WORKSPACE_DATASOURCE_ENABLED = 'IS_GLOBAL_WORKSPACE_DATASOURCE_ENABLED', } 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 c10c30b707..8856a127fc 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 @@ -139,6 +139,7 @@ describe('WorkspaceEntityManager', () => { IS_EMAILING_DOMAIN_ENABLED: false, IS_WORKFLOW_RUN_STOPPAGE_ENABLED: false, IS_DASHBOARD_V2_ENABLED: false, + IS_GLOBAL_WORKSPACE_DATASOURCE_ENABLED: false, }, eventEmitterService: { emitMutationEvent: jest.fn(), @@ -166,6 +167,7 @@ describe('WorkspaceEntityManager', () => { IS_EMAILING_DOMAIN_ENABLED: false, IS_WORKFLOW_RUN_STOPPAGE_ENABLED: false, IS_DASHBOARD_V2_ENABLED: false, + IS_GLOBAL_WORKSPACE_DATASOURCE_ENABLED: false, }, permissionsPerRoleId: {}, } as WorkspaceDataSource; diff --git a/packages/twenty-server/src/engine/twenty-orm/entity-manager/workspace-entity-manager.ts b/packages/twenty-server/src/engine/twenty-orm/entity-manager/workspace-entity-manager.ts index 7e0503a50f..ffe53ff2cf 100644 --- a/packages/twenty-server/src/engine/twenty-orm/entity-manager/workspace-entity-manager.ts +++ b/packages/twenty-server/src/engine/twenty-orm/entity-manager/workspace-entity-manager.ts @@ -49,6 +49,7 @@ import { type DeepPartialWithNestedRelationFields } from 'src/engine/twenty-orm/ import { type QueryDeepPartialEntityWithNestedRelationFields } from 'src/engine/twenty-orm/entity-manager/types/query-deep-partial-entity-with-nested-relation-fields.type'; import { getEntityTarget } from 'src/engine/twenty-orm/entity-manager/utils/get-entity-target'; import { computeTwentyORMException } from 'src/engine/twenty-orm/error-handling/compute-twenty-orm-exception'; +import { type GlobalWorkspaceDataSource } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource'; import { RelationNestedQueries } from 'src/engine/twenty-orm/relation-nested-queries/relation-nested-queries'; import { type OperationType, @@ -76,7 +77,7 @@ export class WorkspaceEntityManager extends EntityManager { constructor( internalContext: WorkspaceInternalContext, - connection: WorkspaceDataSource, + connection: WorkspaceDataSource | GlobalWorkspaceDataSource, queryRunner?: QueryRunner, ) { super(connection, queryRunner); diff --git a/packages/twenty-server/src/engine/twenty-orm/factories/entity-schema.factory.ts b/packages/twenty-server/src/engine/twenty-orm/factories/entity-schema.factory.ts index fe017a5c7e..540455377a 100644 --- a/packages/twenty-server/src/engine/twenty-orm/factories/entity-schema.factory.ts +++ b/packages/twenty-server/src/engine/twenty-orm/factories/entity-schema.factory.ts @@ -8,6 +8,7 @@ import { EntitySchemaColumnFactory } from 'src/engine/twenty-orm/factories/entit import { EntitySchemaRelationFactory } from 'src/engine/twenty-orm/factories/entity-schema-relation.factory'; import { WorkspaceEntitiesStorage } from 'src/engine/twenty-orm/storage/workspace-entities.storage'; import { computeTableName } from 'src/engine/utils/compute-table-name.util'; +import { getWorkspaceSchemaName } from 'src/engine/workspace-datasource/utils/get-workspace-schema-name.util'; @Injectable() export class EntitySchemaFactory { @@ -18,7 +19,6 @@ export class EntitySchemaFactory { async create( workspaceId: string, - _metadataVersion: number, objectMetadata: ObjectMetadataItemWithFieldMaps, objectMetadataMaps: ObjectMetadataMaps, ): Promise { @@ -29,6 +29,8 @@ export class EntitySchemaFactory { objectMetadataMaps, ); + const schemaName = getWorkspaceSchemaName(workspaceId); + const entitySchema = new EntitySchema({ name: objectMetadata.nameSingular, tableName: computeTableName( @@ -37,6 +39,7 @@ export class EntitySchemaFactory { ), columns, relations, + schema: schemaName, }); WorkspaceEntitiesStorage.setEntitySchema( diff --git a/packages/twenty-server/src/engine/twenty-orm/factories/workspace-datasource.factory.ts b/packages/twenty-server/src/engine/twenty-orm/factories/workspace-datasource.factory.ts index de42f166e5..3e72aacb3c 100644 --- a/packages/twenty-server/src/engine/twenty-orm/factories/workspace-datasource.factory.ts +++ b/packages/twenty-server/src/engine/twenty-orm/factories/workspace-datasource.factory.ts @@ -155,7 +155,6 @@ export class WorkspaceDatasourceFactory { .map((objectMetadata) => this.entitySchemaFactory.create( workspaceId, - dataSourceMetadataVersion, objectMetadata, cachedObjectMetadataMaps, ), 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 new file mode 100644 index 0000000000..09f4f41b4a --- /dev/null +++ b/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource.module.ts @@ -0,0 +1,41 @@ +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'; +import { WorkspaceFeatureFlagsMapCacheModule } from 'src/engine/metadata-modules/workspace-feature-flags-map-cache/workspace-feature-flags-map-cache.module'; +import { WorkspaceMetadataCacheModule } from 'src/engine/metadata-modules/workspace-metadata-cache/workspace-metadata-cache.module'; +import { WorkspacePermissionsCacheModule } from 'src/engine/metadata-modules/workspace-permissions-cache/workspace-permissions-cache.module'; +import { EntitySchemaColumnFactory } from 'src/engine/twenty-orm/factories/entity-schema-column.factory'; +import { EntitySchemaRelationFactory } from 'src/engine/twenty-orm/factories/entity-schema-relation.factory'; +import { EntitySchemaFactory } from 'src/engine/twenty-orm/factories/entity-schema.factory'; +import { GlobalWorkspaceOrmManager } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager'; +import { GlobalWorkspaceDataSourceService } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource.service'; +import { WorkspaceCacheStorageModule } from 'src/engine/workspace-cache-storage/workspace-cache-storage.module'; +import { WorkspaceEventEmitterModule } from 'src/engine/workspace-event-emitter/workspace-event-emitter.module'; + +@Global() +@Module({ + imports: [ + TypeOrmModule.forFeature([WorkspaceEntity]), + DataSourceModule, + WorkspaceCacheStorageModule, + WorkspaceMetadataCacheModule, + WorkspacePermissionsCacheModule, + WorkspaceFeatureFlagsMapCacheModule, + FeatureFlagModule, + TwentyConfigModule, + WorkspaceEventEmitterModule, + ], + providers: [ + GlobalWorkspaceDataSourceService, + GlobalWorkspaceOrmManager, + EntitySchemaFactory, + EntitySchemaColumnFactory, + EntitySchemaRelationFactory, + ], + exports: [GlobalWorkspaceDataSourceService, GlobalWorkspaceOrmManager], +}) +export class GlobalWorkspaceDataSourceModule {} 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 new file mode 100644 index 0000000000..99824e4a77 --- /dev/null +++ b/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource.service.ts @@ -0,0 +1,70 @@ +import { + Injectable, + Logger, + OnApplicationShutdown, + OnModuleInit, +} from '@nestjs/common'; + +import { TwentyConfigService } from 'src/engine/core-modules/twenty-config/twenty-config.service'; +import { EntitySchemaFactory } from 'src/engine/twenty-orm/factories/entity-schema.factory'; +import { GlobalWorkspaceDataSource } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource'; +import { WorkspaceEventEmitter } from 'src/engine/workspace-event-emitter/workspace-event-emitter'; + +const TWENTY_MINUTES_IN_MS = 120_000; + +@Injectable() +export class GlobalWorkspaceDataSourceService + implements OnModuleInit, OnApplicationShutdown +{ + private readonly logger = new Logger(GlobalWorkspaceDataSourceService.name); + private globalWorkspaceDataSource: GlobalWorkspaceDataSource | null = null; + + constructor( + private readonly twentyConfigService: TwentyConfigService, + private readonly entitySchemaFactory: EntitySchemaFactory, + private readonly workspaceEventEmitter: WorkspaceEventEmitter, + ) {} + + async onModuleInit(): Promise { + this.globalWorkspaceDataSource = new GlobalWorkspaceDataSource( + { + url: this.twentyConfigService.get('PG_DATABASE_URL'), + type: 'postgres', + logging: this.twentyConfigService.getLoggingConfig(), + entities: [], + ssl: this.twentyConfigService.get('PG_SSL_ALLOW_SELF_SIGNED') + ? { + rejectUnauthorized: false, + } + : undefined, + extra: { + query_timeout: 10000, + idleTimeoutMillis: TWENTY_MINUTES_IN_MS, + max: 4, + allowExitOnIdle: true, + }, + }, + this.workspaceEventEmitter, + this.entitySchemaFactory, + ); + + await this.globalWorkspaceDataSource.initialize(); + } + + public getGlobalWorkspaceDataSource(): GlobalWorkspaceDataSource { + if (!this.globalWorkspaceDataSource) { + throw new Error( + 'GlobalWorkspaceDataSource has not been initialized. Make sure the module has been initialized.', + ); + } + + return this.globalWorkspaceDataSource; + } + + async onApplicationShutdown(): Promise { + if (this.globalWorkspaceDataSource) { + await this.globalWorkspaceDataSource.destroy(); + this.globalWorkspaceDataSource = null; + } + } +} diff --git a/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource.ts b/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource.ts new file mode 100644 index 0000000000..58182579c4 --- /dev/null +++ b/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource.ts @@ -0,0 +1,364 @@ +import { type ObjectsPermissionsByRoleId } from 'twenty-shared/types'; +import { isDefined } from 'twenty-shared/utils'; +import { + DataSource, + type DataSourceOptions, + type EntityMetadata, + type EntitySchema, + type EntityTarget, + type ObjectLiteral, + type QueryRunner, + type ReplicationMode, + type SelectQueryBuilder, +} from 'typeorm'; +import { EntityManagerFactory } from 'typeorm/entity-manager/EntityManagerFactory'; +import { EntitySchemaTransformer } from 'typeorm/entity-schema/EntitySchemaTransformer'; +import { EntityMetadataNotFoundError } from 'typeorm/error/EntityMetadataNotFoundError'; +import { EntityMetadataBuilder } from 'typeorm/metadata-builder/EntityMetadataBuilder'; + +import { type FeatureFlagMap } from 'src/engine/core-modules/feature-flag/interfaces/feature-flag-map.interface'; +import { type WorkspaceInternalContext } from 'src/engine/twenty-orm/interfaces/workspace-internal-context.interface'; + +import { type AuthContext } from 'src/engine/core-modules/auth/types/auth-context.type'; +import { + PermissionsException, + PermissionsExceptionCode, +} from 'src/engine/metadata-modules/permissions/permissions.exception'; +import { type ObjectMetadataMaps } from 'src/engine/metadata-modules/types/object-metadata-maps'; +import { WorkspaceEntityManager } from 'src/engine/twenty-orm/entity-manager/workspace-entity-manager'; +import { type EntitySchemaFactory } from 'src/engine/twenty-orm/factories/entity-schema.factory'; +import { type WorkspaceQueryRunner } from 'src/engine/twenty-orm/query-runner/workspace-query-runner'; +import { type WorkspaceRepository } from 'src/engine/twenty-orm/repository/workspace.repository'; +import { getWorkspaceContext } from 'src/engine/twenty-orm/storage/workspace-context.storage'; +import { type RolePermissionConfig } from 'src/engine/twenty-orm/types/role-permission-config'; +import { type WorkspaceEventEmitter } from 'src/engine/workspace-event-emitter/workspace-event-emitter'; + +type CreateQueryBuilderOptions = { + calledByWorkspaceEntityManager?: boolean; +}; + +const ENTITY_METADATA_CACHE_TTL_MS = 60 * 60 * 1000; // 1 hour + +type CachedEntityMetadata = { + entityMetadataMap: Map, EntityMetadata>; + timestamp: number; +}; + +export class GlobalWorkspaceDataSource extends DataSource { + private entityMetadataCache: Map; + private eventEmitterService: WorkspaceEventEmitter; + private entitySchemaFactory: EntitySchemaFactory; + private _isConstructing = true; + dataSourceWithOverridenCreateQueryBuilder: GlobalWorkspaceDataSource; + + constructor( + options: DataSourceOptions, + eventEmitterService: WorkspaceEventEmitter, + entitySchemaFactory: EntitySchemaFactory, + ) { + super(options); + this.eventEmitterService = eventEmitterService; + this.entitySchemaFactory = entitySchemaFactory; + this.entityMetadataCache = new Map(); + this._isConstructing = false; + + Object.defineProperty(this, 'manager', { + get: () => this.createEntityManager(), + }); + } + + get featureFlagMap(): FeatureFlagMap { + const context = getWorkspaceContext(); + + return context.featureFlagsMap; + } + + get permissionsPerRoleId(): ObjectsPermissionsByRoleId { + const context = getWorkspaceContext(); + + return context.permissionsPerRoleId; + } + + override getRepository( + target: EntityTarget, + permissionOptions?: RolePermissionConfig, + authContext?: AuthContext, + ): WorkspaceRepository { + const manager = this.createEntityManager(); + + return manager.getRepository(target, permissionOptions, authContext); + } + + override findMetadata( + target: EntityTarget, + ): EntityMetadata | undefined { + const context = getWorkspaceContext(); + const { workspaceId, metadataVersion } = context; + const cacheKey = `${workspaceId}-${metadataVersion}`; + + const cachedEntityMetadata = this.getCachedEntityMetadata(cacheKey); + + if (!cachedEntityMetadata) { + return undefined; + } + + return cachedEntityMetadata.entityMetadataMap.get(target); + } + + override getMetadata(target: EntityTarget): EntityMetadata { + const metadata = this.findMetadata(target); + + if (!metadata) { + throw new EntityMetadataNotFoundError(target); + } + + return metadata; + } + + override createEntityManager( + queryRunner?: QueryRunner, + ): WorkspaceEntityManager { + if (this._isConstructing !== false) { + return super.createEntityManager(queryRunner) as WorkspaceEntityManager; + } + + const context = getWorkspaceContext(); + const fullContext: WorkspaceInternalContext = { + ...context, + eventEmitterService: this.eventEmitterService, + }; + + return new WorkspaceEntityManager(fullContext, this, queryRunner); + } + + override createQueryRunner( + mode = 'master' as ReplicationMode, + ): WorkspaceQueryRunner { + const queryRunner = this.driver.createQueryRunner(mode); + const manager = this.createEntityManager(queryRunner); + + Object.assign(queryRunner, { manager: manager }); + + // eslint-disable-next-line @typescript-eslint/no-explicit-any + return queryRunner as any as WorkspaceQueryRunner; + } + + // Do not use, only for specific permission-related purpose + createQueryRunnerForEntityPersistExecutor( + mode = 'master' as ReplicationMode, + ) { + if (this.dataSourceWithOverridenCreateQueryBuilder) { + const queryRunner = this.driver.createQueryRunner(mode); + const manager = new EntityManagerFactory().create( + this.dataSourceWithOverridenCreateQueryBuilder, + queryRunner, + ); + + Object.assign(queryRunner, { manager: manager }); + + return queryRunner; + } + + const dataSourceWithOverridenCreateQueryBuilder = Object.assign( + Object.create(Object.getPrototypeOf(this)), + this, + { + createQueryBuilder: ( + entityOrRunner: EntityTarget | QueryRunner, + alias?: string, + queryRunner?: QueryRunner, + ) => { + if (isDefined(alias) && typeof alias === 'string') { + const entity = entityOrRunner as EntityTarget; + + return this.createQueryBuilder(entity, alias, queryRunner, { + calledByWorkspaceEntityManager: true, + }); + } else { + const runner = entityOrRunner as QueryRunner; + + return this.createQueryBuilder(runner, { + calledByWorkspaceEntityManager: true, + }); + } + }, + }, + ); + const queryRunner = this.driver.createQueryRunner(mode); + const manager = new EntityManagerFactory().create( + dataSourceWithOverridenCreateQueryBuilder, + queryRunner, + ); + + Object.assign(queryRunner, { manager: manager }); + + return queryRunner; + } + + override createQueryBuilder( + entityClass: EntityTarget, + alias: string, + queryRunner?: QueryRunner, + options?: CreateQueryBuilderOptions, + ): SelectQueryBuilder; + + override createQueryBuilder( + queryRunner?: QueryRunner, + options?: CreateQueryBuilderOptions, // eslint-disable-next-line @typescript-eslint/no-explicit-any + ): SelectQueryBuilder; + + // Only callable from workspaceEntityManager to guarantee a permission check was run + override createQueryBuilder( + // eslint-disable-next-line @typescript-eslint/no-explicit-any + queryRunnerOrEntityClass?: QueryRunner | EntityTarget, + aliasOrOptions?: string | CreateQueryBuilderOptions, + queryRunner?: QueryRunner, + options?: CreateQueryBuilderOptions, + // eslint-disable-next-line @typescript-eslint/no-explicit-any + ): SelectQueryBuilder { + let calledByWorkspaceEntityManager; + + const isCalledWithEntityTarget = + isDefined(aliasOrOptions) && typeof aliasOrOptions === 'string'; + + if (isCalledWithEntityTarget) { + calledByWorkspaceEntityManager = options?.calledByWorkspaceEntityManager; + } else { + calledByWorkspaceEntityManager = ( + aliasOrOptions as CreateQueryBuilderOptions + )?.calledByWorkspaceEntityManager; + } + + if (!(calledByWorkspaceEntityManager === true)) { + throw new PermissionsException( + 'Method not allowed because permissions are not implemented at datasource level.', + PermissionsExceptionCode.METHOD_NOT_ALLOWED, + ); + } + + if (isCalledWithEntityTarget) { + // eslint-disable-next-line @typescript-eslint/no-explicit-any + const entityClass = queryRunnerOrEntityClass as EntityTarget; + + return super.createQueryBuilder( + entityClass, + aliasOrOptions as string, + queryRunner, + ); + } else { + const queryRunner = queryRunnerOrEntityClass as QueryRunner; + + return super.createQueryBuilder(queryRunner); + } + } + + // eslint-disable-next-line @typescript-eslint/no-explicit-any + override query( + query: string, + // eslint-disable-next-line @typescript-eslint/no-explicit-any + parameters?: any[], + queryRunner?: QueryRunner, + options?: { + shouldBypassPermissionChecks?: boolean; + }, + ): Promise { + if (!options?.shouldBypassPermissionChecks) { + throw new PermissionsException( + 'Method not allowed because permissions are not implemented at datasource level.', + PermissionsExceptionCode.METHOD_NOT_ALLOWED, + ); + } + + return super.query(query, parameters, queryRunner); + } + + async buildWorkspaceMetadata( + workspaceId: string, + metadataVersion: number, + objectMetadataMaps: ObjectMetadataMaps, + ): Promise { + const cacheKey = `${workspaceId}-${metadataVersion}`; + + if (this.getCachedEntityMetadata(cacheKey)) { + return; + } + + this.clearWorkspaceEntityMetadataCache(workspaceId); + + const entitySchemas = await Promise.all( + Object.values(objectMetadataMaps.byId) + .filter(isDefined) + .map((objectMetadata) => + this.entitySchemaFactory.create( + workspaceId, + objectMetadata, + objectMetadataMaps, + ), + ), + ); + + const entityMetadatas = this.buildMetadatasFromSchemas(entitySchemas); + + const metadataMap = new Map, EntityMetadata>(); + + for (const metadata of entityMetadatas) { + metadataMap.set(metadata.target, metadata); + } + + this.entityMetadataCache.set(cacheKey, { + entityMetadataMap: metadataMap, + timestamp: Date.now(), + }); + } + + private clearWorkspaceEntityMetadataCache(workspaceId: string): void { + for (const key of this.entityMetadataCache.keys()) { + if (key.startsWith(`${workspaceId}-`)) { + this.entityMetadataCache.delete(key); + } + } + } + + private buildMetadatasFromSchemas( + entitySchemas: EntitySchema[], + ): EntityMetadata[] { + const transformer = new EntitySchemaTransformer(); + const metadataArgsStorage = transformer.transform(entitySchemas); + + const entityMetadataBuilder = new EntityMetadataBuilder( + this, + metadataArgsStorage, + ); + + const entityMetadatas = entityMetadataBuilder.build(); + + return entityMetadatas; + } + + hasWorkspaceEntityMetadataCacheForVersion( + workspaceId: string, + metadataVersion: number, + ): boolean { + const cacheKey = `${workspaceId}-${metadataVersion}`; + + return !!this.getCachedEntityMetadata(cacheKey); + } + + private getCachedEntityMetadata( + cacheKey: string, + ): CachedEntityMetadata | undefined { + const cached = this.entityMetadataCache.get(cacheKey); + + if (!cached) { + return undefined; + } + + if (Date.now() - cached.timestamp > ENTITY_METADATA_CACHE_TTL_MS) { + this.entityMetadataCache.delete(cacheKey); + + return undefined; + } + + return cached; + } +} 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 new file mode 100644 index 0000000000..9517d4534b --- /dev/null +++ b/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager.ts @@ -0,0 +1,118 @@ +import { Injectable, type Type } from '@nestjs/common'; + +import { type ObjectLiteral } from 'typeorm'; + +import { WorkspaceFeatureFlagsMapCacheService } from 'src/engine/metadata-modules/workspace-feature-flags-map-cache/workspace-feature-flags-map-cache.service'; +import { WorkspaceMetadataCacheService } from 'src/engine/metadata-modules/workspace-metadata-cache/services/workspace-metadata-cache.service'; +import { WorkspacePermissionsCacheService } from 'src/engine/metadata-modules/workspace-permissions-cache/workspace-permissions-cache.service'; +import { GlobalWorkspaceDataSourceService } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource.service'; +import { type WorkspaceRepository } from 'src/engine/twenty-orm/repository/workspace.repository'; +import { + type WorkspaceContextForStorage, + withWorkspaceContext, +} from 'src/engine/twenty-orm/storage/workspace-context.storage'; +import { type RolePermissionConfig } from 'src/engine/twenty-orm/types/role-permission-config'; +import { convertClassNameToObjectMetadataName } from 'src/engine/workspace-manager/workspace-sync-metadata/utils/convert-class-to-object-metadata-name.util'; + +@Injectable() +export class GlobalWorkspaceOrmManager { + constructor( + private readonly globalWorkspaceDataSourceService: GlobalWorkspaceDataSourceService, + private readonly workspaceMetadataCacheService: WorkspaceMetadataCacheService, + private readonly workspaceFeatureFlagsMapCacheService: WorkspaceFeatureFlagsMapCacheService, + private readonly workspacePermissionsCacheService: WorkspacePermissionsCacheService, + ) {} + + async getRepository( + workspaceId: string, + workspaceEntity: Type, + permissionOptions?: RolePermissionConfig, + ): Promise>; + + async getRepository( + workspaceId: string, + objectMetadataName: string, + permissionOptions?: RolePermissionConfig, + ): Promise>; + + async getRepository( + workspaceId: string, + workspaceEntityOrObjectMetadataName: Type | string, + permissionOptions?: RolePermissionConfig, + ): Promise> { + let objectMetadataName: string; + + if (typeof workspaceEntityOrObjectMetadataName === 'string') { + objectMetadataName = workspaceEntityOrObjectMetadataName; + } else { + objectMetadataName = convertClassNameToObjectMetadataName( + workspaceEntityOrObjectMetadataName.name, + ); + } + + const globalDataSource = + this.globalWorkspaceDataSourceService.getGlobalWorkspaceDataSource(); + + return globalDataSource.getRepository( + objectMetadataName, + permissionOptions, + ); + } + + async getGlobalWorkspaceDataSource() { + return this.globalWorkspaceDataSourceService.getGlobalWorkspaceDataSource(); + } + + async executeInWorkspaceContext( + workspaceId: string, + fn: () => T | Promise, + ): Promise { + const context = await this.loadWorkspaceContext(workspaceId); + const globalDataSource = + this.globalWorkspaceDataSourceService.getGlobalWorkspaceDataSource(); + + if ( + !globalDataSource.hasWorkspaceEntityMetadataCacheForVersion( + workspaceId, + context.metadataVersion, + ) + ) { + await globalDataSource.buildWorkspaceMetadata( + workspaceId, + context.metadataVersion, + context.objectMetadataMaps, + ); + } + + return withWorkspaceContext(context, fn); + } + + private async loadWorkspaceContext( + workspaceId: string, + ): Promise { + const { objectMetadataMaps, metadataVersion } = + await this.workspaceMetadataCacheService.getExistingOrRecomputeMetadataMaps( + { + workspaceId, + }, + ); + + const { data: featureFlagsMap } = + await this.workspaceFeatureFlagsMapCacheService.getWorkspaceFeatureFlagsMapAndVersion( + { workspaceId }, + ); + + const { data: permissionsPerRoleId } = + await this.workspacePermissionsCacheService.getRolesPermissionsFromCache({ + workspaceId, + }); + + return { + workspaceId, + objectMetadataMaps, + metadataVersion, + featureFlagsMap, + permissionsPerRoleId, + }; + } +} diff --git a/packages/twenty-server/src/engine/twenty-orm/storage/workspace-context.storage.ts b/packages/twenty-server/src/engine/twenty-orm/storage/workspace-context.storage.ts new file mode 100644 index 0000000000..fbb924a7cf --- /dev/null +++ b/packages/twenty-server/src/engine/twenty-orm/storage/workspace-context.storage.ts @@ -0,0 +1,42 @@ +import { AsyncLocalStorage } from 'async_hooks'; + +import { type ObjectsPermissionsByRoleId } from 'twenty-shared/types'; + +import { type FeatureFlagKey } from 'src/engine/core-modules/feature-flag/enums/feature-flag-key.enum'; +import { type ObjectMetadataMaps } from 'src/engine/metadata-modules/types/object-metadata-maps'; + +export type WorkspaceContextForStorage = { + workspaceId: string; + objectMetadataMaps: ObjectMetadataMaps; + metadataVersion: number; + featureFlagsMap: Record; + permissionsPerRoleId: ObjectsPermissionsByRoleId; +}; + +export const workspaceContextStorage = + new AsyncLocalStorage(); + +export const getWorkspaceContext = (): WorkspaceContextForStorage => { + const context = workspaceContextStorage.getStore(); + + if (!context) { + throw new Error( + 'Workspace context not set. Operations must be wrapped with withWorkspaceContext()', + ); + } + + return context; +}; + +export const withWorkspaceContext = ( + context: WorkspaceContextForStorage, + fn: () => T | Promise, +): T | Promise => { + return workspaceContextStorage.run(context, fn); +}; + +export const setWorkspaceContext = ( + context: WorkspaceContextForStorage, +): void => { + workspaceContextStorage.enterWith(context); +};