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 7c76c87b4dc..8f45e03ebbd 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 6b7caeebf63..98496dcfc4e 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 1250ed7829e..c4958f66c96 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 88e65e2168a..5583926f4e6 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 9e0046429db..c289d6d6dfb 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 09916cff6ff..0e7d6135d87 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 88d7a6924e0..6a126886bc0 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 ca50eb35870..00000000000 --- 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 2e761f7af1c..8b779d48387 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 745c7754928..5942b997a52 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 d8303c0a2a5..5fc6a6dbbcd 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 047ec83ff50..95b261736cd 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 6acd8e9646d..54b0dfa22ba 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 34eb6ede359..c16384c5556 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 fb92cb6d502..52411a3ee8a 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 263defed78e..ef1fd7a80a6 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 cd7fd615d26..77de1c29614 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 16e6a40a09e..f2044bab1de 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 9e38bf37d2d..086106d0143 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 76d42fbaca2..461eeab92b6 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 6f4938420f2..4f44e9aecd8 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 a37a10ca25a..90b712c8983 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 e2ffd39c47f..e04208f1e2b 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 14636c513ba..6e3edead474 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 67b50607aa2..db52bba3d1f 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 ff1f201ab4b..cbc3020f02f 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 40fd165f43b..31c79f8ccd9 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 47751514098..a9e8f2e8855 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 91367a4538b..2110a0735c9 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 9984e26a54b..56236f0fceb 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 af46be1ab1d..3cfef05f3af 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 a577a8467d2..bee075dbd8b 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 43d93c36b24..427f4f05ee8 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 aeb7c19de17..e717eabfab2 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);