diff --git a/packages/twenty-server/src/modules/messaging/message-cleaner/services/messaging-message-cleaner.service.ts b/packages/twenty-server/src/modules/messaging/message-cleaner/services/messaging-message-cleaner.service.ts index ddb7f038ff..746c341afa 100644 --- a/packages/twenty-server/src/modules/messaging/message-cleaner/services/messaging-message-cleaner.service.ts +++ b/packages/twenty-server/src/modules/messaging/message-cleaner/services/messaging-message-cleaner.service.ts @@ -1,4 +1,4 @@ -import { Injectable } from '@nestjs/common'; +import { Injectable, Logger } from '@nestjs/common'; import { IsNull } from 'typeorm'; @@ -10,6 +10,7 @@ import { deleteUsingPagination } from 'src/modules/messaging/message-cleaner/uti @Injectable() export class MessagingMessageCleanerService { + private readonly logger = new Logger(MessagingMessageCleanerService.name); constructor(private readonly twentyORMManager: TwentyORMManager) {} public async cleanWorkspaceThreads(workspaceId: string) { @@ -57,6 +58,9 @@ export class MessagingMessageCleanerService { workspaceId: string, transactionManager?: WorkspaceEntityManager, ) => { + this.logger.log( + `WorkspaceId: ${workspaceId} Deleting ${ids.length} messages from message cleaner`, + ); await messageRepository.delete(ids, transactionManager); }, transactionManager, 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 efbe01acc9..5f87a93346 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 @@ -91,14 +91,15 @@ export class MessagingMessageListFetchService { (messageList) => messageList.messageExternalIdsToDelete, ); - const isFullSync = messageLists.every( - (messageList) => !isNonEmptyString(messageList.previousSyncCursor), - ); + const isFullSync = + messageLists.every( + (messageList) => !isNonEmptyString(messageList.previousSyncCursor), + ) && isNonEmptyString(messageChannel.syncCursor); let totalMessagesToImportCount = 0; this.logger.log( - `messageChannelId: ${messageChannel.id} Is full sync: ${isFullSync} and message lists length: ${messageLists.length}`, + `messageChannelId: ${messageChannel.id} Is full sync: ${isFullSync} and toImportCount: ${messageExternalIds.length}, toDeleteCount: ${messageExternalIdsToDelete.length}`, ); const messageChannelMessageAssociationRepository = @@ -158,83 +159,12 @@ export class MessagingMessageListFetchService { ); } - const fullSyncMessageChannelMessageAssociationsToDelete = []; - - if (isFullSync) { - const firstMessageChannelMessageAssociation = - await messageChannelMessageAssociationRepository.findOne({ - where: { - messageChannelId: messageChannelWithFreshTokens.id, - }, - order: { - id: 'ASC', - }, - }); - - if (!isDefined(firstMessageChannelMessageAssociation)) { - this.logger.log( - `messageChannelId: ${messageChannel.id} Full sync: No message channel message associations found`, - ); - - return; - } - - this.logger.log( - `messageChannelId: ${messageChannel.id} Full sync: First message channel message association id: ${firstMessageChannelMessageAssociation.id}`, - ); - - let nextFirstBatchMessageChannelMessageAssociationId: - | string - | undefined = firstMessageChannelMessageAssociation.id; - let batchIndex = 0; - - while (isDefined(nextFirstBatchMessageChannelMessageAssociationId)) { - const existingMessageChannelMessageAssociations = - await messageChannelMessageAssociationRepository.find({ - where: { - messageChannelId: messageChannelWithFreshTokens.id, - id: MoreThanOrEqual( - nextFirstBatchMessageChannelMessageAssociationId, - ), - }, - order: { - id: 'ASC', - }, - take: 200, - }); - - if (existingMessageChannelMessageAssociations.length < 200) { - nextFirstBatchMessageChannelMessageAssociationId = undefined; - break; - } - - nextFirstBatchMessageChannelMessageAssociationId = - existingMessageChannelMessageAssociations[ - existingMessageChannelMessageAssociations.length - 1 - ].id; - - batchIndex++; - - const messageChannelMessageAssociationsToDelete = - existingMessageChannelMessageAssociations.filter( - (existingMessageChannelMessageAssociation) => - isDefined( - existingMessageChannelMessageAssociation.messageExternalId, - ) && - !messageExternalIds.includes( - existingMessageChannelMessageAssociation.messageExternalId, - ), - ); - - this.logger.log( - `messageChannelId: ${messageChannel.id} Full sync: Message channel message associations to delete in batch ${batchIndex}: ${messageChannelMessageAssociationsToDelete.length}`, - ); - - fullSyncMessageChannelMessageAssociationsToDelete.push( - ...messageChannelMessageAssociationsToDelete, - ); - } - } + const fullSyncMessageChannelMessageAssociationsToDelete = isFullSync + ? await this.computeFullSyncMessageChannelMessageAssociationsToDelete( + messageChannel, + messageExternalIds, + ) + : []; const allMessageExternalIdsToDelete = [ ...messageExternalIdsToDelete, @@ -308,4 +238,91 @@ export class MessagingMessageListFetchService { ); } } + + private async computeFullSyncMessageChannelMessageAssociationsToDelete( + messageChannel: Pick, + messageExternalIds: string[], + ) { + const messageChannelMessageAssociationRepository = + await this.twentyORMManager.getRepository( + 'messageChannelMessageAssociation', + ); + + const fullSyncMessageChannelMessageAssociationsToDelete = []; + + const firstMessageChannelMessageAssociation = + await messageChannelMessageAssociationRepository.findOne({ + where: { + messageChannelId: messageChannel.id, + }, + order: { + id: 'ASC', + }, + }); + + if (!isDefined(firstMessageChannelMessageAssociation)) { + this.logger.log( + `messageChannelId: ${messageChannel.id} Full sync: No message channel message associations found`, + ); + + return []; + } + + this.logger.log( + `messageChannelId: ${messageChannel.id} Full sync: First message channel message association id: ${firstMessageChannelMessageAssociation.id}`, + ); + + let nextFirstBatchMessageChannelMessageAssociationId: string | undefined = + firstMessageChannelMessageAssociation.id; + let batchIndex = 0; + + while (isDefined(nextFirstBatchMessageChannelMessageAssociationId)) { + const existingMessageChannelMessageAssociations = + await messageChannelMessageAssociationRepository.find({ + where: { + messageChannelId: messageChannel.id, + id: MoreThanOrEqual( + nextFirstBatchMessageChannelMessageAssociationId, + ), + }, + order: { + id: 'ASC', + }, + take: 200, + }); + + const messageChannelMessageAssociationsToDelete = + existingMessageChannelMessageAssociations.filter( + (existingMessageChannelMessageAssociation) => + isDefined( + existingMessageChannelMessageAssociation.messageExternalId, + ) && + !messageExternalIds.includes( + existingMessageChannelMessageAssociation.messageExternalId, + ), + ); + + this.logger.log( + `messageChannelId: ${messageChannel.id} Full sync: Message channel message associations to delete in batch ${batchIndex}: ${messageChannelMessageAssociationsToDelete.length}`, + ); + + fullSyncMessageChannelMessageAssociationsToDelete.push( + ...messageChannelMessageAssociationsToDelete, + ); + + if (existingMessageChannelMessageAssociations.length < 200) { + nextFirstBatchMessageChannelMessageAssociationId = undefined; + break; + } + + nextFirstBatchMessageChannelMessageAssociationId = + existingMessageChannelMessageAssociations[ + existingMessageChannelMessageAssociations.length - 1 + ].id; + + batchIndex++; + } + + return fullSyncMessageChannelMessageAssociationsToDelete; + } } 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 c0267e073c..874b7701c1 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 @@ -123,12 +123,14 @@ export class MessagingMessagesImportService { blocklist.map((blocklistItem) => blocklistItem.handle), ); - await this.saveMessagesAndEnqueueContactCreationService.saveMessagesAndEnqueueContactCreation( - messagesToSave, - messageChannel, - connectedAccountWithFreshTokens, - workspaceId, - ); + if (messagesToSave.length > 0) { + await this.saveMessagesAndEnqueueContactCreationService.saveMessagesAndEnqueueContactCreation( + messagesToSave, + messageChannel, + connectedAccountWithFreshTokens, + workspaceId, + ); + } if ( messageIdsToFetch.length < MESSAGING_GMAIL_USERS_MESSAGES_GET_BATCH_SIZE diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/utils/filter-emails.util.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/utils/filter-emails.util.ts index ffdbef9893..2801e43acd 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/utils/filter-emails.util.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/utils/filter-emails.util.ts @@ -10,17 +10,19 @@ export const filterEmails = ( messages: MessageWithParticipants[], blocklist: string[], ) => { + const messagesWithoutIcsAttachments = filterOutIcsAttachments(messages); + const messagesWithoutBlocklisted = filterOutBlocklistedMessages( [primaryHandle, ...handleAliases], - filterOutIcsAttachments(messages), + messagesWithoutIcsAttachments, blocklist, ); - if (isWorkEmail(primaryHandle)) { - return filterOutInternals(primaryHandle, messagesWithoutBlocklisted); - } + const messagesWithoutInternals = isWorkEmail(primaryHandle) + ? filterOutInternals(primaryHandle, messagesWithoutBlocklisted) + : messagesWithoutBlocklisted; - return messagesWithoutBlocklisted; + return messagesWithoutInternals; }; const filterOutBlocklistedMessages = (