import { Injectable, Logger } from '@nestjs/common'; import { InjectRepository } from '@nestjs/typeorm'; import { isNonEmptyString } from '@sniptt/guards'; import { FeatureFlagKey } from 'twenty-shared/types'; import { isDefined } from 'twenty-shared/utils'; import { validate as uuidValidate } from 'uuid'; import { type FindOneOptions, type FindOptionsWhere, Repository, } from 'typeorm'; import { FeatureFlagService } from 'src/engine/core-modules/feature-flag/services/feature-flag.service'; import { MessageFolderEntity } from 'src/engine/metadata-modules/message-folder/entities/message-folder.entity'; import { type WorkspaceEntityManager } from 'src/engine/twenty-orm/entity-manager/workspace-entity-manager'; import { GlobalWorkspaceOrmManager } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager'; import { type MessageFolderWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; @Injectable() export class MessageFolderDataAccessService { private readonly logger = new Logger(MessageFolderDataAccessService.name); constructor( @InjectRepository(MessageFolderEntity) private readonly coreRepository: Repository, private readonly featureFlagService: FeatureFlagService, private readonly globalWorkspaceOrmManager: GlobalWorkspaceOrmManager, ) {} private async isMigrated(workspaceId: string): Promise { return this.featureFlagService.isFeatureEnabled( FeatureFlagKey.IS_CONNECTED_ACCOUNT_MIGRATED, workspaceId, ); } private async resolveParentFolderIdForCore( workspaceId: string, parentFolderId: string | null, messageChannelId: string | undefined, ): Promise { if (!isNonEmptyString(parentFolderId)) { return null; } if (uuidValidate(parentFolderId)) { return parentFolderId; } if (!isDefined(messageChannelId)) { return null; } const parentFolder = await this.coreRepository.findOne({ where: { workspaceId, messageChannelId, externalId: parentFolderId, }, select: ['id'], }); return parentFolder?.id ?? null; } private async toCore( workspaceId: string, data: Partial, messageChannelId?: string, ): Promise> { const coreData: Record = { ...data, workspaceId }; const channelId = (coreData.messageChannelId as string) ?? messageChannelId; if ('parentFolderId' in coreData) { coreData.parentFolderId = await this.resolveParentFolderIdForCore( workspaceId, coreData.parentFolderId as string | null, channelId, ); } return coreData; } async getWorkspaceRepository(workspaceId: string) { return this.globalWorkspaceOrmManager.getRepository( workspaceId, 'messageFolder', ); } async findOne( workspaceId: string, options: FindOneOptions, ): Promise { if (await this.isMigrated(workspaceId)) { const where = options.where as Record; const coreWhere = Array.isArray(where) ? where.map((whereItem) => ({ ...(whereItem as Record), workspaceId, })) : { ...(where as Record), workspaceId, }; return this.coreRepository.findOne({ ...options, where: coreWhere, } as FindOneOptions) as unknown as Promise; } const workspaceRepository = await this.getWorkspaceRepository(workspaceId); return workspaceRepository.findOne(options); } async find( workspaceId: string, where?: FindOptionsWhere, ): Promise { if (await this.isMigrated(workspaceId)) { return this.coreRepository.find({ where: { ...(where as Record), workspaceId, } as FindOptionsWhere, }) as unknown as Promise; } const workspaceRepository = await this.getWorkspaceRepository(workspaceId); return workspaceRepository.find({ where }); } async save( workspaceId: string, data: Partial, ): Promise { const workspaceRepository = await this.getWorkspaceRepository(workspaceId); const savedData = await workspaceRepository.save(data); if (await this.isMigrated(workspaceId)) { try { const coreData = await this.toCore(workspaceId, savedData); await this.coreRepository.save( coreData as unknown as MessageFolderEntity, ); } catch (error) { this.logger.error( `Failed to dual-write messageFolder to core: ${error}`, ); throw error; } } } async update( workspaceId: string, where: FindOptionsWhere, data: Partial, manager?: WorkspaceEntityManager, ): Promise { const workspaceRepository = await this.getWorkspaceRepository(workspaceId); await workspaceRepository.update(where, data, manager); if (await this.isMigrated(workspaceId)) { try { const coreData = await this.toCore( workspaceId, data, (where as Record).messageChannelId as string, ); await this.coreRepository.update( { ...where, workspaceId } as FindOptionsWhere, coreData, ); } catch (error) { this.logger.error( `Failed to dual-write messageFolder update to core: ${error}`, ); throw error; } } } async delete( workspaceId: string, where: FindOptionsWhere, ): Promise { const workspaceRepository = await this.getWorkspaceRepository(workspaceId); await workspaceRepository.delete(where); if (await this.isMigrated(workspaceId)) { try { await this.coreRepository.delete({ ...where, workspaceId, } as FindOptionsWhere); } catch (error) { this.logger.error( `Failed to dual-write messageFolder delete to core: ${error}`, ); throw error; } } } }