diff --git a/packages/twenty-server/src/engine/core-modules/auth/services/create-message-channel.service.ts b/packages/twenty-server/src/engine/core-modules/auth/services/create-message-channel.service.ts index bd1873ba9f..3c7e893250 100644 --- a/packages/twenty-server/src/engine/core-modules/auth/services/create-message-channel.service.ts +++ b/packages/twenty-server/src/engine/core-modules/auth/services/create-message-channel.service.ts @@ -75,8 +75,10 @@ export class CreateMessageChannelService { if (isDefined(connectedAccount)) { await this.syncMessageFoldersService.syncMessageFolders({ workspaceId, - messageChannelId: newMessageChannel.id, - connectedAccount, + messageChannel: { + ...newMessageChannel, + connectedAccount, + }, manager, }); } diff --git a/packages/twenty-server/src/engine/workspace-manager/workspace-sync-metadata/constants/standard-field-ids.ts b/packages/twenty-server/src/engine/workspace-manager/workspace-sync-metadata/constants/standard-field-ids.ts index af096ce5d7..b2309f2eb9 100644 --- a/packages/twenty-server/src/engine/workspace-manager/workspace-sync-metadata/constants/standard-field-ids.ts +++ b/packages/twenty-server/src/engine/workspace-manager/workspace-sync-metadata/constants/standard-field-ids.ts @@ -248,7 +248,6 @@ export const MESSAGE_CHANNEL_STANDARD_FIELD_IDS = { messageChannelMessageAssociations: '20202020-49b8-4766-88fd-75f1e21b3d5f', messageFolders: '20202020-cc39-4432-9fe8-ec8ab8bbed94', messageFolderImportPolicy: '20202020-cc39-4432-9fe8-ec8ab8bbed95', - syncAllFolders: '20202020-e36b-409c-a9bd-79071596cdc0', pendingGroupEmailsAction: '20202020-17c5-4e9f-bc50-af46a89fdd42', isSyncEnabled: '20202020-d9a6-48e9-990b-b97fdf22e8dd', syncCursor: '20202020-79d1-41cf-b738-bcf5ed61e256', diff --git a/packages/twenty-server/src/modules/messaging/common/standard-objects/message-channel.workspace-entity.ts b/packages/twenty-server/src/modules/messaging/common/standard-objects/message-channel.workspace-entity.ts index f9dfe2ee55..081d0af3f2 100644 --- a/packages/twenty-server/src/modules/messaging/common/standard-objects/message-channel.workspace-entity.ts +++ b/packages/twenty-server/src/modules/messaging/common/standard-objects/message-channel.workspace-entity.ts @@ -257,16 +257,6 @@ export class MessageChannelWorkspaceEntity extends BaseWorkspaceEntity { }) excludeGroupEmails: boolean; - @WorkspaceField({ - standardId: MESSAGE_CHANNEL_STANDARD_FIELD_IDS.syncAllFolders, - type: FieldMetadataType.BOOLEAN, - label: msg`Sync all folders`, - description: msg`Sync all folders including new ones`, - icon: 'IconFolders', - defaultValue: true, - }) - syncAllFolders: boolean; - @WorkspaceField({ standardId: MESSAGE_CHANNEL_STANDARD_FIELD_IDS.pendingGroupEmailsAction, type: FieldMetadataType.SELECT, diff --git a/packages/twenty-server/src/modules/messaging/common/standard-objects/message-folder.workspace-entity.ts b/packages/twenty-server/src/modules/messaging/common/standard-objects/message-folder.workspace-entity.ts index 3459cf3813..1f0c0d8323 100644 --- a/packages/twenty-server/src/modules/messaging/common/standard-objects/message-folder.workspace-entity.ts +++ b/packages/twenty-server/src/modules/messaging/common/standard-objects/message-folder.workspace-entity.ts @@ -21,7 +21,6 @@ import { MessageChannelWorkspaceEntity } from 'src/modules/messaging/common/stan export enum MessageFolderPendingSyncAction { FOLDER_DELETION = 'FOLDER_DELETION', - FOLDER_IMPORT = 'FOLDER_IMPORT', NONE = 'NONE', } @@ -126,16 +125,10 @@ export class MessageFolderWorkspaceEntity extends BaseWorkspaceEntity { position: 0, color: 'red', }, - { - value: MessageFolderPendingSyncAction.FOLDER_IMPORT, - label: 'Folder import', - position: 1, - color: 'green', - }, { value: MessageFolderPendingSyncAction.NONE, label: 'None', - position: 2, + position: 1, color: 'blue', }, ], diff --git a/packages/twenty-server/src/modules/messaging/message-channel-manager/query-hooks/message-channel-query-hook.module.ts b/packages/twenty-server/src/modules/messaging/message-channel-manager/query-hooks/message-channel-query-hook.module.ts index 00e2dfe18d..d0876baaa0 100644 --- a/packages/twenty-server/src/modules/messaging/message-channel-manager/query-hooks/message-channel-query-hook.module.ts +++ b/packages/twenty-server/src/modules/messaging/message-channel-manager/query-hooks/message-channel-query-hook.module.ts @@ -1,8 +1,10 @@ import { Module } from '@nestjs/common'; import { MessageChannelUpdateOnePreQueryHook } from 'src/modules/messaging/message-channel-manager/query-hooks/message-channel-update-one.pre-query.hook'; +import { MessagingImportManagerModule } from 'src/modules/messaging/message-import-manager/messaging-import-manager.module'; @Module({ + imports: [MessagingImportManagerModule], providers: [MessageChannelUpdateOnePreQueryHook], exports: [MessageChannelUpdateOnePreQueryHook], }) diff --git a/packages/twenty-server/src/modules/messaging/message-channel-manager/query-hooks/message-channel-update-one.pre-query.hook.ts b/packages/twenty-server/src/modules/messaging/message-channel-manager/query-hooks/message-channel-update-one.pre-query.hook.ts index 51a8c26fa4..4a788e1a47 100644 --- a/packages/twenty-server/src/modules/messaging/message-channel-manager/query-hooks/message-channel-update-one.pre-query.hook.ts +++ b/packages/twenty-server/src/modules/messaging/message-channel-manager/query-hooks/message-channel-update-one.pre-query.hook.ts @@ -22,6 +22,7 @@ import { MessageFolderPendingSyncAction, type MessageFolderWorkspaceEntity, } from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; +import { MessagingProcessGroupEmailActionsService } from 'src/modules/messaging/message-import-manager/services/messaging-process-group-email-actions.service'; const ONGOING_SYNC_STAGES = [ MessageChannelSyncStage.MESSAGE_LIST_FETCH_ONGOING, @@ -34,6 +35,7 @@ export class MessageChannelUpdateOnePreQueryHook { constructor( private readonly twentyORMGlobalManager: TwentyORMGlobalManager, + private readonly messagingProcessGroupEmailActionsService: MessagingProcessGroupEmailActionsService, ) {} async execute( @@ -101,6 +103,19 @@ export class MessageChannelUpdateOnePreQueryHook ); } + const excludeGroupEmailsChanged = + payload.data.excludeGroupEmails !== messageChannel.excludeGroupEmails; + + if (excludeGroupEmailsChanged) { + await this.messagingProcessGroupEmailActionsService.markMessageChannelAsPendingGroupEmailsAction( + messageChannel, + workspace.id, + payload.data.excludeGroupEmails + ? MessageChannelPendingGroupEmailsAction.GROUP_EMAILS_DELETION + : MessageChannelPendingGroupEmailsAction.GROUP_EMAILS_IMPORT, + ); + } + return payload; } } diff --git a/packages/twenty-server/src/modules/messaging/message-folder-manager/drivers/gmail/gmail-get-all-folders.service.ts b/packages/twenty-server/src/modules/messaging/message-folder-manager/drivers/gmail/gmail-get-all-folders.service.ts index bd2bdb35a8..5d03d1245e 100644 --- a/packages/twenty-server/src/modules/messaging/message-folder-manager/drivers/gmail/gmail-get-all-folders.service.ts +++ b/packages/twenty-server/src/modules/messaging/message-folder-manager/drivers/gmail/gmail-get-all-folders.service.ts @@ -7,8 +7,10 @@ import { import { OAuth2ClientManagerService } from 'src/modules/connected-account/oauth2-client-manager/services/oauth2-client-manager.service'; import { type ConnectedAccountWorkspaceEntity } from 'src/modules/connected-account/standard-objects/connected-account.workspace-entity'; +import { MessageChannelWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-channel.workspace-entity'; import { extractGmailFolderName } from 'src/modules/messaging/message-folder-manager/drivers/gmail/utils/extract-gmail-folder-name.util'; import { getGmailFolderParentId } from 'src/modules/messaging/message-folder-manager/drivers/gmail/utils/get-gmail-folder-parent-id.util'; +import { shouldSyncFolderByDefault } from 'src/modules/messaging/message-folder-manager/utils/should-sync-folder-by-default.util'; import { MESSAGING_GMAIL_DEFAULT_NOT_SYNCED_LABELS } from 'src/modules/messaging/message-import-manager/drivers/gmail/constants/messaging-gmail-default-not-synced-labels'; import { GmailMessageListFetchErrorHandler } from 'src/modules/messaging/message-import-manager/drivers/gmail/services/gmail-message-list-fetch-error-handler.service'; @@ -21,15 +23,15 @@ export class GmailGetAllFoldersService implements MessageFolderDriver { private readonly gmailMessageListFetchErrorHandler: GmailMessageListFetchErrorHandler, ) {} - private isSyncedByDefault(labelId: string): boolean { - return !MESSAGING_GMAIL_DEFAULT_NOT_SYNCED_LABELS.includes(labelId); - } - async getAllMessageFolders( connectedAccount: Pick< ConnectedAccountWorkspaceEntity, 'provider' | 'refreshToken' | 'accessToken' | 'id' | 'handle' >, + messageChannel: Pick< + MessageChannelWorkspaceEntity, + 'messageFolderImportPolicy' + >, ): Promise { try { const oAuth2Client = @@ -70,17 +72,24 @@ export class GmailGetAllFoldersService implements MessageFolderDriver { continue; } + if (MESSAGING_GMAIL_DEFAULT_NOT_SYNCED_LABELS.includes(label.id)) { + continue; + } + const isSentFolder = label.id === 'SENT'; const folderName = extractGmailFolderName(label.name); const parentFolderId = getGmailFolderParentId( label.name, labelNameToIdMap, ); + const isSynced = shouldSyncFolderByDefault( + messageChannel.messageFolderImportPolicy, + ); folders.push({ externalId: label.id, name: folderName, - isSynced: this.isSyncedByDefault(label.id), + isSynced, isSentFolder, parentFolderId, }); diff --git a/packages/twenty-server/src/modules/messaging/message-folder-manager/drivers/imap/imap-get-all-folders.service.ts b/packages/twenty-server/src/modules/messaging/message-folder-manager/drivers/imap/imap-get-all-folders.service.ts index 3fe5df7679..d781b38e69 100644 --- a/packages/twenty-server/src/modules/messaging/message-folder-manager/drivers/imap/imap-get-all-folders.service.ts +++ b/packages/twenty-server/src/modules/messaging/message-folder-manager/drivers/imap/imap-get-all-folders.service.ts @@ -9,10 +9,11 @@ import { } from 'src/modules/messaging/message-folder-manager/interfaces/message-folder-driver.interface'; import { type ConnectedAccountWorkspaceEntity } from 'src/modules/connected-account/standard-objects/connected-account.workspace-entity'; +import { MessageChannelWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-channel.workspace-entity'; +import { shouldCreateFolderByDefault } from 'src/modules/messaging/message-folder-manager/utils/should-create-folder-by-default.util'; +import { shouldSyncFolderByDefault } from 'src/modules/messaging/message-folder-manager/utils/should-sync-folder-by-default.util'; import { ImapClientProvider } from 'src/modules/messaging/message-import-manager/drivers/imap/providers/imap-client.provider'; import { ImapFindSentFolderService } from 'src/modules/messaging/message-import-manager/drivers/imap/services/imap-find-sent-folder.service'; -import { MessageFolderName } from 'src/modules/messaging/message-import-manager/drivers/imap/types/folders'; -import { StandardFolder } from 'src/modules/messaging/message-import-manager/drivers/types/standard-folder'; import { getStandardFolderByRegex } from 'src/modules/messaging/message-import-manager/drivers/utils/get-standard-folder-by-regex'; @Injectable() @@ -29,13 +30,21 @@ export class ImapGetAllFoldersService implements MessageFolderDriver { ConnectedAccountWorkspaceEntity, 'id' | 'provider' | 'connectionParameters' | 'handle' >, + messageChannel: Pick< + MessageChannelWorkspaceEntity, + 'messageFolderImportPolicy' + >, ): Promise { try { const client = await this.imapClientProvider.getClient(connectedAccount); const mailboxList = await client.list(); - const folders = await this.filterAndMapFolders(client, mailboxList); + const folders = await this.filterAndMapFolders( + client, + mailboxList, + messageChannel, + ); await this.imapClientProvider.closeClient(client); @@ -53,6 +62,10 @@ export class ImapGetAllFoldersService implements MessageFolderDriver { private async filterAndMapFolders( client: ImapFlow, mailboxList: ListResponse[], + messageChannel: Pick< + MessageChannelWorkspaceEntity, + 'messageFolderImportPolicy' + >, ): Promise { const folders: MessageFolder[] = []; const pathToExternalIdMap = new Map(); @@ -89,12 +102,14 @@ export class ImapGetAllFoldersService implements MessageFolderDriver { pathToExternalIdMap.set(mailbox.path, externalId); if (this.isValidMailbox(mailbox, folders)) { - const isInbox = await this.isInboxFolder(mailbox); const standardFolder = getStandardFolderByRegex(mailbox.path); - const isSynced = this.shouldSyncByDefault( - mailbox, - standardFolder, - isInbox, + + if (!shouldCreateFolderByDefault(standardFolder)) { + continue; + } + + const isSynced = shouldSyncFolderByDefault( + messageChannel.messageFolderImportPolicy, ); folders.push({ @@ -122,7 +137,7 @@ export class ImapGetAllFoldersService implements MessageFolderDriver { mailbox: ListResponse, existingFolders: MessageFolder[], ): boolean { - if (this.shouldExcludeFolder(mailbox)) { + if (mailbox.flags?.has('\\Noselect')) { return false; } @@ -139,53 +154,6 @@ export class ImapGetAllFoldersService implements MessageFolderDriver { return !isDuplicate; } - private async isInboxFolder(mailbox: ListResponse): Promise { - if ( - mailbox.path.toLowerCase() === MessageFolderName.INBOX || - mailbox.specialUse === '\\Inbox' - ) { - return true; - } - - return false; - } - - private shouldExcludeFolder(mailbox: ListResponse): boolean { - if (mailbox.flags?.has('\\Noselect')) { - return true; - } - - return false; - } - - private shouldSyncByDefault( - mailbox: ListResponse, - standardFolder: StandardFolder | null, - isInbox: boolean, - ): boolean { - if ( - mailbox.specialUse === '\\Drafts' || - mailbox.specialUse === '\\Trash' || - mailbox.specialUse === '\\Junk' - ) { - return false; - } - - if ( - standardFolder === StandardFolder.DRAFTS || - standardFolder === StandardFolder.TRASH || - standardFolder === StandardFolder.JUNK - ) { - return false; - } - - if (isInbox) { - return true; - } - - return false; - } - private async getUidValidity( client: ImapFlow, mailbox: ListResponse, diff --git a/packages/twenty-server/src/modules/messaging/message-folder-manager/drivers/microsoft/microsoft-get-all-folders.service.ts b/packages/twenty-server/src/modules/messaging/message-folder-manager/drivers/microsoft/microsoft-get-all-folders.service.ts index 04066a332d..9d3da2d4e0 100644 --- a/packages/twenty-server/src/modules/messaging/message-folder-manager/drivers/microsoft/microsoft-get-all-folders.service.ts +++ b/packages/twenty-server/src/modules/messaging/message-folder-manager/drivers/microsoft/microsoft-get-all-folders.service.ts @@ -9,8 +9,12 @@ import { import { OAuth2ClientManagerService } from 'src/modules/connected-account/oauth2-client-manager/services/oauth2-client-manager.service'; import { type ConnectedAccountWorkspaceEntity } from 'src/modules/connected-account/standard-objects/connected-account.workspace-entity'; +import { MessageChannelWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-channel.workspace-entity'; +import { shouldCreateFolderByDefault } from 'src/modules/messaging/message-folder-manager/utils/should-create-folder-by-default.util'; +import { shouldSyncFolderByDefault } from 'src/modules/messaging/message-folder-manager/utils/should-sync-folder-by-default.util'; import { MicrosoftMessageListFetchErrorHandler } from 'src/modules/messaging/message-import-manager/drivers/microsoft/services/microsoft-message-list-fetch-error-handler.service'; import { StandardFolder } from 'src/modules/messaging/message-import-manager/drivers/types/standard-folder'; +import { getStandardFolderByRegex } from 'src/modules/messaging/message-import-manager/drivers/utils/get-standard-folder-by-regex'; type MicrosoftGraphFolder = { id: string; @@ -36,6 +40,10 @@ export class MicrosoftGetAllFoldersService implements MessageFolderDriver { ConnectedAccountWorkspaceEntity, 'accessToken' | 'refreshToken' | 'id' | 'handle' | 'provider' >, + messageChannel: Pick< + MessageChannelWorkspaceEntity, + 'messageFolderImportPolicy' + >, ): Promise { try { const microsoftClient = @@ -66,11 +74,18 @@ export class MicrosoftGetAllFoldersService implements MessageFolderDriver { continue; } - const standardFolder = this.getStandardFolderFromWellKnownName( - folder.wellKnownName, - ); + const standardFolder = folder.wellKnownName + ? getStandardFolderByRegex(folder.wellKnownName) + : null; + + if (!shouldCreateFolderByDefault(standardFolder)) { + continue; + } + const isSentFolder = this.isSentFolder(standardFolder); - const isSynced = this.shouldSyncByDefault(standardFolder); + const isSynced = shouldSyncFolderByDefault( + messageChannel.messageFolderImportPolicy, + ); folderInfos.push({ externalId: folder.id, @@ -103,41 +118,6 @@ export class MicrosoftGetAllFoldersService implements MessageFolderDriver { return standardFolder === StandardFolder.SENT; } - private shouldSyncByDefault(standardFolder: StandardFolder | null): boolean { - if ( - standardFolder === StandardFolder.JUNK || - standardFolder === StandardFolder.DRAFTS || - standardFolder === StandardFolder.TRASH - ) { - return false; - } - - return true; - } - - private getStandardFolderFromWellKnownName( - wellKnownName?: string, - ): StandardFolder | null { - if (!isDefined(wellKnownName)) { - return null; - } - - switch (wellKnownName.toLowerCase()) { - case 'inbox': - return StandardFolder.INBOX; - case 'drafts': - return StandardFolder.DRAFTS; - case 'sentitems': - return StandardFolder.SENT; - case 'deleteditems': - return StandardFolder.TRASH; - case 'junkemail': - return StandardFolder.JUNK; - default: - return null; - } - } - /* * All Microsoft folders have a parentFolderId including the standard folders * which point to root node which doesn't exits in the API response. diff --git a/packages/twenty-server/src/modules/messaging/message-folder-manager/interfaces/message-folder-driver.interface.ts b/packages/twenty-server/src/modules/messaging/message-folder-manager/interfaces/message-folder-driver.interface.ts index 7ac6f31953..783a77fa07 100644 --- a/packages/twenty-server/src/modules/messaging/message-folder-manager/interfaces/message-folder-driver.interface.ts +++ b/packages/twenty-server/src/modules/messaging/message-folder-manager/interfaces/message-folder-driver.interface.ts @@ -1,4 +1,5 @@ import { type ConnectedAccountWorkspaceEntity } from 'src/modules/connected-account/standard-objects/connected-account.workspace-entity'; +import { type MessageChannelWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-channel.workspace-entity'; import { type MessageFolderWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; export type MessageFolder = Pick< @@ -17,5 +18,9 @@ export interface MessageFolderDriver { | 'handle' | 'connectionParameters' >, + messageChannel: Pick< + MessageChannelWorkspaceEntity, + 'messageFolderImportPolicy' + >, ): Promise; } diff --git a/packages/twenty-server/src/modules/messaging/message-folder-manager/messaging-folder-sync-manager.module.ts b/packages/twenty-server/src/modules/messaging/message-folder-manager/messaging-folder-sync-manager.module.ts index 9afe07f2a9..d4dd634fbb 100644 --- a/packages/twenty-server/src/modules/messaging/message-folder-manager/messaging-folder-sync-manager.module.ts +++ b/packages/twenty-server/src/modules/messaging/message-folder-manager/messaging-folder-sync-manager.module.ts @@ -6,10 +6,10 @@ import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.ent import { DataSourceModule } from 'src/engine/metadata-modules/data-source/data-source.module'; import { WorkspaceDataSourceModule } from 'src/engine/workspace-datasource/workspace-datasource.module'; import { OAuth2ClientManagerModule } from 'src/modules/connected-account/oauth2-client-manager/oauth2-client-manager.module'; -import { SyncMessageFoldersService } from 'src/modules/messaging/message-folder-manager/services/sync-message-folders.service'; import { GmailGetAllFoldersService } from 'src/modules/messaging/message-folder-manager/drivers/gmail/gmail-get-all-folders.service'; import { ImapGetAllFoldersService } from 'src/modules/messaging/message-folder-manager/drivers/imap/imap-get-all-folders.service'; import { MicrosoftGetAllFoldersService } from 'src/modules/messaging/message-folder-manager/drivers/microsoft/microsoft-get-all-folders.service'; +import { SyncMessageFoldersService } from 'src/modules/messaging/message-folder-manager/services/sync-message-folders.service'; import { MessagingGmailDriverModule } from 'src/modules/messaging/message-import-manager/drivers/gmail/messaging-gmail-driver.module'; import { MessagingIMAPDriverModule } from 'src/modules/messaging/message-import-manager/drivers/imap/messaging-imap-driver.module'; import { MessagingMicrosoftDriverModule } from 'src/modules/messaging/message-import-manager/drivers/microsoft/messaging-microsoft-driver.module'; diff --git a/packages/twenty-server/src/modules/messaging/message-folder-manager/services/sync-message-folders.service.ts b/packages/twenty-server/src/modules/messaging/message-folder-manager/services/sync-message-folders.service.ts index b79acce18d..909e650475 100644 --- a/packages/twenty-server/src/modules/messaging/message-folder-manager/services/sync-message-folders.service.ts +++ b/packages/twenty-server/src/modules/messaging/message-folder-manager/services/sync-message-folders.service.ts @@ -10,7 +10,10 @@ import { type WorkspaceEntityManager } from 'src/engine/twenty-orm/entity-manage import { WorkspaceRepository } from 'src/engine/twenty-orm/repository/workspace.repository'; import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager'; import { type MessageChannelWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-channel.workspace-entity'; -import { type MessageFolderWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; +import { + MessageFolderPendingSyncAction, + type MessageFolderWorkspaceEntity, +} from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; import { GmailGetAllFoldersService } from 'src/modules/messaging/message-folder-manager/drivers/gmail/gmail-get-all-folders.service'; import { ImapGetAllFoldersService } from 'src/modules/messaging/message-folder-manager/drivers/imap/imap-get-all-folders.service'; import { MicrosoftGetAllFoldersService } from 'src/modules/messaging/message-folder-manager/drivers/microsoft/microsoft-get-all-folders.service'; @@ -18,8 +21,10 @@ import { MessageFolderName } from 'src/modules/messaging/message-import-manager/ type SyncMessageFoldersInput = { workspaceId: string; - messageChannelId: string; - connectedAccount: MessageChannelWorkspaceEntity['connectedAccount']; + messageChannel: Pick< + MessageChannelWorkspaceEntity, + 'messageFolderImportPolicy' | 'connectedAccount' | 'id' + >; manager: WorkspaceEntityManager; }; @@ -52,13 +57,16 @@ export class SyncMessageFoldersService { ) {} async syncMessageFolders(input: SyncMessageFoldersInput): Promise { - const { workspaceId, messageChannelId, connectedAccount, manager } = input; + const { workspaceId, messageChannel, manager } = input; - const folders = await this.discoverAllFolders(connectedAccount); + const folders = await this.discoverAllFolders( + messageChannel.connectedAccount, + messageChannel, + ); await this.upsertDiscoveredFolders({ workspaceId, - messageChannelId, + messageChannelId: messageChannel.id, folders, manager, }); @@ -88,7 +96,7 @@ export class SyncMessageFoldersService { const inserts: MessageFolderToInsert[] = []; const updates: [string, MessageFolderToUpdate][] = []; - const deletes: string[] = []; + const foldersToMarkForDeletion: string[] = []; const discoveredExternalIds = new Set( folders @@ -101,7 +109,7 @@ export class SyncMessageFoldersService { existingFolder.externalId && !discoveredExternalIds.has(existingFolder.externalId) ) { - deletes.push(existingFolder.id); + foldersToMarkForDeletion.push(existingFolder.id); } } @@ -150,26 +158,41 @@ export class SyncMessageFoldersService { ); } - if (deletes.length > 0) { - await messageFolderRepository.delete(deletes, manager); + if (foldersToMarkForDeletion.length > 0) { + await messageFolderRepository.updateMany( + foldersToMarkForDeletion.map((id) => ({ + criteria: id, + partialEntity: { + pendingSyncAction: MessageFolderPendingSyncAction.FOLDER_DELETION, + }, + })), + manager, + ); } } async discoverAllFolders( connectedAccount: MessageChannelWorkspaceEntity['connectedAccount'], + messageChannel: Pick< + MessageChannelWorkspaceEntity, + 'messageFolderImportPolicy' + >, ): Promise { switch (connectedAccount.provider) { case ConnectedAccountProvider.GOOGLE: return await this.gmailGetAllFoldersService.getAllMessageFolders( connectedAccount, + messageChannel, ); case ConnectedAccountProvider.MICROSOFT: return await this.microsoftGetAllFoldersService.getAllMessageFolders( connectedAccount, + messageChannel, ); case ConnectedAccountProvider.IMAP_SMTP_CALDAV: return await this.imapGetAllFoldersService.getAllMessageFolders( connectedAccount, + messageChannel, ); default: throw new Error( diff --git a/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/MESSAGING_FOLDER_MANAGER_ALWAYS_EXCLUDED_FOLDERS.ts b/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/MESSAGING_FOLDER_MANAGER_ALWAYS_EXCLUDED_FOLDERS.ts new file mode 100644 index 0000000000..3e70413fd1 --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/MESSAGING_FOLDER_MANAGER_ALWAYS_EXCLUDED_FOLDERS.ts @@ -0,0 +1,7 @@ +import { StandardFolder } from 'src/modules/messaging/message-import-manager/drivers/types/standard-folder'; + +export const MESSAGING_FOLDER_MANAGER_ALWAYS_EXCLUDED_FOLDERS = [ + StandardFolder.DRAFTS, + StandardFolder.TRASH, + StandardFolder.JUNK, +]; diff --git a/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/__tests__/should-create-folder-by-default.util.spec.ts b/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/__tests__/should-create-folder-by-default.util.spec.ts new file mode 100644 index 0000000000..b10370f793 --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/__tests__/should-create-folder-by-default.util.spec.ts @@ -0,0 +1,20 @@ +import { shouldCreateFolderByDefault } from 'src/modules/messaging/message-folder-manager/utils/should-create-folder-by-default.util'; +import { StandardFolder } from 'src/modules/messaging/message-import-manager/drivers/types/standard-folder'; + +describe('shouldCreateFolderByDefault', () => { + it('should allow creating user folders', () => { + expect(shouldCreateFolderByDefault(StandardFolder.INBOX)).toBe(true); + expect(shouldCreateFolderByDefault(StandardFolder.SENT)).toBe(true); + }); + + it('should allow creating custom folders', () => { + expect(shouldCreateFolderByDefault(null)).toBe(true); + expect(shouldCreateFolderByDefault(undefined)).toBe(true); + }); + + it('should prevent creating system-excluded folders', () => { + expect(shouldCreateFolderByDefault(StandardFolder.DRAFTS)).toBe(false); + expect(shouldCreateFolderByDefault(StandardFolder.TRASH)).toBe(false); + expect(shouldCreateFolderByDefault(StandardFolder.JUNK)).toBe(false); + }); +}); diff --git a/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/__tests__/should-sync-folder-by-default.util.spec.ts b/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/__tests__/should-sync-folder-by-default.util.spec.ts new file mode 100644 index 0000000000..bb783720ed --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/__tests__/should-sync-folder-by-default.util.spec.ts @@ -0,0 +1,24 @@ +import { MessageFolderImportPolicy } from 'src/modules/messaging/common/standard-objects/message-channel.workspace-entity'; +import { shouldSyncFolderByDefault } from 'src/modules/messaging/message-folder-manager/utils/should-sync-folder-by-default.util'; + +describe('shouldSyncFolderByDefault', () => { + describe('when messageFolderImportPolicy is SELECTED_FOLDERS', () => { + it('should return false for all folders', () => { + const result = shouldSyncFolderByDefault( + MessageFolderImportPolicy.SELECTED_FOLDERS, + ); + + expect(result).toBe(false); + }); + }); + + describe('when messageFolderImportPolicy is ALL_FOLDERS', () => { + it('should return true for all folders', () => { + const result = shouldSyncFolderByDefault( + MessageFolderImportPolicy.ALL_FOLDERS, + ); + + expect(result).toBe(true); + }); + }); +}); diff --git a/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/should-create-folder-by-default.util.ts b/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/should-create-folder-by-default.util.ts new file mode 100644 index 0000000000..a5139b9c7d --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/should-create-folder-by-default.util.ts @@ -0,0 +1,19 @@ +import { isDefined } from 'twenty-shared/utils'; + +import { MESSAGING_FOLDER_MANAGER_ALWAYS_EXCLUDED_FOLDERS } from 'src/modules/messaging/message-folder-manager/utils/MESSAGING_FOLDER_MANAGER_ALWAYS_EXCLUDED_FOLDERS'; +import { type StandardFolder } from 'src/modules/messaging/message-import-manager/drivers/types/standard-folder'; + +export const shouldCreateFolderByDefault = ( + standardFolder?: StandardFolder | null, +): boolean => { + if ( + isDefined(standardFolder) && + MESSAGING_FOLDER_MANAGER_ALWAYS_EXCLUDED_FOLDERS.includes( + standardFolder as StandardFolder, + ) + ) { + return false; + } + + return true; +}; diff --git a/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/should-sync-folder-by-default.util.ts b/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/should-sync-folder-by-default.util.ts new file mode 100644 index 0000000000..99b61431e6 --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/should-sync-folder-by-default.util.ts @@ -0,0 +1,7 @@ +import { MessageFolderImportPolicy } from 'src/modules/messaging/common/standard-objects/message-channel.workspace-entity'; + +export const shouldSyncFolderByDefault = ( + messageFolderImportPolicy: MessageFolderImportPolicy, +): boolean => { + return messageFolderImportPolicy === MessageFolderImportPolicy.ALL_FOLDERS; +}; diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/crons/commands/messaging-process-folder-actions.cron.command.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/crons/commands/messaging-process-folder-actions.cron.command.ts new file mode 100644 index 0000000000..84f97a2433 --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/crons/commands/messaging-process-folder-actions.cron.command.ts @@ -0,0 +1,35 @@ +import { Command, CommandRunner } from 'nest-commander'; + +import { InjectMessageQueue } from 'src/engine/core-modules/message-queue/decorators/message-queue.decorator'; +import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants'; +import { MessageQueueService } from 'src/engine/core-modules/message-queue/services/message-queue.service'; +import { + MESSAGING_PROCESS_FOLDER_ACTIONS_CRON_PATTERN, + MessagingProcessFolderActionsCronJob, +} from 'src/modules/messaging/message-import-manager/crons/jobs/messaging-process-folder-actions.cron.job'; + +@Command({ + name: 'cron:messaging:process-folder-actions', + description: + 'Starts a cron job to process pending folder actions (deletion) for message channels', +}) +export class MessagingProcessFolderActionsCronCommand extends CommandRunner { + constructor( + @InjectMessageQueue(MessageQueue.cronQueue) + private readonly messageQueueService: MessageQueueService, + ) { + super(); + } + + async run(): Promise { + await this.messageQueueService.addCron({ + jobName: MessagingProcessFolderActionsCronJob.name, + data: undefined, + options: { + repeat: { + pattern: MESSAGING_PROCESS_FOLDER_ACTIONS_CRON_PATTERN, + }, + }, + }); + } +} diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/crons/commands/messaging-process-group-email-actions.cron.command.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/crons/commands/messaging-process-group-email-actions.cron.command.ts new file mode 100644 index 0000000000..53015e3f1d --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/crons/commands/messaging-process-group-email-actions.cron.command.ts @@ -0,0 +1,35 @@ +import { Command, CommandRunner } from 'nest-commander'; + +import { InjectMessageQueue } from 'src/engine/core-modules/message-queue/decorators/message-queue.decorator'; +import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants'; +import { MessageQueueService } from 'src/engine/core-modules/message-queue/services/message-queue.service'; +import { + MESSAGING_PROCESS_GROUP_EMAIL_ACTIONS_CRON_PATTERN, + MessagingProcessGroupEmailActionsCronJob, +} from 'src/modules/messaging/message-import-manager/crons/jobs/messaging-process-group-email-actions.cron.job'; + +@Command({ + name: 'cron:messaging:process-group-email-actions', + description: + 'Starts a cron job to process pending group email actions (deletion or import) for message channels', +}) +export class MessagingProcessGroupEmailActionsCronCommand extends CommandRunner { + constructor( + @InjectMessageQueue(MessageQueue.cronQueue) + private readonly messageQueueService: MessageQueueService, + ) { + super(); + } + + async run(): Promise { + await this.messageQueueService.addCron({ + jobName: MessagingProcessGroupEmailActionsCronJob.name, + data: undefined, + options: { + repeat: { + pattern: MESSAGING_PROCESS_GROUP_EMAIL_ACTIONS_CRON_PATTERN, + }, + }, + }); + } +} diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/crons/jobs/messaging-process-folder-actions.cron.job.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/crons/jobs/messaging-process-folder-actions.cron.job.ts new file mode 100644 index 0000000000..95298dc3bd --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/crons/jobs/messaging-process-folder-actions.cron.job.ts @@ -0,0 +1,76 @@ +import { InjectDataSource, InjectRepository } from '@nestjs/typeorm'; + +import { WorkspaceActivationStatus } from 'twenty-shared/workspace'; +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'; +import { InjectMessageQueue } from 'src/engine/core-modules/message-queue/decorators/message-queue.decorator'; +import { Process } from 'src/engine/core-modules/message-queue/decorators/process.decorator'; +import { Processor } from 'src/engine/core-modules/message-queue/decorators/processor.decorator'; +import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants'; +import { MessageQueueService } from 'src/engine/core-modules/message-queue/services/message-queue.service'; +import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity'; +import { getWorkspaceSchemaName } from 'src/engine/workspace-datasource/utils/get-workspace-schema-name.util'; +import { MessageFolderPendingSyncAction } from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; +import { + MessagingProcessFolderActionsJob, + type MessagingProcessFolderActionsJobData, +} from 'src/modules/messaging/message-import-manager/jobs/messaging-process-folder-actions.job'; + +export const MESSAGING_PROCESS_FOLDER_ACTIONS_CRON_PATTERN = '*/15 * * * *'; + +@Processor(MessageQueue.cronQueue) +export class MessagingProcessFolderActionsCronJob { + constructor( + @InjectRepository(WorkspaceEntity) + private readonly workspaceRepository: Repository, + @InjectMessageQueue(MessageQueue.messagingQueue) + private readonly messageQueueService: MessageQueueService, + @InjectDataSource() + private readonly coreDataSource: DataSource, + private readonly exceptionHandlerService: ExceptionHandlerService, + ) {} + + @Process(MessagingProcessFolderActionsCronJob.name) + @SentryCronMonitor( + MessagingProcessFolderActionsCronJob.name, + MESSAGING_PROCESS_FOLDER_ACTIONS_CRON_PATTERN, + ) + async handle(): Promise { + const activeWorkspaces = await this.workspaceRepository.find({ + where: { + activationStatus: WorkspaceActivationStatus.ACTIVE, + }, + }); + + for (const activeWorkspace of activeWorkspaces) { + try { + const schemaName = getWorkspaceSchemaName(activeWorkspace.id); + + const messageChannels = await this.coreDataSource.query( + `SELECT DISTINCT mc.id + FROM ${schemaName}."messageChannel" mc + INNER JOIN ${schemaName}."messageFolder" mf ON mf."messageChannelId" = mc.id + WHERE mf."pendingSyncAction" = '${MessageFolderPendingSyncAction.FOLDER_DELETION}'`, + ); + + for (const messageChannel of messageChannels) { + await this.messageQueueService.add( + MessagingProcessFolderActionsJob.name, + { + workspaceId: activeWorkspace.id, + messageChannelId: messageChannel.id, + }, + ); + } + } catch (error) { + this.exceptionHandlerService.captureExceptions([error], { + workspace: { + id: activeWorkspace.id, + }, + }); + } + } + } +} diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/crons/jobs/messaging-process-group-email-actions.cron.job.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/crons/jobs/messaging-process-group-email-actions.cron.job.ts new file mode 100644 index 0000000000..2e636f3a38 --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/crons/jobs/messaging-process-group-email-actions.cron.job.ts @@ -0,0 +1,73 @@ +import { InjectDataSource, InjectRepository } from '@nestjs/typeorm'; + +import { WorkspaceActivationStatus } from 'twenty-shared/workspace'; +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'; +import { InjectMessageQueue } from 'src/engine/core-modules/message-queue/decorators/message-queue.decorator'; +import { Process } from 'src/engine/core-modules/message-queue/decorators/process.decorator'; +import { Processor } from 'src/engine/core-modules/message-queue/decorators/processor.decorator'; +import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants'; +import { MessageQueueService } from 'src/engine/core-modules/message-queue/services/message-queue.service'; +import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity'; +import { getWorkspaceSchemaName } from 'src/engine/workspace-datasource/utils/get-workspace-schema-name.util'; +import { MessageChannelPendingGroupEmailsAction } from 'src/modules/messaging/common/standard-objects/message-channel.workspace-entity'; +import { + MessagingProcessGroupEmailActionsJob, + type MessagingProcessGroupEmailActionsJobData, +} from 'src/modules/messaging/message-import-manager/jobs/messaging-process-group-email-actions.job'; + +export const MESSAGING_PROCESS_GROUP_EMAIL_ACTIONS_CRON_PATTERN = '0 */2 * * *'; + +@Processor(MessageQueue.cronQueue) +export class MessagingProcessGroupEmailActionsCronJob { + constructor( + @InjectRepository(WorkspaceEntity) + private readonly workspaceRepository: Repository, + @InjectMessageQueue(MessageQueue.messagingQueue) + private readonly messageQueueService: MessageQueueService, + @InjectDataSource() + private readonly coreDataSource: DataSource, + private readonly exceptionHandlerService: ExceptionHandlerService, + ) {} + + @Process(MessagingProcessGroupEmailActionsCronJob.name) + @SentryCronMonitor( + MessagingProcessGroupEmailActionsCronJob.name, + MESSAGING_PROCESS_GROUP_EMAIL_ACTIONS_CRON_PATTERN, + ) + async handle(): Promise { + const activeWorkspaces = await this.workspaceRepository.find({ + where: { + activationStatus: WorkspaceActivationStatus.ACTIVE, + }, + }); + + for (const activeWorkspace of activeWorkspaces) { + try { + const schemaName = getWorkspaceSchemaName(activeWorkspace.id); + + const messageChannels = await this.coreDataSource.query( + `SELECT id FROM ${schemaName}."messageChannel" WHERE "pendingGroupEmailsAction" IN ('${MessageChannelPendingGroupEmailsAction.GROUP_EMAILS_DELETION}', '${MessageChannelPendingGroupEmailsAction.GROUP_EMAILS_IMPORT}')`, + ); + + for (const messageChannel of messageChannels) { + await this.messageQueueService.add( + MessagingProcessGroupEmailActionsJob.name, + { + workspaceId: activeWorkspace.id, + messageChannelId: messageChannel.id, + }, + ); + } + } catch (error) { + this.exceptionHandlerService.captureExceptions([error], { + workspace: { + id: activeWorkspace.id, + }, + }); + } + } + } +} diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/jobs/messaging-process-folder-actions.job.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/jobs/messaging-process-folder-actions.job.ts new file mode 100644 index 0000000000..1fe47a92e2 --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/jobs/messaging-process-folder-actions.job.ts @@ -0,0 +1,92 @@ +import { Logger, Scope } from '@nestjs/common'; + +import { Process } from 'src/engine/core-modules/message-queue/decorators/process.decorator'; +import { Processor } from 'src/engine/core-modules/message-queue/decorators/processor.decorator'; +import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants'; +import { TwentyORMManager } from 'src/engine/twenty-orm/twenty-orm.manager'; +import { type MessageChannelWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-channel.workspace-entity'; +import { + MessageFolderPendingSyncAction, + type MessageFolderWorkspaceEntity, +} from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; +import { MessagingProcessFolderActionsService } from 'src/modules/messaging/message-import-manager/services/messaging-process-folder-actions.service'; + +export type MessagingProcessFolderActionsJobData = { + workspaceId: string; + messageChannelId: string; +}; + +@Processor({ + queueName: MessageQueue.messagingQueue, + scope: Scope.REQUEST, +}) +export class MessagingProcessFolderActionsJob { + private readonly logger = new Logger(MessagingProcessFolderActionsJob.name); + + constructor( + private readonly twentyORMManager: TwentyORMManager, + private readonly messagingProcessFolderActionsService: MessagingProcessFolderActionsService, + ) {} + + @Process(MessagingProcessFolderActionsJob.name) + async handle(data: MessagingProcessFolderActionsJobData): Promise { + const { workspaceId, messageChannelId } = data; + + this.logger.log( + `Processing pending folder actions for message channel ${messageChannelId} in workspace ${workspaceId}`, + ); + + const messageChannelRepository = + await this.twentyORMManager.getRepository( + 'messageChannel', + ); + + const messageChannel = await messageChannelRepository.findOne({ + where: { + id: messageChannelId, + }, + }); + + if (!messageChannel) { + this.logger.warn( + `Message channel ${messageChannelId} not found in workspace ${workspaceId}`, + ); + + return; + } + + const messageFolderRepository = + await this.twentyORMManager.getRepository( + 'messageFolder', + ); + + const messageFolders = await messageFolderRepository.find({ + where: { + messageChannelId: messageChannel.id, + pendingSyncAction: MessageFolderPendingSyncAction.FOLDER_DELETION, + }, + }); + + if (messageFolders.length === 0) { + this.logger.log( + `Message channel ${messageChannelId} has no folders with pending deletion actions, skipping`, + ); + + return; + } + + try { + await this.messagingProcessFolderActionsService.processFolderActions( + messageChannel, + messageFolders, + workspaceId, + ); + } catch (error) { + this.logger.error( + `Error processing folder actions for message channel ${messageChannelId} in workspace ${workspaceId}: ${error.message}`, + error.stack, + ); + throw error; + } + } +} diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/jobs/messaging-process-group-email-actions.job.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/jobs/messaging-process-group-email-actions.job.ts new file mode 100644 index 0000000000..0c6250afda --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/jobs/messaging-process-group-email-actions.job.ts @@ -0,0 +1,84 @@ +import { Logger, Scope } from '@nestjs/common'; + +import { Process } from 'src/engine/core-modules/message-queue/decorators/process.decorator'; +import { Processor } from 'src/engine/core-modules/message-queue/decorators/processor.decorator'; +import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants'; +import { TwentyORMManager } from 'src/engine/twenty-orm/twenty-orm.manager'; +import { + MessageChannelPendingGroupEmailsAction, + type MessageChannelWorkspaceEntity, +} from 'src/modules/messaging/common/standard-objects/message-channel.workspace-entity'; +import { MessagingProcessGroupEmailActionsService } from 'src/modules/messaging/message-import-manager/services/messaging-process-group-email-actions.service'; + +export type MessagingProcessGroupEmailActionsJobData = { + workspaceId: string; + messageChannelId: string; +}; + +@Processor({ + queueName: MessageQueue.messagingQueue, + scope: Scope.REQUEST, +}) +export class MessagingProcessGroupEmailActionsJob { + private readonly logger = new Logger( + MessagingProcessGroupEmailActionsJob.name, + ); + + constructor( + private readonly twentyORMManager: TwentyORMManager, + private readonly messagingProcessGroupEmailActionsService: MessagingProcessGroupEmailActionsService, + ) {} + + @Process(MessagingProcessGroupEmailActionsJob.name) + async handle(data: MessagingProcessGroupEmailActionsJobData): Promise { + const { workspaceId, messageChannelId } = data; + + this.logger.log( + `Processing pending group email action for message channel ${messageChannelId} in workspace ${workspaceId}`, + ); + + const messageChannelRepository = + await this.twentyORMManager.getRepository( + 'messageChannel', + ); + + const messageChannel = await messageChannelRepository.findOne({ + where: { + id: messageChannelId, + }, + }); + + if (!messageChannel) { + this.logger.warn( + `Message channel ${messageChannelId} not found in workspace ${workspaceId}`, + ); + + return; + } + + if ( + messageChannel.pendingGroupEmailsAction === + MessageChannelPendingGroupEmailsAction.NONE || + !messageChannel.pendingGroupEmailsAction + ) { + this.logger.log( + `Message channel ${messageChannelId} no longer has a pending action, skipping`, + ); + + return; + } + + try { + await this.messagingProcessGroupEmailActionsService.processGroupEmailActions( + messageChannel, + workspaceId, + ); + } catch (error) { + this.logger.error( + `Error processing group email actions for message channel ${messageChannelId} in workspace ${workspaceId}: ${error.message}`, + error.stack, + ); + throw error; + } + } +} diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/messaging-import-manager.module.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/messaging-import-manager.module.ts index 5736b5c404..344e8f3acd 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/messaging-import-manager.module.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/messaging-import-manager.module.ts @@ -18,10 +18,14 @@ import { MessagingSingleMessageImportCommand } from 'src/modules/messaging/messa import { MessagingMessageListFetchCronCommand } from 'src/modules/messaging/message-import-manager/crons/commands/messaging-message-list-fetch.cron.command'; import { MessagingMessagesImportCronCommand } from 'src/modules/messaging/message-import-manager/crons/commands/messaging-messages-import.cron.command'; import { MessagingOngoingStaleCronCommand } from 'src/modules/messaging/message-import-manager/crons/commands/messaging-ongoing-stale.cron.command'; +import { MessagingProcessFolderActionsCronCommand } from 'src/modules/messaging/message-import-manager/crons/commands/messaging-process-folder-actions.cron.command'; +import { MessagingProcessGroupEmailActionsCronCommand } from 'src/modules/messaging/message-import-manager/crons/commands/messaging-process-group-email-actions.cron.command'; import { MessagingRelaunchFailedMessageChannelsCronCommand } from 'src/modules/messaging/message-import-manager/crons/commands/messaging-relaunch-failed-message-channels.cron.command'; import { MessagingMessageListFetchCronJob } from 'src/modules/messaging/message-import-manager/crons/jobs/messaging-message-list-fetch.cron.job'; import { MessagingMessagesImportCronJob } from 'src/modules/messaging/message-import-manager/crons/jobs/messaging-messages-import.cron.job'; import { MessagingOngoingStaleCronJob } from 'src/modules/messaging/message-import-manager/crons/jobs/messaging-ongoing-stale.cron.job'; +import { MessagingProcessFolderActionsCronJob } from 'src/modules/messaging/message-import-manager/crons/jobs/messaging-process-folder-actions.cron.job'; +import { MessagingProcessGroupEmailActionsCronJob } from 'src/modules/messaging/message-import-manager/crons/jobs/messaging-process-group-email-actions.cron.job'; import { MessagingRelaunchFailedMessageChannelsCronJob } from 'src/modules/messaging/message-import-manager/crons/jobs/messaging-relaunch-failed-message-channels.cron.job'; import { MessagingGmailDriverModule } from 'src/modules/messaging/message-import-manager/drivers/gmail/messaging-gmail-driver.module'; import { MessagingIMAPDriverModule } from 'src/modules/messaging/message-import-manager/drivers/imap/messaging-imap-driver.module'; @@ -32,17 +36,23 @@ import { MessagingCleanCacheJob } from 'src/modules/messaging/message-import-man import { MessagingMessageListFetchJob } from 'src/modules/messaging/message-import-manager/jobs/messaging-message-list-fetch.job'; import { MessagingMessagesImportJob } from 'src/modules/messaging/message-import-manager/jobs/messaging-messages-import.job'; import { MessagingOngoingStaleJob } from 'src/modules/messaging/message-import-manager/jobs/messaging-ongoing-stale.job'; +import { MessagingProcessFolderActionsJob } from 'src/modules/messaging/message-import-manager/jobs/messaging-process-folder-actions.job'; +import { MessagingProcessGroupEmailActionsJob } from 'src/modules/messaging/message-import-manager/jobs/messaging-process-group-email-actions.job'; import { MessagingRelaunchFailedMessageChannelJob } from 'src/modules/messaging/message-import-manager/jobs/messaging-relaunch-failed-message-channel.job'; import { MessagingMessageImportManagerMessageChannelListener } from 'src/modules/messaging/message-import-manager/listeners/messaging-import-manager-message-channel.listener'; import { MessagingAccountAuthenticationService } from 'src/modules/messaging/message-import-manager/services/messaging-account-authentication.service'; import { MessagingClearCursorsModule } from 'src/modules/messaging/message-import-manager/services/messaging-clear-cursors.module'; import { MessagingCursorService } from 'src/modules/messaging/message-import-manager/services/messaging-cursor.service'; +import { MessagingDeleteFolderMessagesService } from 'src/modules/messaging/message-import-manager/services/messaging-delete-folder-messages.service'; +import { MessagingDeleteGroupEmailMessagesService } from 'src/modules/messaging/message-import-manager/services/messaging-delete-group-email-messages.service'; import { MessagingGetMessageListService } from 'src/modules/messaging/message-import-manager/services/messaging-get-message-list.service'; import { MessagingGetMessagesService } from 'src/modules/messaging/message-import-manager/services/messaging-get-messages.service'; import { MessageImportExceptionHandlerService } from 'src/modules/messaging/message-import-manager/services/messaging-import-exception-handler.service'; import { MessagingMessageListFetchService } from 'src/modules/messaging/message-import-manager/services/messaging-message-list-fetch.service'; import { MessagingMessageService } from 'src/modules/messaging/message-import-manager/services/messaging-message.service'; import { MessagingMessagesImportService } from 'src/modules/messaging/message-import-manager/services/messaging-messages-import.service'; +import { MessagingProcessFolderActionsService } from 'src/modules/messaging/message-import-manager/services/messaging-process-folder-actions.service'; +import { MessagingProcessGroupEmailActionsService } from 'src/modules/messaging/message-import-manager/services/messaging-process-group-email-actions.service'; import { MessagingSaveMessagesAndEnqueueContactCreationService } from 'src/modules/messaging/message-import-manager/services/messaging-save-messages-and-enqueue-contact-creation.service'; import { MessagingSendMessageService } from 'src/modules/messaging/message-import-manager/services/messaging-send-message.service'; import { MessageParticipantManagerModule } from 'src/modules/messaging/message-participant-manager/message-participant-manager.module'; @@ -76,15 +86,21 @@ import { MessagingMonitoringModule } from 'src/modules/messaging/monitoring/mess MessagingMessageListFetchCronCommand, MessagingMessagesImportCronCommand, MessagingOngoingStaleCronCommand, + MessagingProcessFolderActionsCronCommand, + MessagingProcessGroupEmailActionsCronCommand, MessagingRelaunchFailedMessageChannelsCronCommand, MessagingSingleMessageImportCommand, MessagingMessageListFetchJob, MessagingMessagesImportJob, MessagingOngoingStaleJob, + MessagingProcessFolderActionsJob, + MessagingProcessGroupEmailActionsJob, MessagingRelaunchFailedMessageChannelJob, MessagingMessageListFetchCronJob, MessagingMessagesImportCronJob, MessagingOngoingStaleCronJob, + MessagingProcessFolderActionsCronJob, + MessagingProcessGroupEmailActionsCronJob, MessagingRelaunchFailedMessageChannelsCronJob, MessagingAddSingleMessageToCacheForImportJob, MessagingMessageImportManagerMessageChannelListener, @@ -99,6 +115,10 @@ import { MessagingMonitoringModule } from 'src/modules/messaging/monitoring/mess MessagingCursorService, MessagingSendMessageService, MessagingAccountAuthenticationService, + MessagingProcessFolderActionsService, + MessagingProcessGroupEmailActionsService, + MessagingDeleteFolderMessagesService, + MessagingDeleteGroupEmailMessagesService, ], exports: [ MessagingSendMessageService, @@ -106,6 +126,7 @@ import { MessagingMonitoringModule } from 'src/modules/messaging/monitoring/mess MessagingMessagesImportCronCommand, MessagingOngoingStaleCronCommand, MessagingRelaunchFailedMessageChannelsCronCommand, + MessagingProcessGroupEmailActionsService, ], }) export class MessagingImportManagerModule {} diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-delete-folder-messages.service.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-delete-folder-messages.service.ts new file mode 100644 index 0000000000..173ab17246 --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-delete-folder-messages.service.ts @@ -0,0 +1,77 @@ +import { Injectable, Logger } from '@nestjs/common'; + +import chunk from 'lodash.chunk'; +import { isDefined } from 'twenty-shared/utils'; + +import { type MessageChannelWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-channel.workspace-entity'; +import { type MessageFolderWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; +import { MessagingMessageCleanerService } from 'src/modules/messaging/message-cleaner/services/messaging-message-cleaner.service'; +import { MessagingGetMessageListService } from 'src/modules/messaging/message-import-manager/services/messaging-get-message-list.service'; + +@Injectable() +export class MessagingDeleteFolderMessagesService { + private readonly logger = new Logger( + MessagingDeleteFolderMessagesService.name, + ); + + constructor( + private readonly messagingMessageCleanerService: MessagingMessageCleanerService, + private readonly messagingGetMessageListService: MessagingGetMessageListService, + ) {} + + async deleteFolderMessages( + workspaceId: string, + messageChannel: MessageChannelWorkspaceEntity, + messageFolder: MessageFolderWorkspaceEntity, + ): Promise { + this.logger.log( + `WorkspaceId: ${workspaceId}, MessageChannelId: ${messageChannel.id}, FolderId: ${messageFolder.id} - Deleting messages from folder: ${messageFolder.name}`, + ); + + const messageLists = + await this.messagingGetMessageListService.getMessageLists( + messageChannel, + [messageFolder], + ); + + let totalDeletedCount = 0; + + for (const messageList of messageLists) { + const { messageExternalIds } = messageList; + + if (messageExternalIds.length === 0) { + continue; + } + + const messageExternalIdsChunks = chunk(messageExternalIds, 200); + + for (const messageExternalIdsChunk of messageExternalIdsChunks) { + const validExternalIds = messageExternalIdsChunk.filter(isDefined); + + if (validExternalIds.length === 0) { + continue; + } + + await this.messagingMessageCleanerService.deleteMessagesChannelMessageAssociationsAndRelatedOrphans( + { + workspaceId, + messageExternalIds: validExternalIds, + messageChannelId: messageChannel.id, + }, + ); + + totalDeletedCount += validExternalIds.length; + + this.logger.log( + `WorkspaceId: ${workspaceId}, MessageChannelId: ${messageChannel.id}, FolderId: ${messageFolder.id} - Processed ${validExternalIds.length} message deletions`, + ); + } + } + + this.logger.log( + `WorkspaceId: ${workspaceId}, MessageChannelId: ${messageChannel.id}, FolderId: ${messageFolder.id} - Completed deleting ${totalDeletedCount} messages from folder: ${messageFolder.name}`, + ); + + return totalDeletedCount; + } +} diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-delete-group-email-messages.service.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-delete-group-email-messages.service.ts new file mode 100644 index 0000000000..a2a6209a1c --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-delete-group-email-messages.service.ts @@ -0,0 +1,126 @@ +import { Injectable, Logger } from '@nestjs/common'; + +import chunk from 'lodash.chunk'; +import { isDefined } from 'twenty-shared/utils'; + +import { TwentyORMGlobalManager } from 'src/engine/twenty-orm/twenty-orm-global.manager'; +import { MessageChannelMessageAssociationWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-channel-message-association.workspace-entity'; +import { MessagingMessageCleanerService } from 'src/modules/messaging/message-cleaner/services/messaging-message-cleaner.service'; +import { isGroupEmail } from 'src/modules/messaging/message-import-manager/utils/is-group-email'; + +const MESSAGE_CHANNEL_MESSAGE_ASSOCIATION_BATCH_SIZE = 500; + +type MessageBatchRawResult = { + messageId: string; + messageExternalId: string; + participantHandle: string; +}; + +@Injectable() +export class MessagingDeleteGroupEmailMessagesService { + private readonly logger = new Logger( + MessagingDeleteGroupEmailMessagesService.name, + ); + + constructor( + private readonly twentyORMGlobalManager: TwentyORMGlobalManager, + private readonly messagingMessageCleanerService: MessagingMessageCleanerService, + ) {} + + async deleteGroupEmailMessages( + workspaceId: string, + messageChannelId: string, + ): Promise { + this.logger.log( + `WorkspaceId: ${workspaceId}, MessageChannelId: ${messageChannelId} - Deleting messages from group email addresses`, + ); + + const messageChannelMessageAssociationRepository = + await this.twentyORMGlobalManager.getRepositoryForWorkspace( + workspaceId, + 'messageChannelMessageAssociation', + ); + + let offset = 0; + let totalDeletedCount = 0; + + while (true) { + const batch = await messageChannelMessageAssociationRepository + .createQueryBuilder('mcma') + .select('mcma.messageId', 'messageId') + .addSelect('mcma.messageExternalId', 'messageExternalId') + .addSelect('participant.handle', 'participantHandle') + .innerJoin('mcma.message', 'message') + .innerJoin( + 'message.messageParticipants', + 'participant', + 'participant.role = :role', + { role: 'from' }, + ) + .where('mcma.messageChannelId = :messageChannelId', { + messageChannelId, + }) + .skip(offset) + .take(MESSAGE_CHANNEL_MESSAGE_ASSOCIATION_BATCH_SIZE) + .getRawMany(); + + if (batch.length === 0) { + break; + } + + const groupEmailRecords = batch.filter( + (record) => + isDefined(record.participantHandle) && + isGroupEmail(record.participantHandle), + ); + + if (groupEmailRecords.length > 0) { + const uniqueMessageIds = new Set( + groupEmailRecords.map((r) => r.messageId), + ); + + const messageExternalIdsToDelete = batch + .filter((record) => uniqueMessageIds.has(record.messageId)) + .map((record) => record.messageExternalId) + .filter(isDefined); + + if (messageExternalIdsToDelete.length > 0) { + const messageExternalIdsChunks = chunk( + messageExternalIdsToDelete, + 200, + ); + + for (const messageExternalIdsChunk of messageExternalIdsChunks) { + await this.messagingMessageCleanerService.deleteMessagesChannelMessageAssociationsAndRelatedOrphans( + { + workspaceId, + messageExternalIds: messageExternalIdsChunk, + messageChannelId, + }, + ); + + totalDeletedCount += messageExternalIdsChunk.length; + + this.logger.log( + `WorkspaceId: ${workspaceId}, MessageChannelId: ${messageChannelId} - Deleted ${messageExternalIdsChunk.length} group email messages`, + ); + } + } + } + + if (batch.length < MESSAGE_CHANNEL_MESSAGE_ASSOCIATION_BATCH_SIZE) { + break; + } + + if (groupEmailRecords.length === 0) { + offset += MESSAGE_CHANNEL_MESSAGE_ASSOCIATION_BATCH_SIZE; + } + } + + this.logger.log( + `WorkspaceId: ${workspaceId}, MessageChannelId: ${messageChannelId} - Completed deleting ${totalDeletedCount} group email messages`, + ); + + return totalDeletedCount; + } +} diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-message-list-fetch.service.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-message-list-fetch.service.ts index 8909124b1a..bd3e32e3b0 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-message-list-fetch.service.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-message-list-fetch.service.ts @@ -12,10 +12,14 @@ import { TwentyORMManager } from 'src/engine/twenty-orm/twenty-orm.manager'; import { MessageChannelSyncStatusService } from 'src/modules/messaging/common/services/message-channel-sync-status.service'; import { type MessageChannelMessageAssociationWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-channel-message-association.workspace-entity'; import { + MessageChannelPendingGroupEmailsAction, MessageChannelSyncStage, MessageChannelWorkspaceEntity, } from 'src/modules/messaging/common/standard-objects/message-channel.workspace-entity'; -import { MessageFolderWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; +import { + MessageFolderPendingSyncAction, + MessageFolderWorkspaceEntity, +} from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; import { MessagingMessageCleanerService } from 'src/modules/messaging/message-cleaner/services/messaging-message-cleaner.service'; import { SyncMessageFoldersService } from 'src/modules/messaging/message-folder-manager/services/sync-message-folders.service'; import { MessagingAccountAuthenticationService } from 'src/modules/messaging/message-import-manager/services/messaging-account-authentication.service'; @@ -51,6 +55,19 @@ export class MessagingMessageListFetchService { workspaceId: string, ) { try { + if ( + messageChannel.pendingGroupEmailsAction === + MessageChannelPendingGroupEmailsAction.GROUP_EMAILS_DELETION || + messageChannel.pendingGroupEmailsAction === + MessageChannelPendingGroupEmailsAction.GROUP_EMAILS_IMPORT + ) { + this.logger.log( + `messageChannelId: ${messageChannel.id} Skipping message list fetch due to pending group emails action: ${messageChannel.pendingGroupEmailsAction}`, + ); + + return; + } + await this.messageChannelSyncStatusService.markAsMessagesListFetchOngoing( [messageChannel.id], ); @@ -81,8 +98,7 @@ export class MessagingMessageListFetchService { await this.syncMessageFoldersService.syncMessageFolders({ workspaceId, - messageChannelId: messageChannelWithFreshTokens.id, - connectedAccount: messageChannelWithFreshTokens.connectedAccount, + messageChannel: messageChannelWithFreshTokens, manager: datasource.manager, }); @@ -94,6 +110,7 @@ export class MessagingMessageListFetchService { const messageFolders = await messageFolderRepository.find({ where: { messageChannelId: messageChannel.id, + pendingSyncAction: MessageFolderPendingSyncAction.NONE, }, }); diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-messages-import.service.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-messages-import.service.ts index 661c10628c..57ffaa8e41 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-messages-import.service.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-messages-import.service.ts @@ -119,6 +119,7 @@ export class MessagingMessagesImportService { [...connectedAccountWithFreshTokens.handleAliases.split(',')], allMessages, blocklist.map((blocklistItem) => blocklistItem.handle), + messageChannel.excludeGroupEmails, ); if (messagesToSave.length > 0) { diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-process-folder-actions.service.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-process-folder-actions.service.ts new file mode 100644 index 0000000000..5966bbb01a --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-process-folder-actions.service.ts @@ -0,0 +1,124 @@ +import { Injectable, Logger } from '@nestjs/common'; + +import { isDefined } from 'twenty-shared/utils'; +import { In } from 'typeorm'; + +import { type WorkspaceEntityManager } from 'src/engine/twenty-orm/entity-manager/workspace-entity-manager'; +import { TwentyORMManager } from 'src/engine/twenty-orm/twenty-orm.manager'; +import { type MessageChannelWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-channel.workspace-entity'; +import { + MessageFolderPendingSyncAction, + MessageFolderWorkspaceEntity, +} from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; +import { MessagingDeleteFolderMessagesService } from 'src/modules/messaging/message-import-manager/services/messaging-delete-folder-messages.service'; + +@Injectable() +export class MessagingProcessFolderActionsService { + private readonly logger = new Logger( + MessagingProcessFolderActionsService.name, + ); + + constructor( + private readonly twentyORMManager: TwentyORMManager, + private readonly messagingDeleteFolderMessagesService: MessagingDeleteFolderMessagesService, + ) {} + + async processFolderActions( + messageChannel: MessageChannelWorkspaceEntity, + messageFolders: MessageFolderWorkspaceEntity[], + workspaceId: string, + ): Promise { + const foldersWithPendingActions = messageFolders.filter( + (folder) => + isDefined(folder.pendingSyncAction) && + folder.pendingSyncAction !== MessageFolderPendingSyncAction.NONE, + ); + + if (foldersWithPendingActions.length === 0) { + return; + } + + this.logger.log( + `WorkspaceId: ${workspaceId}, MessageChannelId: ${messageChannel.id} - Processing ${foldersWithPendingActions.length} folders with pending actions`, + ); + + const folderIdsToDelete: string[] = []; + const processedFolderIds: string[] = []; + const failedFolderIds: Array<{ folderId: string; error: Error }> = []; + + for (const folder of foldersWithPendingActions) { + try { + this.logger.log( + `WorkspaceId: ${workspaceId}, MessageChannelId: ${messageChannel.id}, FolderId: ${folder.id} - Processing folder action: ${folder.pendingSyncAction}`, + ); + + if ( + folder.pendingSyncAction === + MessageFolderPendingSyncAction.FOLDER_DELETION + ) { + await this.messagingDeleteFolderMessagesService.deleteFolderMessages( + workspaceId, + messageChannel, + folder, + ); + + folderIdsToDelete.push(folder.id); + + this.logger.log( + `WorkspaceId: ${workspaceId}, MessageChannelId: ${messageChannel.id}, FolderId: ${folder.id} - Completed FOLDER_DELETION action`, + ); + } + + processedFolderIds.push(folder.id); + } catch (error) { + this.logger.error( + `WorkspaceId: ${workspaceId}, MessageChannelId: ${messageChannel.id}, FolderId: ${folder.id} - Error processing folder action: ${error.message}`, + error.stack, + ); + failedFolderIds.push({ folderId: folder.id, error }); + } + } + + if (failedFolderIds.length > 0) { + this.logger.warn( + `WorkspaceId: ${workspaceId}, MessageChannelId: ${messageChannel.id} - Failed to process ${failedFolderIds.length} folders. They will be retried on next sync.`, + ); + } + + if (processedFolderIds.length > 0 || folderIdsToDelete.length > 0) { + const workspaceDataSource = await this.twentyORMManager.getDatasource(); + + await workspaceDataSource?.transaction( + async (transactionManager: WorkspaceEntityManager) => { + const messageFolderRepository = + await this.twentyORMManager.getRepository( + 'messageFolder', + ); + + if (processedFolderIds.length > 0) { + await messageFolderRepository.update( + { id: In(processedFolderIds) }, + { pendingSyncAction: MessageFolderPendingSyncAction.NONE }, + transactionManager, + ); + + this.logger.log( + `WorkspaceId: ${workspaceId}, MessageChannelId: ${messageChannel.id} - Reset pendingSyncAction to NONE for ${processedFolderIds.length} folders`, + ); + } + + if (folderIdsToDelete.length > 0) { + await messageFolderRepository.delete( + { id: In(folderIdsToDelete) }, + transactionManager, + ); + + this.logger.log( + `WorkspaceId: ${workspaceId}, MessageChannelId: ${messageChannel.id} - Deleted ${folderIdsToDelete.length} folders`, + ); + } + }, + ); + } + } +} diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-process-group-email-actions.service.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-process-group-email-actions.service.ts new file mode 100644 index 0000000000..f2a18c3e9c --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-process-group-email-actions.service.ts @@ -0,0 +1,157 @@ +import { Injectable, Logger } from '@nestjs/common'; + +import { isDefined } from 'twenty-shared/utils'; + +import { type WorkspaceEntityManager } from 'src/engine/twenty-orm/entity-manager/workspace-entity-manager'; +import { TwentyORMManager } from 'src/engine/twenty-orm/twenty-orm.manager'; +import { MessageChannelSyncStatusService } from 'src/modules/messaging/common/services/message-channel-sync-status.service'; +import { + MessageChannelPendingGroupEmailsAction, + MessageChannelWorkspaceEntity, +} from 'src/modules/messaging/common/standard-objects/message-channel.workspace-entity'; +import { MessagingClearCursorsService } from 'src/modules/messaging/message-import-manager/services/messaging-clear-cursors.service'; +import { MessagingDeleteGroupEmailMessagesService } from 'src/modules/messaging/message-import-manager/services/messaging-delete-group-email-messages.service'; + +@Injectable() +export class MessagingProcessGroupEmailActionsService { + private readonly logger = new Logger( + MessagingProcessGroupEmailActionsService.name, + ); + + constructor( + private readonly twentyORMManager: TwentyORMManager, + private readonly messagingDeleteGroupEmailMessagesService: MessagingDeleteGroupEmailMessagesService, + private readonly messagingClearCursorsService: MessagingClearCursorsService, + private readonly messageChannelSyncStatusService: MessageChannelSyncStatusService, + ) {} + + async markMessageChannelAsPendingGroupEmailsAction( + messageChannel: MessageChannelWorkspaceEntity, + workspaceId: string, + pendingGroupEmailsAction: MessageChannelPendingGroupEmailsAction, + ): Promise { + const messageChannelRepository = + await this.twentyORMManager.getRepository( + 'messageChannel', + ); + + await messageChannelRepository.update( + { id: messageChannel.id }, + { pendingGroupEmailsAction }, + ); + + this.logger.log( + `WorkspaceId: ${workspaceId}, MessageChannelId: ${messageChannel.id} - Marked message channel as pending group emails action: ${pendingGroupEmailsAction}`, + ); + } + + async processGroupEmailActions( + messageChannel: MessageChannelWorkspaceEntity, + workspaceId: string, + ): Promise { + const { pendingGroupEmailsAction } = messageChannel; + + if ( + !isDefined(pendingGroupEmailsAction) || + pendingGroupEmailsAction === MessageChannelPendingGroupEmailsAction.NONE + ) { + return; + } + + this.logger.log( + `WorkspaceId: ${workspaceId}, MessageChannelId: ${messageChannel.id} - Processing group email action: ${pendingGroupEmailsAction}`, + ); + + const workspaceDataSource = await this.twentyORMManager.getDatasource(); + + await workspaceDataSource?.transaction( + async (transactionManager: WorkspaceEntityManager) => { + try { + const messageChannelRepository = + await this.twentyORMManager.getRepository( + 'messageChannel', + ); + + switch (pendingGroupEmailsAction) { + case MessageChannelPendingGroupEmailsAction.GROUP_EMAILS_DELETION: + await this.handleGroupEmailsDeletion( + workspaceId, + messageChannel.id, + transactionManager, + ); + break; + case MessageChannelPendingGroupEmailsAction.GROUP_EMAILS_IMPORT: + await this.handleGroupEmailsImport( + workspaceId, + messageChannel.id, + transactionManager, + ); + break; + } + + await messageChannelRepository.update( + { id: messageChannel.id }, + { + pendingGroupEmailsAction: + MessageChannelPendingGroupEmailsAction.NONE, + }, + transactionManager, + ); + + this.logger.log( + `WorkspaceId: ${workspaceId}, MessageChannelId: ${messageChannel.id} - Reset pendingGroupEmailsAction to NONE`, + ); + } catch (error) { + this.logger.error( + `WorkspaceId: ${workspaceId}, MessageChannelId: ${messageChannel.id} - Error processing group email action: ${error.message}`, + error.stack, + ); + throw error; + } + }, + ); + + await this.messageChannelSyncStatusService.scheduleMessageListFetch([ + messageChannel.id, + ]); + + this.logger.log( + `WorkspaceId: ${workspaceId}, MessageChannelId: ${messageChannel.id} - Scheduled message list fetch after processing group email action`, + ); + } + + private async handleGroupEmailsDeletion( + workspaceId: string, + messageChannelId: string, + transactionManager: WorkspaceEntityManager, + ): Promise { + await this.messagingDeleteGroupEmailMessagesService.deleteGroupEmailMessages( + workspaceId, + messageChannelId, + ); + + await this.messagingClearCursorsService.clearAllMessageChannelCursors( + messageChannelId, + transactionManager, + ); + + this.logger.log( + `WorkspaceId: ${workspaceId}, MessageChannelId: ${messageChannelId} - Completed GROUP_EMAILS_DELETION action`, + ); + } + + private async handleGroupEmailsImport( + workspaceId: string, + messageChannelId: string, + transactionManager: WorkspaceEntityManager, + ): Promise { + await this.messagingClearCursorsService.clearAllMessageChannelCursors( + messageChannelId, + transactionManager, + ); + + this.logger.log( + `WorkspaceId: ${workspaceId}, MessageChannelId: ${messageChannelId} - Completed GROUP_EMAILS_IMPORT action`, + ); + } +} diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-save-messages-and-enqueue-contact-creation.service.spec.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-save-messages-and-enqueue-contact-creation.service.spec.ts index 80bacc6226..5112e34238 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-save-messages-and-enqueue-contact-creation.service.spec.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-save-messages-and-enqueue-contact-creation.service.spec.ts @@ -228,36 +228,6 @@ describe('MessagingSaveMessagesAndEnqueueContactCreationService', () => { ); }); - it('should not create group emails contacts', async () => { - await service.saveMessagesAndEnqueueContactCreation( - [ - { - ...mockMessages[0], - participants: [ - { - role: 'from', - handle: 'contact@group.com', - displayName: 'participant that is the Connected Account', - }, - ], - }, - ], - mockMessageChannel, - mockConnectedAccount, - workspaceId, - ); - - expect(messageQueueService.add).toHaveBeenCalledWith( - CreateCompanyAndContactJob.name, - { - workspaceId, - connectedAccount: mockConnectedAccount, - source: FieldActorSource.EMAIL, - contactsToCreate: [], - }, - ); - }); - it('should not create personal emails contacts', async () => { await service.saveMessagesAndEnqueueContactCreation( [ diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-save-messages-and-enqueue-contact-creation.service.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-save-messages-and-enqueue-contact-creation.service.ts index cd95942b50..e8d01e502e 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-save-messages-and-enqueue-contact-creation.service.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-save-messages-and-enqueue-contact-creation.service.ts @@ -22,7 +22,6 @@ import { } from 'src/modules/messaging/message-import-manager/drivers/gmail/types/gmail-message.type'; import { MessagingMessageService } from 'src/modules/messaging/message-import-manager/services/messaging-message.service'; import { type MessageWithParticipants } from 'src/modules/messaging/message-import-manager/types/message'; -import { isGroupEmail } from 'src/modules/messaging/message-import-manager/utils/is-group-email'; import { MessagingMessageParticipantService } from 'src/modules/messaging/message-participant-manager/services/messaging-message-participant.service'; import { isWorkEmail } from 'src/utils/is-work-email'; @@ -79,15 +78,10 @@ export class MessagingSaveMessagesAndEnqueueContactCreationService { messageChannel.excludeNonProfessionalEmails && !isWorkEmail(participant.handle); - const isExcludedByGroupEmails = - messageChannel.excludeGroupEmails && - isGroupEmail(participant.handle); - const shouldCreateContact = !!participant.handle && !isParticipantConnectedAccount && !isExcludedByNonProfessionalEmails && - !isExcludedByGroupEmails && (messageChannel.contactAutoCreationPolicy === MessageChannelContactAutoCreationPolicy.SENT_AND_RECEIVED || (messageChannel.contactAutoCreationPolicy ===