From dec2239ae730210223b01e2afcb412ab7244a5a8 Mon Sep 17 00:00:00 2001 From: Weiko Date: Thu, 28 Aug 2025 13:21:26 +0200 Subject: [PATCH] Remove typeorm service (#14116) ## Context To simplify the way we inject our default datasource, I've recently removed the token injection that was confusion since we only had once configured on the module level. Now I'm removing TypeORM service which allows us to instantiate a new Datasource with the same parameters as the default one, it was redundant and confusing. --- ...enqueued-status-to-workflow-run.command.ts | 13 ++-- .../1-1/1-1-fix-schema-array-type.command.ts | 12 ++- ...ueued-status-to-workflow-run-v2.command.ts | 13 ++-- ...ep-ids-to-workflow-runs-trigger.command.ts | 17 ++--- ...column-type-in-workspace-schema.command.ts | 13 ++-- ...flow-versions-and-workflow-runs.command.ts | 15 ++-- .../src/database/typeorm/typeorm.module.ts | 12 +-- .../src/database/typeorm/typeorm.service.ts | 73 ------------------ .../user-workspace.service.spec.ts | 12 --- .../engine/core-modules/user/user.module.ts | 8 +- .../object-metadata.service.ts | 19 ++--- .../remote-server/remote-server.service.ts | 2 +- .../distant-table/distant-table.service.ts | 12 ++- .../foreign-table/foreign-table.service.ts | 10 ++- .../remote-table/remote-table.service.ts | 11 ++- .../utils/get-remote-table-local-name.util.ts | 11 ++- .../workspace-datasource.service.ts | 36 ++++----- .../workspace-manager.service.spec.ts | 9 +++ .../dev-seeder-permissions.service.ts | 75 ++++++++----------- .../data/services/dev-seeder-data.service.ts | 16 ++-- .../services/dev-seeder-metadata.service.ts | 19 +++-- .../dev-seeder/services/dev-seeder.service.ts | 15 ++-- .../prefill-core-views.ts | 4 +- .../standard-objects-prefill-data.ts | 57 +++++++------- .../services/database-structure.service.ts | 24 +++--- .../object-metadata-health.service.ts | 14 ++-- .../workspace-manager.service.ts | 16 ++-- .../workspace-migration-runner.service.ts | 17 ++--- .../calendar-event-list-fetch.cron.job.ts | 13 ++-- .../jobs/calendar-events-import.cron.job.ts | 12 ++- .../messaging-message-list-fetch.cron.job.ts | 13 ++-- .../messaging-messages-import.cron.job.ts | 13 ++-- .../workflow-clean-workflow-runs.cron.job.ts | 13 ++-- .../crons/jobs/cron-trigger.cron.job.ts | 13 ++-- .../test/integration/utils/setup-test.ts | 7 +- 35 files changed, 240 insertions(+), 399 deletions(-) delete mode 100644 packages/twenty-server/src/database/typeorm/typeorm.service.ts diff --git a/packages/twenty-server/src/database/commands/upgrade-version-command/1-1/1-1-add-enqueued-status-to-workflow-run.command.ts b/packages/twenty-server/src/database/commands/upgrade-version-command/1-1/1-1-add-enqueued-status-to-workflow-run.command.ts index 7c76c87b4d..8f45e03ebb 100644 --- a/packages/twenty-server/src/database/commands/upgrade-version-command/1-1/1-1-add-enqueued-status-to-workflow-run.command.ts +++ b/packages/twenty-server/src/database/commands/upgrade-version-command/1-1/1-1-add-enqueued-status-to-workflow-run.command.ts @@ -1,7 +1,7 @@ -import { InjectRepository } from '@nestjs/typeorm'; +import { InjectDataSource, InjectRepository } from '@nestjs/typeorm'; import { Command } from 'nest-commander'; -import { Repository } from 'typeorm'; +import { DataSource, Repository } from 'typeorm'; import { ActiveOrSuspendedWorkspacesMigrationCommandRunner, @@ -11,7 +11,6 @@ import { Workspace } from 'src/engine/core-modules/workspace/workspace.entity'; import { FieldMetadataEntity } from 'src/engine/metadata-modules/field-metadata/field-metadata.entity'; import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager'; import { getWorkspaceSchemaName } from 'src/engine/workspace-datasource/utils/get-workspace-schema-name.util'; -import { WorkspaceDataSourceService } from 'src/engine/workspace-datasource/workspace-datasource.service'; import { WORKFLOW_RUN_STANDARD_FIELD_IDS } from 'src/engine/workspace-manager/workspace-sync-metadata/constants/standard-field-ids'; import { WorkflowRunStatus } from 'src/modules/workflow/common/standard-objects/workflow-run.workspace-entity'; @@ -26,7 +25,8 @@ export class AddEnqueuedStatusToWorkflowRunCommand extends ActiveOrSuspendedWork protected readonly twentyORMGlobalManager: TwentyORMGlobalManager, @InjectRepository(FieldMetadataEntity) private readonly fieldMetadataRepository: Repository, - private readonly workspaceDataSourceService: WorkspaceDataSourceService, + @InjectDataSource() + private readonly coreDataSource: DataSource, ) { super(workspaceRepository, twentyORMGlobalManager); } @@ -89,16 +89,13 @@ export class AddEnqueuedStatusToWorkflowRunCommand extends ActiveOrSuspendedWork const schemaName = getWorkspaceSchemaName(workspaceId); - const mainDataSource = - await this.workspaceDataSourceService.connectToMainDataSource(); - if (options.dryRun) { this.logger.log( `Would try to add enqueued status to workflow run status enum for workspace ${workspaceId}`, ); } else { try { - await mainDataSource.query( + await this.coreDataSource.query( `ALTER TYPE ${schemaName}."workflowRun_status_enum" ADD VALUE 'ENQUEUED'`, ); this.logger.log( diff --git a/packages/twenty-server/src/database/commands/upgrade-version-command/1-1/1-1-fix-schema-array-type.command.ts b/packages/twenty-server/src/database/commands/upgrade-version-command/1-1/1-1-fix-schema-array-type.command.ts index 6b7caeebf6..98496dcfc4 100644 --- a/packages/twenty-server/src/database/commands/upgrade-version-command/1-1/1-1-fix-schema-array-type.command.ts +++ b/packages/twenty-server/src/database/commands/upgrade-version-command/1-1/1-1-fix-schema-array-type.command.ts @@ -1,14 +1,13 @@ -import { InjectRepository } from '@nestjs/typeorm'; +import { InjectDataSource, InjectRepository } from '@nestjs/typeorm'; import { Command } from 'nest-commander'; import { FieldMetadataType } from 'twenty-shared/types'; -import { Repository } from 'typeorm'; +import { DataSource, Repository } from 'typeorm'; import { ActiveOrSuspendedWorkspacesMigrationCommandRunner, type RunOnWorkspaceArgs, } from 'src/database/commands/command-runners/active-or-suspended-workspaces-migration.command-runner'; -import { TypeORMService } from 'src/database/typeorm/typeorm.service'; import { Workspace } from 'src/engine/core-modules/workspace/workspace.entity'; import { FieldMetadataEntity } from 'src/engine/metadata-modules/field-metadata/field-metadata.entity'; import { computeColumnName } from 'src/engine/metadata-modules/field-metadata/utils/compute-column-name.util'; @@ -27,7 +26,8 @@ export class FixSchemaArrayTypeCommand extends ActiveOrSuspendedWorkspacesMigrat protected readonly workspaceRepository: Repository, protected readonly twentyORMGlobalManager: TwentyORMGlobalManager, private readonly databaseStructureService: DatabaseStructureService, - private readonly typeORMService: TypeORMService, + @InjectDataSource() + private readonly coreDataSource: DataSource, @InjectRepository(FieldMetadataEntity) private readonly fieldMetadataRepository: Repository, ) { @@ -82,9 +82,7 @@ export class FixSchemaArrayTypeCommand extends ActiveOrSuspendedWorkspacesMigrat `Altering column ${schemaName}.${tableName}.${columnName} to type text[] (was ${dbColumn.dataType})`, ); if (!options.dryRun) { - const queryRunner = this.typeORMService - .getMainDataSource() - .createQueryRunner(); + const queryRunner = this.coreDataSource.createQueryRunner(); await queryRunner.connect(); try { diff --git a/packages/twenty-server/src/database/commands/upgrade-version-command/1-2/1-2-add-enqueued-status-to-workflow-run-v2.command.ts b/packages/twenty-server/src/database/commands/upgrade-version-command/1-2/1-2-add-enqueued-status-to-workflow-run-v2.command.ts index 1250ed7829..c4958f66c9 100644 --- a/packages/twenty-server/src/database/commands/upgrade-version-command/1-2/1-2-add-enqueued-status-to-workflow-run-v2.command.ts +++ b/packages/twenty-server/src/database/commands/upgrade-version-command/1-2/1-2-add-enqueued-status-to-workflow-run-v2.command.ts @@ -1,7 +1,7 @@ -import { InjectRepository } from '@nestjs/typeorm'; +import { InjectDataSource, InjectRepository } from '@nestjs/typeorm'; import { Command } from 'nest-commander'; -import { Repository } from 'typeorm'; +import { DataSource, Repository } from 'typeorm'; import { ActiveOrSuspendedWorkspacesMigrationCommandRunner, @@ -12,7 +12,6 @@ import { FieldMetadataEntity } from 'src/engine/metadata-modules/field-metadata/ import { ObjectMetadataEntity } from 'src/engine/metadata-modules/object-metadata/object-metadata.entity'; import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager'; import { getWorkspaceSchemaName } from 'src/engine/workspace-datasource/utils/get-workspace-schema-name.util'; -import { WorkspaceDataSourceService } from 'src/engine/workspace-datasource/workspace-datasource.service'; import { WORKFLOW_RUN_STANDARD_FIELD_IDS } from 'src/engine/workspace-manager/workspace-sync-metadata/constants/standard-field-ids'; import { STANDARD_OBJECT_IDS } from 'src/engine/workspace-manager/workspace-sync-metadata/constants/standard-object-ids'; import { WorkflowRunStatus } from 'src/modules/workflow/common/standard-objects/workflow-run.workspace-entity'; @@ -30,7 +29,8 @@ export class AddEnqueuedStatusToWorkflowRunV2Command extends ActiveOrSuspendedWo private readonly objectMetadataRepository: Repository, @InjectRepository(FieldMetadataEntity) private readonly fieldMetadataRepository: Repository, - private readonly workspaceDataSourceService: WorkspaceDataSourceService, + @InjectDataSource() + private readonly coreDataSource: DataSource, ) { super(workspaceRepository, twentyORMGlobalManager); } @@ -110,16 +110,13 @@ export class AddEnqueuedStatusToWorkflowRunV2Command extends ActiveOrSuspendedWo const schemaName = getWorkspaceSchemaName(workspaceId); - const mainDataSource = - await this.workspaceDataSourceService.connectToMainDataSource(); - if (options.dryRun) { this.logger.log( `Would try to add enqueued status to workflow run status enum for workspace ${workspaceId}`, ); } else { try { - await mainDataSource.query( + await this.coreDataSource.query( `ALTER TYPE ${schemaName}."workflowRun_status_enum" ADD VALUE 'ENQUEUED'`, ); this.logger.log( diff --git a/packages/twenty-server/src/database/commands/upgrade-version-command/1-3/1-3-add-next-step-ids-to-workflow-runs-trigger.command.ts b/packages/twenty-server/src/database/commands/upgrade-version-command/1-3/1-3-add-next-step-ids-to-workflow-runs-trigger.command.ts index 88e65e2168..5583926f4e 100644 --- a/packages/twenty-server/src/database/commands/upgrade-version-command/1-3/1-3-add-next-step-ids-to-workflow-runs-trigger.command.ts +++ b/packages/twenty-server/src/database/commands/upgrade-version-command/1-3/1-3-add-next-step-ids-to-workflow-runs-trigger.command.ts @@ -1,8 +1,8 @@ -import { InjectRepository } from '@nestjs/typeorm'; +import { InjectDataSource, InjectRepository } from '@nestjs/typeorm'; import { Command } from 'nest-commander'; -import { Repository } from 'typeorm'; import { isDefined } from 'twenty-shared/utils'; +import { DataSource, Repository } from 'typeorm'; import { ActiveOrSuspendedWorkspacesMigrationCommandRunner, @@ -10,9 +10,8 @@ import { } from 'src/database/commands/command-runners/active-or-suspended-workspaces-migration.command-runner'; import { Workspace } from 'src/engine/core-modules/workspace/workspace.entity'; import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager'; -import { type WorkflowAction } from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type'; -import { WorkspaceDataSourceService } from 'src/engine/workspace-datasource/workspace-datasource.service'; import { getWorkspaceSchemaName } from 'src/engine/workspace-datasource/utils/get-workspace-schema-name.util'; +import { type WorkflowAction } from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type'; @Command({ name: 'upgrade:1-3:add-next-step-ids-to-workflow-runs-trigger', @@ -22,7 +21,8 @@ export class AddNextStepIdsToWorkflowRunsTrigger extends ActiveOrSuspendedWorksp constructor( @InjectRepository(Workspace) protected readonly workspaceRepository: Repository, - private readonly workspaceDataSourceService: WorkspaceDataSourceService, + @InjectDataSource() + private readonly coreDataSource: DataSource, protected readonly twentyORMGlobalManager: TwentyORMGlobalManager, ) { super(workspaceRepository, twentyORMGlobalManager); @@ -31,12 +31,9 @@ export class AddNextStepIdsToWorkflowRunsTrigger extends ActiveOrSuspendedWorksp override async runOnWorkspace({ workspaceId, }: RunOnWorkspaceArgs): Promise { - const mainDataSource = - await this.workspaceDataSourceService.connectToMainDataSource(); - const schemaName = getWorkspaceSchemaName(workspaceId); - const workflowRuns = await mainDataSource.query( + const workflowRuns = await this.coreDataSource.query( `SELECT id, state FROM ${schemaName}."workflowRun"`, ); @@ -63,7 +60,7 @@ export class AddNextStepIdsToWorkflowRunsTrigger extends ActiveOrSuspendedWorksp }, }; - await mainDataSource.query( + await this.coreDataSource.query( `UPDATE ${schemaName}."workflowRun" SET state = $1::jsonb WHERE id = $2;`, [updatedState, workflowRun.id], ); diff --git a/packages/twenty-server/src/database/commands/upgrade-version-command/1-3/1-3-update-timestamp-column-type-in-workspace-schema.command.ts b/packages/twenty-server/src/database/commands/upgrade-version-command/1-3/1-3-update-timestamp-column-type-in-workspace-schema.command.ts index 9e0046429d..c289d6d6df 100644 --- a/packages/twenty-server/src/database/commands/upgrade-version-command/1-3/1-3-update-timestamp-column-type-in-workspace-schema.command.ts +++ b/packages/twenty-server/src/database/commands/upgrade-version-command/1-3/1-3-update-timestamp-column-type-in-workspace-schema.command.ts @@ -1,8 +1,8 @@ -import { InjectRepository } from '@nestjs/typeorm'; +import { InjectDataSource, InjectRepository } from '@nestjs/typeorm'; import { Command } from 'nest-commander'; import { FieldMetadataType } from 'twenty-shared/types'; -import { Repository } from 'typeorm'; +import { DataSource, Repository } from 'typeorm'; import { ActiveOrSuspendedWorkspacesMigrationCommandRunner, @@ -13,7 +13,6 @@ import { FieldMetadataEntity } from 'src/engine/metadata-modules/field-metadata/ import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.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'; -import { WorkspaceDataSourceService } from 'src/engine/workspace-datasource/workspace-datasource.service'; @Command({ name: 'upgrade:1-3:update-timestamp-column-type-in-workspace-schema', @@ -24,7 +23,8 @@ export class UpdateTimestampColumnTypeInWorkspaceSchemaCommand extends ActiveOrS constructor( @InjectRepository(Workspace) protected readonly workspaceRepository: Repository, - private readonly workspaceDataSourceService: WorkspaceDataSourceService, + @InjectDataSource() + private readonly coreDataSource: DataSource, protected readonly twentyORMGlobalManager: TwentyORMGlobalManager, @InjectRepository(FieldMetadataEntity) private readonly fieldMetadataRepository: Repository, @@ -43,16 +43,13 @@ export class UpdateTimestampColumnTypeInWorkspaceSchemaCommand extends ActiveOrS relations: ['object'], }); - const mainDataSource = - await this.workspaceDataSourceService.connectToMainDataSource(); - const schemaName = getWorkspaceSchemaName(workspaceId); for (const fieldMetadataItem of dateTimeFieldMetadataItems) { this.logger.log( `Updating column type for ${fieldMetadataItem.name} in ${schemaName}."${computeObjectTargetTable(fieldMetadataItem.object)}"`, ); - await mainDataSource.query( + await this.coreDataSource.query( `ALTER TABLE ${schemaName}."${computeObjectTargetTable(fieldMetadataItem.object)}" ALTER COLUMN "${fieldMetadataItem.name}" TYPE timestamptz(3);`, ); diff --git a/packages/twenty-server/src/database/commands/upgrade-version-command/1-5/1-5-add-positions-to-workflow-versions-and-workflow-runs.command.ts b/packages/twenty-server/src/database/commands/upgrade-version-command/1-5/1-5-add-positions-to-workflow-versions-and-workflow-runs.command.ts index 09916cff6f..0e7d6135d8 100644 --- a/packages/twenty-server/src/database/commands/upgrade-version-command/1-5/1-5-add-positions-to-workflow-versions-and-workflow-runs.command.ts +++ b/packages/twenty-server/src/database/commands/upgrade-version-command/1-5/1-5-add-positions-to-workflow-versions-and-workflow-runs.command.ts @@ -1,9 +1,9 @@ -import { InjectRepository } from '@nestjs/typeorm'; +import { InjectDataSource, InjectRepository } from '@nestjs/typeorm'; import Dagre from '@dagrejs/dagre'; import { Command, Option } from 'nest-commander'; import { isDefined } from 'twenty-shared/utils'; -import { Repository } from 'typeorm'; +import { DataSource, Repository } from 'typeorm'; import { v4 } from 'uuid'; import { @@ -14,7 +14,6 @@ import { import { Workspace } from 'src/engine/core-modules/workspace/workspace.entity'; import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager'; import { getWorkspaceSchemaName } from 'src/engine/workspace-datasource/utils/get-workspace-schema-name.util'; -import { WorkspaceDataSourceService } from 'src/engine/workspace-datasource/workspace-datasource.service'; import { type WorkflowVersionWorkspaceEntity } from 'src/modules/workflow/common/standard-objects/workflow-version.workspace-entity'; import { type WorkflowAction } from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type'; import { type WorkflowTrigger } from 'src/modules/workflow/workflow-trigger/types/workflow-trigger.type'; @@ -46,7 +45,8 @@ export class AddPositionsToWorkflowVersionsAndWorkflowRuns extends ActiveOrSuspe constructor( @InjectRepository(Workspace) protected readonly workspaceRepository: Repository, - private readonly workspaceDataSourceService: WorkspaceDataSourceService, + @InjectDataSource() + private readonly coreDataSource: DataSource, protected readonly twentyORMGlobalManager: TwentyORMGlobalManager, ) { super(workspaceRepository, twentyORMGlobalManager); @@ -142,12 +142,9 @@ export class AddPositionsToWorkflowVersionsAndWorkflowRuns extends ActiveOrSuspe }: { workspaceId: string; }) { - const mainDataSource = - await this.workspaceDataSourceService.connectToMainDataSource(); - const schemaName = getWorkspaceSchemaName(workspaceId); - const workflowRuns = await mainDataSource.query( + const workflowRuns = await this.coreDataSource.query( `SELECT id, state FROM ${schemaName}."workflowRun"`, ); @@ -168,7 +165,7 @@ export class AddPositionsToWorkflowVersionsAndWorkflowRuns extends ActiveOrSuspe }, }; - await mainDataSource.query( + await this.coreDataSource.query( `UPDATE ${schemaName}."workflowRun" SET state = $1::jsonb WHERE id = $2`, [updatedState, workflowRun.id], ); diff --git a/packages/twenty-server/src/database/typeorm/typeorm.module.ts b/packages/twenty-server/src/database/typeorm/typeorm.module.ts index 88d7a6924e..6a126886bc 100644 --- a/packages/twenty-server/src/database/typeorm/typeorm.module.ts +++ b/packages/twenty-server/src/database/typeorm/typeorm.module.ts @@ -2,16 +2,10 @@ import { Module } from '@nestjs/common'; import { TypeOrmModule } from '@nestjs/typeorm'; import { typeORMCoreModuleOptions } from 'src/database/typeorm/core/core.datasource'; -import { TwentyConfigModule } from 'src/engine/core-modules/twenty-config/twenty-config.module'; - -import { TypeORMService } from './typeorm.service'; @Module({ - imports: [ - TwentyConfigModule, - TypeOrmModule.forRoot(typeORMCoreModuleOptions), - ], - providers: [TypeORMService], - exports: [TypeORMService], + imports: [TypeOrmModule.forRoot(typeORMCoreModuleOptions)], + providers: [], + exports: [], }) export class TypeORMModule {} diff --git a/packages/twenty-server/src/database/typeorm/typeorm.service.ts b/packages/twenty-server/src/database/typeorm/typeorm.service.ts deleted file mode 100644 index ca50eb3587..0000000000 --- a/packages/twenty-server/src/database/typeorm/typeorm.service.ts +++ /dev/null @@ -1,73 +0,0 @@ -import { - Injectable, - Logger, - type OnModuleDestroy, - type OnModuleInit, -} from '@nestjs/common'; - -import { DataSource } from 'typeorm'; - -import { TwentyConfigService } from 'src/engine/core-modules/twenty-config/twenty-config.service'; - -@Injectable() -export class TypeORMService implements OnModuleInit, OnModuleDestroy { - private mainDataSource: DataSource; - private readonly logger = new Logger(TypeORMService.name); - - constructor(private readonly twentyConfigService: TwentyConfigService) { - const isJest = process.argv.some((arg) => arg.includes('jest')); - - this.mainDataSource = new DataSource({ - url: twentyConfigService.get('PG_DATABASE_URL'), - type: 'postgres', - logging: twentyConfigService.getLoggingConfig(), - schema: 'core', - entities: [ - `${isJest ? '' : 'dist/'}src/engine/core-modules/**/*.entity{.ts,.js}`, - `${isJest ? '' : 'dist/'}src/engine/metadata-modules/**/*.entity{.ts,.js}`, - ], - metadataTableName: '_typeorm_generated_columns_and_materialized_views', - ssl: twentyConfigService.get('PG_SSL_ALLOW_SELF_SIGNED') - ? { - rejectUnauthorized: false, - } - : undefined, - extra: { - query_timeout: 10000, - }, - }); - } - - public getMainDataSource(): DataSource { - return this.mainDataSource; - } - - public async createSchema(schemaName: string): Promise { - const queryRunner = this.mainDataSource.createQueryRunner(); - - await queryRunner.createSchema(schemaName, true); - - await queryRunner.release(); - - return schemaName; - } - - public async deleteSchema(schemaName: string) { - const queryRunner = this.mainDataSource.createQueryRunner(); - - await queryRunner.dropSchema(schemaName, true, true); - - await queryRunner.release(); - } - - async onModuleInit() { - // Init main data source "default" schema - await this.mainDataSource.initialize(); - } - - async onModuleDestroy() { - // Destroy main data source "default" schema - this.logger.log('Destroying main data source'); - await this.mainDataSource.destroy(); - } -} diff --git a/packages/twenty-server/src/engine/core-modules/user-workspace/user-workspace.service.spec.ts b/packages/twenty-server/src/engine/core-modules/user-workspace/user-workspace.service.spec.ts index 2e761f7af1..8b779d4838 100644 --- a/packages/twenty-server/src/engine/core-modules/user-workspace/user-workspace.service.spec.ts +++ b/packages/twenty-server/src/engine/core-modules/user-workspace/user-workspace.service.spec.ts @@ -5,7 +5,6 @@ import { type DataSource, type Repository } from 'typeorm'; import { FileFolder } from 'src/engine/core-modules/file/interfaces/file-folder.interface'; -import { TypeORMService } from 'src/database/typeorm/typeorm.service'; import { type ApprovedAccessDomain } from 'src/engine/core-modules/approved-access-domain/approved-access-domain.entity'; import { ApprovedAccessDomainService } from 'src/engine/core-modules/approved-access-domain/services/approved-access-domain.service'; import { AuthException } from 'src/engine/core-modules/auth/auth.exception'; @@ -33,7 +32,6 @@ describe('UserWorkspaceService', () => { let service: UserWorkspaceService; let userWorkspaceRepository: Repository; let userRepository: Repository; - let typeORMService: TypeORMService; let workspaceInvitationService: WorkspaceInvitationService; let approvedAccessDomainService: ApprovedAccessDomainService; let twentyORMGlobalManager: TwentyORMGlobalManager; @@ -75,12 +73,6 @@ describe('UserWorkspaceService', () => { getLastDataSourceMetadataFromWorkspaceIdOrFail: jest.fn(), }, }, - { - provide: TypeORMService, - useValue: { - getMainDataSource: jest.fn(), - }, - }, { provide: WorkspaceInvitationService, useValue: { @@ -142,7 +134,6 @@ describe('UserWorkspaceService', () => { fileService = module.get(FileService); userWorkspaceRepository = module.get(getRepositoryToken(UserWorkspace)); userRepository = module.get(getRepositoryToken(User)); - typeORMService = module.get(TypeORMService); workspaceInvitationService = module.get( WorkspaceInvitationService, ); @@ -346,9 +337,6 @@ describe('UserWorkspaceService', () => { find: jest.fn().mockResolvedValue(workspaceMember), }; - jest - .spyOn(typeORMService, 'getMainDataSource') - .mockReturnValue(mainDataSource); jest .spyOn(mainDataSource, 'query') .mockResolvedValueOnce(undefined) diff --git a/packages/twenty-server/src/engine/core-modules/user/user.module.ts b/packages/twenty-server/src/engine/core-modules/user/user.module.ts index 745c775492..5942b997a5 100644 --- a/packages/twenty-server/src/engine/core-modules/user/user.module.ts +++ b/packages/twenty-server/src/engine/core-modules/user/user.module.ts @@ -5,7 +5,6 @@ import { NestjsQueryGraphQLModule } from '@ptc-org/nestjs-query-graphql'; import { NestjsQueryTypeOrmModule } from '@ptc-org/nestjs-query-typeorm'; import { TypeORMModule } from 'src/database/typeorm/typeorm.module'; -import { TypeORMService } from 'src/database/typeorm/typeorm.service'; import { AuditModule } from 'src/engine/core-modules/audit/audit.module'; import { FeatureFlagModule } from 'src/engine/core-modules/feature-flag/feature-flag.module'; import { FileUploadModule } from 'src/engine/core-modules/file/file-upload/file-upload.module'; @@ -53,11 +52,6 @@ import { UserService } from './services/user.service'; UserWorkspaceModule, ], exports: [UserService, WorkspaceMemberTranspiler], - providers: [ - UserService, - UserResolver, - TypeORMService, - WorkspaceMemberTranspiler, - ], + providers: [UserService, UserResolver, WorkspaceMemberTranspiler], }) export class UserModule {} diff --git a/packages/twenty-server/src/engine/metadata-modules/object-metadata/object-metadata.service.ts b/packages/twenty-server/src/engine/metadata-modules/object-metadata/object-metadata.service.ts index d8303c0a2a..5fc6a6dbbc 100644 --- a/packages/twenty-server/src/engine/metadata-modules/object-metadata/object-metadata.service.ts +++ b/packages/twenty-server/src/engine/metadata-modules/object-metadata/object-metadata.service.ts @@ -1,11 +1,12 @@ import { Injectable } from '@nestjs/common'; -import { InjectRepository } from '@nestjs/typeorm'; +import { InjectDataSource, InjectRepository } from '@nestjs/typeorm'; import { type Query, type QueryOptions } from '@ptc-org/nestjs-query-core'; import { TypeOrmQueryService } from '@ptc-org/nestjs-query-typeorm'; import { FieldMetadataType } from 'twenty-shared/types'; import { capitalize, isDefined } from 'twenty-shared/utils'; import { + DataSource, In, Repository, type FindManyOptions, @@ -47,7 +48,6 @@ import { WorkspaceMetadataVersionService } from 'src/engine/metadata-modules/wor import { WorkspacePermissionsCacheService } from 'src/engine/metadata-modules/workspace-permissions-cache/workspace-permissions-cache.service'; import { computeObjectTargetTable } from 'src/engine/utils/compute-object-target-table.util'; import { isFieldMetadataEntityOfType } from 'src/engine/utils/is-field-metadata-of-type.util'; -import { WorkspaceDataSourceService } from 'src/engine/workspace-datasource/workspace-datasource.service'; import { WorkspaceMigrationRunnerService } from 'src/engine/workspace-manager/workspace-migration-runner/workspace-migration-runner.service'; import { CUSTOM_OBJECT_STANDARD_FIELD_IDS } from 'src/engine/workspace-manager/workspace-sync-metadata/constants/standard-field-ids'; import { isSearchableFieldType } from 'src/engine/workspace-manager/workspace-sync-metadata/utils/is-searchable-field.util'; @@ -72,7 +72,8 @@ export class ObjectMetadataService extends TypeOrmQueryService { - const mainDataSource = - await this.workspaceDataSourceService.connectToMainDataSource(); - const queryRunner = mainDataSource.createQueryRunner(); + const queryRunner = this.coreDataSource.createQueryRunner(); await queryRunner.connect(); await queryRunner.startTransaction(); @@ -467,9 +464,7 @@ export class ObjectMetadataService extends TypeOrmQueryService { diff --git a/packages/twenty-server/src/engine/metadata-modules/remote-server/remote-table/distant-table/distant-table.service.ts b/packages/twenty-server/src/engine/metadata-modules/remote-server/remote-table/distant-table/distant-table.service.ts index 047ec83ff5..95b261736c 100644 --- a/packages/twenty-server/src/engine/metadata-modules/remote-server/remote-table/distant-table/distant-table.service.ts +++ b/packages/twenty-server/src/engine/metadata-modules/remote-server/remote-table/distant-table/distant-table.service.ts @@ -1,6 +1,7 @@ import { Injectable } from '@nestjs/common'; +import { InjectDataSource } from '@nestjs/typeorm'; -import { type EntityManager } from 'typeorm'; +import { DataSource, type EntityManager } from 'typeorm'; import { v4 } from 'uuid'; import { @@ -15,12 +16,12 @@ import { type DistantTables } from 'src/engine/metadata-modules/remote-server/re import { STRIPE_DISTANT_TABLES } from 'src/engine/metadata-modules/remote-server/remote-table/distant-table/utils/stripe-distant-tables.util'; import { type PostgresTableSchemaColumn } from 'src/engine/metadata-modules/remote-server/types/postgres-table-schema-column'; import { isQueryTimeoutError } from 'src/engine/utils/query-timeout.util'; -import { WorkspaceDataSourceService } from 'src/engine/workspace-datasource/workspace-datasource.service'; @Injectable() export class DistantTableService { constructor( - private readonly workspaceDataSourceService: WorkspaceDataSourceService, + @InjectDataSource() + private readonly coreDataSource: DataSource, ) {} public async fetchDistantTables( @@ -68,11 +69,8 @@ export class DistantTableService { const tmpSchemaId = v4(); const tmpSchemaName = `${workspaceId}_${remoteServer.id}_${tmpSchemaId}`; - const mainDataSource = - await this.workspaceDataSourceService.connectToMainDataSource(); - try { - const distantTables = await mainDataSource.transaction( + const distantTables = await this.coreDataSource.transaction( async (entityManager: EntityManager) => { await entityManager.query(`CREATE SCHEMA "${tmpSchemaName}"`); diff --git a/packages/twenty-server/src/engine/metadata-modules/remote-server/remote-table/foreign-table/foreign-table.service.ts b/packages/twenty-server/src/engine/metadata-modules/remote-server/remote-table/foreign-table/foreign-table.service.ts index 6acd8e9646..54b0dfa22b 100644 --- a/packages/twenty-server/src/engine/metadata-modules/remote-server/remote-table/foreign-table/foreign-table.service.ts +++ b/packages/twenty-server/src/engine/metadata-modules/remote-server/remote-table/foreign-table/foreign-table.service.ts @@ -1,4 +1,7 @@ import { Injectable } from '@nestjs/common'; +import { InjectDataSource } from '@nestjs/typeorm'; + +import { DataSource } from 'typeorm'; import { type RemoteServerEntity, @@ -31,18 +34,17 @@ export class ForeignTableService { private readonly workspaceMigrationRunnerService: WorkspaceMigrationRunnerService, private readonly workspaceDataSourceService: WorkspaceDataSourceService, private readonly workspaceMetadataVersionService: WorkspaceMetadataVersionService, + @InjectDataSource() + private readonly coreDataSource: DataSource, ) {} public async fetchForeignTableNamesWithinWorkspace( _workspaceId: string, foreignDataWrapperId: string, ): Promise { - const mainDataSource = - await this.workspaceDataSourceService.connectToMainDataSource(); - return ( ( - await mainDataSource.query( + await this.coreDataSource.query( `SELECT foreign_table_name, foreign_server_name FROM information_schema.foreign_tables WHERE foreign_server_name = $1`, [foreignDataWrapperId], ) diff --git a/packages/twenty-server/src/engine/metadata-modules/remote-server/remote-table/remote-table.service.ts b/packages/twenty-server/src/engine/metadata-modules/remote-server/remote-table/remote-table.service.ts index 34eb6ede35..c16384c555 100644 --- a/packages/twenty-server/src/engine/metadata-modules/remote-server/remote-table/remote-table.service.ts +++ b/packages/twenty-server/src/engine/metadata-modules/remote-server/remote-table/remote-table.service.ts @@ -1,9 +1,9 @@ import { Logger } from '@nestjs/common'; -import { InjectRepository } from '@nestjs/typeorm'; +import { InjectDataSource, InjectRepository } from '@nestjs/typeorm'; import isEmpty from 'lodash.isempty'; import { plural } from 'pluralize'; -import { Repository } from 'typeorm'; +import { DataSource, Repository } from 'typeorm'; import { DataSourceService } from 'src/engine/metadata-modules/data-source/data-source.service'; import { type CreateFieldInput } from 'src/engine/metadata-modules/field-metadata/dtos/create-field.input'; @@ -63,6 +63,8 @@ export class RemoteTableService { private readonly foreignTableService: ForeignTableService, private readonly workspaceDataSourceService: WorkspaceDataSourceService, private readonly remoteTableSchemaUpdateService: RemoteTableSchemaUpdateService, + @InjectDataSource() + private readonly coreDataSource: DataSource, ) {} public async findDistantTablesWithStatus( @@ -182,14 +184,11 @@ export class RemoteTableService { workspaceId, ); - const mainDataSource = - await this.workspaceDataSourceService.connectToMainDataSource(); - const { baseName: localTableBaseName, suffix: localTableSuffix } = await getRemoteTableLocalName( input.name, dataSourceMetatada.schema, - mainDataSource, + this.coreDataSource, ); const localTableName = localTableSuffix diff --git a/packages/twenty-server/src/engine/metadata-modules/remote-server/remote-table/utils/get-remote-table-local-name.util.ts b/packages/twenty-server/src/engine/metadata-modules/remote-server/remote-table/utils/get-remote-table-local-name.util.ts index fb92cb6d50..52411a3ee8 100644 --- a/packages/twenty-server/src/engine/metadata-modules/remote-server/remote-table/utils/get-remote-table-local-name.util.ts +++ b/packages/twenty-server/src/engine/metadata-modules/remote-server/remote-table/utils/get-remote-table-local-name.util.ts @@ -17,11 +17,10 @@ type RemoteTableLocalName = { const isNameAvailable = async ( tableName: string, workspaceSchemaName: string, - workspaceDataSource: DataSource, + coreDataSource: DataSource, ) => { - // TO DO workspaceDataSource.query method is not allowed, this will throw const numberOfTablesWithSameName = +( - await workspaceDataSource.query( + await coreDataSource.query( `SELECT count(table_name) FROM information_schema.tables WHERE table_name LIKE '${tableName}' AND table_schema IN ('core', '${workspaceSchemaName}')`, ) )[0].count; @@ -32,13 +31,13 @@ const isNameAvailable = async ( export const getRemoteTableLocalName = async ( distantTableName: string, workspaceSchemaName: string, - workspaceDataSource: DataSource, + coreDataSource: DataSource, ): Promise => { const baseName = singular(camelCase(distantTableName)); const isBaseNameValid = await isNameAvailable( baseName, workspaceSchemaName, - workspaceDataSource, + coreDataSource, ); if (isBaseNameValid) { @@ -50,7 +49,7 @@ export const getRemoteTableLocalName = async ( const isNameWithSuffixValid = await isNameAvailable( name, workspaceSchemaName, - workspaceDataSource, + coreDataSource, ); if (isNameWithSuffixValid) { diff --git a/packages/twenty-server/src/engine/workspace-datasource/workspace-datasource.service.ts b/packages/twenty-server/src/engine/workspace-datasource/workspace-datasource.service.ts index 263defed78..ef1fd7a80a 100644 --- a/packages/twenty-server/src/engine/workspace-datasource/workspace-datasource.service.ts +++ b/packages/twenty-server/src/engine/workspace-datasource/workspace-datasource.service.ts @@ -1,8 +1,8 @@ import { Injectable } from '@nestjs/common'; +import { InjectDataSource } from '@nestjs/typeorm'; import { type DataSource, type EntityManager } from 'typeorm'; -import { TypeORMService } from 'src/database/typeorm/typeorm.service'; import { DataSourceService } from 'src/engine/metadata-modules/data-source/data-source.service'; import { PermissionsException, @@ -14,26 +14,10 @@ import { getWorkspaceSchemaName } from 'src/engine/workspace-datasource/utils/ge export class WorkspaceDataSourceService { constructor( private readonly dataSourceService: DataSourceService, - private readonly typeormService: TypeORMService, + @InjectDataSource() + private readonly coreDataSource: DataSource, ) {} - /** - * - * Connect to the workspace data source - * - * @param workspaceId - * @returns - */ - public async connectToMainDataSource(): Promise { - const dataSource = this.typeormService.getMainDataSource(); - - if (!dataSource) { - throw new Error(`Could not connect to workspace data source`); - } - - return dataSource; - } - public async checkSchemaExists(workspaceId: string) { const dataSource = await this.dataSourceService.getDataSourcesMetadataFromWorkspaceId( @@ -53,7 +37,13 @@ export class WorkspaceDataSourceService { public async createWorkspaceDBSchema(workspaceId: string): Promise { const schemaName = getWorkspaceSchemaName(workspaceId); - return await this.typeormService.createSchema(schemaName); + const queryRunner = this.coreDataSource.createQueryRunner(); + + await queryRunner.createSchema(schemaName, true); + + await queryRunner.release(); + + return schemaName; } /** @@ -66,7 +56,11 @@ export class WorkspaceDataSourceService { public async deleteWorkspaceDBSchema(workspaceId: string): Promise { const schemaName = getWorkspaceSchemaName(workspaceId); - return await this.typeormService.deleteSchema(schemaName); + const queryRunner = this.coreDataSource.createQueryRunner(); + + await queryRunner.dropSchema(schemaName, true); + + await queryRunner.release(); } public async executeRawQuery( diff --git a/packages/twenty-server/src/engine/workspace-manager/__tests__/workspace-manager.service.spec.ts b/packages/twenty-server/src/engine/workspace-manager/__tests__/workspace-manager.service.spec.ts index cd7fd615d2..77de1c2961 100644 --- a/packages/twenty-server/src/engine/workspace-manager/__tests__/workspace-manager.service.spec.ts +++ b/packages/twenty-server/src/engine/workspace-manager/__tests__/workspace-manager.service.spec.ts @@ -19,6 +19,7 @@ import { RoleService } from 'src/engine/metadata-modules/role/role.service'; import { UserRoleService } from 'src/engine/metadata-modules/user-role/user-role.service'; import { WorkspaceMigrationEntity } from 'src/engine/metadata-modules/workspace-migration/workspace-migration.entity'; import { WorkspaceMigrationService } from 'src/engine/metadata-modules/workspace-migration/workspace-migration.service'; +import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager'; import { WorkspaceDataSourceService } from 'src/engine/workspace-datasource/workspace-datasource.service'; import { WorkspaceManagerService } from 'src/engine/workspace-manager/workspace-manager.service'; import { WorkspaceSyncMetadataService } from 'src/engine/workspace-manager/workspace-sync-metadata/workspace-sync-metadata.service'; @@ -124,6 +125,14 @@ describe('WorkspaceManagerService', () => { .mockResolvedValue({ id: 'mock-agent-id' }), }, }, + { + provide: TwentyORMGlobalManager, + useValue: { + getDataSourceForWorkspace: jest.fn().mockResolvedValue({ + transaction: jest.fn(), + }), + }, + }, ], }).compile(); diff --git a/packages/twenty-server/src/engine/workspace-manager/dev-seeder/core/services/dev-seeder-permissions.service.ts b/packages/twenty-server/src/engine/workspace-manager/dev-seeder/core/services/dev-seeder-permissions.service.ts index 16e6a40a09..f2044bab1d 100644 --- a/packages/twenty-server/src/engine/workspace-manager/dev-seeder/core/services/dev-seeder-permissions.service.ts +++ b/packages/twenty-server/src/engine/workspace-manager/dev-seeder/core/services/dev-seeder-permissions.service.ts @@ -1,10 +1,9 @@ import { Injectable, Logger } from '@nestjs/common'; -import { InjectRepository } from '@nestjs/typeorm'; +import { InjectDataSource, InjectRepository } from '@nestjs/typeorm'; import { WorkspaceActivationStatus } from 'twenty-shared/workspace'; -import { Repository } from 'typeorm'; +import { DataSource, Repository } from 'typeorm'; -import { TypeORMService } from 'src/database/typeorm/typeorm.service'; import { Workspace } from 'src/engine/core-modules/workspace/workspace.entity'; import { ObjectMetadataEntity } from 'src/engine/metadata-modules/object-metadata/object-metadata.entity'; import { FieldPermissionService } from 'src/engine/metadata-modules/object-permission/field-permission/field-permission.service'; @@ -33,9 +32,10 @@ export class DevSeederPermissionsService { private readonly objectMetadataRepository: Repository, @InjectRepository(RoleEntity) private readonly roleRepository: Repository, - private readonly typeORMService: TypeORMService, private readonly workspacePermissionsCacheService: WorkspacePermissionsCacheService, private readonly fieldPermissionService: FieldPermissionService, + @InjectDataSource() + private readonly coreDataSource: DataSource, ) {} public async initPermissions(workspaceId: string) { @@ -52,39 +52,33 @@ export class DevSeederPermissionsService { ); } - const dataSource = this.typeORMService.getMainDataSource(); - - if (dataSource) { - try { - await dataSource - .createQueryBuilder() - .insert() - .into('core.roleTargets', ['roleId', 'apiKeyId', 'workspaceId']) - .orIgnore() - .values([ - { - roleId: adminRole.id, - apiKeyId: API_KEY_DATA_SEED_IDS.ID_1, - workspaceId: workspaceId, - }, - ]) - .execute(); - - await this.workspacePermissionsCacheService.recomputeApiKeyRoleMapCache( + try { + await this.coreDataSource + .createQueryBuilder() + .insert() + .into('core.roleTargets', ['roleId', 'apiKeyId', 'workspaceId']) + .orIgnore() + .values([ { - workspaceId, + roleId: adminRole.id, + apiKeyId: API_KEY_DATA_SEED_IDS.ID_1, + workspaceId: workspaceId, }, - ); - await this.workspacePermissionsCacheService.recomputeUserWorkspaceRoleMapCache( - { - workspaceId, - }, - ); - } catch (error) { - this.logger.error( - `Could not assign role to test API key: ${error.message}`, - ); - } + ]) + .execute(); + + await this.workspacePermissionsCacheService.recomputeApiKeyRoleMapCache({ + workspaceId, + }); + await this.workspacePermissionsCacheService.recomputeUserWorkspaceRoleMapCache( + { + workspaceId, + }, + ); + } catch (error) { + this.logger.error( + `Could not assign role to test API key: ${error.message}`, + ); } let adminUserWorkspaceId: string | undefined; @@ -137,13 +131,10 @@ export class DevSeederPermissionsService { workspaceId, }); - await this.typeORMService - .getMainDataSource() - ?.getRepository(Workspace) - .update(workspaceId, { - defaultRoleId: memberRole.id, - activationStatus: WorkspaceActivationStatus.ACTIVE, - }); + await this.coreDataSource.getRepository(Workspace).update(workspaceId, { + defaultRoleId: memberRole.id, + activationStatus: WorkspaceActivationStatus.ACTIVE, + }); if (memberUserWorkspaceIds) { for (const memberUserWorkspaceId of memberUserWorkspaceIds) { diff --git a/packages/twenty-server/src/engine/workspace-manager/dev-seeder/data/services/dev-seeder-data.service.ts b/packages/twenty-server/src/engine/workspace-manager/dev-seeder/data/services/dev-seeder-data.service.ts index 9e38bf37d2..086106d014 100644 --- a/packages/twenty-server/src/engine/workspace-manager/dev-seeder/data/services/dev-seeder-data.service.ts +++ b/packages/twenty-server/src/engine/workspace-manager/dev-seeder/data/services/dev-seeder-data.service.ts @@ -1,10 +1,12 @@ import { Injectable } from '@nestjs/common'; +import { InjectDataSource } from '@nestjs/typeorm'; + +import { DataSource } from 'typeorm'; import { ObjectMetadataService } from 'src/engine/metadata-modules/object-metadata/object-metadata.service'; import { type WorkspaceEntityManager } from 'src/engine/twenty-orm/entity-manager/workspace-entity-manager'; import { computeTableName } from 'src/engine/utils/compute-table-name.util'; import { shouldSeedWorkspaceFavorite } from 'src/engine/utils/should-seed-workspace-favorite'; -import { WorkspaceDataSourceService } from 'src/engine/workspace-datasource/workspace-datasource.service'; import { CALENDAR_CHANNEL_DATA_SEED_COLUMNS, CALENDAR_CHANNEL_DATA_SEEDS, @@ -196,7 +198,8 @@ const RECORD_SEEDS_CONFIGS = [ @Injectable() export class DevSeederDataService { constructor( - private readonly workspaceDataSourceService: WorkspaceDataSourceService, + @InjectDataSource() + private readonly coreDataSource: DataSource, private readonly objectMetadataService: ObjectMetadataService, private readonly timelineActivitySeederService: TimelineActivitySeederService, ) {} @@ -208,17 +211,10 @@ export class DevSeederDataService { schemaName: string; workspaceId: string; }) { - const mainDataSource = - await this.workspaceDataSourceService.connectToMainDataSource(); - - if (!mainDataSource) { - throw new Error('Could not connect to main data source'); - } - const objectMetadataItems = await this.objectMetadataService.findManyWithinWorkspace(workspaceId); - await mainDataSource.transaction( + await this.coreDataSource.transaction( async (entityManager: WorkspaceEntityManager) => { for (const recordSeedsConfig of RECORD_SEEDS_CONFIGS) { const objectMetadata = objectMetadataItems.find( diff --git a/packages/twenty-server/src/engine/workspace-manager/dev-seeder/metadata/services/dev-seeder-metadata.service.ts b/packages/twenty-server/src/engine/workspace-manager/dev-seeder/metadata/services/dev-seeder-metadata.service.ts index 76d42fbaca..461eeab92b 100644 --- a/packages/twenty-server/src/engine/workspace-manager/dev-seeder/metadata/services/dev-seeder-metadata.service.ts +++ b/packages/twenty-server/src/engine/workspace-manager/dev-seeder/metadata/services/dev-seeder-metadata.service.ts @@ -1,8 +1,8 @@ import { Injectable } from '@nestjs/common'; +import { InjectDataSource } from '@nestjs/typeorm'; -import { isDefined } from 'class-validator'; +import { DataSource } from 'typeorm'; -import { TypeORMService } from 'src/database/typeorm/typeorm.service'; import { type DataSourceEntity } from 'src/engine/metadata-modules/data-source/data-source.entity'; import { FieldMetadataService } from 'src/engine/metadata-modules/field-metadata/services/field-metadata.service'; import { ObjectMetadataService } from 'src/engine/metadata-modules/object-metadata/object-metadata.service'; @@ -26,7 +26,8 @@ export class DevSeederMetadataService { constructor( private readonly objectMetadataService: ObjectMetadataService, private readonly fieldMetadataService: FieldMetadataService, - private readonly typeORMService: TypeORMService, + @InjectDataSource() + private readonly coreDataSource: DataSource, ) {} private readonly workspaceConfigs: Record< @@ -152,15 +153,13 @@ export class DevSeederMetadataService { } private async seedCoreViews(workspaceId: string): Promise { - const mainDataSource = this.typeORMService.getMainDataSource(); - - if (!isDefined(mainDataSource)) { - throw new Error('Could not connect to main data source'); - } - const createdObjectMetadata = await this.objectMetadataService.findManyWithinWorkspace(workspaceId); - await prefillCoreViews(mainDataSource, workspaceId, createdObjectMetadata); + await prefillCoreViews( + this.coreDataSource, + workspaceId, + createdObjectMetadata, + ); } } diff --git a/packages/twenty-server/src/engine/workspace-manager/dev-seeder/services/dev-seeder.service.ts b/packages/twenty-server/src/engine/workspace-manager/dev-seeder/services/dev-seeder.service.ts index 6f4938420f..4f44e9aecd 100644 --- a/packages/twenty-server/src/engine/workspace-manager/dev-seeder/services/dev-seeder.service.ts +++ b/packages/twenty-server/src/engine/workspace-manager/dev-seeder/services/dev-seeder.service.ts @@ -1,6 +1,8 @@ import { Injectable } from '@nestjs/common'; +import { InjectDataSource } from '@nestjs/typeorm'; + +import { DataSource } from 'typeorm'; -import { TypeORMService } from 'src/database/typeorm/typeorm.service'; import { FeatureFlagService } from 'src/engine/core-modules/feature-flag/services/feature-flag.service'; import { TwentyConfigService } from 'src/engine/core-modules/twenty-config/twenty-config.service'; import { DataSourceService } from 'src/engine/metadata-modules/data-source/data-source.service'; @@ -15,7 +17,6 @@ import { WorkspaceSyncMetadataService } from 'src/engine/workspace-manager/works @Injectable() export class DevSeederService { constructor( - private readonly typeORMService: TypeORMService, private readonly workspaceCacheStorageService: WorkspaceCacheStorageService, private readonly twentyConfigService: TwentyConfigService, private readonly workspaceDataSourceService: WorkspaceDataSourceService, @@ -25,20 +26,16 @@ export class DevSeederService { private readonly devSeederMetadataService: DevSeederMetadataService, private readonly devSeederPermissionsService: DevSeederPermissionsService, private readonly devSeederDataService: DevSeederDataService, + @InjectDataSource() + private readonly coreDataSource: DataSource, ) {} public async seedDev(workspaceId: string): Promise { - const mainDataSource = this.typeORMService.getMainDataSource(); - - if (!mainDataSource) { - throw new Error('Could not connect to workspace data source'); - } - const isBillingEnabled = this.twentyConfigService.get('IS_BILLING_ENABLED'); const appVersion = this.twentyConfigService.get('APP_VERSION'); await seedCoreSchema({ - dataSource: mainDataSource, + dataSource: this.coreDataSource, workspaceId, seedBilling: isBillingEnabled, appVersion, diff --git a/packages/twenty-server/src/engine/workspace-manager/standard-objects-prefill-data/prefill-core-views.ts b/packages/twenty-server/src/engine/workspace-manager/standard-objects-prefill-data/prefill-core-views.ts index a37a10ca25..90b712c898 100644 --- a/packages/twenty-server/src/engine/workspace-manager/standard-objects-prefill-data/prefill-core-views.ts +++ b/packages/twenty-server/src/engine/workspace-manager/standard-objects-prefill-data/prefill-core-views.ts @@ -29,7 +29,7 @@ import { ViewOpenRecordInType } from 'src/modules/view/standard-objects/view.wor import { convertViewFilterOperandToCoreOperand } from 'src/modules/view/utils/convert-view-filter-operand-to-core-operand.util'; export const prefillCoreViews = async ( - dataSource: DataSource, + coreDataSource: DataSource, workspaceId: string, objectMetadataItems: ObjectMetadataEntity[], featureFlags?: Record, @@ -52,7 +52,7 @@ export const prefillCoreViews = async ( views.push(dashboardsAllView(objectMetadataItems, true)); } - const queryRunner = dataSource.createQueryRunner(); + const queryRunner = coreDataSource.createQueryRunner(); await queryRunner.connect(); diff --git a/packages/twenty-server/src/engine/workspace-manager/standard-objects-prefill-data/standard-objects-prefill-data.ts b/packages/twenty-server/src/engine/workspace-manager/standard-objects-prefill-data/standard-objects-prefill-data.ts index e2ffd39c47..e04208f1e2 100644 --- a/packages/twenty-server/src/engine/workspace-manager/standard-objects-prefill-data/standard-objects-prefill-data.ts +++ b/packages/twenty-server/src/engine/workspace-manager/standard-objects-prefill-data/standard-objects-prefill-data.ts @@ -1,6 +1,5 @@ -import { type DataSource } from 'typeorm'; - import { type ObjectMetadataEntity } from 'src/engine/metadata-modules/object-metadata/object-metadata.entity'; +import { type WorkspaceDataSource } from 'src/engine/twenty-orm/datasource/workspace.datasource'; import { type WorkspaceEntityManager } from 'src/engine/twenty-orm/entity-manager/workspace-entity-manager'; import { shouldSeedWorkspaceFavorite } from 'src/engine/utils/should-seed-workspace-favorite'; import { prefillCompanies } from 'src/engine/workspace-manager/standard-objects-prefill-data/prefill-companies'; @@ -10,38 +9,40 @@ import { prefillWorkflows } from 'src/engine/workspace-manager/standard-objects- import { prefillWorkspaceFavorites } from 'src/engine/workspace-manager/standard-objects-prefill-data/prefill-workspace-favorites'; export const standardObjectsPrefillData = async ( - mainDataSource: DataSource, + workspaceDataSource: WorkspaceDataSource, schemaName: string, objectMetadataItems: ObjectMetadataEntity[], featureFlags?: Record, ) => { - mainDataSource.transaction(async (entityManager: WorkspaceEntityManager) => { - await prefillCompanies(entityManager, schemaName); + workspaceDataSource.transaction( + async (entityManager: WorkspaceEntityManager) => { + await prefillCompanies(entityManager, schemaName); - await prefillPeople(entityManager, schemaName); + await prefillPeople(entityManager, schemaName); - await prefillWorkflows(entityManager, schemaName, objectMetadataItems); + await prefillWorkflows(entityManager, schemaName, objectMetadataItems); - const viewDefinitionsWithId = await prefillViews( - entityManager, - schemaName, - objectMetadataItems, - featureFlags, - ); + const viewDefinitionsWithId = await prefillViews( + entityManager, + schemaName, + objectMetadataItems, + featureFlags, + ); - await prefillWorkspaceFavorites( - viewDefinitionsWithId - .filter( - (view) => - view.key === 'INDEX' && - shouldSeedWorkspaceFavorite( - view.objectMetadataId, - objectMetadataItems, - ), - ) - .map((view) => view.id), - entityManager, - schemaName, - ); - }); + await prefillWorkspaceFavorites( + viewDefinitionsWithId + .filter( + (view) => + view.key === 'INDEX' && + shouldSeedWorkspaceFavorite( + view.objectMetadataId, + objectMetadataItems, + ), + ) + .map((view) => view.id), + entityManager, + schemaName, + ); + }, + ); }; diff --git a/packages/twenty-server/src/engine/workspace-manager/workspace-health/services/database-structure.service.ts b/packages/twenty-server/src/engine/workspace-manager/workspace-health/services/database-structure.service.ts index 14636c513b..6e3edead47 100644 --- a/packages/twenty-server/src/engine/workspace-manager/workspace-health/services/database-structure.service.ts +++ b/packages/twenty-server/src/engine/workspace-manager/workspace-health/services/database-structure.service.ts @@ -1,8 +1,9 @@ import { Injectable } from '@nestjs/common'; +import { InjectDataSource } from '@nestjs/typeorm'; -import { type ColumnType } from 'typeorm'; -import { type ColumnMetadata } from 'typeorm/metadata/ColumnMetadata'; import { FieldMetadataType } from 'twenty-shared/types'; +import { DataSource, type ColumnType } from 'typeorm'; +import { type ColumnMetadata } from 'typeorm/metadata/ColumnMetadata'; import { type FieldMetadataDefaultValue, @@ -13,7 +14,6 @@ import { type WorkspaceTableStructureResult, } from 'src/engine/workspace-manager/workspace-health/interfaces/workspace-table-definition.interface'; -import { TypeORMService } from 'src/database/typeorm/typeorm.service'; import { compositeTypeDefinitions } from 'src/engine/metadata-modules/field-metadata/composite-types'; import { type FieldMetadataDefaultValueFunctionNames } from 'src/engine/metadata-modules/field-metadata/dtos/default-value.input'; import { type FieldMetadataEntity } from 'src/engine/metadata-modules/field-metadata/field-metadata.entity'; @@ -26,14 +26,16 @@ import { isRelationFieldMetadataType } from 'src/engine/utils/is-relation-field- @Injectable() export class DatabaseStructureService { - constructor(private readonly typeORMService: TypeORMService) {} + constructor( + @InjectDataSource() + private readonly coreDataSource: DataSource, + ) {} async getWorkspaceTableColumns( schemaName: string, tableName: string, ): Promise { - const mainDataSource = this.typeORMService.getMainDataSource(); - const results = await mainDataSource.query< + const results = await this.coreDataSource.query< WorkspaceTableStructureResult[] >(` WITH foreign_keys AS ( @@ -150,8 +152,6 @@ export class DatabaseStructureService { } getPostgresDataTypes(fieldMetadata: FieldMetadataEntity): string[] { - const mainDataSource = this.typeORMService.getMainDataSource(); - const normalizer = ( type: FieldMetadataType, isArray: boolean | undefined, @@ -166,7 +166,7 @@ export class DatabaseStructureService { return `${objectName}_${columnName}_enum${isArray ? '[]' : ''}`; } - return mainDataSource.driver.normalizeType({ + return this.coreDataSource.driver.normalizeType({ type: typeORMType, }); }; @@ -201,7 +201,6 @@ export class DatabaseStructureService { getFieldMetadataTypeFromPostgresDataType( postgresDataType: string, ): FieldMetadataType | null { - const mainDataSource = this.typeORMService.getMainDataSource(); const types = Object.values(FieldMetadataType).filter((type) => { // We're skipping composite and relation types, as they're not directly mapped to a column type if (isCompositeFieldMetadataType(type)) { @@ -219,7 +218,7 @@ export class DatabaseStructureService { const typeORMType = fieldMetadataTypeToColumnType( FieldMetadataType[type], ) as ColumnType; - const dataType = mainDataSource.driver.normalizeType({ + const dataType = this.coreDataSource.driver.normalizeType({ type: typeORMType, }); @@ -248,7 +247,6 @@ export class DatabaseStructureService { | null, ) => { const typeORMType = fieldMetadataTypeToColumnType(type) as ColumnType; - const mainDataSource = this.typeORMService.getMainDataSource(); // eslint-disable-next-line @typescript-eslint/no-explicit-any let value: any = @@ -282,7 +280,7 @@ export class DatabaseStructureService { value = value.replace(/^'/, '').replace(/'$/, ''); } - return mainDataSource.driver.normalizeDefault({ + return this.coreDataSource.driver.normalizeDefault({ type: typeORMType, default: value, isArray: false, diff --git a/packages/twenty-server/src/engine/workspace-manager/workspace-health/services/object-metadata-health.service.ts b/packages/twenty-server/src/engine/workspace-manager/workspace-health/services/object-metadata-health.service.ts index 67b50607aa..db52bba3d1 100644 --- a/packages/twenty-server/src/engine/workspace-manager/workspace-health/services/object-metadata-health.service.ts +++ b/packages/twenty-server/src/engine/workspace-manager/workspace-health/services/object-metadata-health.service.ts @@ -1,4 +1,7 @@ import { Injectable } from '@nestjs/common'; +import { InjectDataSource } from '@nestjs/typeorm'; + +import { DataSource } from 'typeorm'; import { type WorkspaceHealthIssue, @@ -7,13 +10,15 @@ import { import { type WorkspaceHealthOptions } from 'src/engine/workspace-manager/workspace-health/interfaces/workspace-health-options.interface'; import { type ObjectMetadataEntity } from 'src/engine/metadata-modules/object-metadata/object-metadata.entity'; -import { validName } from 'src/engine/workspace-manager/workspace-health/utils/valid-name.util'; -import { TypeORMService } from 'src/database/typeorm/typeorm.service'; import { computeObjectTargetTable } from 'src/engine/utils/compute-object-target-table.util'; +import { validName } from 'src/engine/workspace-manager/workspace-health/utils/valid-name.util'; @Injectable() export class ObjectMetadataHealthService { - constructor(private readonly typeORMService: TypeORMService) {} + constructor( + @InjectDataSource() + private readonly coreDataSource: DataSource, + ) {} async healthCheck( schemaName: string, @@ -50,11 +55,10 @@ export class ObjectMetadataHealthService { schemaName: string, objectMetadata: ObjectMetadataEntity, ): Promise { - const mainDataSource = this.typeORMService.getMainDataSource(); const issues: WorkspaceHealthIssue[] = []; // Check if the table exist in database - const tableExist = await mainDataSource.query( + const tableExist = await this.coreDataSource.query( `SELECT EXISTS (SELECT FROM information_schema.tables WHERE table_schema = '${schemaName}' AND table_name = '${computeObjectTargetTable(objectMetadata)}')`, ); diff --git a/packages/twenty-server/src/engine/workspace-manager/workspace-manager.service.ts b/packages/twenty-server/src/engine/workspace-manager/workspace-manager.service.ts index ff1f201ab4..cbc3020f02 100644 --- a/packages/twenty-server/src/engine/workspace-manager/workspace-manager.service.ts +++ b/packages/twenty-server/src/engine/workspace-manager/workspace-manager.service.ts @@ -17,6 +17,7 @@ import { RoleEntity } from 'src/engine/metadata-modules/role/role.entity'; import { RoleService } from 'src/engine/metadata-modules/role/role.service'; import { UserRoleService } from 'src/engine/metadata-modules/user-role/user-role.service'; import { WorkspaceMigrationService } from 'src/engine/metadata-modules/workspace-migration/workspace-migration.service'; +import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager'; import { WorkspaceDataSourceService } from 'src/engine/workspace-datasource/workspace-datasource.service'; import { prefillCoreViews } from 'src/engine/workspace-manager/standard-objects-prefill-data/prefill-core-views'; import { standardObjectsPrefillData } from 'src/engine/workspace-manager/standard-objects-prefill-data/standard-objects-prefill-data'; @@ -47,6 +48,7 @@ export class WorkspaceManagerService { @InjectRepository(RoleTargetsEntity) private readonly roleTargetsRepository: Repository, private readonly agentService: AgentService, + protected readonly twentyORMGlobalManager: TwentyORMGlobalManager, ) {} public async init({ @@ -123,18 +125,16 @@ export class WorkspaceManagerService { workspaceId: string, featureFlags: Record, ) { - const mainDataSource = - await this.workspaceDataSourceService.connectToMainDataSource(); - - if (!mainDataSource) { - throw new Error('Could not connect to main data source'); - } + const workspaceDataSource = + await this.twentyORMGlobalManager.getDataSourceForWorkspace({ + workspaceId, + }); const createdObjectMetadata = await this.objectMetadataService.findManyWithinWorkspace(workspaceId); await standardObjectsPrefillData( - mainDataSource, + workspaceDataSource, dataSourceMetadata.schema, createdObjectMetadata, featureFlags, @@ -144,7 +144,7 @@ export class WorkspaceManagerService { this.logger.log(`Prefilling core views for workspace ${workspaceId}`); await prefillCoreViews( - mainDataSource, + workspaceDataSource, workspaceId, createdObjectMetadata, featureFlags, diff --git a/packages/twenty-server/src/engine/workspace-manager/workspace-migration-runner/workspace-migration-runner.service.ts b/packages/twenty-server/src/engine/workspace-manager/workspace-migration-runner/workspace-migration-runner.service.ts index 40fd165f43..31c79f8ccd 100644 --- a/packages/twenty-server/src/engine/workspace-manager/workspace-migration-runner/workspace-migration-runner.service.ts +++ b/packages/twenty-server/src/engine/workspace-manager/workspace-migration-runner/workspace-migration-runner.service.ts @@ -1,8 +1,9 @@ import { Injectable, Logger } from '@nestjs/common'; +import { InjectDataSource } from '@nestjs/typeorm'; import { t } from '@lingui/core/macro'; import { isDefined } from 'twenty-shared/utils'; -import { type QueryRunner, Table, type TableColumn } from 'typeorm'; +import { DataSource, type QueryRunner, Table, type TableColumn } from 'typeorm'; import { IndexMetadataException, @@ -21,7 +22,6 @@ import { } from 'src/engine/metadata-modules/workspace-migration/workspace-migration.entity'; import { WorkspaceMigrationService } from 'src/engine/metadata-modules/workspace-migration/workspace-migration.service'; import { getWorkspaceSchemaName } from 'src/engine/workspace-datasource/utils/get-workspace-schema-name.util'; -import { WorkspaceDataSourceService } from 'src/engine/workspace-datasource/workspace-datasource.service'; import { WorkspaceMigrationColumnService } from 'src/engine/workspace-manager/workspace-migration-runner/services/workspace-migration-column.service'; import { type PostgresQueryRunner } from 'src/engine/workspace-manager/workspace-migration-runner/types/postgres-query-runner.type'; import { tableDefaultColumns } from 'src/engine/workspace-manager/workspace-migration-runner/utils/table-default-column.util'; @@ -33,7 +33,8 @@ export class WorkspaceMigrationRunnerService { private readonly logger = new Logger(WorkspaceMigrationRunnerService.name); constructor( - private readonly workspaceDataSourceService: WorkspaceDataSourceService, + @InjectDataSource() + private readonly coreDataSource: DataSource, private readonly workspaceMigrationService: WorkspaceMigrationService, private readonly workspaceMigrationColumnService: WorkspaceMigrationColumnService, ) {} @@ -88,15 +89,7 @@ export class WorkspaceMigrationRunnerService { public async executeMigrationFromPendingMigrations( workspaceId: string, ): Promise { - const mainDataSource = - await this.workspaceDataSourceService.connectToMainDataSource(); - - if (!mainDataSource) { - throw new Error('Main data source not found'); - } - - const queryRunner = - mainDataSource.createQueryRunner() as PostgresQueryRunner; + const queryRunner = this.coreDataSource.createQueryRunner(); await queryRunner.connect(); await queryRunner.startTransaction(); diff --git a/packages/twenty-server/src/modules/calendar/calendar-event-import-manager/crons/jobs/calendar-event-list-fetch.cron.job.ts b/packages/twenty-server/src/modules/calendar/calendar-event-import-manager/crons/jobs/calendar-event-list-fetch.cron.job.ts index 4775151409..a9e8f2e885 100644 --- a/packages/twenty-server/src/modules/calendar/calendar-event-import-manager/crons/jobs/calendar-event-list-fetch.cron.job.ts +++ b/packages/twenty-server/src/modules/calendar/calendar-event-import-manager/crons/jobs/calendar-event-list-fetch.cron.job.ts @@ -1,7 +1,7 @@ -import { InjectRepository } from '@nestjs/typeorm'; +import { InjectDataSource, InjectRepository } from '@nestjs/typeorm'; import { WorkspaceActivationStatus } from 'twenty-shared/workspace'; -import { Repository } from 'typeorm'; +import { DataSource, Repository } from 'typeorm'; import { SentryCronMonitor } from 'src/engine/core-modules/cron/sentry-cron-monitor.decorator'; import { ExceptionHandlerService } from 'src/engine/core-modules/exception-handler/exception-handler.service'; @@ -12,7 +12,6 @@ import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queu import { MessageQueueService } from 'src/engine/core-modules/message-queue/services/message-queue.service'; import { Workspace } from 'src/engine/core-modules/workspace/workspace.entity'; import { getWorkspaceSchemaName } from 'src/engine/workspace-datasource/utils/get-workspace-schema-name.util'; -import { WorkspaceDataSourceService } from 'src/engine/workspace-datasource/workspace-datasource.service'; import { CalendarEventListFetchJob, type CalendarEventListFetchJobData, @@ -31,7 +30,8 @@ export class CalendarEventListFetchCronJob { @InjectMessageQueue(MessageQueue.calendarQueue) private readonly messageQueueService: MessageQueueService, private readonly exceptionHandlerService: ExceptionHandlerService, - private readonly workspaceDataSourceService: WorkspaceDataSourceService, + @InjectDataSource() + private readonly coreDataSource: DataSource, ) {} @Process(CalendarEventListFetchCronJob.name) @@ -46,14 +46,11 @@ export class CalendarEventListFetchCronJob { }, }); - const mainDataSource = - await this.workspaceDataSourceService.connectToMainDataSource(); - for (const activeWorkspace of activeWorkspaces) { try { const schemaName = getWorkspaceSchemaName(activeWorkspace.id); - const calendarChannels = await mainDataSource.query( + const calendarChannels = await this.coreDataSource.query( `SELECT * FROM ${schemaName}."calendarChannel" WHERE "isSyncEnabled" = true AND "syncStage" IN ('${CalendarChannelSyncStage.FULL_CALENDAR_EVENT_LIST_FETCH_PENDING}', '${CalendarChannelSyncStage.PARTIAL_CALENDAR_EVENT_LIST_FETCH_PENDING}')`, ); diff --git a/packages/twenty-server/src/modules/calendar/calendar-event-import-manager/crons/jobs/calendar-events-import.cron.job.ts b/packages/twenty-server/src/modules/calendar/calendar-event-import-manager/crons/jobs/calendar-events-import.cron.job.ts index 91367a4538..2110a0735c 100644 --- a/packages/twenty-server/src/modules/calendar/calendar-event-import-manager/crons/jobs/calendar-events-import.cron.job.ts +++ b/packages/twenty-server/src/modules/calendar/calendar-event-import-manager/crons/jobs/calendar-events-import.cron.job.ts @@ -1,7 +1,7 @@ -import { InjectRepository } from '@nestjs/typeorm'; +import { InjectDataSource, InjectRepository } from '@nestjs/typeorm'; import { WorkspaceActivationStatus } from 'twenty-shared/workspace'; -import { Repository } from 'typeorm'; +import { DataSource, Repository } from 'typeorm'; import { SentryCronMonitor } from 'src/engine/core-modules/cron/sentry-cron-monitor.decorator'; import { ExceptionHandlerService } from 'src/engine/core-modules/exception-handler/exception-handler.service'; @@ -12,7 +12,6 @@ import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queu import { MessageQueueService } from 'src/engine/core-modules/message-queue/services/message-queue.service'; import { Workspace } from 'src/engine/core-modules/workspace/workspace.entity'; import { getWorkspaceSchemaName } from 'src/engine/workspace-datasource/utils/get-workspace-schema-name.util'; -import { WorkspaceDataSourceService } from 'src/engine/workspace-datasource/workspace-datasource.service'; import { type CalendarEventListFetchJobData } from 'src/modules/calendar/calendar-event-import-manager/jobs/calendar-event-list-fetch.job'; import { CalendarEventsImportJob } from 'src/modules/calendar/calendar-event-import-manager/jobs/calendar-events-import.job'; import { CalendarChannelSyncStage } from 'src/modules/calendar/common/standard-objects/calendar-channel.workspace-entity'; @@ -28,7 +27,8 @@ export class CalendarEventsImportCronJob { private readonly workspaceRepository: Repository, @InjectMessageQueue(MessageQueue.calendarQueue) private readonly messageQueueService: MessageQueueService, - private readonly workspaceDataSourceService: WorkspaceDataSourceService, + @InjectDataSource() + private readonly coreDataSource: DataSource, private readonly exceptionHandlerService: ExceptionHandlerService, ) {} @@ -43,14 +43,12 @@ export class CalendarEventsImportCronJob { activationStatus: WorkspaceActivationStatus.ACTIVE, }, }); - const mainDataSource = - await this.workspaceDataSourceService.connectToMainDataSource(); for (const activeWorkspace of activeWorkspaces) { try { const schemaName = getWorkspaceSchemaName(activeWorkspace.id); - const calendarChannels = await mainDataSource.query( + const calendarChannels = await this.coreDataSource.query( `SELECT * FROM ${schemaName}."calendarChannel" WHERE "isSyncEnabled" = true AND "syncStage" = '${CalendarChannelSyncStage.CALENDAR_EVENTS_IMPORT_PENDING}'`, ); diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/crons/jobs/messaging-message-list-fetch.cron.job.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/crons/jobs/messaging-message-list-fetch.cron.job.ts index 9984e26a54..56236f0fce 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/crons/jobs/messaging-message-list-fetch.cron.job.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/crons/jobs/messaging-message-list-fetch.cron.job.ts @@ -1,7 +1,7 @@ -import { InjectRepository } from '@nestjs/typeorm'; +import { InjectDataSource, InjectRepository } from '@nestjs/typeorm'; import { WorkspaceActivationStatus } from 'twenty-shared/workspace'; -import { Repository } from 'typeorm'; +import { DataSource, Repository } from 'typeorm'; import { SentryCronMonitor } from 'src/engine/core-modules/cron/sentry-cron-monitor.decorator'; import { ExceptionHandlerService } from 'src/engine/core-modules/exception-handler/exception-handler.service'; @@ -12,7 +12,6 @@ import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queu import { MessageQueueService } from 'src/engine/core-modules/message-queue/services/message-queue.service'; import { Workspace } from 'src/engine/core-modules/workspace/workspace.entity'; import { getWorkspaceSchemaName } from 'src/engine/workspace-datasource/utils/get-workspace-schema-name.util'; -import { WorkspaceDataSourceService } from 'src/engine/workspace-datasource/workspace-datasource.service'; import { MessageChannelSyncStage } from 'src/modules/messaging/common/standard-objects/message-channel.workspace-entity'; import { MessagingMessageListFetchJob, @@ -28,7 +27,8 @@ export class MessagingMessageListFetchCronJob { private readonly workspaceRepository: Repository, @InjectMessageQueue(MessageQueue.messagingQueue) private readonly messageQueueService: MessageQueueService, - private readonly workspaceDataSourceService: WorkspaceDataSourceService, + @InjectDataSource() + private readonly coreDataSource: DataSource, private readonly exceptionHandlerService: ExceptionHandlerService, ) {} @@ -44,15 +44,12 @@ export class MessagingMessageListFetchCronJob { }, }); - const mainDataSource = - await this.workspaceDataSourceService.connectToMainDataSource(); - for (const activeWorkspace of activeWorkspaces) { try { const schemaName = getWorkspaceSchemaName(activeWorkspace.id); // TODO: deprecate looking for FULL_MESSAGE_LIST_FETCH_PENDING as we introduce MESSAGE_LIST_FETCH_PENDING - const messageChannels = await mainDataSource.query( + const messageChannels = await this.coreDataSource.query( `SELECT * FROM ${schemaName}."messageChannel" WHERE "isSyncEnabled" = true AND "syncStage" IN ('${MessageChannelSyncStage.PARTIAL_MESSAGE_LIST_FETCH_PENDING}', '${MessageChannelSyncStage.FULL_MESSAGE_LIST_FETCH_PENDING}')`, ); diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/crons/jobs/messaging-messages-import.cron.job.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/crons/jobs/messaging-messages-import.cron.job.ts index af46be1ab1..3cfef05f3a 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/crons/jobs/messaging-messages-import.cron.job.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/crons/jobs/messaging-messages-import.cron.job.ts @@ -1,8 +1,8 @@ -import { InjectRepository } from '@nestjs/typeorm'; +import { InjectDataSource, InjectRepository } from '@nestjs/typeorm'; import { isDefined } from 'twenty-shared/utils'; import { WorkspaceActivationStatus } from 'twenty-shared/workspace'; -import { Repository } from 'typeorm'; +import { DataSource, Repository } from 'typeorm'; import { SentryCronMonitor } from 'src/engine/core-modules/cron/sentry-cron-monitor.decorator'; import { ExceptionHandlerService } from 'src/engine/core-modules/exception-handler/exception-handler.service'; @@ -17,7 +17,6 @@ import { DataSourceExceptionCode, } from 'src/engine/metadata-modules/data-source/data-source.exception'; import { getWorkspaceSchemaName } from 'src/engine/workspace-datasource/utils/get-workspace-schema-name.util'; -import { WorkspaceDataSourceService } from 'src/engine/workspace-datasource/workspace-datasource.service'; import { MessageChannelSyncStage } from 'src/modules/messaging/common/standard-objects/message-channel.workspace-entity'; import { MessagingMessagesImportJob, @@ -34,7 +33,8 @@ export class MessagingMessagesImportCronJob { @InjectMessageQueue(MessageQueue.messagingQueue) private readonly messageQueueService: MessageQueueService, private readonly exceptionHandlerService: ExceptionHandlerService, - private readonly workspaceDataSourceService: WorkspaceDataSourceService, + @InjectDataSource() + private readonly coreDataSource: DataSource, ) {} @Process(MessagingMessagesImportCronJob.name) @@ -49,14 +49,11 @@ export class MessagingMessagesImportCronJob { }, }); - const mainDataSource = - await this.workspaceDataSourceService.connectToMainDataSource(); - for (const activeWorkspace of activeWorkspaces) { try { const schemaName = getWorkspaceSchemaName(activeWorkspace.id); - const messageChannels = await mainDataSource.query( + const messageChannels = await this.coreDataSource.query( `SELECT * FROM ${schemaName}."messageChannel" WHERE "isSyncEnabled" = true AND "syncStage" = '${MessageChannelSyncStage.MESSAGES_IMPORT_PENDING}'`, ); diff --git a/packages/twenty-server/src/modules/workflow/workflow-runner/workflow-run-queue/cron/jobs/workflow-clean-workflow-runs.cron.job.ts b/packages/twenty-server/src/modules/workflow/workflow-runner/workflow-run-queue/cron/jobs/workflow-clean-workflow-runs.cron.job.ts index a577a8467d..bee075dbd8 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-runner/workflow-run-queue/cron/jobs/workflow-clean-workflow-runs.cron.job.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-runner/workflow-run-queue/cron/jobs/workflow-clean-workflow-runs.cron.job.ts @@ -1,8 +1,8 @@ import { Logger } from '@nestjs/common'; -import { InjectRepository } from '@nestjs/typeorm'; +import { InjectDataSource, InjectRepository } from '@nestjs/typeorm'; import { WorkspaceActivationStatus } from 'twenty-shared/workspace'; -import { Repository } from 'typeorm'; +import { DataSource, Repository } from 'typeorm'; import { SentryCronMonitor } from 'src/engine/core-modules/cron/sentry-cron-monitor.decorator'; import { Process } from 'src/engine/core-modules/message-queue/decorators/process.decorator'; @@ -11,7 +11,6 @@ import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queu import { Workspace } from 'src/engine/core-modules/workspace/workspace.entity'; import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager'; import { getWorkspaceSchemaName } from 'src/engine/workspace-datasource/utils/get-workspace-schema-name.util'; -import { WorkspaceDataSourceService } from 'src/engine/workspace-datasource/workspace-datasource.service'; import { WorkflowRunStatus, WorkflowRunWorkspaceEntity, @@ -28,8 +27,9 @@ export class WorkflowCleanWorkflowRunsJob { constructor( @InjectRepository(Workspace) private readonly workspaceRepository: Repository, - private readonly workspaceDataSourceService: WorkspaceDataSourceService, private readonly twentyORMGlobalManager: TwentyORMGlobalManager, + @InjectDataSource() + private readonly coreDataSource: DataSource, ) {} @Process(WorkflowCleanWorkflowRunsJob.name) @@ -44,13 +44,10 @@ export class WorkflowCleanWorkflowRunsJob { }, }); - const mainDataSource = - await this.workspaceDataSourceService.connectToMainDataSource(); - for (const activeWorkspace of activeWorkspaces) { const schemaName = getWorkspaceSchemaName(activeWorkspace.id); - const workflowRunsToDelete = await mainDataSource.query( + const workflowRunsToDelete = await this.coreDataSource.query( ` WITH ranked_runs AS ( SELECT id, diff --git a/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/crons/jobs/cron-trigger.cron.job.ts b/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/crons/jobs/cron-trigger.cron.job.ts index 43d93c36b2..427f4f05ee 100644 --- a/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/crons/jobs/cron-trigger.cron.job.ts +++ b/packages/twenty-server/src/modules/workflow/workflow-trigger/automated-trigger/crons/jobs/cron-trigger.cron.job.ts @@ -1,8 +1,8 @@ -import { InjectRepository } from '@nestjs/typeorm'; +import { InjectDataSource, InjectRepository } from '@nestjs/typeorm'; import { isDefined } from 'twenty-shared/utils'; import { WorkspaceActivationStatus } from 'twenty-shared/workspace'; -import { Repository } from 'typeorm'; +import { DataSource, Repository } from 'typeorm'; import { SentryCronMonitor } from 'src/engine/core-modules/cron/sentry-cron-monitor.decorator'; import { ExceptionHandlerService } from 'src/engine/core-modules/exception-handler/exception-handler.service'; @@ -13,7 +13,6 @@ import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queu import { MessageQueueService } from 'src/engine/core-modules/message-queue/services/message-queue.service'; import { Workspace } from 'src/engine/core-modules/workspace/workspace.entity'; import { getWorkspaceSchemaName } from 'src/engine/workspace-datasource/utils/get-workspace-schema-name.util'; -import { WorkspaceDataSourceService } from 'src/engine/workspace-datasource/workspace-datasource.service'; import { AutomatedTriggerType } from 'src/modules/workflow/common/standard-objects/workflow-automated-trigger.workspace-entity'; import { type CronTriggerSettings } from 'src/modules/workflow/workflow-trigger/automated-trigger/constants/automated-trigger-settings'; import { shouldRunNow } from 'src/modules/workflow/workflow-trigger/automated-trigger/crons/utils/should-run-now.utils'; @@ -27,7 +26,8 @@ export const CRON_TRIGGER_CRON_PATTERN = '* * * * *'; @Processor(MessageQueue.cronQueue) export class CronTriggerCronJob { constructor( - private readonly workspaceDataSourceService: WorkspaceDataSourceService, + @InjectDataSource() + private readonly coreDataSource: DataSource, @InjectRepository(Workspace) private readonly workspaceRepository: Repository, @InjectMessageQueue(MessageQueue.workflowQueue) @@ -46,14 +46,11 @@ export class CronTriggerCronJob { const now = new Date(); - const mainDataSource = - await this.workspaceDataSourceService.connectToMainDataSource(); - for (const activeWorkspace of activeWorkspaces) { try { const schemaName = getWorkspaceSchemaName(activeWorkspace.id); - const workflowAutomatedCronTriggers = await mainDataSource.query( + const workflowAutomatedCronTriggers = await this.coreDataSource.query( `SELECT * FROM ${schemaName}."workflowAutomatedTrigger" WHERE type = '${AutomatedTriggerType.CRON}'`, ); diff --git a/packages/twenty-server/test/integration/utils/setup-test.ts b/packages/twenty-server/test/integration/utils/setup-test.ts index aeb7c19de1..e717eabfab 100644 --- a/packages/twenty-server/test/integration/utils/setup-test.ts +++ b/packages/twenty-server/test/integration/utils/setup-test.ts @@ -1,10 +1,9 @@ import { type JestConfigWithTsJest } from 'ts-jest'; import 'tsconfig-paths/register'; -import { rawDataSource } from 'src/database/typeorm/raw/raw.datasource'; -import { TypeORMService } from 'src/database/typeorm/typeorm.service'; -import { DataSourceService } from 'src/engine/metadata-modules/data-source/data-source.service'; import { DataSeedWorkspaceCommand } from 'src/database/commands/data-seed-dev-workspace.command'; +import { rawDataSource } from 'src/database/typeorm/raw/raw.datasource'; +import { DataSourceService } from 'src/engine/metadata-modules/data-source/data-source.service'; import { createApp } from './create-app'; @@ -25,8 +24,6 @@ export default async (_, projectConfig: JestConfigWithTsJest) => { // @ts-expect-error legacy noImplicitAny global.testDataSource = rawDataSource; // @ts-expect-error legacy noImplicitAny - global.typeOrmService = app.get(TypeORMService); - // @ts-expect-error legacy noImplicitAny global.dataSourceService = app.get(DataSourceService); // @ts-expect-error legacy noImplicitAny global.dataSeedWorkspaceCommand = app.get(DataSeedWorkspaceCommand);