Global workspace datasource poc (#15744)

## 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<ObjectLiteral>): 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
```
This commit is contained in:
Weiko
2025-11-10 17:32:49 +01:00
committed by GitHub
parent 7a1e699fc8
commit 02fda92a93
14 changed files with 745 additions and 17 deletions
@@ -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',
@@ -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',
@@ -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<Args>,
authContext: WorkspaceAuthContext,
queryRunnerContext: CommonBaseQueryRunnerContext,
commonQueryParser: GraphqlQueryParser,
isGlobalDatasourceEnabled: boolean,
): Promise<Output> {
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<Omit<CommonExtendedQueryRunnerContext, 'commonQueryParser'>> {
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;
@@ -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,
@@ -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',
}
@@ -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;
@@ -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);
@@ -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<EntitySchema> {
@@ -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(
@@ -155,7 +155,6 @@ export class WorkspaceDatasourceFactory {
.map((objectMetadata) =>
this.entitySchemaFactory.create(
workspaceId,
dataSourceMetadataVersion,
objectMetadata,
cachedObjectMetadataMaps,
),
@@ -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 {}
@@ -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<void> {
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<void> {
if (this.globalWorkspaceDataSource) {
await this.globalWorkspaceDataSource.destroy();
this.globalWorkspaceDataSource = null;
}
}
}
@@ -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<EntityTarget<ObjectLiteral>, EntityMetadata>;
timestamp: number;
};
export class GlobalWorkspaceDataSource extends DataSource {
private entityMetadataCache: Map<string, CachedEntityMetadata>;
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<Entity extends ObjectLiteral>(
target: EntityTarget<Entity>,
permissionOptions?: RolePermissionConfig,
authContext?: AuthContext,
): WorkspaceRepository<Entity> {
const manager = this.createEntityManager();
return manager.getRepository(target, permissionOptions, authContext);
}
override findMetadata(
target: EntityTarget<ObjectLiteral>,
): 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<ObjectLiteral>): 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<ObjectLiteral> | QueryRunner,
alias?: string,
queryRunner?: QueryRunner,
) => {
if (isDefined(alias) && typeof alias === 'string') {
const entity = entityOrRunner as EntityTarget<ObjectLiteral>;
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<Entity extends ObjectLiteral>(
entityClass: EntityTarget<Entity>,
alias: string,
queryRunner?: QueryRunner,
options?: CreateQueryBuilderOptions,
): SelectQueryBuilder<Entity>;
override createQueryBuilder(
queryRunner?: QueryRunner,
options?: CreateQueryBuilderOptions, // eslint-disable-next-line @typescript-eslint/no-explicit-any
): SelectQueryBuilder<any>;
// 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<any>,
aliasOrOptions?: string | CreateQueryBuilderOptions,
queryRunner?: QueryRunner,
options?: CreateQueryBuilderOptions,
// eslint-disable-next-line @typescript-eslint/no-explicit-any
): SelectQueryBuilder<any> {
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<any>;
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<T = any>(
query: string,
// eslint-disable-next-line @typescript-eslint/no-explicit-any
parameters?: any[],
queryRunner?: QueryRunner,
options?: {
shouldBypassPermissionChecks?: boolean;
},
): Promise<T> {
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<void> {
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<EntityTarget<ObjectLiteral>, 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;
}
}
@@ -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<T extends ObjectLiteral>(
workspaceId: string,
workspaceEntity: Type<T>,
permissionOptions?: RolePermissionConfig,
): Promise<WorkspaceRepository<T>>;
async getRepository<T extends ObjectLiteral>(
workspaceId: string,
objectMetadataName: string,
permissionOptions?: RolePermissionConfig,
): Promise<WorkspaceRepository<T>>;
async getRepository<T extends ObjectLiteral>(
workspaceId: string,
workspaceEntityOrObjectMetadataName: Type<T> | string,
permissionOptions?: RolePermissionConfig,
): Promise<WorkspaceRepository<T>> {
let objectMetadataName: string;
if (typeof workspaceEntityOrObjectMetadataName === 'string') {
objectMetadataName = workspaceEntityOrObjectMetadataName;
} else {
objectMetadataName = convertClassNameToObjectMetadataName(
workspaceEntityOrObjectMetadataName.name,
);
}
const globalDataSource =
this.globalWorkspaceDataSourceService.getGlobalWorkspaceDataSource();
return globalDataSource.getRepository<T>(
objectMetadataName,
permissionOptions,
);
}
async getGlobalWorkspaceDataSource() {
return this.globalWorkspaceDataSourceService.getGlobalWorkspaceDataSource();
}
async executeInWorkspaceContext<T>(
workspaceId: string,
fn: () => T | Promise<T>,
): Promise<T> {
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<WorkspaceContextForStorage> {
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,
};
}
}
@@ -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<FeatureFlagKey, boolean>;
permissionsPerRoleId: ObjectsPermissionsByRoleId;
};
export const workspaceContextStorage =
new AsyncLocalStorage<WorkspaceContextForStorage>();
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 = <T>(
context: WorkspaceContextForStorage,
fn: () => T | Promise<T>,
): T | Promise<T> => {
return workspaceContextStorage.run(context, fn);
};
export const setWorkspaceContext = (
context: WorkspaceContextForStorage,
): void => {
workspaceContextStorage.enterWith(context);
};