From 083a038f8af603a48ca9fd27bb20ff07f332d2ce Mon Sep 17 00:00:00 2001 From: Weiko Date: Thu, 11 Dec 2025 18:27:49 +0100 Subject: [PATCH] Deprecate workspace datasoure (#16507) --- .../__tests__/upgrade.command-runner.spec.ts | 18 +- .../workspaces-migration.command-runner.ts | 17 +- ...mpty-string-null-in-text-fields.command.ts | 7 +- .../common-base-query-runner.service.ts | 4 +- ...common-create-many-query-runner.service.ts | 4 +- ...mmon-extended-query-runner-context.type.ts | 4 +- .../process-nested-relations-v2.helper.ts | 8 +- .../process-nested-relations.helper.ts | 6 +- .../auth/services/google-apis.service.spec.ts | 4 +- .../auth/services/google-apis.service.ts | 4 +- .../services/microsoft-apis.service.spec.ts | 4 +- .../auth/services/microsoft-apis.service.ts | 4 +- .../datasource/workspace.datasource.ts | 258 -------- .../workspace-entity-manager.spec.ts | 70 ++- .../workspace-entity-manager.ts | 48 +- .../src/engine/twenty-orm/factories/index.ts | 2 - .../factories/workspace-datasource.factory.ts | 353 ----------- .../global-workspace-datasource.service.ts | 10 +- .../global-workspace-datasource.ts | 9 +- .../global-workspace-orm.manager.ts | 84 +-- .../workspace-datasource.interface.ts | 49 -- .../__tests__/pg-shared-pool.service.spec.ts | 317 ---------- .../pg-shared-pool/pg-shared-pool.module.ts | 39 -- .../pg-shared-pool/pg-shared-pool.service.ts | 560 ------------------ .../twenty-orm/twenty-orm-global.manager.ts | 61 -- .../twenty-orm.module-definition.ts | 10 - .../engine/twenty-orm/twenty-orm.module.ts | 8 +- .../services/cleaner.workspace-service.ts | 4 - .../services/calendar-save-events.service.ts | 4 +- .../imap-smtp-caldav-apis.service.spec.ts | 4 +- .../services/imap-smtp-caldav-apis.service.ts | 4 +- .../services/dashboard-duplication.service.ts | 122 ++-- .../messaging-message-cleaner.service.ts | 4 +- ...ssaging-message-list-fetch.service.spec.ts | 7 +- .../messaging-message-list-fetch.service.ts | 4 +- ...essaging-process-folder-actions.service.ts | 4 +- ...ing-process-group-email-actions.service.ts | 4 +- ...d-enqueue-contact-creation.service.spec.ts | 2 +- ...es-and-enqueue-contact-creation.service.ts | 4 +- 39 files changed, 200 insertions(+), 1929 deletions(-) delete mode 100644 packages/twenty-server/src/engine/twenty-orm/datasource/workspace.datasource.ts delete mode 100644 packages/twenty-server/src/engine/twenty-orm/factories/workspace-datasource.factory.ts delete mode 100644 packages/twenty-server/src/engine/twenty-orm/interfaces/workspace-datasource.interface.ts delete mode 100644 packages/twenty-server/src/engine/twenty-orm/pg-shared-pool/__tests__/pg-shared-pool.service.spec.ts delete mode 100644 packages/twenty-server/src/engine/twenty-orm/pg-shared-pool/pg-shared-pool.module.ts delete mode 100644 packages/twenty-server/src/engine/twenty-orm/pg-shared-pool/pg-shared-pool.service.ts delete mode 100644 packages/twenty-server/src/engine/twenty-orm/twenty-orm-global.manager.ts delete mode 100644 packages/twenty-server/src/engine/twenty-orm/twenty-orm.module-definition.ts diff --git a/packages/twenty-server/src/database/commands/command-runners/__tests__/upgrade.command-runner.spec.ts b/packages/twenty-server/src/database/commands/command-runners/__tests__/upgrade.command-runner.spec.ts index f651bec5521..69144983924 100644 --- a/packages/twenty-server/src/database/commands/command-runners/__tests__/upgrade.command-runner.spec.ts +++ b/packages/twenty-server/src/database/commands/command-runners/__tests__/upgrade.command-runner.spec.ts @@ -140,7 +140,6 @@ describe('UpgradeCommandRunner', () => { let upgradeCommandRunner: BasicUpgradeCommandRunner; let workspaceRepository: Repository; let runCoreMigrationsSpy: jest.SpyInstance; - let globalWorkspaceOrmManagerSpy: GlobalWorkspaceOrmManager; type BuildModuleAndSetupSpiesArgs = { numberOfWorkspace?: number; @@ -184,9 +183,6 @@ describe('UpgradeCommandRunner', () => { workspaceRepository = module.get>( getRepositoryToken(WorkspaceEntity), ); - globalWorkspaceOrmManagerSpy = module.get( - GlobalWorkspaceOrmManager, - ); }; it('should ignore and list as succesfull upgrade on workspace with higher version', async () => { @@ -214,10 +210,9 @@ describe('UpgradeCommandRunner', () => { expect(successReport.length).toBe(1); expect(failReport.length).toBe(0); - [ - globalWorkspaceOrmManagerSpy.destroyDataSourceForWorkspace, - upgradeCommandRunner.runOnWorkspace, - ].forEach((fn) => expect(fn).toHaveBeenCalledTimes(1)); + [upgradeCommandRunner.runOnWorkspace].forEach((fn) => + expect(fn).toHaveBeenCalledTimes(1), + ); [workspaceRepository.update].forEach((fn) => expect(fn).not.toHaveBeenCalled(), @@ -239,10 +234,9 @@ describe('UpgradeCommandRunner', () => { // @ts-expect-error legacy noImplicitAny await upgradeCommandRunner.run(passedParams, options); - [ - upgradeCommandRunner.runOnWorkspace, - globalWorkspaceOrmManagerSpy.destroyDataSourceForWorkspace, - ].forEach((fn) => expect(fn).toHaveBeenCalledTimes(numberOfWorkspace)); + [upgradeCommandRunner.runOnWorkspace].forEach((fn) => + expect(fn).toHaveBeenCalledTimes(numberOfWorkspace), + ); expect(workspaceRepository.update).toHaveBeenNthCalledWith( numberOfWorkspace, { id: expect.any(String) }, diff --git a/packages/twenty-server/src/database/commands/command-runners/workspaces-migration.command-runner.ts b/packages/twenty-server/src/database/commands/command-runners/workspaces-migration.command-runner.ts index 2749619e1a2..f5a8baa0db1 100644 --- a/packages/twenty-server/src/database/commands/command-runners/workspaces-migration.command-runner.ts +++ b/packages/twenty-server/src/database/commands/command-runners/workspaces-migration.command-runner.ts @@ -4,11 +4,10 @@ import { isDefined } from 'twenty-shared/utils'; import { WorkspaceActivationStatus } from 'twenty-shared/workspace'; import { In, MoreThanOrEqual, type Repository } from 'typeorm'; -import { type WorkspaceDataSourceInterface } from 'src/engine/twenty-orm/interfaces/workspace-datasource.interface'; - import { MigrationCommandRunner } from 'src/database/commands/command-runners/migration.command-runner'; import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity'; import { DataSourceService } from 'src/engine/metadata-modules/data-source/data-source.service'; +import { GlobalWorkspaceDataSource } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource'; import { type GlobalWorkspaceOrmManager } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager'; import { buildSystemAuthContext } from 'src/engine/twenty-orm/utils/build-system-auth-context.util'; @@ -23,7 +22,7 @@ export type WorkspacesMigrationCommandOptions = { export type RunOnWorkspaceArgs = { options: WorkspacesMigrationCommandOptions; workspaceId: string; - dataSource?: WorkspaceDataSourceInterface; + dataSource?: GlobalWorkspaceDataSource; index: number; total: number; }; @@ -151,9 +150,7 @@ export abstract class WorkspacesMigrationCommandRunner< ); const dataSource = isDefined(workspaceHasDataSource) - ? await this.globalWorkspaceOrmManager.getDataSourceForWorkspace( - workspaceId, - ) + ? await this.globalWorkspaceOrmManager.getGlobalWorkspaceDataSource() : undefined; await this.runOnWorkspace({ @@ -178,14 +175,6 @@ export abstract class WorkspacesMigrationCommandRunner< chalk.red(`Error in workspace ${workspaceId}: ${error.message}`), ); } - - try { - await this.globalWorkspaceOrmManager.destroyDataSourceForWorkspace( - workspaceId, - ); - } catch (error) { - this.logger.error(error); - } } this.migrationReport.fail.forEach(({ error, workspaceId }) => diff --git a/packages/twenty-server/src/database/commands/upgrade-version-command/1-13/1-13-clean-empty-string-null-in-text-fields.command.ts b/packages/twenty-server/src/database/commands/upgrade-version-command/1-13/1-13-clean-empty-string-null-in-text-fields.command.ts index 5f153365b91..e5ba84922be 100644 --- a/packages/twenty-server/src/database/commands/upgrade-version-command/1-13/1-13-clean-empty-string-null-in-text-fields.command.ts +++ b/packages/twenty-server/src/database/commands/upgrade-version-command/1-13/1-13-clean-empty-string-null-in-text-fields.command.ts @@ -6,14 +6,13 @@ import { FieldMetadataType } from 'twenty-shared/types'; import { isDefined } from 'twenty-shared/utils'; import { Repository } from 'typeorm'; -import { type WorkspaceDataSourceInterface } from 'src/engine/twenty-orm/interfaces/workspace-datasource.interface'; - import { ActiveOrSuspendedWorkspacesMigrationCommandRunner } from 'src/database/commands/command-runners/active-or-suspended-workspaces-migration.command-runner'; import { RunOnWorkspaceArgs } from 'src/database/commands/command-runners/workspaces-migration.command-runner'; import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity'; import { DataSourceService } from 'src/engine/metadata-modules/data-source/data-source.service'; import { FieldMetadataEntity } from 'src/engine/metadata-modules/field-metadata/field-metadata.entity'; import { ObjectMetadataEntity } from 'src/engine/metadata-modules/object-metadata/object-metadata.entity'; +import { GlobalWorkspaceDataSource } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource'; import { GlobalWorkspaceOrmManager } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager'; import { computeObjectTargetTable } from 'src/engine/utils/compute-object-target-table.util'; import { getWorkspaceSchemaName } from 'src/engine/workspace-datasource/utils/get-workspace-schema-name.util'; @@ -114,7 +113,7 @@ export class CleanEmptyStringNullInTextFieldsCommand extends ActiveOrSuspendedWo objectMetadataItem: ObjectMetadataEntity, tableName: string, schemaName: string, - dataSource: WorkspaceDataSourceInterface, + dataSource: GlobalWorkspaceDataSource, isDryRun: boolean, ): Promise { const textFields = objectMetadataItem.fields.filter( @@ -155,7 +154,7 @@ export class CleanEmptyStringNullInTextFieldsCommand extends ActiveOrSuspendedWo objectMetadataItem: ObjectMetadataEntity, tableName: string, schemaName: string, - dataSource: WorkspaceDataSourceInterface, + dataSource: GlobalWorkspaceDataSource, isDryRun: boolean, ): Promise { const nameField = objectMetadataItem.fields.find( diff --git a/packages/twenty-server/src/engine/api/common/common-query-runners/common-base-query-runner.service.ts b/packages/twenty-server/src/engine/api/common/common-query-runners/common-base-query-runner.service.ts index 5b2b15a0666..17c80a335dd 100644 --- a/packages/twenty-server/src/engine/api/common/common-query-runners/common-base-query-runner.service.ts +++ b/packages/twenty-server/src/engine/api/common/common-query-runners/common-base-query-runner.service.ts @@ -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, }; diff --git a/packages/twenty-server/src/engine/api/common/common-query-runners/common-create-many-query-runner/common-create-many-query-runner.service.ts b/packages/twenty-server/src/engine/api/common/common-query-runners/common-create-many-query-runner/common-create-many-query-runner.service.ts index e0e5a41e1ac..0e05b92f3ba 100644 --- a/packages/twenty-server/src/engine/api/common/common-query-runners/common-create-many-query-runner/common-create-many-query-runner.service.ts +++ b/packages/twenty-server/src/engine/api/common/common-query-runners/common-create-many-query-runner/common-create-many-query-runner.service.ts @@ -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; flatFieldMetadataMaps: FlatEntityMaps; authContext: AuthContext; - workspaceDataSource: WorkspaceDataSource; + workspaceDataSource: GlobalWorkspaceDataSource; rolePermissionConfig?: RolePermissionConfig; }): Promise { if (!args.selectedFieldsResult.relations) { diff --git a/packages/twenty-server/src/engine/api/common/types/common-extended-query-runner-context.type.ts b/packages/twenty-server/src/engine/api/common/types/common-extended-query-runner-context.type.ts index c8116998a34..b1765695df9 100644 --- a/packages/twenty-server/src/engine/api/common/types/common-extended-query-runner-context.type.ts +++ b/packages/twenty-server/src/engine/api/common/types/common-extended-query-runner-context.type.ts @@ -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; commonQueryParser: GraphqlQueryParser; - workspaceDataSource: WorkspaceDataSource; + workspaceDataSource: GlobalWorkspaceDataSource; }; diff --git a/packages/twenty-server/src/engine/api/graphql/graphql-query-runner/helpers/process-nested-relations-v2.helper.ts b/packages/twenty-server/src/engine/api/graphql/graphql-query-runner/helpers/process-nested-relations-v2.helper.ts index f4f6017ee01..02f82840818 100644 --- a/packages/twenty-server/src/engine/api/graphql/graphql-query-runner/helpers/process-nested-relations-v2.helper.ts +++ b/packages/twenty-server/src/engine/api/graphql/graphql-query-runner/helpers/process-nested-relations-v2.helper.ts @@ -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; limit: number; authContext: AuthContext; - workspaceDataSource: WorkspaceDataSource; + workspaceDataSource: GlobalWorkspaceDataSource; rolePermissionConfig?: RolePermissionConfig; // eslint-disable-next-line @typescript-eslint/no-explicit-any selectedFields: Record; @@ -111,7 +111,7 @@ export class ProcessNestedRelationsV2Helper { aggregate: Record; limit: number; authContext: AuthContext; - workspaceDataSource: WorkspaceDataSource; + workspaceDataSource: GlobalWorkspaceDataSource; rolePermissionConfig?: RolePermissionConfig; selectedFields: Record; }): Promise { diff --git a/packages/twenty-server/src/engine/api/graphql/graphql-query-runner/helpers/process-nested-relations.helper.ts b/packages/twenty-server/src/engine/api/graphql/graphql-query-runner/helpers/process-nested-relations.helper.ts index a048a06a3f5..48cef4944eb 100644 --- a/packages/twenty-server/src/engine/api/graphql/graphql-query-runner/helpers/process-nested-relations.helper.ts +++ b/packages/twenty-server/src/engine/api/graphql/graphql-query-runner/helpers/process-nested-relations.helper.ts @@ -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; limit: number; authContext: AuthContext; - workspaceDataSource: WorkspaceDataSource; + workspaceDataSource: GlobalWorkspaceDataSource; rolePermissionConfig?: RolePermissionConfig; // eslint-disable-next-line @typescript-eslint/no-explicit-any selectedFields: Record; diff --git a/packages/twenty-server/src/engine/core-modules/auth/services/google-apis.service.spec.ts b/packages/twenty-server/src/engine/core-modules/auth/services/google-apis.service.spec.ts index c0152d1fb26..8be530dc6bf 100644 --- a/packages/twenty-server/src/engine/core-modules/auth/services/google-apis.service.spec.ts +++ b/packages/twenty-server/src/engine/core-modules/auth/services/google-apis.service.spec.ts @@ -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()), diff --git a/packages/twenty-server/src/engine/core-modules/auth/services/google-apis.service.ts b/packages/twenty-server/src/engine/core-modules/auth/services/google-apis.service.ts index 8ca046e6568..c0441da7383 100644 --- a/packages/twenty-server/src/engine/core-modules/auth/services/google-apis.service.ts +++ b/packages/twenty-server/src/engine/core-modules/auth/services/google-apis.service.ts @@ -126,9 +126,7 @@ export class GoogleAPIsService { ); const workspaceDataSource = - await this.globalWorkspaceOrmManager.getDataSourceForWorkspace( - workspaceId, - ); + await this.globalWorkspaceOrmManager.getGlobalWorkspaceDataSource(); await workspaceDataSource.transaction( async (manager: WorkspaceEntityManager) => { diff --git a/packages/twenty-server/src/engine/core-modules/auth/services/microsoft-apis.service.spec.ts b/packages/twenty-server/src/engine/core-modules/auth/services/microsoft-apis.service.spec.ts index ab2c1295983..7c27cdb3c77 100644 --- a/packages/twenty-server/src/engine/core-modules/auth/services/microsoft-apis.service.spec.ts +++ b/packages/twenty-server/src/engine/core-modules/auth/services/microsoft-apis.service.spec.ts @@ -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()), diff --git a/packages/twenty-server/src/engine/core-modules/auth/services/microsoft-apis.service.ts b/packages/twenty-server/src/engine/core-modules/auth/services/microsoft-apis.service.ts index b1214972c50..08ab5fbfc17 100644 --- a/packages/twenty-server/src/engine/core-modules/auth/services/microsoft-apis.service.ts +++ b/packages/twenty-server/src/engine/core-modules/auth/services/microsoft-apis.service.ts @@ -107,9 +107,7 @@ export class MicrosoftAPIsService { ); const workspaceDataSource = - await this.globalWorkspaceOrmManager.getDataSourceForWorkspace( - workspaceId, - ); + await this.globalWorkspaceOrmManager.getGlobalWorkspaceDataSource(); await workspaceDataSource.transaction( async (manager: WorkspaceEntityManager) => { diff --git a/packages/twenty-server/src/engine/twenty-orm/datasource/workspace.datasource.ts b/packages/twenty-server/src/engine/twenty-orm/datasource/workspace.datasource.ts deleted file mode 100644 index a21dc4f984e..00000000000 --- a/packages/twenty-server/src/engine/twenty-orm/datasource/workspace.datasource.ts +++ /dev/null @@ -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( - target: EntityTarget, - permissionOptions?: RolePermissionConfig, - authContext?: AuthContext, - ): WorkspaceRepository { - 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 | QueryRunner, - alias?: string, - queryRunner?: QueryRunner, - ) => { - if (isDefined(alias) && typeof alias === 'string') { - const entity = entityOrRunner as EntityTarget; - - return this.createQueryBuilder(entity, alias, queryRunner, { - calledByWorkspaceEntityManager: true, - }); - } else { - const runner = entityOrRunner as QueryRunner; - - return this.createQueryBuilder(runner, { - calledByWorkspaceEntityManager: true, - }); - } - }, - }, - ); - const queryRunner = this.driver.createQueryRunner(mode); - const manager = new EntityManagerFactory().create( - dataSourceWithOverridenCreateQueryBuilder, - queryRunner, - ); - - Object.assign(queryRunner, { manager: manager }); - - return queryRunner; - } - - override createQueryBuilder( - entityClass: EntityTarget, - alias: string, - queryRunner?: QueryRunner, - options?: CreateQueryBuilderOptions, - ): SelectQueryBuilder; - - override createQueryBuilder( - queryRunner?: QueryRunner, - options?: CreateQueryBuilderOptions, // eslint-disable-next-line @typescript-eslint/no-explicit-any - ): SelectQueryBuilder; - - // Only callable from workspaceEntityManager to guarantee a permission check was run - override createQueryBuilder( - // eslint-disable-next-line @typescript-eslint/no-explicit-any - queryRunnerOrEntityClass?: QueryRunner | EntityTarget, - aliasOrOptions?: string | CreateQueryBuilderOptions, - queryRunner?: QueryRunner, - options?: CreateQueryBuilderOptions, - // eslint-disable-next-line @typescript-eslint/no-explicit-any - ): SelectQueryBuilder { - let calledByWorkspaceEntityManager; - - const isCalledWithEntityTarget = - isDefined(aliasOrOptions) && typeof aliasOrOptions === 'string'; - - if (isCalledWithEntityTarget) { - calledByWorkspaceEntityManager = options?.calledByWorkspaceEntityManager; - } else { - calledByWorkspaceEntityManager = ( - aliasOrOptions as CreateQueryBuilderOptions - )?.calledByWorkspaceEntityManager; - } - - if (!(calledByWorkspaceEntityManager === true)) { - throw new PermissionsException( - 'Method not allowed because permissions are not implemented at datasource level.', - PermissionsExceptionCode.METHOD_NOT_ALLOWED, - ); - } - - if (isCalledWithEntityTarget) { - // eslint-disable-next-line @typescript-eslint/no-explicit-any - const entityClass = queryRunnerOrEntityClass as EntityTarget; - - return super.createQueryBuilder( - entityClass, - aliasOrOptions as string, - queryRunner, - ); - } else { - const queryRunner = queryRunnerOrEntityClass as QueryRunner; - - return super.createQueryBuilder(queryRunner); - } - } - - // eslint-disable-next-line @typescript-eslint/no-explicit-any - override query( - query: string, - // eslint-disable-next-line @typescript-eslint/no-explicit-any - parameters?: any[], - queryRunner?: QueryRunner, - options?: { - shouldBypassPermissionChecks?: boolean; - }, - ): Promise { - if (!options?.shouldBypassPermissionChecks) { - throw new PermissionsException( - 'Method not allowed because permissions are not implemented at datasource level.', - PermissionsExceptionCode.METHOD_NOT_ALLOWED, - ); - } - - return super.query(query, parameters, queryRunner); - } - - 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 { - 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(); - } - } -} diff --git a/packages/twenty-server/src/engine/twenty-orm/entity-manager/workspace-entity-manager.spec.ts b/packages/twenty-server/src/engine/twenty-orm/entity-manager/workspace-entity-manager.spec.ts index 9028bafd0f8..5f2d5341247 100644 --- a/packages/twenty-server/src/engine/twenty-orm/entity-manager/workspace-entity-manager.spec.ts +++ b/packages/twenty-server/src/engine/twenty-orm/entity-manager/workspace-entity-manager.spec.ts @@ -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', diff --git a/packages/twenty-server/src/engine/twenty-orm/entity-manager/workspace-entity-manager.ts b/packages/twenty-server/src/engine/twenty-orm/entity-manager/workspace-entity-manager.ts index 6f630b0b7ea..ff512e94ca4 100644 --- a/packages/twenty-server/src/engine/twenty-orm/entity-manager/workspace-entity-manager.ts +++ b/packages/twenty-server/src/engine/twenty-orm/entity-manager/workspace-entity-manager.ts @@ -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>; - 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, + target: EntityTarget, ): string { return this.connection.getMetadata(target).name; } diff --git a/packages/twenty-server/src/engine/twenty-orm/factories/index.ts b/packages/twenty-server/src/engine/twenty-orm/factories/index.ts index b07da0d7bc5..951debc166e 100644 --- a/packages/twenty-server/src/engine/twenty-orm/factories/index.ts +++ b/packages/twenty-server/src/engine/twenty-orm/factories/index.ts @@ -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, ]; diff --git a/packages/twenty-server/src/engine/twenty-orm/factories/workspace-datasource.factory.ts b/packages/twenty-server/src/engine/twenty-orm/factories/workspace-datasource.factory.ts deleted file mode 100644 index e1bf939dfe6..00000000000 --- a/packages/twenty-server/src/engine/twenty-orm/factories/workspace-datasource.factory.ts +++ /dev/null @@ -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(); - - 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, - private readonly workspaceEventEmitter: WorkspaceEventEmitter, - ) {} - - private async safelyDestroyDataSource( - dataSource: WorkspaceDataSource, - ): Promise { - 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 { - 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({ - 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 { - 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 { - 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 { - 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}`, - ); - } - } -} diff --git a/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource.service.ts b/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource.service.ts index df70e4b94c3..f011a1718a5 100644 --- a/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource.service.ts +++ b/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource.service.ts @@ -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, diff --git a/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource.ts b/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource.ts index ecc167c1c15..40ec2d44061 100644 --- a/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource.ts +++ b/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-datasource.ts @@ -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( diff --git a/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager.ts b/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager.ts index 2355a2822ea..926116b3426 100644 --- a/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager.ts +++ b/packages/twenty-server/src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager.ts @@ -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 { - const { featureFlagsMap } = await this.workspaceCacheService.getOrRecompute( - workspaceId, - ['featureFlagsMap'], - ); - - return ( - featureFlagsMap[FeatureFlagKey.IS_GLOBAL_WORKSPACE_DATASOURCE_ENABLED] === - true - ); - } - async getRepository( workspaceId: string, workspaceEntity: Type, @@ -51,62 +36,29 @@ export class GlobalWorkspaceOrmManager { ): Promise>; async getRepository( - workspaceId: string, + _workspaceId: string, workspaceEntityOrObjectMetadataName: Type | string, permissionOptions?: RolePermissionConfig, ): Promise> { - 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( - objectMetadataName, - permissionOptions, - ); - } + let objectMetadataName: string; if (typeof workspaceEntityOrObjectMetadataName === 'string') { - return this.twentyORMGlobalManager.getRepositoryForWorkspace( - workspaceId, - workspaceEntityOrObjectMetadataName, - permissionOptions, + objectMetadataName = workspaceEntityOrObjectMetadataName; + } else { + objectMetadataName = convertClassNameToObjectMetadataName( + workspaceEntityOrObjectMetadataName.name, ); } - return this.twentyORMGlobalManager.getRepositoryForWorkspace( - workspaceId, - workspaceEntityOrObjectMetadataName, + const globalDataSource = await this.getGlobalWorkspaceDataSource(); + + return globalDataSource.getRepository( + objectMetadataName, permissionOptions, ); } - async getDataSourceForWorkspace( - workspaceId: string, - ): Promise { - const isGlobalFlow = await this.isGlobalDataSourceFlow(workspaceId); - - if (isGlobalFlow) { - return this.globalWorkspaceDataSourceService.getGlobalWorkspaceDataSource(); - } - - return this.twentyORMGlobalManager.getDataSourceForWorkspace({ - workspaceId, - }); - } - - async getGlobalWorkspaceDataSource() { + async getGlobalWorkspaceDataSource(): Promise { return this.globalWorkspaceDataSourceService.getGlobalWorkspaceDataSource(); } @@ -119,18 +71,6 @@ export class GlobalWorkspaceOrmManager { return withWorkspaceContext(context, fn); } - async destroyDataSourceForWorkspace(workspaceId: string): Promise { - const isGlobalFlow = await this.isGlobalDataSourceFlow(workspaceId); - - if (isGlobalFlow) { - return; - } - - await this.twentyORMGlobalManager.destroyDataSourceForWorkspace( - workspaceId, - ); - } - private async loadWorkspaceContext( authContext: WorkspaceAuthContext, ): Promise { diff --git a/packages/twenty-server/src/engine/twenty-orm/interfaces/workspace-datasource.interface.ts b/packages/twenty-server/src/engine/twenty-orm/interfaces/workspace-datasource.interface.ts deleted file mode 100644 index 01e20edbc24..00000000000 --- a/packages/twenty-server/src/engine/twenty-orm/interfaces/workspace-datasource.interface.ts +++ /dev/null @@ -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( - target: EntityTarget, - permissionOptions?: RolePermissionConfig, - authContext?: AuthContext, - ): WorkspaceRepository; - - createQueryRunner(mode?: ReplicationMode): WorkspaceQueryRunner; - - createEntityManager(queryRunner?: QueryRunner): WorkspaceEntityManager; - - createQueryRunnerForEntityPersistExecutor( - mode?: ReplicationMode, - ): QueryRunner; - - transaction( - runInTransaction: (entityManager: EntityManager) => Promise, - ): Promise; - - query( - query: string, - parameters?: unknown[], - queryRunner?: QueryRunner, - options?: { shouldBypassPermissionChecks?: boolean }, - ): Promise; -} diff --git a/packages/twenty-server/src/engine/twenty-orm/pg-shared-pool/__tests__/pg-shared-pool.service.spec.ts b/packages/twenty-server/src/engine/twenty-orm/pg-shared-pool/__tests__/pg-shared-pool.service.spec.ts deleted file mode 100644 index 31468a858fe..00000000000 --- a/packages/twenty-server/src/engine/twenty-orm/pg-shared-pool/__tests__/pg-shared-pool.service.spec.ts +++ /dev/null @@ -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; - - const configValues: Record = { - 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); - configService = module.get(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); - - 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(); - }); - }); -}); diff --git a/packages/twenty-server/src/engine/twenty-orm/pg-shared-pool/pg-shared-pool.module.ts b/packages/twenty-server/src/engine/twenty-orm/pg-shared-pool/pg-shared-pool.module.ts deleted file mode 100644 index 7a45a53a366..00000000000 --- a/packages/twenty-server/src/engine/twenty-orm/pg-shared-pool/pg-shared-pool.module.ts +++ /dev/null @@ -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(); - } -} diff --git a/packages/twenty-server/src/engine/twenty-orm/pg-shared-pool/pg-shared-pool.service.ts b/packages/twenty-server/src/engine/twenty-orm/pg-shared-pool/pg-shared-pool.service.ts deleted file mode 100644 index 479e616bd6f..00000000000 --- a/packages/twenty-server/src/engine/twenty-orm/pg-shared-pool/pg-shared-pool.service.ts +++ /dev/null @@ -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; - _idle?: Array; - _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(); - - 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 | 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 { - 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 { - 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 { - if (this.poolsMap.size === 0) { - return; - } - - const closePromises: Promise[] = []; - - 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 => { - 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((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.', - ); - } -} diff --git a/packages/twenty-server/src/engine/twenty-orm/twenty-orm-global.manager.ts b/packages/twenty-server/src/engine/twenty-orm/twenty-orm-global.manager.ts deleted file mode 100644 index 8f72b71fb31..00000000000 --- a/packages/twenty-server/src/engine/twenty-orm/twenty-orm-global.manager.ts +++ /dev/null @@ -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( - workspaceId: string, - workspaceEntity: Type, - options?: RolePermissionConfig, - ): Promise>; - - async getRepositoryForWorkspace( - workspaceId: string, - objectMetadataName: string, - options?: RolePermissionConfig, - ): Promise>; - - async getRepositoryForWorkspace( - workspaceId: string, - workspaceEntityOrObjectMetadataName: Type | string, - options?: RolePermissionConfig, - ): Promise> { - 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( - objectMetadataName, - options, - ); - - return repository; - } - - async getDataSourceForWorkspace({ workspaceId }: { workspaceId: string }) { - return await this.workspaceDataSourceFactory.create(workspaceId); - } - - async destroyDataSourceForWorkspace(workspaceId: string) { - await this.workspaceDataSourceFactory.destroy(workspaceId); - } -} diff --git a/packages/twenty-server/src/engine/twenty-orm/twenty-orm.module-definition.ts b/packages/twenty-server/src/engine/twenty-orm/twenty-orm.module-definition.ts deleted file mode 100644 index 0cc798e7481..00000000000 --- a/packages/twenty-server/src/engine/twenty-orm/twenty-orm.module-definition.ts +++ /dev/null @@ -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().build(); diff --git a/packages/twenty-server/src/engine/twenty-orm/twenty-orm.module.ts b/packages/twenty-server/src/engine/twenty-orm/twenty-orm.module.ts index a5ab96929a9..d95c5db711f 100644 --- a/packages/twenty-server/src/engine/twenty-orm/twenty-orm.module.ts +++ b/packages/twenty-server/src/engine/twenty-orm/twenty-orm.module.ts @@ -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 {} diff --git a/packages/twenty-server/src/engine/workspace-manager/workspace-cleaner/services/cleaner.workspace-service.ts b/packages/twenty-server/src/engine/workspace-manager/workspace-cleaner/services/cleaner.workspace-service.ts index a965d8b5569..63f3e07723a 100644 --- a/packages/twenty-server/src/engine/workspace-manager/workspace-cleaner/services/cleaner.workspace-service.ts +++ b/packages/twenty-server/src/engine/workspace-manager/workspace-cleaner/services/cleaner.workspace-service.ts @@ -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!`, diff --git a/packages/twenty-server/src/modules/calendar/calendar-event-import-manager/services/calendar-save-events.service.ts b/packages/twenty-server/src/modules/calendar/calendar-event-import-manager/services/calendar-save-events.service.ts index 1e237fb0da9..e3b99055931 100644 --- a/packages/twenty-server/src/modules/calendar/calendar-event-import-manager/services/calendar-save-events.service.ts +++ b/packages/twenty-server/src/modules/calendar/calendar-event-import-manager/services/calendar-save-events.service.ts @@ -72,9 +72,7 @@ export class CalendarSaveEventsService { ); const workspaceDataSource = - await this.globalWorkspaceOrmManager.getDataSourceForWorkspace( - workspaceId, - ); + await this.globalWorkspaceOrmManager.getGlobalWorkspaceDataSource(); await workspaceDataSource.transaction( async (transactionManager: WorkspaceEntityManager) => { diff --git a/packages/twenty-server/src/modules/connected-account/services/imap-smtp-caldav-apis.service.spec.ts b/packages/twenty-server/src/modules/connected-account/services/imap-smtp-caldav-apis.service.spec.ts index b35f28dfd5a..9503f6f3b18 100644 --- a/packages/twenty-server/src/modules/connected-account/services/imap-smtp-caldav-apis.service.spec.ts +++ b/packages/twenty-server/src/modules/connected-account/services/imap-smtp-caldav-apis.service.spec.ts @@ -64,9 +64,9 @@ describe('ImapSmtpCalDavAPIService', () => { return {}; }), - getDataSourceForWorkspace: jest + getGlobalWorkspaceDataSource: jest .fn() - .mockImplementation(() => mockWorkspaceDataSource), + .mockResolvedValue(mockWorkspaceDataSource), executeInWorkspaceContext: jest .fn() diff --git a/packages/twenty-server/src/modules/connected-account/services/imap-smtp-caldav-apis.service.ts b/packages/twenty-server/src/modules/connected-account/services/imap-smtp-caldav-apis.service.ts index 9860252497e..074c12bf082 100644 --- a/packages/twenty-server/src/modules/connected-account/services/imap-smtp-caldav-apis.service.ts +++ b/packages/twenty-server/src/modules/connected-account/services/imap-smtp-caldav-apis.service.ts @@ -89,9 +89,7 @@ export class ImapSmtpCalDavAPIService { existingAccount?.id ?? connectedAccountId ?? v4(); const workspaceDataSource = - await this.globalWorkspaceOrmManager.getDataSourceForWorkspace( - workspaceId, - ); + await this.globalWorkspaceOrmManager.getGlobalWorkspaceDataSource(); const existingMessageChannel = existingAccount ? await messageChannelRepository.findOne({ diff --git a/packages/twenty-server/src/modules/dashboard/services/dashboard-duplication.service.ts b/packages/twenty-server/src/modules/dashboard/services/dashboard-duplication.service.ts index 9bce2d515ae..6702beb9061 100644 --- a/packages/twenty-server/src/modules/dashboard/services/dashboard-duplication.service.ts +++ b/packages/twenty-server/src/modules/dashboard/services/dashboard-duplication.service.ts @@ -3,7 +3,9 @@ import { Injectable, Logger } from '@nestjs/common'; import { appendCopySuffix, isDefined } from 'twenty-shared/utils'; import { PageLayoutDuplicationService } from 'src/engine/metadata-modules/page-layout/services/page-layout-duplication.service'; -import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager'; +import { GlobalWorkspaceOrmManager } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager'; +import { WorkspaceRepository } from 'src/engine/twenty-orm/repository/workspace.repository'; +import { buildSystemAuthContext } from 'src/engine/twenty-orm/utils/build-system-auth-context.util'; import { DuplicatedDashboardDTO } from 'src/modules/dashboard/dtos/duplicated-dashboard.dto'; import { DashboardException, @@ -19,82 +21,84 @@ export class DashboardDuplicationService { constructor( private readonly pageLayoutDuplicationService: PageLayoutDuplicationService, - private readonly twentyORMGlobalManager: TwentyORMGlobalManager, + private readonly globalWorkspaceOrmManager: GlobalWorkspaceOrmManager, ) {} async duplicateDashboard( dashboardId: string, workspaceId: string, ): Promise { - const dashboardRepository = - await this.twentyORMGlobalManager.getRepositoryForWorkspace( - workspaceId, - 'dashboard', - { shouldBypassPermissionChecks: true }, - ); + return this.globalWorkspaceOrmManager.executeInWorkspaceContext( + buildSystemAuthContext(workspaceId), + async () => { + const dashboardRepository = + await this.globalWorkspaceOrmManager.getRepository( + workspaceId, + 'dashboard', + { shouldBypassPermissionChecks: true }, + ); - const originalDashboard = await dashboardRepository.findOne({ - where: { id: dashboardId }, - }); + const originalDashboard = await dashboardRepository.findOne({ + where: { id: dashboardId }, + }); - if (!isDefined(originalDashboard)) { - throw new DashboardException( - generateDashboardExceptionMessage( - DashboardExceptionMessageKey.DASHBOARD_NOT_FOUND, - dashboardId, - ), - DashboardExceptionCode.DASHBOARD_NOT_FOUND, - ); - } + if (!isDefined(originalDashboard)) { + throw new DashboardException( + generateDashboardExceptionMessage( + DashboardExceptionMessageKey.DASHBOARD_NOT_FOUND, + dashboardId, + ), + DashboardExceptionCode.DASHBOARD_NOT_FOUND, + ); + } - if (!isDefined(originalDashboard.pageLayoutId)) { - throw new DashboardException( - generateDashboardExceptionMessage( - DashboardExceptionMessageKey.PAGE_LAYOUT_NOT_FOUND, - dashboardId, - ), - DashboardExceptionCode.PAGE_LAYOUT_NOT_FOUND, - ); - } + if (!isDefined(originalDashboard.pageLayoutId)) { + throw new DashboardException( + generateDashboardExceptionMessage( + DashboardExceptionMessageKey.PAGE_LAYOUT_NOT_FOUND, + dashboardId, + ), + DashboardExceptionCode.PAGE_LAYOUT_NOT_FOUND, + ); + } - try { - const newPageLayout = await this.pageLayoutDuplicationService.duplicate({ - pageLayoutId: originalDashboard.pageLayoutId, - workspaceId, - }); + try { + const newPageLayout = + await this.pageLayoutDuplicationService.duplicate({ + pageLayoutId: originalDashboard.pageLayoutId, + workspaceId, + }); - const newDashboard = await this.createDuplicatedDashboard( - originalDashboard, - newPageLayout.id, - dashboardRepository, - ); + const newDashboard = await this.createDuplicatedDashboard( + originalDashboard, + newPageLayout.id, + dashboardRepository, + ); - return { - id: newDashboard.id, - title: newDashboard.title, - pageLayoutId: newDashboard.pageLayoutId, - position: newDashboard.position, - createdAt: newDashboard.createdAt, - updatedAt: newDashboard.updatedAt, - }; - } catch (error) { - this.logger.error( - `Failed to duplicate dashboard ${dashboardId}: ${error.message}`, - error.stack, - ); + return { + id: newDashboard.id, + title: newDashboard.title, + pageLayoutId: newDashboard.pageLayoutId, + position: newDashboard.position, + createdAt: newDashboard.createdAt, + updatedAt: newDashboard.updatedAt, + }; + } catch (error) { + this.logger.error( + `Failed to duplicate dashboard ${dashboardId}: ${error.message}`, + error.stack, + ); - throw error; - } + throw error; + } + }, + ); } private async createDuplicatedDashboard( originalDashboard: DashboardWorkspaceEntity, newPageLayoutId: string, - dashboardRepository: Awaited< - ReturnType< - typeof this.twentyORMGlobalManager.getRepositoryForWorkspace - > - >, + dashboardRepository: WorkspaceRepository, ): Promise { const newTitle = appendCopySuffix(originalDashboard.title ?? ''); diff --git a/packages/twenty-server/src/modules/messaging/message-cleaner/services/messaging-message-cleaner.service.ts b/packages/twenty-server/src/modules/messaging/message-cleaner/services/messaging-message-cleaner.service.ts index cdf545c182e..39661ad09d0 100644 --- a/packages/twenty-server/src/modules/messaging/message-cleaner/services/messaging-message-cleaner.service.ts +++ b/packages/twenty-server/src/modules/messaging/message-cleaner/services/messaging-message-cleaner.service.ts @@ -142,9 +142,7 @@ export class MessagingMessageCleanerService { ); const workspaceDataSource = - await this.globalWorkspaceOrmManager.getDataSourceForWorkspace( - workspaceId, - ); + await this.globalWorkspaceOrmManager.getGlobalWorkspaceDataSource(); await workspaceDataSource.transaction( async (transactionManager: WorkspaceEntityManager) => { diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/services/__tests__/messaging-message-list-fetch.service.spec.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/services/__tests__/messaging-message-list-fetch.service.spec.ts index 7eb0162aea9..ca4b3586011 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/services/__tests__/messaging-message-list-fetch.service.spec.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/services/__tests__/messaging-message-list-fetch.service.spec.ts @@ -198,7 +198,7 @@ describe('MessagingMessageListFetchService', () => { { provide: GlobalWorkspaceOrmManager, useValue: { - getDataSourceForWorkspace: jest.fn().mockResolvedValue({ + getGlobalWorkspaceDataSource: jest.fn().mockResolvedValue({ manager: {}, }), executeInWorkspaceContext: jest @@ -206,6 +206,11 @@ describe('MessagingMessageListFetchService', () => { .mockImplementation((_authContext: any, fn: () => any) => fn()), getRepository: jest.fn().mockImplementation((workspaceId, name) => { + if (name === 'messageChannel') { + return { + findOne: jest.fn().mockResolvedValue(undefined), + }; + } if (name === 'messageChannelMessageAssociation') { return mockMessageChannelMessageAssociationRepository; } diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-message-list-fetch.service.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-message-list-fetch.service.ts index 0e255466b9b..25cf9ae9765 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-message-list-fetch.service.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-message-list-fetch.service.ts @@ -128,9 +128,7 @@ export class MessagingMessageListFetchService { }; const datasource = - await this.globalWorkspaceOrmManager.getDataSourceForWorkspace( - workspaceId, - ); + await this.globalWorkspaceOrmManager.getGlobalWorkspaceDataSource(); await this.syncMessageFoldersService.syncMessageFolders({ workspaceId, diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-process-folder-actions.service.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-process-folder-actions.service.ts index 66204478f7d..092cefb7186 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-process-folder-actions.service.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-process-folder-actions.service.ts @@ -93,9 +93,7 @@ export class MessagingProcessFolderActionsService { authContext, async () => { const workspaceDataSource = - await this.globalWorkspaceOrmManager.getDataSourceForWorkspace( - workspaceId, - ); + await this.globalWorkspaceOrmManager.getGlobalWorkspaceDataSource(); await workspaceDataSource?.transaction( async (transactionManager: WorkspaceEntityManager) => { diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-process-group-email-actions.service.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-process-group-email-actions.service.ts index c93129dc4b4..db8cf2bf64f 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-process-group-email-actions.service.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-process-group-email-actions.service.ts @@ -74,9 +74,7 @@ export class MessagingProcessGroupEmailActionsService { authContext, async () => { const workspaceDataSource = - await this.globalWorkspaceOrmManager.getDataSourceForWorkspace( - workspaceId, - ); + await this.globalWorkspaceOrmManager.getGlobalWorkspaceDataSource(); await workspaceDataSource?.transaction( async (transactionManager: WorkspaceEntityManager) => { diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-save-messages-and-enqueue-contact-creation.service.spec.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-save-messages-and-enqueue-contact-creation.service.spec.ts index fd890b6fae4..400b46f47d8 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-save-messages-and-enqueue-contact-creation.service.spec.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-save-messages-and-enqueue-contact-creation.service.spec.ts @@ -154,7 +154,7 @@ describe('MessagingSaveMessagesAndEnqueueContactCreationService', () => { { provide: GlobalWorkspaceOrmManager, useValue: { - getDataSourceForWorkspace: jest + getGlobalWorkspaceDataSource: jest .fn() .mockResolvedValue(datasourceInstance), executeInWorkspaceContext: jest diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-save-messages-and-enqueue-contact-creation.service.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-save-messages-and-enqueue-contact-creation.service.ts index 04e4dd9c483..0927f2cba8a 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-save-messages-and-enqueue-contact-creation.service.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-save-messages-and-enqueue-contact-creation.service.ts @@ -50,9 +50,7 @@ export class MessagingSaveMessagesAndEnqueueContactCreationService { authContext, async () => { const workspaceDataSource = - await this.globalWorkspaceOrmManager.getDataSourceForWorkspace( - workspaceId, - ); + await this.globalWorkspaceOrmManager.getGlobalWorkspaceDataSource(); return workspaceDataSource?.transaction( async (transactionManager: WorkspaceEntityManager) => {