diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/messaging-imap-driver.module.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/messaging-imap-driver.module.ts index b37e40f486f..bca9eb25fd5 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/messaging-imap-driver.module.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/messaging-imap-driver.module.ts @@ -15,7 +15,8 @@ import { ImapFindSentFolderService } from 'src/modules/messaging/message-import- import { ImapGetMessageListService } from 'src/modules/messaging/message-import-manager/drivers/imap/services/imap-get-message-list.service'; import { ImapGetMessagesService } from 'src/modules/messaging/message-import-manager/drivers/imap/services/imap-get-messages.service'; import { ImapHandleErrorService } from 'src/modules/messaging/message-import-manager/drivers/imap/services/imap-handle-error.service'; -import { ImapMessageLocatorService } from 'src/modules/messaging/message-import-manager/drivers/imap/services/imap-message-locator.service'; +import { ImapIncrementalSyncService } from 'src/modules/messaging/message-import-manager/drivers/imap/services/imap-incremental-sync.service'; +import { ImapMessageFetcherService } from 'src/modules/messaging/message-import-manager/drivers/imap/services/imap-message-fetcher.service'; import { ImapMessageProcessorService } from 'src/modules/messaging/message-import-manager/drivers/imap/services/imap-message-processor.service'; import { MessageParticipantManagerModule } from 'src/modules/messaging/message-participant-manager/message-participant-manager.module'; @@ -36,7 +37,8 @@ import { MessageParticipantManagerModule } from 'src/modules/messaging/message-p ImapGetMessagesService, ImapGetMessageListService, ImapHandleErrorService, - ImapMessageLocatorService, + ImapIncrementalSyncService, + ImapMessageFetcherService, ImapMessageProcessorService, ImapFindSentFolderService, ], diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/services/imap-fetch-by-batch.service.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/services/imap-fetch-by-batch.service.ts index 1c601f9ce50..1ca9008411f 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/services/imap-fetch-by-batch.service.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/services/imap-fetch-by-batch.service.ts @@ -2,7 +2,6 @@ import { Injectable, Logger } from '@nestjs/common'; import { type ConnectedAccountWorkspaceEntity } from 'src/modules/connected-account/standard-objects/connected-account.workspace-entity'; import { ImapClientProvider } from 'src/modules/messaging/message-import-manager/drivers/imap/providers/imap-client.provider'; -import { ImapMessageLocatorService } from 'src/modules/messaging/message-import-manager/drivers/imap/services/imap-message-locator.service'; import { ImapMessageProcessorService, type MessageFetchResult, @@ -14,7 +13,7 @@ type ConnectedAccount = Pick< >; type FetchAllResult = { - messageIdsByBatch: string[][]; + uidsByBatch: number[][]; batchResults: MessageFetchResult[][]; }; @@ -24,48 +23,42 @@ export class ImapFetchByBatchService { constructor( private readonly imapClientProvider: ImapClientProvider, - private readonly imapMessageLocatorService: ImapMessageLocatorService, private readonly imapMessageProcessorService: ImapMessageProcessorService, ) {} async fetchAllByBatches( - messageIds: string[], + uids: number[], connectedAccount: ConnectedAccount, + folder: string, ): Promise { const batchLimit = 20; const batchResults: MessageFetchResult[][] = []; - const messageIdsByBatch: string[][] = []; + const uidsByBatch: number[][] = []; this.logger.log( - `Starting optimized batch fetch for ${messageIds.length} messages`, + `Starting optimized batch fetch for ${uids.length} messages from folder ${folder}`, ); const client = await this.imapClientProvider.getClient(connectedAccount); try { - const messageLocations = - await this.imapMessageLocatorService.locateAllMessages( - messageIds, - client, - ); + for (let i = 0; i < uids.length; i += batchLimit) { + const batchUids = uids.slice(i, i + batchLimit); - for (let i = 0; i < messageIds.length; i += batchLimit) { - const batchMessageIds = messageIds.slice(i, i + batchLimit); - - messageIdsByBatch.push(batchMessageIds); + uidsByBatch.push(batchUids); try { const batchResult = - await this.imapMessageProcessorService.processMessagesByIds( - batchMessageIds, - messageLocations, + await this.imapMessageProcessorService.processMessagesByUidsInFolder( + batchUids, + folder, client, ); batchResults.push(batchResult); this.logger.log( - `Fetched batch ${Math.floor(i / batchLimit) + 1}/${Math.ceil(messageIds.length / batchLimit)} (${batchMessageIds.length} messages)`, + `Fetched batch ${Math.floor(i / batchLimit) + 1}/${Math.ceil(uids.length / batchLimit)} (${batchUids.length} messages)`, ); } catch (error) { this.logger.error( @@ -74,7 +67,8 @@ export class ImapFetchByBatchService { const errorResults = this.imapMessageProcessorService.createErrorResults( - batchMessageIds, + batchUids, + folder, error as Error, ); @@ -83,7 +77,7 @@ export class ImapFetchByBatchService { } return { - messageIdsByBatch, + uidsByBatch, batchResults, }; } finally { diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/services/imap-get-message-list.service.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/services/imap-get-message-list.service.ts index 9651a6a5cd1..f997d27e1ab 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/services/imap-get-message-list.service.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/services/imap-get-message-list.service.ts @@ -1,12 +1,19 @@ import { Injectable, Logger } from '@nestjs/common'; -import { FetchQueryObject, type ImapFlow } from 'imapflow'; +import { type ImapFlow } from 'imapflow'; import { type MessageFolderWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; 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 { ImapHandleErrorService } from 'src/modules/messaging/message-import-manager/drivers/imap/services/imap-handle-error.service'; +import { ImapIncrementalSyncService } from 'src/modules/messaging/message-import-manager/drivers/imap/services/imap-incremental-sync.service'; import { MessageFolderName } from 'src/modules/messaging/message-import-manager/drivers/imap/types/folders'; +import { createSyncCursor } from 'src/modules/messaging/message-import-manager/drivers/imap/utils/create-sync-cursor.util'; +import { extractMailboxState } from 'src/modules/messaging/message-import-manager/drivers/imap/utils/extract-mailbox-state.util'; +import { + ImapSyncCursor, + parseSyncCursor, +} from 'src/modules/messaging/message-import-manager/drivers/imap/utils/parse-sync-cursor.util'; import { type GetMessageListsArgs } from 'src/modules/messaging/message-import-manager/types/get-message-lists-args.type'; import { type GetMessageListsResponse, @@ -20,6 +27,7 @@ export class ImapGetMessageListService { constructor( private readonly imapClientProvider: ImapClientProvider, private readonly imapFindSentFolderService: ImapFindSentFolderService, + private readonly imapIncrementalSyncService: ImapIncrementalSyncService, private readonly imapHandleErrorService: ImapHandleErrorService, ) {} @@ -38,7 +46,9 @@ export class ImapGetMessageListService { const folderName = await this.getFolderName(client, folder.name); if (!folderName) { - this.logger.warn(`No folder name found for folder: ${folder.name}`); + this.logger.warn( + `No IMAP folder found for message folder: ${folder.name}`, + ); continue; } @@ -96,24 +106,26 @@ export class ImapGetMessageListService { folder: string, messageFolder: Pick, ): Promise { - const messages = await this.getMessagesFromFolder( - client, - folder, - messageFolder.syncCursor, + const { messages, messageExternalUidsToDelete, syncCursor } = + await this.getMessagesFromFolder( + client, + folder, + messageFolder.syncCursor, + ); + + messages.sort((a, b) => b.uid - a.uid); + + const messageExternalIds = messages.map( + (message) => `${folder}:${message.uid.toString()}`, ); - messages.sort((a, b) => parseInt(b.uid) - parseInt(a.uid)); - - const messageExternalIds = messages.map((message) => message.id); - - const nextSyncCursor = - messages.length > 0 ? messages[0].uid : messageFolder.syncCursor || ''; - return { messageExternalIds, - nextSyncCursor, + nextSyncCursor: JSON.stringify(syncCursor), previousSyncCursor: messageFolder.syncCursor || '', - messageExternalIdsToDelete: [], + messageExternalIdsToDelete: messageExternalUidsToDelete.map((uid) => + uid.toString(), + ), folderId: undefined, }; } @@ -146,63 +158,52 @@ export class ImapGetMessageListService { client: ImapFlow, folder: string, cursor?: string, - ): Promise<{ id: string; uid: string }[]> { + ): Promise<{ + messages: { uid: number }[]; + messageExternalUidsToDelete: number[]; + syncCursor: ImapSyncCursor; + }> { let lock; try { lock = await client.getMailboxLock(folder); + const mailbox = client.mailbox!; - const supportsUidPlus = client.capabilities.has('UIDPLUS'); - let sequence = '1:*'; - - if (cursor && supportsUidPlus) { - const cursorUid = parseInt(cursor); - - if (!isNaN(cursorUid)) { - sequence = `${cursorUid + 1}:*`; - } + if (typeof mailbox === 'boolean') { + throw new Error(`Invalid mailbox state for folder ${folder}`); } - const messages: { id: string; uid: string }[] = []; + const mailboxState = extractMailboxState(mailbox); + const previousCursor = parseSyncCursor(cursor); - this.logger.log( - `Fetching from folder: ${folder} with sequence: ${sequence} (UIDPLUS: ${supportsUidPlus})`, + const { messages, messageExternalUidsToDelete } = + await this.imapIncrementalSyncService.syncMessages( + client, + previousCursor, + mailboxState, + folder, + ); + + const newSyncCursor = createSyncCursor( + messages, + previousCursor, + mailboxState, ); - const fetchOptions: FetchQueryObject = { - envelope: true, + return { + messages, + messageExternalUidsToDelete, + syncCursor: newSyncCursor, }; - - if (supportsUidPlus) { - fetchOptions.uid = true; - } - - for await (const message of client.fetch(sequence, fetchOptions)) { - if (message.envelope?.messageId) { - messages.push({ - id: message.envelope.messageId, - uid: - supportsUidPlus && message.uid - ? message.uid.toString() - : message.seq?.toString() || '0', - }); - } - } - - this.logger.log(`Found ${messages.length} messages in folder: ${folder}`); - - return messages; - } catch (error) { + } catch (err) { this.logger.error( - `Error fetching from folder ${folder}: ${error.message}`, - error.stack, + `Error fetching from folder ${folder}: ${err.message}`, + err.stack, ); - return []; + throw err; } finally { - if (lock) { - lock.release(); - } + if (lock) lock.release(); } } } diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/services/imap-get-messages.service.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/services/imap-get-messages.service.ts index 29a2f062f10..c613ba6875e 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/services/imap-get-messages.service.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/services/imap-get-messages.service.ts @@ -8,6 +8,7 @@ import { computeMessageDirection } from 'src/modules/messaging/message-import-ma import { ImapFetchByBatchService } from 'src/modules/messaging/message-import-manager/drivers/imap/services/imap-fetch-by-batch.service'; import { type MessageFetchResult } from 'src/modules/messaging/message-import-manager/drivers/imap/services/imap-message-processor.service'; import { extractTextWithoutReplyQuotations } from 'src/modules/messaging/message-import-manager/drivers/imap/utils/extract-message-text.util'; +import { parseMessageId } from 'src/modules/messaging/message-import-manager/drivers/imap/utils/parse-message-id.util'; import { type EmailAddress } from 'src/modules/messaging/message-import-manager/types/email-address'; import { type MessageWithParticipants } from 'src/modules/messaging/message-import-manager/types/message'; import { formatAddressObjectAsParticipants } from 'src/modules/messaging/message-import-manager/utils/format-address-object-as-participants.util'; @@ -34,19 +35,55 @@ export class ImapGetMessagesService { return []; } - const { batchResults } = await this.fetchByBatchService.fetchAllByBatches( - messageIds, - connectedAccount, - ); + const folderToUidsMap = this.groupMessageIdsByFolder(messageIds); - this.logger.log(`IMAP fetch completed`); + const allMessages: MessageWithParticipants[] = []; - const messages = this.formatBatchResponsesAsMessages( - batchResults, - connectedAccount, - ); + for (const [folder, uids] of folderToUidsMap.entries()) { + if (!uids.length) { + continue; + } - return messages; + const { batchResults } = await this.fetchByBatchService.fetchAllByBatches( + uids, + connectedAccount, + folder, + ); + + this.logger.log(`IMAP fetch completed for folder: ${folder}`); + + const messages = this.formatBatchResponsesAsMessages( + batchResults, + connectedAccount, + folder, + ); + + allMessages.push(...messages); + } + + return allMessages; + } + + private groupMessageIdsByFolder(messageIds: string[]): Map { + const folderToUidsMap = new Map(); + + for (const messageId of messageIds) { + const parsedMessageId = parseMessageId(messageId); + + if (!parsedMessageId) { + this.logger.warn(`Invalid messageId format: ${messageId}`); + continue; + } + + const { folder, uid } = parsedMessageId; + + if (!folderToUidsMap.has(folder)) { + folderToUidsMap.set(folder, []); + } + folderToUidsMap.get(folder)!.push(uid); + } + + return folderToUidsMap; } public formatBatchResponsesAsMessages( @@ -55,9 +92,14 @@ export class ImapGetMessagesService { ConnectedAccountWorkspaceEntity, 'handle' | 'handleAliases' >, + folder: string, ): MessageWithParticipants[] { return batchResults.flatMap((batchResult) => { - return this.formatBatchResponseAsMessages(batchResult, connectedAccount); + return this.formatBatchResponseAsMessages( + batchResult, + connectedAccount, + folder, + ); }); } @@ -67,11 +109,12 @@ export class ImapGetMessagesService { ConnectedAccountWorkspaceEntity, 'handle' | 'handleAliases' >, + folder: string, ): MessageWithParticipants[] { const messages = batchResults.map((result) => { if (!result.parsed) { this.logger.warn( - `Message ${result.messageId} could not be parsed - likely not found in current folders`, + `Message UID ${result.uid} could not be parsed - likely not found in current folders`, ); return undefined; @@ -79,8 +122,9 @@ export class ImapGetMessagesService { return this.createMessageFromParsedMail( result.parsed, - result.messageId, + result.uid.toString(), connectedAccount, + folder, ); }); @@ -95,11 +139,12 @@ export class ImapGetMessagesService { private createMessageFromParsedMail( parsed: ParsedMail, - messageId: string, + uid: string, connectedAccount: Pick< ConnectedAccountWorkspaceEntity, 'handle' | 'handleAliases' >, + folder: string, ): MessageWithParticipants { const participants = this.extractAllParticipants(parsed); const attachments = this.extractAttachments(parsed); @@ -120,9 +165,9 @@ export class ImapGetMessagesService { const subject = sanitizeString(parsed.subject || ''); return { - externalId: messageId, - messageThreadExternalId: threadId || messageId, - headerMessageId: parsed.messageId || messageId, + externalId: `${folder}:${uid}`, + messageThreadExternalId: threadId || parsed.messageId || uid, + headerMessageId: parsed.messageId || uid, subject: subject, text: text, receivedAt: parsed.date || new Date(), diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/services/imap-incremental-sync.service.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/services/imap-incremental-sync.service.ts new file mode 100644 index 00000000000..4e21946c21a --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/services/imap-incremental-sync.service.ts @@ -0,0 +1,111 @@ +import { Injectable, Logger } from '@nestjs/common'; + +import { type ImapFlow } from 'imapflow'; + +import { canUseQresync } from 'src/modules/messaging/message-import-manager/drivers/imap/utils/can-use-qresync.util'; +import { type MailboxState } from 'src/modules/messaging/message-import-manager/drivers/imap/utils/extract-mailbox-state.util'; +import { type ImapSyncCursor } from 'src/modules/messaging/message-import-manager/drivers/imap/utils/parse-sync-cursor.util'; + +import { ImapMessageFetcherService } from './imap-message-fetcher.service'; + +type SyncStrategyResult = { + messages: { uid: number }[]; + messageExternalUidsToDelete: number[]; +}; + +@Injectable() +export class ImapIncrementalSyncService { + private readonly logger = new Logger(ImapIncrementalSyncService.name); + + constructor( + private readonly imapMessageFetcherService: ImapMessageFetcherService, + ) {} + + public async syncMessages( + client: ImapFlow, + previousCursor: ImapSyncCursor | null, + mailboxState: MailboxState, + folder: string, + ): Promise { + const messageExternalUidsToDelete = await this.checkUidValidityChange( + client, + previousCursor, + mailboxState, + folder, + ); + + const messages = await this.selectSyncStrategy( + client, + previousCursor, + mailboxState, + folder, + ); + + return { + messages, + messageExternalUidsToDelete, + }; + } + + private async checkUidValidityChange( + client: ImapFlow, + previousCursor: ImapSyncCursor | null, + mailboxState: MailboxState, + folder: string, + ): Promise { + const lastUidValidity = previousCursor?.uidValidity ?? 0; + const { uidValidity } = mailboxState; + + if (lastUidValidity !== 0 && lastUidValidity !== uidValidity) { + this.logger.log( + `UID validity changed from ${lastUidValidity} to ${uidValidity} in ${folder}. Full resync required.`, + ); + + return this.imapMessageFetcherService.getAllMessageUids(client); + } + + return []; + } + + private async selectSyncStrategy( + client: ImapFlow, + previousCursor: ImapSyncCursor | null, + mailboxState: MailboxState, + folder: string, + ): Promise<{ uid: number }[]> { + const lastSeenUid = previousCursor?.highestUid ?? 0; + const supportsQresync = client.capabilities.has('QRESYNC'); + const { maxUid } = mailboxState; + + if (canUseQresync(supportsQresync, previousCursor, mailboxState)) { + this.logger.log(`Using QRESYNC for folder ${folder}`); + const lastModSeq = BigInt(previousCursor!.modSeq!); + + try { + return await this.imapMessageFetcherService.getMessagesWithQresync( + client, + lastSeenUid, + lastModSeq, + ); + } catch (error) { + this.logger.warn( + `QRESYNC failed for folder ${folder}, falling back to UID search: ${error.message}`, + ); + + return this.imapMessageFetcherService.getMessagesWithUidSearch( + client, + lastSeenUid, + maxUid, + ); + } + } + + this.logger.log(`Using standard UID search for folder ${folder}`); + + return this.imapMessageFetcherService.getMessagesWithUidSearch( + client, + lastSeenUid, + maxUid, + ); + } +} diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/services/imap-message-fetcher.service.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/services/imap-message-fetcher.service.ts new file mode 100644 index 00000000000..b6eb0803dfe --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/services/imap-message-fetcher.service.ts @@ -0,0 +1,114 @@ +import { Injectable, Logger } from '@nestjs/common'; + +import { type ImapFlow } from 'imapflow'; + +@Injectable() +export class ImapMessageFetcherService { + private readonly logger = new Logger(ImapMessageFetcherService.name); + + public async getAllMessageUids(client: ImapFlow): Promise { + try { + const uids: number[] = []; + + for await (const msg of client.fetch('1:*', {}, { uid: true })) { + if (msg.uid) { + uids.push(msg.uid); + } + } + + return uids; + } catch (err) { + this.logger.error(`Error getting all message UIDs: ${err.message}`); + + return []; + } + } + + public async getMessagesWithUidSearch( + client: ImapFlow, + lastSeenUid: number, + maxUid: number, + ): Promise<{ uid: number }[]> { + try { + let allUids = await client.search({ all: true }, { uid: true }); + + if (!Array.isArray(allUids)) allUids = []; + + const wantedUids = allUids.filter((u) => u > lastSeenUid && u <= maxUid); + + if (wantedUids.length === 0) { + this.logger.log( + `No new messages. lastSeenUid=${lastSeenUid}, maxUid=${maxUid}`, + ); + + return []; + } + + const messages: { uid: number }[] = []; + + this.logger.log( + `Fetching ${wantedUids.length} messages, UIDs ${wantedUids[0]}..${ + wantedUids[wantedUids.length - 1] + }`, + ); + + for await (const msg of client.fetch( + wantedUids, + {}, + { + uid: true, + }, + )) { + if (msg.uid) { + messages.push({ uid: msg.uid }); + } + } + + return messages; + } catch (err) { + this.logger.error(`Error with UID search: ${err.message}`); + throw err; + } + } + + public async getMessagesWithQresync( + client: ImapFlow, + lastSeenUid: number, + lastModSeq: bigint, + ): Promise<{ uid: number }[]> { + try { + const vanished = await client.search( + { + modseq: lastModSeq + BigInt(1), + uid: `${lastSeenUid + 1}:*`, + }, + { uid: true }, + ); + + const messages: { uid: number }[] = []; + + if (vanished && Array.isArray(vanished) && vanished.length > 0) { + this.logger.log( + `QRESYNC: Fetching ${vanished.length} new/modified messages`, + ); + + for await (const msg of client.fetch( + vanished, + {}, + { + uid: true, + }, + )) { + if (msg.uid) { + messages.push({ uid: msg.uid }); + } + } + } + + return messages; + } catch (err) { + this.logger.error(`Error with QRESYNC: ${err.message}`); + throw err; + } + } +} diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/services/imap-message-locator.service.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/services/imap-message-locator.service.ts deleted file mode 100644 index c3049e21032..00000000000 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/services/imap-message-locator.service.ts +++ /dev/null @@ -1,134 +0,0 @@ -import { Injectable, Logger } from '@nestjs/common'; - -import { FetchMessageObject, type ImapFlow } from 'imapflow'; - -import { ImapFindSentFolderService } from 'src/modules/messaging/message-import-manager/drivers/imap/services/imap-find-sent-folder.service'; - -export type MessageLocation = { - messageId: string; - uid: number; - folder: string; -}; - -@Injectable() -export class ImapMessageLocatorService { - private readonly logger = new Logger(ImapMessageLocatorService.name); - private static readonly FETCH_BATCH_SIZE = 50; - - constructor( - private readonly imapFindSentFolderService: ImapFindSentFolderService, - ) {} - - async locateAllMessages( - messageIds: string[], - client: ImapFlow, - ): Promise> { - const locations = new Map(); - const folders = await this.getFoldersToSearch(client); - const messageIdSet = new Set(messageIds); - - for (const folder of folders) { - await this.searchFolderForMessages( - folder, - client, - messageIdSet, - locations, - ); - } - - return locations; - } - - private async searchFolderForMessages( - folder: string, - client: ImapFlow, - messageIdSet: Set, - locations: Map, - ): Promise { - let lock; - - try { - lock = await client.getMailboxLock(folder); - const uids = await client.search({ all: true }); - - await this.processBatchedMessages( - uids, - folder, - client, - messageIdSet, - locations, - ); - } catch (error) { - this.logger.warn(`Error searching folder ${folder}: ${error.message}`); - } finally { - lock?.release(); - } - } - - private async processBatchedMessages( - uids: number[], - folder: string, - client: ImapFlow, - messageIdSet: Set, - locations: Map, - ): Promise { - const batches = this.chunkArray( - uids, - ImapMessageLocatorService.FETCH_BATCH_SIZE, - ); - - for (const batchUids of batches) { - const fetchResults = client.fetch(batchUids.join(','), { - envelope: true, - }); - - for await (const message of fetchResults) { - this.processMessage(message, folder, messageIdSet, locations); - } - } - } - - private processMessage( - message: FetchMessageObject, - folder: string, - messageIdSet: Set, - locations: Map, - ): void { - const envelopeMessageId = message.envelope?.messageId; - - if (envelopeMessageId && messageIdSet.has(envelopeMessageId)) { - locations.set(envelopeMessageId, { - messageId: envelopeMessageId, - uid: message.uid, - folder, - }); - } - } - - private async getFoldersToSearch(client: ImapFlow): Promise { - const folders = ['INBOX']; - - try { - const sentFolder = - await this.imapFindSentFolderService.findSentFolder(client); - - if (sentFolder && sentFolder !== 'INBOX') { - folders.push(sentFolder); - } - } catch (error) { - this.logger.warn(`Failed to find sent folder: ${error.message}`); - } - - return folders; - } - - private chunkArray(array: T[], chunkSize: number): T[][] { - const chunks: T[][] = []; - - for (let i = 0; i < array.length; i += chunkSize) { - chunks.push(array.slice(i, i + chunkSize)); - } - - return chunks; - } -} diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/services/imap-message-processor.service.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/services/imap-message-processor.service.ts index aee01252fa0..aedfed9c2ed 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/services/imap-message-processor.service.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/services/imap-message-processor.service.ts @@ -4,10 +4,9 @@ import { type FetchMessageObject, type ImapFlow } from 'imapflow'; import { type ParsedMail, simpleParser } from 'mailparser'; import { ImapHandleErrorService } from 'src/modules/messaging/message-import-manager/drivers/imap/services/imap-handle-error.service'; -import { type MessageLocation } from 'src/modules/messaging/message-import-manager/drivers/imap/services/imap-message-locator.service'; export type MessageFetchResult = { - messageId: string; + uid: number; parsed: ParsedMail | null; processingTimeMs?: number; }; @@ -20,71 +19,20 @@ export class ImapMessageProcessorService { private readonly imapHandleErrorService: ImapHandleErrorService, ) {} - async processMessagesByIds( - messageIds: string[], - messageLocations: Map, + async processMessagesByUidsInFolder( + uids: number[], + folder: string, client: ImapFlow, ): Promise { - if (!messageIds.length) { + if (!uids.length) { return []; } - const results: MessageFetchResult[] = []; - - const messagesByFolder = new Map(); - const notFoundIds: string[] = []; - - for (const messageId of messageIds) { - const location = messageLocations.get(messageId); - - if (location) { - const locations = messagesByFolder.get(location.folder) || []; - - locations.push(location); - messagesByFolder.set(location.folder, locations); - } else { - notFoundIds.push(messageId); - } - } - - const fetchPromises = Array.from(messagesByFolder.entries()).map( - ([folder, locations]) => - this.fetchMessagesFromFolder(locations, client, folder), - ); - - const folderResults = await Promise.allSettled(fetchPromises); - - for (const result of folderResults) { - if (result.status === 'fulfilled') { - results.push(...result.value); - } else { - this.logger.error(`Folder batch fetch failed: ${result.reason}`); - } - } - - for (const messageId of notFoundIds) { - results.push({ - messageId, - parsed: null, - processingTimeMs: 0, - }); - } - - return results; - } - - private async fetchMessagesFromFolder( - messageLocations: MessageLocation[], - client: ImapFlow, - folder: string, - ): Promise { - if (!messageLocations.length) return []; - try { const lock = await client.getMailboxLock(folder); try { - return await this.fetchMessagesWithUids(messageLocations, client); + return await this.fetchMessages(uids, client); } finally { lock.release(); } @@ -93,28 +41,30 @@ export class ImapMessageProcessorService { `Failed to fetch messages from folder ${folder}: ${error.message}`, ); - return messageLocations.map((location) => - this.createErrorResult(location.messageId, error as Error, Date.now()), + return uids.map((uid) => + this.createErrorResult(uid, error as Error, Date.now()), ); } } - private async fetchMessagesWithUids( - messageLocations: MessageLocation[], + private async fetchMessages( + uids: number[], client: ImapFlow, ): Promise { const startTime = Date.now(); const results: MessageFetchResult[] = []; try { - const uids = messageLocations.map((loc) => loc.uid.toString()); const uidSet = uids.join(','); - const fetchResults = client.fetch(uidSet, { - uid: true, - source: true, - envelope: true, - }); + const fetchResults = client.fetch( + uidSet, + { + uid: true, + source: true, + }, + { uid: true }, + ); const messagesData = new Map(); @@ -122,12 +72,12 @@ export class ImapMessageProcessorService { messagesData.set(message.uid, message); } - for (const location of messageLocations) { - const messageData = messagesData.get(location.uid); + for (const uid of uids) { + const messageData = messagesData.get(uid); if (messageData) { const result = await this.processMessageData( - location.messageId, + uid, messageData, startTime, ); @@ -135,7 +85,7 @@ export class ImapMessageProcessorService { results.push(result); } else { results.push({ - messageId: location.messageId, + uid, parsed: null, processingTimeMs: Date.now() - startTime, }); @@ -144,8 +94,8 @@ export class ImapMessageProcessorService { } catch (error) { this.logger.error(`Batch fetch failed: ${error.message}`); - return messageLocations.map((location) => - this.createErrorResult(location.messageId, error as Error, startTime), + return uids.map((uid) => + this.createErrorResult(uid, error as Error, startTime), ); } @@ -153,7 +103,7 @@ export class ImapMessageProcessorService { } private async processMessageData( - messageId: string, + uid: number, messageData: FetchMessageObject, startTime: number, ): Promise { @@ -161,77 +111,73 @@ export class ImapMessageProcessorService { const rawContent = messageData.source?.toString() || ''; if (!rawContent) { - this.logger.debug(`No source content for message ${messageId}`); + this.logger.debug(`No source content for message UID ${uid}`); return { - messageId, + uid, parsed: null, processingTimeMs: Date.now() - startTime, }; } - const parsed = await this.parseMessage(rawContent, messageId); + const parsed = await this.parseMessage(rawContent, uid); const processingTime = Date.now() - startTime; - this.logger.debug( - `Processed message ${messageId} in ${processingTime}ms`, - ); + this.logger.debug(`Processed message UID ${uid} in ${processingTime}ms`); return { - messageId, + uid, parsed, processingTimeMs: processingTime, }; } catch (error) { - return this.createErrorResult(messageId, error as Error, startTime); + return this.createErrorResult(uid, error as Error, startTime); } } private async parseMessage( rawContent: string, - messageId: string, + uid: number, ): Promise { try { return await simpleParser(rawContent); } catch (error) { - this.logger.error( - `Failed to parse message ${messageId}: ${error.message}`, - ); + this.logger.error(`Failed to parse message UID ${uid}: ${error.message}`); throw error; } } createErrorResult( - messageId: string, + uid: number, error: Error, startTime: number, ): MessageFetchResult { const processingTime = Date.now() - startTime; - this.logger.error(`Failed to fetch message ${messageId}: ${error.message}`); - - this.imapHandleErrorService.handleImapMessagesImportError(error, messageId); + this.logger.error(`Failed to fetch message UID ${uid}: ${error.message}`); return { - messageId, + uid, parsed: null, processingTimeMs: processingTime, }; } - createErrorResults(messageIds: string[], error: Error): MessageFetchResult[] { - return messageIds.map((messageId) => { - this.logger.error( - `Failed to fetch message ${messageId}: ${error.message}`, - ); + createErrorResults( + uids: number[], + folder: string, + error: Error, + ): MessageFetchResult[] { + return uids.map((uid) => { + this.logger.error(`Failed to fetch message UID ${uid}: ${error.message}`); this.imapHandleErrorService.handleImapMessagesImportError( error, - messageId, + `${folder}:${uid}`, ); return { - messageId, + uid, parsed: null, }; }); diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/utils/__tests__/parse-message-id.util.spec.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/utils/__tests__/parse-message-id.util.spec.ts new file mode 100644 index 00000000000..768feb4afb4 --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/utils/__tests__/parse-message-id.util.spec.ts @@ -0,0 +1,52 @@ +import { parseMessageId } from 'src/modules/messaging/message-import-manager/drivers/imap/utils/parse-message-id.util'; + +describe('parseMessageId', () => { + it('parses a valid message id with folder and uid', () => { + const input = 'INBOX:12345'; + const result = parseMessageId(input); + + expect(result).toEqual({ folder: 'INBOX', uid: 12345 }); + }); + + it('parses a folder name with colons', () => { + const input = 'Sent:Items:6789'; + const result = parseMessageId(input); + + expect(result).toEqual({ folder: 'Sent:Items', uid: 6789 }); + }); + + it('returns null for non-numeric uid', () => { + const input = 'INBOX:abc'; + const result = parseMessageId(input); + + expect(result).toBeNull(); + }); + + it('returns null for empty string', () => { + const input = ''; + const result = parseMessageId(input); + + expect(result).toBeNull(); + }); + + it('returns null for missing folder', () => { + const input = ':123'; + const result = parseMessageId(input); + + expect(result).toBeNull(); + }); + + it('returns null for missing uid', () => { + const input = 'INBOX:'; + const result = parseMessageId(input); + + expect(result).toBeNull(); + }); + + it('parses folder with spaces', () => { + const input = 'My Folder:42'; + const result = parseMessageId(input); + + expect(result).toEqual({ folder: 'My Folder', uid: 42 }); + }); +}); diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/utils/can-use-qresync.util.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/utils/can-use-qresync.util.ts new file mode 100644 index 00000000000..bd9f43c8a9c --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/utils/can-use-qresync.util.ts @@ -0,0 +1,21 @@ +import { type MailboxState } from './extract-mailbox-state.util'; +import { type ImapSyncCursor } from './parse-sync-cursor.util'; + +export const canUseQresync = ( + supportsQresync: boolean, + previousCursor: ImapSyncCursor | null, + mailboxState: MailboxState, +): boolean => { + const lastModSeq = previousCursor?.modSeq + ? BigInt(previousCursor.modSeq) + : undefined; + const lastUidValidity = previousCursor?.uidValidity ?? 0; + const { uidValidity, highestModSeq } = mailboxState; + + return ( + supportsQresync && + lastModSeq !== undefined && + highestModSeq !== undefined && + lastUidValidity === uidValidity + ); +}; diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/utils/create-sync-cursor.util.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/utils/create-sync-cursor.util.ts new file mode 100644 index 00000000000..1711052a416 --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/utils/create-sync-cursor.util.ts @@ -0,0 +1,22 @@ +import { type MailboxState } from './extract-mailbox-state.util'; +import { type ImapSyncCursor } from './parse-sync-cursor.util'; + +export const createSyncCursor = ( + messages: { uid: number }[], + previousCursor: ImapSyncCursor | null, + mailboxState: MailboxState, +): ImapSyncCursor => { + const { uidValidity, highestModSeq } = mailboxState; + const lastSeenUid = previousCursor?.highestUid ?? 0; + + const highestUid = + messages.length > 0 + ? Math.max(...messages.map((message) => message.uid)) + : lastSeenUid; + + return { + highestUid, + uidValidity, + ...(highestModSeq ? { modSeq: highestModSeq.toString() } : {}), + }; +}; diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/utils/extract-mailbox-state.util.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/utils/extract-mailbox-state.util.ts new file mode 100644 index 00000000000..f009478a6a5 --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/utils/extract-mailbox-state.util.ts @@ -0,0 +1,25 @@ +import { type ImapFlow } from 'imapflow'; + +export type MailboxState = { + uidValidity: number; + uidNext: number; + maxUid: number; + highestModSeq?: bigint; +}; + +export const extractMailboxState = ( + mailbox: NonNullable, +): MailboxState => { + if (typeof mailbox === 'boolean') { + throw new Error('Invalid mailbox state'); + } + + const uidNext = Number(mailbox.uidNext ?? 1); + + return { + uidValidity: Number(mailbox.uidValidity ?? 0), + uidNext, + maxUid: Math.max(0, uidNext - 1), + highestModSeq: mailbox.highestModseq, + }; +}; diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/utils/parse-message-id.util.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/utils/parse-message-id.util.ts new file mode 100644 index 00000000000..afcd04661d4 --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/utils/parse-message-id.util.ts @@ -0,0 +1,22 @@ +export type ParsedMessageId = { + folder: string; + uid: number; +}; + +export function parseMessageId(messageId: string): ParsedMessageId | null { + const regex = /^(.+):(\d+)$/; + const match = regex.exec(messageId); + + if (!match) { + return null; + } + + const [, folder, uidStr] = match; + const uid = Number(uidStr); + + if (!Number.isInteger(uid)) { + return null; + } + + return { folder, uid }; +} diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/utils/parse-sync-cursor.util.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/utils/parse-sync-cursor.util.ts new file mode 100644 index 00000000000..386158fbf9f --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/utils/parse-sync-cursor.util.ts @@ -0,0 +1,19 @@ +import { isDefined } from 'twenty-shared/utils'; + +export type ImapSyncCursor = { + highestUid: number; + uidValidity: number; + modSeq?: string; +}; + +export const parseSyncCursor = (cursor?: string): ImapSyncCursor | null => { + if (!isDefined(cursor)) { + return null; + } + + try { + return JSON.parse(cursor) as ImapSyncCursor; + } catch { + return null; + } +};