Deprecate workspace datasoure (#16507)

This commit is contained in:
Weiko
2025-12-11 18:27:49 +01:00
committed by GitHub
parent 3dd2684254
commit 083a038f8a
39 changed files with 200 additions and 1929 deletions
@@ -48,7 +48,6 @@ import {
} from 'src/engine/metadata-modules/permissions/permissions.exception';
import { PermissionsService } from 'src/engine/metadata-modules/permissions/permissions.service';
import { UserRoleService } from 'src/engine/metadata-modules/user-role/user-role.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 { WorkspaceCacheService } from 'src/engine/workspace-cache/services/workspace-cache.service';
@@ -341,8 +340,7 @@ export abstract class CommonBaseQueryRunnerService<
return {
...queryRunnerContext,
authContext,
workspaceDataSource:
globalWorkspaceDataSource as unknown as WorkspaceDataSource,
workspaceDataSource: globalWorkspaceDataSource,
rolePermissionConfig,
repository,
};
@@ -36,7 +36,7 @@ import { type FlatFieldMetadata } from 'src/engine/metadata-modules/flat-field-m
import { buildFieldMapsFromFlatObjectMetadata } from 'src/engine/metadata-modules/flat-field-metadata/utils/build-field-maps-from-flat-object-metadata.util';
import { type FlatObjectMetadata } from 'src/engine/metadata-modules/flat-object-metadata/types/flat-object-metadata.type';
import { assertMutationNotOnRemoteObject } from 'src/engine/metadata-modules/object-metadata/utils/assert-mutation-not-on-remote-object.util';
import { WorkspaceDataSource } from 'src/engine/twenty-orm/datasource/workspace.datasource';
import { GlobalWorkspaceDataSource } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource';
import { WorkspaceRepository } from 'src/engine/twenty-orm/repository/workspace.repository';
import { RolePermissionConfig } from 'src/engine/twenty-orm/types/role-permission-config';
@@ -117,7 +117,7 @@ export class CommonCreateManyQueryRunnerService extends CommonBaseQueryRunnerSer
flatObjectMetadataMaps: FlatEntityMaps<FlatObjectMetadata>;
flatFieldMetadataMaps: FlatEntityMaps<FlatFieldMetadata>;
authContext: AuthContext;
workspaceDataSource: WorkspaceDataSource;
workspaceDataSource: GlobalWorkspaceDataSource;
rolePermissionConfig?: RolePermissionConfig;
}): Promise<void> {
if (!args.selectedFieldsResult.relations) {
@@ -4,7 +4,7 @@ import { type WorkspaceAuthContext } from 'src/engine/api/common/interfaces/work
import { type CommonBaseQueryRunnerContext } from 'src/engine/api/common/types/common-base-query-runner-context.type';
import { type GraphqlQueryParser } from 'src/engine/api/graphql/graphql-query-runner/graphql-query-parsers/graphql-query.parser';
import { type WorkspaceDataSource } from 'src/engine/twenty-orm/datasource/workspace.datasource';
import { type GlobalWorkspaceDataSource } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource';
import { type WorkspaceRepository } from 'src/engine/twenty-orm/repository/workspace.repository';
import { type RolePermissionConfig } from 'src/engine/twenty-orm/types/role-permission-config';
@@ -16,5 +16,5 @@ export type CommonExtendedQueryRunnerContext = Omit<
rolePermissionConfig: RolePermissionConfig;
repository: WorkspaceRepository<ObjectLiteral>;
commonQueryParser: GraphqlQueryParser;
workspaceDataSource: WorkspaceDataSource;
workspaceDataSource: GlobalWorkspaceDataSource;
};
@@ -13,7 +13,7 @@ import { ProcessAggregateHelper } from 'src/engine/api/graphql/graphql-query-run
import { buildColumnsToSelect } from 'src/engine/api/graphql/graphql-query-runner/utils/build-columns-to-select';
import { getTargetObjectMetadataOrThrow } from 'src/engine/api/graphql/graphql-query-runner/utils/get-target-object-metadata.util';
import { type AggregationField } from 'src/engine/api/graphql/workspace-schema-builder/utils/get-available-aggregations-from-object-fields.util';
import { type AuthContext } from 'src/engine/core-modules/auth/types/auth-context.type';
import { AuthContext } from 'src/engine/core-modules/auth/types/auth-context.type';
import { FlatEntityMaps } from 'src/engine/metadata-modules/flat-entity/types/flat-entity-maps.type';
import { findFlatEntityByIdInFlatEntityMaps } from 'src/engine/metadata-modules/flat-entity/utils/find-flat-entity-by-id-in-flat-entity-maps.util';
import { FlatFieldMetadata } from 'src/engine/metadata-modules/flat-field-metadata/types/flat-field-metadata.type';
@@ -22,7 +22,7 @@ import {
type FieldMapsForObject,
} from 'src/engine/metadata-modules/flat-field-metadata/utils/build-field-maps-from-flat-object-metadata.util';
import { FlatObjectMetadata } from 'src/engine/metadata-modules/flat-object-metadata/types/flat-object-metadata.type';
import { type WorkspaceDataSource } from 'src/engine/twenty-orm/datasource/workspace.datasource';
import { GlobalWorkspaceDataSource } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource';
import { type WorkspaceSelectQueryBuilder } from 'src/engine/twenty-orm/repository/workspace-select-query-builder';
import { type RolePermissionConfig } from 'src/engine/twenty-orm/types/role-permission-config';
import { isFieldMetadataEntityOfType } from 'src/engine/utils/is-field-metadata-of-type.util';
@@ -55,7 +55,7 @@ export class ProcessNestedRelationsV2Helper {
aggregate?: Record<string, AggregationField>;
limit: number;
authContext: AuthContext;
workspaceDataSource: WorkspaceDataSource;
workspaceDataSource: GlobalWorkspaceDataSource;
rolePermissionConfig?: RolePermissionConfig;
// eslint-disable-next-line @typescript-eslint/no-explicit-any
selectedFields: Record<string, any>;
@@ -111,7 +111,7 @@ export class ProcessNestedRelationsV2Helper {
aggregate: Record<string, AggregationField>;
limit: number;
authContext: AuthContext;
workspaceDataSource: WorkspaceDataSource;
workspaceDataSource: GlobalWorkspaceDataSource;
rolePermissionConfig?: RolePermissionConfig;
selectedFields: Record<string, unknown>;
}): Promise<void> {
@@ -5,11 +5,11 @@ import { type FindOptionsRelations, type ObjectLiteral } from 'typeorm';
import { ProcessNestedRelationsV2Helper } from 'src/engine/api/graphql/graphql-query-runner/helpers/process-nested-relations-v2.helper';
import { type AggregationField } from 'src/engine/api/graphql/workspace-schema-builder/utils/get-available-aggregations-from-object-fields.util';
import { type AuthContext } from 'src/engine/core-modules/auth/types/auth-context.type';
import { AuthContext } from 'src/engine/core-modules/auth/types/auth-context.type';
import { FlatEntityMaps } from 'src/engine/metadata-modules/flat-entity/types/flat-entity-maps.type';
import { FlatFieldMetadata } from 'src/engine/metadata-modules/flat-field-metadata/types/flat-field-metadata.type';
import { FlatObjectMetadata } from 'src/engine/metadata-modules/flat-object-metadata/types/flat-object-metadata.type';
import { type WorkspaceDataSource } from 'src/engine/twenty-orm/datasource/workspace.datasource';
import { GlobalWorkspaceDataSource } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource';
import { RolePermissionConfig } from 'src/engine/twenty-orm/types/role-permission-config';
@Injectable()
@@ -42,7 +42,7 @@ export class ProcessNestedRelationsHelper {
aggregate?: Record<string, AggregationField>;
limit: number;
authContext: AuthContext;
workspaceDataSource: WorkspaceDataSource;
workspaceDataSource: GlobalWorkspaceDataSource;
rolePermissionConfig?: RolePermissionConfig;
// eslint-disable-next-line @typescript-eslint/no-explicit-any
selectedFields: Record<string, any>;
@@ -93,9 +93,9 @@ describe('GoogleAPIsService', () => {
return {};
}),
getDataSourceForWorkspace: jest
getGlobalWorkspaceDataSource: jest
.fn()
.mockImplementation(() => mockWorkspaceDataSource),
.mockResolvedValue(mockWorkspaceDataSource),
executeInWorkspaceContext: jest
.fn()
.mockImplementation((_authContext: any, fn: () => any) => fn()),
@@ -126,9 +126,7 @@ export class GoogleAPIsService {
);
const workspaceDataSource =
await this.globalWorkspaceOrmManager.getDataSourceForWorkspace(
workspaceId,
);
await this.globalWorkspaceOrmManager.getGlobalWorkspaceDataSource();
await workspaceDataSource.transaction(
async (manager: WorkspaceEntityManager) => {
@@ -92,9 +92,9 @@ describe('MicrosoftAPIsService', () => {
return {};
}),
getDataSourceForWorkspace: jest
getGlobalWorkspaceDataSource: jest
.fn()
.mockImplementation(() => mockWorkspaceDataSource),
.mockResolvedValue(mockWorkspaceDataSource),
executeInWorkspaceContext: jest
.fn()
.mockImplementation((_authContext: any, fn: () => any) => fn()),
@@ -107,9 +107,7 @@ export class MicrosoftAPIsService {
);
const workspaceDataSource =
await this.globalWorkspaceOrmManager.getDataSourceForWorkspace(
workspaceId,
);
await this.globalWorkspaceOrmManager.getGlobalWorkspaceDataSource();
await workspaceDataSource.transaction(
async (manager: WorkspaceEntityManager) => {
@@ -1,258 +0,0 @@
import { type Entity } from '@microsoft/microsoft-graph-types';
import { isDefined } from 'class-validator';
import { type ObjectsPermissionsByRoleId } from 'twenty-shared/types';
import {
DataSource,
type DataSourceOptions,
type EntityTarget,
type ObjectLiteral,
type QueryRunner,
type ReplicationMode,
type SelectQueryBuilder,
} from 'typeorm';
import { EntityManagerFactory } from 'typeorm/entity-manager/EntityManagerFactory';
import { type FeatureFlagMap } from 'src/engine/core-modules/feature-flag/interfaces/feature-flag-map.interface';
import { type WorkspaceDataSourceInterface } from 'src/engine/twenty-orm/interfaces/workspace-datasource.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 { WorkspaceEntityManager } from 'src/engine/twenty-orm/entity-manager/workspace-entity-manager';
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 { type RolePermissionConfig } from 'src/engine/twenty-orm/types/role-permission-config';
type CreateQueryBuilderOptions = {
calledByWorkspaceEntityManager?: boolean;
};
export class WorkspaceDataSource
extends DataSource
implements WorkspaceDataSourceInterface
{
readonly isGlobalFlow = false;
readonly internalContext: WorkspaceInternalContext;
readonly manager: WorkspaceEntityManager;
featureFlagMapVersion: string;
featureFlagMap: FeatureFlagMap;
rolesPermissionsVersion: string;
permissionsPerRoleId: ObjectsPermissionsByRoleId;
dataSourceWithOverridenCreateQueryBuilder: WorkspaceDataSource;
isPoolSharingEnabled: boolean;
constructor(
internalContext: WorkspaceInternalContext,
options: DataSourceOptions,
featureFlagMapVersion: string,
featureFlagMap: FeatureFlagMap,
rolesPermissionsVersion: string,
permissionsPerRoleId: ObjectsPermissionsByRoleId,
isPoolSharingEnabled: boolean,
) {
super(options);
this.internalContext = internalContext;
this.featureFlagMap = featureFlagMap;
this.featureFlagMapVersion = featureFlagMapVersion;
// Recreate manager after internalContext has been initialized
this.manager = this.createEntityManager();
this.rolesPermissionsVersion = rolesPermissionsVersion;
this.permissionsPerRoleId = permissionsPerRoleId;
this.isPoolSharingEnabled = isPoolSharingEnabled;
}
override getRepository<Entity extends ObjectLiteral>(
target: EntityTarget<Entity>,
permissionOptions?: RolePermissionConfig,
authContext?: AuthContext,
): WorkspaceRepository<Entity> {
return this.manager.getRepository(target, permissionOptions, authContext);
}
override createEntityManager(
queryRunner?: QueryRunner,
): WorkspaceEntityManager {
return new WorkspaceEntityManager(this.internalContext, 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<Entity> | QueryRunner,
alias?: string,
queryRunner?: QueryRunner,
) => {
if (isDefined(alias) && typeof alias === 'string') {
const entity = entityOrRunner as EntityTarget<Entity>;
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);
}
setRolesPermissionsVersion(rolesPermissionsVersion: string) {
this.rolesPermissionsVersion = rolesPermissionsVersion;
}
setRolesPermissions(permissionsPerRoleId: ObjectsPermissionsByRoleId) {
this.permissionsPerRoleId = permissionsPerRoleId;
}
setFeatureFlagMap(featureFlagMap: FeatureFlagMap) {
this.featureFlagMap = featureFlagMap;
}
setFeatureFlagMapVersion(featureFlagMapVersion: string) {
this.featureFlagMapVersion = featureFlagMapVersion;
}
override async destroy(): Promise<void> {
if (this.isPoolSharingEnabled) {
// eslint-disable-next-line no-console
console.log(
`PromiseMemoizer Event: A WorkspaceDataSource for workspace ${this.internalContext.workspaceId} is being cleared. Actual pool closure managed by PgPoolSharedService. Not calling dataSource.destroy().`,
);
// We should NOT call dataSource.destroy() here, because that would end
// the shared pool, potentially affecting other active users of that pool.
// The PgPoolSharedService is responsible for the lifecycle of shared pools.
return Promise.resolve();
} else {
// eslint-disable-next-line no-console
console.log(
`PromiseMemoizer Event: A WorkspaceDataSource for workspace ${this.internalContext.workspaceId} is being cleared. Calling safelyDestroyDataSource.`,
);
return super.destroy();
}
}
}
@@ -6,13 +6,19 @@ import { EntityManager } from 'typeorm';
import { EntityPersistExecutor } from 'typeorm/persistence/EntityPersistExecutor';
import { PlainObjectToDatabaseEntityTransformer } from 'typeorm/query-builder/transformer/PlainObjectToDatabaseEntityTransformer';
import { type WorkspaceAuthContext } from 'src/engine/api/common/interfaces/workspace-auth-context.interface';
import { type WorkspaceInternalContext } from 'src/engine/twenty-orm/interfaces/workspace-internal-context.interface';
import { type FlatEntityMaps } from 'src/engine/metadata-modules/flat-entity/types/flat-entity-maps.type';
import { type FlatFieldMetadata } from 'src/engine/metadata-modules/flat-field-metadata/types/flat-field-metadata.type';
import { type FlatObjectMetadata } from 'src/engine/metadata-modules/flat-object-metadata/types/flat-object-metadata.type';
import { type WorkspaceDataSource } from 'src/engine/twenty-orm/datasource/workspace.datasource';
import { type GlobalWorkspaceDataSource } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource';
import { validateOperationIsPermittedOrThrow } from 'src/engine/twenty-orm/repository/permissions.utils';
import {
setWorkspaceContext,
withWorkspaceContext,
type ORMWorkspaceContext,
} from 'src/engine/twenty-orm/storage/orm-workspace-context.storage';
import { getObjectMetadataFromEntityTarget } from 'src/engine/twenty-orm/utils/get-object-metadata-from-entity-target.util';
import { WorkspaceEntityManager } from './workspace-entity-manager';
@@ -77,12 +83,13 @@ jest.mock('../repository/workspace-select-query-builder', () => ({
describe('WorkspaceEntityManager', () => {
let entityManager: WorkspaceEntityManager;
let mockInternalContext: WorkspaceInternalContext;
let mockDataSource: WorkspaceDataSource;
let mockDataSource: GlobalWorkspaceDataSource;
let mockPermissionOptions: {
shouldBypassPermissionChecks: boolean;
objectRecordsPermissions?: ObjectsPermissions;
};
let mockInternalContext: WorkspaceInternalContext;
let mockWorkspaceContext: ORMWorkspaceContext;
beforeEach(() => {
const mockFlatObjectMetadata: FlatObjectMetadata = {
@@ -234,7 +241,8 @@ describe('WorkspaceEntityManager', () => {
IS_GLOBAL_WORKSPACE_DATASOURCE_ENABLED: false,
},
permissionsPerRoleId: {},
} as WorkspaceDataSource;
eventEmitterService: mockInternalContext.eventEmitterService,
} as GlobalWorkspaceDataSource;
mockPermissionOptions = {
shouldBypassPermissionChecks: false,
@@ -249,6 +257,30 @@ describe('WorkspaceEntityManager', () => {
},
};
const mockAuthContext = {
user: { id: 'user-id' },
workspace: { id: 'test-workspace-id' },
workspaceMemberId: 'workspace-member-id',
userWorkspaceId: 'user-workspace-id',
apiKey: null,
} as unknown as WorkspaceAuthContext;
mockWorkspaceContext = {
authContext: mockAuthContext,
flatObjectMetadataMaps,
flatFieldMetadataMaps,
flatIndexMaps: mockInternalContext.flatIndexMaps,
objectIdByNameSingular: mockInternalContext.objectIdByNameSingular,
featureFlagsMap: mockInternalContext.featureFlagsMap,
permissionsPerRoleId: mockDataSource.permissionsPerRoleId,
entityMetadatas: [],
userWorkspaceRoleMap: {
'user-workspace-id': 'role-id',
},
};
setWorkspaceContext(mockWorkspaceContext);
// Mock TypeORM connection methods
const mockWorkspaceDataSource = {
getMetadata: jest.fn().mockReturnValue({
@@ -258,6 +290,7 @@ describe('WorkspaceEntityManager', () => {
findInheritanceMetadata: jest.fn(),
findColumnWithPropertyPath: jest.fn(),
}),
eventEmitterService: mockInternalContext.eventEmitterService,
createQueryBuilder: jest.fn().mockReturnValue({
delete: jest.fn().mockReturnThis(),
from: jest.fn().mockReturnThis(),
@@ -290,10 +323,7 @@ describe('WorkspaceEntityManager', () => {
}),
};
entityManager = new WorkspaceEntityManager(
mockInternalContext,
mockDataSource,
);
entityManager = new WorkspaceEntityManager(mockDataSource);
Object.defineProperty(entityManager, 'connection', {
get: () => mockWorkspaceDataSource,
@@ -362,7 +392,9 @@ describe('WorkspaceEntityManager', () => {
describe('Query Method', () => {
it('should call validatePermissions and validateOperationIsPermittedOrThrow for find', async () => {
await entityManager.find('test-entity', {}, mockPermissionOptions);
await withWorkspaceContext(mockWorkspaceContext, () =>
entityManager.find('test-entity', {}, mockPermissionOptions),
);
expect(entityManager.createQueryBuilder).toHaveBeenCalledWith(
'test-entity',
@@ -380,11 +412,13 @@ describe('WorkspaceEntityManager', () => {
describe('Save Methods', () => {
it('should call validatePermissions and validateOperationIsPermittedOrThrow for save', async () => {
await entityManager.save(
'test-entity',
{},
{ reload: false },
mockPermissionOptions,
await withWorkspaceContext(mockWorkspaceContext, () =>
entityManager.save(
'test-entity',
{},
{ reload: false },
mockPermissionOptions,
),
);
expect(entityManager['validatePermissions']).toHaveBeenCalledWith({
target: 'test-entity',
@@ -409,7 +443,9 @@ describe('WorkspaceEntityManager', () => {
describe('Update Methods', () => {
it('should call createQueryBuilder with permissionOptions for update', async () => {
await entityManager.update('test-entity', {}, {}, mockPermissionOptions);
await withWorkspaceContext(mockWorkspaceContext, () =>
entityManager.update('test-entity', {}, {}, mockPermissionOptions),
);
expect(entityManager['createQueryBuilder']).toHaveBeenCalledWith(
'test-entity',
undefined,
@@ -421,7 +457,9 @@ describe('WorkspaceEntityManager', () => {
describe('Other Methods', () => {
it('should call validatePermissions and validateOperationIsPermittedOrThrow for clear', async () => {
await entityManager.clear('test-entity', mockPermissionOptions);
await withWorkspaceContext(mockWorkspaceContext, () =>
entityManager.clear('test-entity', mockPermissionOptions),
);
expect(entityManager['validatePermissions']).toHaveBeenCalledWith({
target: 'test-entity',
operationType: 'delete',
@@ -34,7 +34,6 @@ import { type UpsertOptions } from 'typeorm/repository/UpsertOptions';
import { InstanceChecker } from 'typeorm/util/InstanceChecker';
import { type FeatureFlagMap } from 'src/engine/core-modules/feature-flag/interfaces/feature-flag-map.interface';
import { type WorkspaceDataSourceInterface } from 'src/engine/twenty-orm/interfaces/workspace-datasource.interface';
import { type WorkspaceInternalContext } from 'src/engine/twenty-orm/interfaces/workspace-internal-context.interface';
import { DatabaseEventAction } from 'src/engine/api/graphql/graphql-query-runner/enums/database-event-action';
@@ -45,7 +44,6 @@ import {
PermissionsException,
PermissionsExceptionCode,
} from 'src/engine/metadata-modules/permissions/permissions.exception';
import { type WorkspaceDataSource } from 'src/engine/twenty-orm/datasource/workspace.datasource';
import { type DeepPartialWithNestedRelationFields } from 'src/engine/twenty-orm/entity-manager/types/deep-partial-entity-with-nested-relation-fields.type';
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';
@@ -73,52 +71,34 @@ type PermissionOptions = {
};
export class WorkspaceEntityManager extends EntityManager {
private readonly _legacyInternalContext?: WorkspaceInternalContext;
// eslint-disable-next-line @typescript-eslint/no-explicit-any
readonly repositories: Map<string, Repository<any>>;
declare connection: WorkspaceDataSource;
declare connection: GlobalWorkspaceDataSource;
constructor(
legacyInternalContext: WorkspaceInternalContext | undefined,
connection: WorkspaceDataSource | GlobalWorkspaceDataSource,
connection: GlobalWorkspaceDataSource,
queryRunner?: QueryRunner,
) {
super(connection, queryRunner);
this.repositories = new Map();
this._legacyInternalContext = legacyInternalContext;
}
private get eventEmitterService(): WorkspaceEventEmitter {
const isGlobalFlow = (this.connection as WorkspaceDataSourceInterface)
.isGlobalFlow;
if (isGlobalFlow) {
return (this.connection as unknown as GlobalWorkspaceDataSource)
.eventEmitterService;
}
return this._legacyInternalContext!.eventEmitterService;
return this.connection.eventEmitterService;
}
get internalContext(): WorkspaceInternalContext {
const isGlobalFlow = (this.connection as WorkspaceDataSourceInterface)
.isGlobalFlow;
const context = getWorkspaceContext();
if (isGlobalFlow) {
const context = getWorkspaceContext();
return {
workspaceId: context.authContext.workspace.id,
flatObjectMetadataMaps: context.flatObjectMetadataMaps,
flatFieldMetadataMaps: context.flatFieldMetadataMaps,
flatIndexMaps: context.flatIndexMaps,
objectIdByNameSingular: context.objectIdByNameSingular,
featureFlagsMap: context.featureFlagsMap,
eventEmitterService: this.eventEmitterService,
};
}
return this._legacyInternalContext!;
return {
workspaceId: context.authContext.workspace.id,
flatObjectMetadataMaps: context.flatObjectMetadataMaps,
flatFieldMetadataMaps: context.flatFieldMetadataMaps,
flatIndexMaps: context.flatIndexMaps,
objectIdByNameSingular: context.objectIdByNameSingular,
featureFlagsMap: context.featureFlagsMap,
eventEmitterService: this.eventEmitterService,
};
}
getFeatureFlagMap(): FeatureFlagMap {
@@ -487,7 +467,7 @@ export class WorkspaceEntityManager extends EntityManager {
}
private extractTargetNameSingularFromEntityTarget(
target: EntityTarget<unknown>,
target: EntityTarget<ObjectLiteral>,
): string {
return this.connection.getMetadata(target).name;
}
@@ -1,11 +1,9 @@
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 { WorkspaceDatasourceFactory } from 'src/engine/twenty-orm/factories/workspace-datasource.factory';
export const entitySchemaFactories = [
EntitySchemaColumnFactory,
EntitySchemaRelationFactory,
EntitySchemaFactory,
WorkspaceDatasourceFactory,
];
@@ -1,353 +0,0 @@
import { Injectable, Logger } from '@nestjs/common';
import { InjectRepository } from '@nestjs/typeorm';
import crypto from 'crypto';
import { type ObjectsPermissionsByRoleId } from 'twenty-shared/types';
import { isDefined } from 'twenty-shared/utils';
import { EntitySchema, Repository } from 'typeorm';
import { type FeatureFlagMap } from 'src/engine/core-modules/feature-flag/interfaces/feature-flag-map.interface';
import { TwentyConfigService } from 'src/engine/core-modules/twenty-config/twenty-config.service';
import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity';
import { DataSourceService } from 'src/engine/metadata-modules/data-source/data-source.service';
import { WorkspaceManyOrAllFlatEntityMapsCacheService } from 'src/engine/metadata-modules/flat-entity/services/workspace-many-or-all-flat-entity-maps-cache.service';
import { buildObjectIdByNameMaps } from 'src/engine/metadata-modules/flat-object-metadata/utils/build-object-id-by-name-maps.util';
import { WorkspaceDataSource } from 'src/engine/twenty-orm/datasource/workspace.datasource';
import {
TwentyORMException,
TwentyORMExceptionCode,
} from 'src/engine/twenty-orm/exceptions/twenty-orm.exception';
import { EntitySchemaFactory } from 'src/engine/twenty-orm/factories/entity-schema.factory';
import { PromiseMemoizer } from 'src/engine/twenty-orm/storage/promise-memoizer.storage';
import { type CacheKey } from 'src/engine/twenty-orm/storage/types/cache-key.type';
import { WorkspaceCacheStorageService } from 'src/engine/workspace-cache-storage/workspace-cache-storage.service';
import { WorkspaceCacheService } from 'src/engine/workspace-cache/services/workspace-cache.service';
import { WorkspaceEventEmitter } from 'src/engine/workspace-event-emitter/workspace-event-emitter';
const TWENTY_MINUTES_IN_MS = 120_000;
@Injectable()
export class WorkspaceDatasourceFactory {
private readonly logger = new Logger(WorkspaceDatasourceFactory.name);
private promiseMemoizer = new PromiseMemoizer<WorkspaceDataSource>();
constructor(
private readonly dataSourceService: DataSourceService,
private readonly twentyConfigService: TwentyConfigService,
private readonly workspaceCacheStorageService: WorkspaceCacheStorageService,
private readonly workspaceManyOrAllFlatEntityMapsCacheService: WorkspaceManyOrAllFlatEntityMapsCacheService,
private readonly entitySchemaFactory: EntitySchemaFactory,
private readonly workspaceCacheService: WorkspaceCacheService,
@InjectRepository(WorkspaceEntity)
private readonly workspaceRepository: Repository<WorkspaceEntity>,
private readonly workspaceEventEmitter: WorkspaceEventEmitter,
) {}
private async safelyDestroyDataSource(
dataSource: WorkspaceDataSource,
): Promise<void> {
try {
await dataSource.destroy();
} catch (error) {
// Ignore known race-condition errors to prevent noise during shutdown
if (
error.message === 'Called end on pool more than once' ||
error.message?.includes(
'pool is draining and cannot accommodate new clients',
)
) {
this.logger.debug(
`Ignoring pool error during cleanup: ${error.message}`,
);
return;
}
throw error;
}
}
public async create(workspaceId: string): Promise<WorkspaceDataSource> {
const dataSourceMetadataVersion =
await this.getWorkspaceMetadataVersionFromCacheOrFromDB(workspaceId);
const {
featureFlagsMap: cachedFeatureFlagMap,
rolesPermissions: cachedRolesPermissions,
} = await this.workspaceCacheService.getOrRecompute(workspaceId, [
'featureFlagsMap',
'rolesPermissions',
]);
const cachedFeatureFlagMapVersion = crypto
.createHash('sha256')
.update(JSON.stringify(cachedFeatureFlagMap))
.digest('hex');
const cachedRolesPermissionsVersion = crypto
.createHash('sha256')
.update(JSON.stringify(cachedRolesPermissions))
.digest('hex');
const cacheKey: CacheKey = `${workspaceId}-${dataSourceMetadataVersion}`;
const workspaceDataSource =
await this.promiseMemoizer.memoizePromiseAndExecute(
cacheKey,
async () => {
const dataSourceMetadata =
await this.dataSourceService.getLastDataSourceMetadataFromWorkspaceId(
workspaceId,
);
if (!dataSourceMetadata) {
throw new TwentyORMException(
`Workspace Schema not found for workspace ${workspaceId}`,
TwentyORMExceptionCode.WORKSPACE_SCHEMA_NOT_FOUND,
);
}
const cachedEntitySchemaOptions =
await this.workspaceCacheStorageService.getORMEntitySchema(
workspaceId,
dataSourceMetadataVersion,
);
let cachedEntitySchemas: EntitySchema[];
const {
flatObjectMetadataMaps,
flatFieldMetadataMaps,
flatIndexMaps,
} =
await this.workspaceManyOrAllFlatEntityMapsCacheService.getOrRecomputeManyOrAllFlatEntityMaps(
{
workspaceId,
flatMapsKeys: [
'flatObjectMetadataMaps',
'flatFieldMetadataMaps',
'flatIndexMaps',
],
},
);
const { idByNameSingular: objectIdByNameSingular } =
buildObjectIdByNameMaps(flatObjectMetadataMaps);
const metadataVersionForFinalUpToDateCheck =
await this.workspaceCacheStorageService.getMetadataVersion(
workspaceId,
);
if (
metadataVersionForFinalUpToDateCheck !== dataSourceMetadataVersion
) {
throw new TwentyORMException(
`Workspace metadata version mismatch detected for workspace ${workspaceId}. Latest version: ${metadataVersionForFinalUpToDateCheck}. Built version: ${dataSourceMetadataVersion}`,
TwentyORMExceptionCode.METADATA_VERSION_MISMATCH,
);
}
if (cachedEntitySchemaOptions) {
cachedEntitySchemas = cachedEntitySchemaOptions.map(
(option) => new EntitySchema(option),
);
} else {
const entitySchemas = Object.values(flatObjectMetadataMaps.byId)
.filter(isDefined)
.map((flatObjectMetadata) =>
this.entitySchemaFactory.create(
workspaceId,
flatObjectMetadata,
flatObjectMetadataMaps,
flatFieldMetadataMaps,
),
);
await this.workspaceCacheStorageService.setORMEntitySchema(
workspaceId,
dataSourceMetadataVersion,
entitySchemas.map((entitySchema) => entitySchema.options),
);
cachedEntitySchemas = entitySchemas;
}
const workspaceDataSource = new WorkspaceDataSource(
{
workspaceId,
flatObjectMetadataMaps,
flatFieldMetadataMaps,
flatIndexMaps,
objectIdByNameSingular,
featureFlagsMap: cachedFeatureFlagMap,
eventEmitterService: this.workspaceEventEmitter,
},
{
url:
dataSourceMetadata.url ??
this.twentyConfigService.get('PG_DATABASE_URL'),
type: 'postgres',
logging: this.twentyConfigService.getLoggingConfig(),
schema: dataSourceMetadata.schema,
entities: cachedEntitySchemas,
ssl: this.twentyConfigService.get('PG_SSL_ALLOW_SELF_SIGNED')
? {
rejectUnauthorized: false,
}
: undefined,
extra: {
query_timeout: 10000,
// https://node-postgres.com/apis/pool
// TypeORM doesn't allow sharing connection pools between data sources
// So we keep a small pool open for longer if connection pooling patch isn't enabled
// TODO: Probably not needed anymore when connection pooling patch is enabled
idleTimeoutMillis: TWENTY_MINUTES_IN_MS,
max: 4,
allowExitOnIdle: true,
},
},
cachedFeatureFlagMapVersion,
cachedFeatureFlagMap,
cachedRolesPermissionsVersion,
cachedRolesPermissions,
this.twentyConfigService.get('PG_ENABLE_POOL_SHARING'),
);
await workspaceDataSource.initialize();
return workspaceDataSource;
},
this.safelyDestroyDataSource.bind(this),
);
if (!workspaceDataSource) {
throw new Error(`Failed to create WorkspaceDataSource for ${cacheKey}`);
}
await this.updateWorkspaceDataSourceRolesPermissionsIfNeeded({
workspaceDataSource,
cachedRolesPermissionsVersion,
cachedRolesPermissions,
});
await this.updateWorkspaceDataSourceFeatureFlagsMapIfNeeded({
workspaceDataSource,
cachedFeatureFlagMapVersion,
cachedFeatureFlagMap,
});
return workspaceDataSource;
}
private updateWorkspaceDataSourceIfNeeded<T>({
workspaceDataSource,
currentVersion,
newVersion,
newData,
setData,
setVersion,
}: {
workspaceDataSource: WorkspaceDataSource;
currentVersion: string | undefined;
newVersion: string | undefined;
newData: T | undefined;
setData: (data: T) => void;
setVersion: (version: string) => void;
}): void {
if (
isDefined(newVersion) &&
isDefined(newData) &&
currentVersion !== newVersion
) {
workspaceDataSource.manager.repositories.clear();
setData(newData);
setVersion(newVersion);
}
}
private async updateWorkspaceDataSourceRolesPermissionsIfNeeded({
workspaceDataSource,
cachedRolesPermissionsVersion,
cachedRolesPermissions,
}: {
workspaceDataSource: WorkspaceDataSource;
cachedRolesPermissionsVersion: string;
cachedRolesPermissions: ObjectsPermissionsByRoleId;
}): Promise<void> {
this.updateWorkspaceDataSourceIfNeeded({
workspaceDataSource,
currentVersion: workspaceDataSource.rolesPermissionsVersion,
newVersion: cachedRolesPermissionsVersion,
newData: cachedRolesPermissions,
setData: (data) => workspaceDataSource.setRolesPermissions(data),
setVersion: (version) =>
workspaceDataSource.setRolesPermissionsVersion(version),
});
}
private async updateWorkspaceDataSourceFeatureFlagsMapIfNeeded({
workspaceDataSource,
cachedFeatureFlagMapVersion,
cachedFeatureFlagMap,
}: {
workspaceDataSource: WorkspaceDataSource;
cachedFeatureFlagMapVersion: string | undefined;
cachedFeatureFlagMap: FeatureFlagMap | undefined;
}): Promise<void> {
this.updateWorkspaceDataSourceIfNeeded({
workspaceDataSource,
currentVersion: workspaceDataSource.featureFlagMapVersion,
newVersion: cachedFeatureFlagMapVersion,
newData: cachedFeatureFlagMap,
setData: (data) => workspaceDataSource.setFeatureFlagMap(data),
setVersion: (version) =>
workspaceDataSource.setFeatureFlagMapVersion(version),
});
}
private async getWorkspaceMetadataVersionFromCacheOrFromDB(
workspaceId: string,
): Promise<number> {
const latestWorkspaceMetadataVersion =
await this.workspaceCacheStorageService.getMetadataVersion(workspaceId);
if (isDefined(latestWorkspaceMetadataVersion)) {
return latestWorkspaceMetadataVersion;
}
const workspace = await this.workspaceRepository.findOne({
where: { id: workspaceId },
});
if (!workspace) {
throw new TwentyORMException(
`Workspace not found for workspace ${workspaceId}`,
TwentyORMExceptionCode.WORKSPACE_NOT_FOUND,
);
}
await this.workspaceCacheStorageService.setMetadataVersion(
workspaceId,
workspace.metadataVersion,
);
return workspace.metadataVersion;
}
public async destroy(workspaceId: string) {
try {
await this.promiseMemoizer.clearKeys(
`${workspaceId}-`,
this.safelyDestroyDataSource.bind(this),
);
} catch (error) {
// Log and swallow any errors during cleanup to prevent crashes
this.logger.warn(
`Error cleaning up datasources for workspace ${workspaceId}: ${error.message}`,
);
}
}
}
@@ -33,11 +33,15 @@ export class GlobalWorkspaceDataSourceService
rejectUnauthorized: false,
}
: undefined,
poolSize: this.twentyConfigService.get('PG_POOL_MAX_CONNECTIONS'),
extra: {
query_timeout: 10000, // 10 seconds,
idleTimeoutMillis: 120_000, // 2 minutes,
max: 4,
allowExitOnIdle: true,
idleTimeoutMillis: this.twentyConfigService.get(
'PG_POOL_IDLE_TIMEOUT_MS',
),
allowExitOnIdle: this.twentyConfigService.get(
'PG_POOL_ALLOW_EXIT_ON_IDLE',
),
},
},
this.workspaceEventEmitter,
@@ -15,7 +15,6 @@ import { EntityMetadataNotFoundError } from 'typeorm/error/EntityMetadataNotFoun
import { type WorkspaceAuthContext } from 'src/engine/api/common/interfaces/workspace-auth-context.interface';
import { type FeatureFlagMap } from 'src/engine/core-modules/feature-flag/interfaces/feature-flag-map.interface';
import { type WorkspaceDataSourceInterface } from 'src/engine/twenty-orm/interfaces/workspace-datasource.interface';
import {
PermissionsException,
@@ -32,11 +31,7 @@ type CreateQueryBuilderOptions = {
calledByWorkspaceEntityManager?: boolean;
};
export class GlobalWorkspaceDataSource
extends DataSource
implements WorkspaceDataSourceInterface
{
readonly isGlobalFlow = true;
export class GlobalWorkspaceDataSource extends DataSource {
readonly eventEmitterService: WorkspaceEventEmitter;
private _isConstructing = true;
dataSourceWithOverridenCreateQueryBuilder: GlobalWorkspaceDataSource;
@@ -107,7 +102,7 @@ export class GlobalWorkspaceDataSource
return super.createEntityManager(queryRunner) as WorkspaceEntityManager;
}
return new WorkspaceEntityManager(undefined, this, queryRunner);
return new WorkspaceEntityManager(this, queryRunner);
}
override createQueryRunner(
@@ -3,17 +3,15 @@ import { Injectable, type Type } from '@nestjs/common';
import { type ObjectLiteral } from 'typeorm';
import { type WorkspaceAuthContext } from 'src/engine/api/common/interfaces/workspace-auth-context.interface';
import type { WorkspaceDataSourceInterface } from 'src/engine/twenty-orm/interfaces/workspace-datasource.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';
import type { WorkspaceRepository } from 'src/engine/twenty-orm/repository/workspace.repository';
import {
type ORMWorkspaceContext,
withWorkspaceContext,
} from 'src/engine/twenty-orm/storage/orm-workspace-context.storage';
import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager';
import type { RolePermissionConfig } from 'src/engine/twenty-orm/types/role-permission-config';
import { WorkspaceCacheService } from 'src/engine/workspace-cache/services/workspace-cache.service';
import { convertClassNameToObjectMetadataName } from 'src/engine/workspace-manager/workspace-sync-metadata/utils/convert-class-to-object-metadata-name.util';
@@ -23,21 +21,8 @@ export class GlobalWorkspaceOrmManager {
constructor(
private readonly globalWorkspaceDataSourceService: GlobalWorkspaceDataSourceService,
private readonly workspaceCacheService: WorkspaceCacheService,
private readonly twentyORMGlobalManager: TwentyORMGlobalManager,
) {}
async isGlobalDataSourceFlow(workspaceId: string): Promise<boolean> {
const { featureFlagsMap } = await this.workspaceCacheService.getOrRecompute(
workspaceId,
['featureFlagsMap'],
);
return (
featureFlagsMap[FeatureFlagKey.IS_GLOBAL_WORKSPACE_DATASOURCE_ENABLED] ===
true
);
}
async getRepository<T extends ObjectLiteral>(
workspaceId: string,
workspaceEntity: Type<T>,
@@ -51,62 +36,29 @@ export class GlobalWorkspaceOrmManager {
): Promise<WorkspaceRepository<T>>;
async getRepository<T extends ObjectLiteral>(
workspaceId: string,
_workspaceId: string,
workspaceEntityOrObjectMetadataName: Type<T> | string,
permissionOptions?: RolePermissionConfig,
): Promise<WorkspaceRepository<T>> {
const isGlobalFlow = await this.isGlobalDataSourceFlow(workspaceId);
if (isGlobalFlow) {
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,
);
}
let objectMetadataName: string;
if (typeof workspaceEntityOrObjectMetadataName === 'string') {
return this.twentyORMGlobalManager.getRepositoryForWorkspace<T>(
workspaceId,
workspaceEntityOrObjectMetadataName,
permissionOptions,
objectMetadataName = workspaceEntityOrObjectMetadataName;
} else {
objectMetadataName = convertClassNameToObjectMetadataName(
workspaceEntityOrObjectMetadataName.name,
);
}
return this.twentyORMGlobalManager.getRepositoryForWorkspace<T>(
workspaceId,
workspaceEntityOrObjectMetadataName,
const globalDataSource = await this.getGlobalWorkspaceDataSource();
return globalDataSource.getRepository<T>(
objectMetadataName,
permissionOptions,
);
}
async getDataSourceForWorkspace(
workspaceId: string,
): Promise<WorkspaceDataSourceInterface> {
const isGlobalFlow = await this.isGlobalDataSourceFlow(workspaceId);
if (isGlobalFlow) {
return this.globalWorkspaceDataSourceService.getGlobalWorkspaceDataSource();
}
return this.twentyORMGlobalManager.getDataSourceForWorkspace({
workspaceId,
});
}
async getGlobalWorkspaceDataSource() {
async getGlobalWorkspaceDataSource(): Promise<GlobalWorkspaceDataSource> {
return this.globalWorkspaceDataSourceService.getGlobalWorkspaceDataSource();
}
@@ -119,18 +71,6 @@ export class GlobalWorkspaceOrmManager {
return withWorkspaceContext(context, fn);
}
async destroyDataSourceForWorkspace(workspaceId: string): Promise<void> {
const isGlobalFlow = await this.isGlobalDataSourceFlow(workspaceId);
if (isGlobalFlow) {
return;
}
await this.twentyORMGlobalManager.destroyDataSourceForWorkspace(
workspaceId,
);
}
private async loadWorkspaceContext(
authContext: WorkspaceAuthContext,
): Promise<ORMWorkspaceContext> {
@@ -1,49 +0,0 @@
import { type ObjectsPermissionsByRoleId } from 'twenty-shared/types';
import {
type EntityManager,
type EntityTarget,
type ObjectLiteral,
type QueryRunner,
type ReplicationMode,
} from 'typeorm';
import { type FeatureFlagMap } from 'src/engine/core-modules/feature-flag/interfaces/feature-flag-map.interface';
import { type AuthContext } from 'src/engine/core-modules/auth/types/auth-context.type';
import { type WorkspaceEntityManager } from 'src/engine/twenty-orm/entity-manager/workspace-entity-manager';
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 { type RolePermissionConfig } from 'src/engine/twenty-orm/types/role-permission-config';
export interface WorkspaceDataSourceInterface {
readonly isGlobalFlow: boolean;
readonly manager: EntityManager;
featureFlagMap: FeatureFlagMap;
permissionsPerRoleId: ObjectsPermissionsByRoleId;
getRepository<Entity extends ObjectLiteral>(
target: EntityTarget<Entity>,
permissionOptions?: RolePermissionConfig,
authContext?: AuthContext,
): WorkspaceRepository<Entity>;
createQueryRunner(mode?: ReplicationMode): WorkspaceQueryRunner;
createEntityManager(queryRunner?: QueryRunner): WorkspaceEntityManager;
createQueryRunnerForEntityPersistExecutor(
mode?: ReplicationMode,
): QueryRunner;
transaction<T>(
runInTransaction: (entityManager: EntityManager) => Promise<T>,
): Promise<T>;
query<T = unknown>(
query: string,
parameters?: unknown[],
queryRunner?: QueryRunner,
options?: { shouldBypassPermissionChecks?: boolean },
): Promise<T>;
}
@@ -1,317 +0,0 @@
import { type Logger, type LogLevel } from '@nestjs/common';
import { Test, type TestingModule } from '@nestjs/testing';
import { Pool } from 'pg';
import { TwentyConfigService } from 'src/engine/core-modules/twenty-config/twenty-config.service';
import { PgPoolSharedService } from 'src/engine/twenty-orm/pg-shared-pool/pg-shared-pool.service';
type ConfigKey =
| 'PG_ENABLE_POOL_SHARING'
| 'PG_POOL_MAX_CONNECTIONS'
| 'LOG_LEVELS';
type ConfigValue = boolean | number | LogLevel[] | string;
interface PoolWithEndTracker extends Pool {
__hasEnded?: boolean;
}
jest.mock('pg', () => {
const mockPool = jest.fn().mockImplementation(() => ({
on: jest.fn(),
end: jest.fn().mockImplementation((callback) => {
if (callback) callback();
return Promise.resolve();
}),
_clients: [],
_idle: [],
_pendingQueue: { length: 0 },
}));
mockPool.prototype = {
on: jest.fn(),
end: jest.fn(),
};
return {
Pool: mockPool,
};
});
describe('PgPoolSharedService', () => {
let service: PgPoolSharedService;
let configService: TwentyConfigService;
let mockLogger: Partial<Logger>;
const configValues: Record<ConfigKey, ConfigValue> = {
PG_ENABLE_POOL_SHARING: true,
PG_POOL_MAX_CONNECTIONS: 10,
LOG_LEVELS: ['error', 'warn'],
};
beforeEach(async () => {
mockLogger = {
log: jest.fn(),
debug: jest.fn(),
warn: jest.fn(),
error: jest.fn(),
};
const module: TestingModule = await Test.createTestingModule({
providers: [
PgPoolSharedService,
{
provide: TwentyConfigService,
useValue: {
get: jest
.fn()
.mockImplementation(
(key: string) => configValues[key as ConfigKey],
),
},
},
],
}).compile();
service = module.get<PgPoolSharedService>(PgPoolSharedService);
configService = module.get<TwentyConfigService>(TwentyConfigService);
Object.defineProperty(service, 'logger', {
value: mockLogger,
writable: true,
});
});
afterEach(() => {
if (service && typeof service.unpatchForTesting === 'function') {
service.unpatchForTesting();
}
jest.clearAllMocks();
});
it('should be defined', () => {
expect(service).toBeDefined();
});
describe('initialize', () => {
it('should initialize and patch pg Pool when enabled', () => {
service.initialize();
expect(mockLogger.log).toHaveBeenCalledWith(
expect.stringContaining(
'Pool sharing will use max 10 connections per pool',
),
);
expect(mockLogger.log).toHaveBeenCalledWith(
expect.stringContaining('Pg pool sharing initialized'),
);
});
it('should not initialize when pool sharing is disabled', () => {
jest.spyOn(configService, 'get').mockImplementation((key: string) => {
if (key === 'PG_ENABLE_POOL_SHARING') return false;
return configValues[key as ConfigKey];
});
service.initialize();
expect(mockLogger.log).toHaveBeenCalledWith(
'Pg pool sharing is disabled by configuration',
);
expect(mockLogger.log).not.toHaveBeenCalledWith(
expect.stringContaining('Pg pool sharing initialized'),
);
});
it('should not initialize twice', () => {
service.initialize();
jest.clearAllMocks();
service.initialize();
expect(mockLogger.debug).toHaveBeenCalledWith(
'Pg pool sharing already initialized, skipping',
);
});
});
describe('pool sharing functionality', () => {
beforeEach(() => {
service.initialize();
});
afterEach(() => {
jest.clearAllMocks();
});
it('should reuse pools with identical connection parameters', () => {
const pool1 = new Pool({
host: 'localhost',
port: 5432,
database: 'testdb',
user: 'testuser',
});
const pool2 = new Pool({
host: 'localhost',
port: 5432,
database: 'testdb',
user: 'testuser',
});
expect(pool1).toBe(pool2);
const poolsMap = service.getPoolsMapForTesting();
expect(poolsMap?.size).toBe(1);
});
it('should create separate pools for different connection parameters', () => {
const pool1 = new Pool({
host: 'localhost',
port: 5432,
database: 'db1',
user: 'user1',
});
const pool2 = new Pool({
host: 'localhost',
port: 5432,
database: 'db2',
user: 'user1',
});
expect(pool1).not.toBe(pool2);
const poolsMap = service.getPoolsMapForTesting();
expect(poolsMap?.size).toBe(2);
});
it('should remove pools from cache when they are ended', async () => {
const pool = new Pool({
host: 'localhost',
database: 'testdb',
});
const poolsMapBefore = service.getPoolsMapForTesting();
expect(poolsMapBefore?.size).toBe(1);
await pool.end();
const poolsMapAfter = service.getPoolsMapForTesting();
expect(poolsMapAfter?.size).toBe(0);
expect(mockLogger.log).toHaveBeenCalledWith(
expect.stringContaining('pg Pool for key'),
);
expect(mockLogger.log).toHaveBeenCalledWith(
expect.stringContaining('has been closed'),
);
});
it('should handle calling end() multiple times on the same pool', async () => {
const pool = new Pool({
host: 'localhost',
database: 'testdb',
});
await pool.end();
expect(service.getPoolsMapForTesting()?.size).toBe(0);
expect((pool as PoolWithEndTracker).__hasEnded).toBe(true);
await pool.end();
expect(mockLogger.debug).toHaveBeenCalledWith(
expect.stringContaining('Ignoring duplicate end() call'),
);
const closeMessageCalls = (mockLogger.log as jest.Mock).mock.calls.filter(
(call: any[]) => call[0].includes('has been closed'),
);
expect(closeMessageCalls.length).toBe(1);
expect(mockLogger.log).toHaveBeenCalledWith(
expect.stringContaining('has been closed'),
);
});
});
describe('debug logging', () => {
it('should enable debug logging when debug log level is set', async () => {
jest
.spyOn(configService, 'get')
.mockImplementation((key: string): ConfigValue => {
if (key === 'LOG_LEVELS') return ['debug', 'error', 'warn'];
return configValues[key as ConfigKey];
});
const module = await Test.createTestingModule({
providers: [
PgPoolSharedService,
{
provide: TwentyConfigService,
useValue: configService,
},
],
}).compile();
const debugService = module.get<PgPoolSharedService>(PgPoolSharedService);
Object.defineProperty(debugService, 'logger', {
value: mockLogger,
writable: true,
});
const spyInterval = jest.spyOn(global, 'setInterval');
debugService.initialize();
expect(spyInterval).toHaveBeenCalled();
expect(mockLogger.debug).toHaveBeenCalledWith(
expect.stringContaining('Pool statistics logging enabled'),
);
});
});
describe('logPoolStats', () => {
it('should log pool statistics correctly', () => {
service.initialize();
new Pool({ host: 'localhost', database: 'testdb' });
jest.clearAllMocks();
service.logPoolStats();
expect(mockLogger.debug).toHaveBeenCalledWith(
'=== PostgreSQL Connection Pool Stats ===',
);
expect(mockLogger.debug).toHaveBeenCalledWith(
expect.stringContaining('Total pools: 1'),
);
});
});
describe('onApplicationShutdown', () => {
it('should clear interval on shutdown', () => {
const clearIntervalSpy = jest.spyOn(global, 'clearInterval');
service['logStatsInterval'] = setInterval(() => {}, 1000);
service.onApplicationShutdown();
expect(clearIntervalSpy).toHaveBeenCalled();
expect(service['logStatsInterval']).toBeNull();
});
});
});
@@ -1,39 +0,0 @@
import {
Global,
Logger,
Module,
type OnApplicationShutdown,
type OnModuleInit,
} from '@nestjs/common';
import { TwentyConfigModule } from 'src/engine/core-modules/twenty-config/twenty-config.module';
import { PgPoolSharedService } from 'src/engine/twenty-orm/pg-shared-pool/pg-shared-pool.service';
/**
* Module that initializes the shared pg pool at application bootstrap
*/
@Global()
@Module({
imports: [TwentyConfigModule],
providers: [PgPoolSharedService],
exports: [PgPoolSharedService],
})
export class PgPoolSharedModule implements OnModuleInit, OnApplicationShutdown {
constructor(private readonly pgPoolSharedService: PgPoolSharedService) {}
private readonly logger = new Logger(PgPoolSharedModule.name);
/**
* Initialize the pool sharing service when the module is initialized
*/
async onModuleInit() {
await this.pgPoolSharedService.initialize();
}
/**
* Clean up any resources when the application shuts down
*/
async onApplicationShutdown() {
this.logger.log('Shutting down PgPoolSharedModule');
await this.pgPoolSharedService.onApplicationShutdown();
}
}
@@ -1,560 +0,0 @@
import { Injectable, Logger } from '@nestjs/common';
import pg, { type Pool, type PoolConfig } from 'pg';
import { isDefined } from 'twenty-shared/utils';
import { TwentyConfigService } from 'src/engine/core-modules/twenty-config/twenty-config.service';
interface PgWithPatchSymbol {
Pool: typeof Pool;
[key: symbol]: boolean;
}
interface SSLConfig {
rejectUnauthorized?: boolean;
[key: string]: unknown;
}
interface PoolWithEndTracker extends Pool {
__hasEnded?: boolean;
}
interface ExtendedPoolConfig extends PoolConfig {
allowExitOnIdle?: boolean;
}
interface PoolInternalStats {
_clients?: Array<unknown>;
_idle?: Array<unknown>;
_pendingQueue?: {
length: number;
};
}
/**
* Service that manages shared pg connection pools across tenants.
* It patches the pg.Pool constructor to return cached instances for
* identical connection parameters.
*/
@Injectable()
export class PgPoolSharedService {
private readonly logger = new Logger('PgPoolSharedService');
private initialized = false;
private readonly PATCH_SYMBOL = Symbol.for('@@twenty/pg-shared-pool');
private isDebugEnabled = false;
private logStatsInterval: NodeJS.Timeout | null = null;
// Internal pool cache - exposed for testing only
private poolsMap = new Map<string, Pool>();
private capturedOriginalPoolValue: typeof Pool | null = null;
private originalPgPoolDescriptor: PropertyDescriptor | undefined = undefined;
constructor(private readonly configService: TwentyConfigService) {
this.detectDebugMode();
}
private detectDebugMode(): void {
const logLevels = this.configService.get('LOG_LEVELS');
this.isDebugEnabled =
Array.isArray(logLevels) &&
logLevels.some((level) => level === 'debug' || level === 'verbose');
}
/**
* Provides access to the internal pools map for testing purposes
*/
getPoolsMapForTesting(): Map<string, Pool> | null {
if (!this.initialized) {
return null;
}
return this.poolsMap;
}
/**
* Applies the pg.Pool patch to enable connection pool sharing.
* Safe to call multiple times - will only apply the patch once.
*/
async initialize(): Promise<void> {
if (this.initialized) {
this.logger.debug('Pg pool sharing already initialized, skipping');
return;
}
if (!this.isPoolSharingEnabled()) {
this.logger.log('Pg pool sharing is disabled by configuration');
return;
}
this.logPoolConfiguration();
this.patchPgPool();
if (!this.isPatchSuccessful()) {
this.logger.error(
'Failed to patch pg.Pool. PgPoolSharedService will not be active.',
);
this.initialized = false;
return;
}
this.initialized = true;
this.logger.log(
'Pg pool sharing initialized - pools will be shared across tenants',
);
if (this.isDebugEnabled) {
this.startPoolStatsLogging();
}
}
private isPoolSharingEnabled(): boolean {
return !!this.configService.get('PG_ENABLE_POOL_SHARING');
}
private logPoolConfiguration(): void {
const maxConnections = this.configService.get('PG_POOL_MAX_CONNECTIONS');
const idleTimeoutMs = this.configService.get('PG_POOL_IDLE_TIMEOUT_MS');
const allowExitOnIdle = this.configService.get(
'PG_POOL_ALLOW_EXIT_ON_IDLE',
);
this.logger.log(
`Pool sharing will use max ${maxConnections} connections per pool with ${idleTimeoutMs}ms idle timeout and allowExitOnIdle=${allowExitOnIdle}`,
);
}
private isPatchSuccessful(): boolean {
const pgWithSymbol = pg as PgWithPatchSymbol;
return !!pgWithSymbol[this.PATCH_SYMBOL];
}
/**
* Stops the periodic logging of pool statistics.
* Call this during application shutdown.
*/
async onApplicationShutdown(): Promise<void> {
this.stopStatsLogging();
this.logger.log('onApplicationShutdown called in PgPoolSharedService');
await this.closeAllPools();
}
private stopStatsLogging(): void {
if (!this.logStatsInterval) {
return;
}
clearInterval(this.logStatsInterval);
this.logStatsInterval = null;
}
private async closeAllPools(): Promise<void> {
if (this.poolsMap.size === 0) {
return;
}
const closePromises: Promise<void>[] = [];
for (const [key, pool] of this.poolsMap.entries()) {
closePromises.push(
pool
.end()
.catch((err) => {
if (err?.message !== 'Called end on pool more than once') {
this.logger.debug(
`Pool[${key}] error during shutdown: ${err.message}`,
);
}
})
.then(() => {
this.logger.debug(
`Pool[${key}] closed during application shutdown`,
);
}),
);
}
this.logger.debug('Attempting to close all pg pools...');
await Promise.allSettled(closePromises);
this.logger.debug('All pg pools closure attempts completed.');
}
/**
* Logs detailed statistics about all connection pools
*/
logPoolStats(): void {
if (!this.initialized || this.poolsMap.size === 0) {
this.logger.debug('No active pg pools to log stats for');
return;
}
let totalActive = 0;
let totalIdle = 0;
let totalPoolSize = 0;
let totalQueueSize = 0;
this.logger.debug('=== PostgreSQL Connection Pool Stats ===');
for (const [key, pool] of this.poolsMap.entries()) {
const stats = this.collectPoolStats(key, pool);
totalActive += stats.active;
totalIdle += stats.idle;
totalPoolSize += stats.poolSize;
totalQueueSize += stats.queueSize;
}
this.logTotalStats(totalActive, totalIdle, totalPoolSize, totalQueueSize);
}
private collectPoolStats(key: string, pool: Pool) {
const poolStats = pool as PoolInternalStats;
const active =
(poolStats._clients?.length || 0) - (poolStats._idle?.length || 0);
const idle = poolStats._idle?.length || 0;
const poolSize = poolStats._clients?.length || 0;
const queueSize = poolStats._pendingQueue?.length || 0;
this.logger.debug(
`Pool [${key}]: active=${active}, idle=${idle}, size=${poolSize}, queue=${queueSize}`,
);
return { active, idle, poolSize, queueSize };
}
private logTotalStats(
totalActive: number,
totalIdle: number,
totalPoolSize: number,
totalQueueSize: number,
): void {
this.logger.debug(
`Total pools: ${this.poolsMap.size}, active=${totalActive}, idle=${totalIdle}, ` +
`total connections=${totalPoolSize}, queued requests=${totalQueueSize}`,
);
this.logger.debug('=========================================');
}
/**
* Starts periodically logging pool statistics if debug is enabled
*/
private startPoolStatsLogging(): void {
this.logPoolStats();
this.logStatsInterval = setInterval(() => {
this.logPoolStats();
}, 30000);
this.logger.debug('Pool statistics logging enabled (30s interval)');
}
/**
* Patches the pg module's Pool constructor to provide shared instances
* across all tenant workspaces.
*/
private patchPgPool(): void {
const pgWithSymbol = pg as PgWithPatchSymbol;
if (this.isAlreadyPatched(pgWithSymbol)) {
return;
}
if (!this.captureOriginalPool(pgWithSymbol)) {
return;
}
this.applyPatchWithSharedPool(pgWithSymbol);
}
private isAlreadyPatched(pgWithSymbol: PgWithPatchSymbol): boolean {
if (pgWithSymbol[this.PATCH_SYMBOL]) {
this.logger.debug(
'pg.Pool is already patched. Skipping patch operation for this instance.',
);
return true;
}
return false;
}
private captureOriginalPool(pgWithSymbol: PgWithPatchSymbol): boolean {
this.originalPgPoolDescriptor = Object.getOwnPropertyDescriptor(
pgWithSymbol,
'Pool',
);
if (
!this.originalPgPoolDescriptor ||
typeof this.originalPgPoolDescriptor.value !== 'function'
) {
this.logger.error(
'Could not get original pg.Pool constructor or descriptor is invalid. Aborting patch.',
);
return false;
}
this.capturedOriginalPoolValue = this.originalPgPoolDescriptor
.value as typeof Pool;
return true;
}
private applyPatchWithSharedPool(pgWithSymbol: PgWithPatchSymbol): void {
const OriginalPool = this.capturedOriginalPoolValue as typeof Pool;
const SharedPool = this.createSharedPoolConstructor(OriginalPool);
// Preserve prototype chain for instanceof checks
SharedPool.prototype = Object.create(OriginalPool.prototype);
SharedPool.prototype.constructor = SharedPool;
// Replace the original Pool with our patched version
Object.defineProperty(pgWithSymbol, 'Pool', {
value: SharedPool as unknown as typeof Pool,
writable: true,
configurable: true,
enumerable: this.originalPgPoolDescriptor?.enumerable,
});
pgWithSymbol[this.PATCH_SYMBOL] = true;
this.logger.log('pg.Pool patched successfully by this service instance.');
}
private createSharedPoolConstructor(OriginalPool: typeof Pool) {
const maxConnections = this.configService.get('PG_POOL_MAX_CONNECTIONS');
const idleTimeoutMs = this.configService.get('PG_POOL_IDLE_TIMEOUT_MS');
const allowExitOnIdle = this.configService.get(
'PG_POOL_ALLOW_EXIT_ON_IDLE',
);
// Store references to service functions/properties that we need in our constructor
const buildPoolKey = this.buildPoolKey.bind(this);
const poolsMap = this.poolsMap;
const logger = this.logger;
const isDebugEnabled = this.isDebugEnabled;
const setupPoolEvents = this.setupPoolEvents.bind(this);
const replacePoolEndMethod = this.replacePoolEndMethod.bind(this);
// Define a proper constructor function that can be used with "new"
function SharedPool(this: Pool, config?: PoolConfig): Pool {
// When called as a function (without new), make sure to return a new instance
if (!(this instanceof SharedPool)) {
// @ts-expect-error We know this works at runtime
return new SharedPool(config);
}
const poolConfig = config
? ({ ...config } as ExtendedPoolConfig)
: ({} as ExtendedPoolConfig);
if (maxConnections) {
poolConfig.max = maxConnections;
}
if (idleTimeoutMs) {
poolConfig.idleTimeoutMillis = idleTimeoutMs;
}
if (allowExitOnIdle) {
poolConfig.allowExitOnIdle = allowExitOnIdle;
}
const key = buildPoolKey(poolConfig);
const existing = poolsMap.get(key);
if (existing) {
if (isDebugEnabled) {
logger.debug(`Reusing existing pg Pool for key "${key}"`);
}
return existing;
}
const pool = new OriginalPool(poolConfig);
poolsMap.set(key, pool);
logger.log(
`Created new shared pg Pool for key "${key}" with ${poolConfig.max ?? 'default'} max connections and ${poolConfig.idleTimeoutMillis ?? 'default'} ms idle timeout. Total pools: ${poolsMap.size}`,
);
if (isDebugEnabled) {
setupPoolEvents(pool, key);
}
replacePoolEndMethod(pool, key);
return pool;
}
return SharedPool;
}
private setupPoolEvents(pool: Pool, key: string): void {
pool.on('connect', () => {
this.logger.debug(`Pool[${key}]: New connection established`);
});
pool.on('acquire', () => {
this.logger.debug(`Pool[${key}]: Client acquired from pool`);
});
pool.on('remove', () => {
this.logger.debug(`Pool[${key}]: Connection removed from pool`);
});
pool.on('error', (err) => {
this.logger.warn(`Pool[${key}]: Connection error: ${err.message}`);
});
}
private replacePoolEndMethod(pool: Pool, key: string): void {
const originalEnd = pool.end.bind(pool) as {
(callback?: (err?: Error) => void): void;
};
(pool as PoolWithEndTracker).end = (
callback?: (err?: Error) => void,
): Promise<void> => {
if ((pool as PoolWithEndTracker).__hasEnded) {
if (callback) {
callback();
}
this.logger.debug(`Ignoring duplicate end() call for pool "${key}"`);
return Promise.resolve();
}
// Mark this pool as ended to prevent subsequent calls
(pool as PoolWithEndTracker).__hasEnded = true;
this.poolsMap.delete(key);
this.logger.log(
`pg Pool for key "${key}" has been closed. Remaining pools: ${this.poolsMap.size}`,
);
return new Promise<void>((resolve, reject) => {
originalEnd((err) => {
if (err) {
// If error is about duplicate end, suppress it
if (err.message === 'Called end on pool more than once') {
if (callback) callback();
resolve();
return;
}
if (callback) callback(err);
reject(err);
return;
}
if (callback) callback();
resolve();
});
});
};
}
/**
* Builds a unique key for a pool configuration to identify identical connections
*/
private buildPoolKey(config: PoolConfig = {}): string {
// We identify pools only by parameters that open a *physical* connection.
// `search_path`/schema is not included because it is changed at session level.
const {
host = 'localhost',
port = 5432,
user = 'postgres',
database = '',
ssl,
} = config;
// Note: SSL object can contain certificates, so only stringify relevant
// properties that influence connection reuse.
const sslKey = isDefined(ssl)
? typeof ssl === 'object'
? JSON.stringify({
rejectUnauthorized: (ssl as SSLConfig).rejectUnauthorized,
})
: String(ssl)
: 'no-ssl';
return [host, port, user, database, sslKey].join('|');
}
/**
* Resets the pg.Pool patch and clears service state. For testing purposes only.
*/
public unpatchForTesting(): void {
this.logger.debug('Attempting to unpatch pg.Pool for testing...');
const pgWithSymbol = pg as PgWithPatchSymbol;
if (!pgWithSymbol[this.PATCH_SYMBOL]) {
this.logger.debug(
'pg.Pool was not patched by this instance or PATCH_SYMBOL not found, no unpatch needed from this instance.',
);
this.resetStateForTesting();
return;
}
this.restoreOriginalPool(pgWithSymbol);
delete pgWithSymbol[this.PATCH_SYMBOL];
this.logger.debug('PATCH_SYMBOL removed from pg module.');
this.resetStateForTesting();
}
private restoreOriginalPool(pgWithSymbol: PgWithPatchSymbol): void {
if (this.originalPgPoolDescriptor) {
Object.defineProperty(
pgWithSymbol,
'Pool',
this.originalPgPoolDescriptor,
);
this.logger.debug('pg.Pool restored using original property descriptor.');
return;
}
if (this.capturedOriginalPoolValue) {
// Fallback if descriptor wasn't captured
pgWithSymbol.Pool = this.capturedOriginalPoolValue;
this.logger.warn(
'pg.Pool restored using captured value (descriptor method preferred).',
);
return;
}
this.logger.error(
'Cannot unpatch pg.Pool: no original Pool reference or descriptor was captured by this instance.',
);
}
private resetStateForTesting(): void {
this.initialized = false;
this.poolsMap.clear();
this.capturedOriginalPoolValue = null;
this.originalPgPoolDescriptor = undefined;
if (this.logStatsInterval) {
clearInterval(this.logStatsInterval);
this.logStatsInterval = null;
}
this.logger.debug(
'Service instance state (initialized, poolsMap, captured originals, timers) reset for testing.',
);
}
}
@@ -1,61 +0,0 @@
import { Injectable, type Type } from '@nestjs/common';
import { type ObjectLiteral } from 'typeorm';
import { WorkspaceDatasourceFactory } from 'src/engine/twenty-orm/factories/workspace-datasource.factory';
import { type WorkspaceRepository } from 'src/engine/twenty-orm/repository/workspace.repository';
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 TwentyORMGlobalManager {
constructor(
private readonly workspaceDataSourceFactory: WorkspaceDatasourceFactory,
) {}
async getRepositoryForWorkspace<T extends ObjectLiteral>(
workspaceId: string,
workspaceEntity: Type<T>,
options?: RolePermissionConfig,
): Promise<WorkspaceRepository<T>>;
async getRepositoryForWorkspace<T extends ObjectLiteral>(
workspaceId: string,
objectMetadataName: string,
options?: RolePermissionConfig,
): Promise<WorkspaceRepository<T>>;
async getRepositoryForWorkspace<T extends ObjectLiteral>(
workspaceId: string,
workspaceEntityOrObjectMetadataName: Type<T> | string,
options?: RolePermissionConfig,
): Promise<WorkspaceRepository<T>> {
let objectMetadataName: string;
if (typeof workspaceEntityOrObjectMetadataName === 'string') {
objectMetadataName = workspaceEntityOrObjectMetadataName;
} else {
objectMetadataName = convertClassNameToObjectMetadataName(
workspaceEntityOrObjectMetadataName.name,
);
}
const workspaceDataSource =
await this.workspaceDataSourceFactory.create(workspaceId);
const repository = workspaceDataSource.getRepository<T>(
objectMetadataName,
options,
);
return repository;
}
async getDataSourceForWorkspace({ workspaceId }: { workspaceId: string }) {
return await this.workspaceDataSourceFactory.create(workspaceId);
}
async destroyDataSourceForWorkspace(workspaceId: string) {
await this.workspaceDataSourceFactory.destroy(workspaceId);
}
}
@@ -1,10 +0,0 @@
import { ConfigurableModuleBuilder } from '@nestjs/common';
import { type TwentyORMOptions } from './interfaces/twenty-orm-options.interface';
export const {
ConfigurableModuleClass,
MODULE_OPTIONS_TOKEN,
OPTIONS_TYPE,
ASYNC_OPTIONS_TYPE,
} = new ConfigurableModuleBuilder<TwentyORMOptions>().build();
@@ -12,12 +12,9 @@ import { RoleTargetEntity } from 'src/engine/metadata-modules/role-target/role-t
import { WorkspaceFeatureFlagsMapCacheModule } from 'src/engine/metadata-modules/workspace-feature-flags-map-cache/workspace-feature-flags-map-cache.module';
import { entitySchemaFactories } from 'src/engine/twenty-orm/factories';
import { EntitySchemaFactory } from 'src/engine/twenty-orm/factories/entity-schema.factory';
import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager';
import { WorkspaceCacheStorageModule } from 'src/engine/workspace-cache-storage/workspace-cache-storage.module';
import { WorkspaceCacheModule } from 'src/engine/workspace-cache/workspace-cache.module';
import { PgPoolSharedModule } from './pg-shared-pool/pg-shared-pool.module';
@Global()
@Module({
imports: [
@@ -33,10 +30,9 @@ import { PgPoolSharedModule } from './pg-shared-pool/pg-shared-pool.module';
WorkspaceFeatureFlagsMapCacheModule,
FeatureFlagModule,
TwentyConfigModule,
PgPoolSharedModule,
WorkspaceCacheModule,
],
providers: [...entitySchemaFactories, TwentyORMGlobalManager],
exports: [EntitySchemaFactory, TwentyORMGlobalManager, PgPoolSharedModule],
providers: [...entitySchemaFactories],
exports: [EntitySchemaFactory],
})
export class TwentyORMModule {}
@@ -400,10 +400,6 @@ export class CleanerWorkspaceService {
`Error while processing workspace ${workspace.id} ${workspace.displayName}: ${error}`,
);
}
await this.globalWorkspaceOrmManager.destroyDataSourceForWorkspace(
workspace.id,
);
}
this.logger.log(
`${dryRun ? 'DRY RUN - ' : ''}batchWarnOrCleanSuspendedWorkspaces done!`,