From 734b0bd2d243dc0619605e329c6eb97c6c2cd1b9 Mon Sep 17 00:00:00 2001 From: neo773 <62795688+neo773@users.noreply.github.com> Date: Mon, 25 Aug 2025 20:12:07 +0530 Subject: [PATCH] IMAP Refactor (#14053) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Removed `IMAPMessageLocator` service which relied on `Message-ID` header lookups. We now fetch messages directly by UID, which should eliminate most of the issues we were hitting. This had a challenge, `messageIdsToFetch` were flattened and UIDs are only unique within their own folder. To preserve context without modifying the core message module, I landed on using composite keys for IMAP message IDs. Each UID is now prefixed with its folder name, and we parse it back when needed to recover the folder. This feels like a reasonable trade-off since it resolves a lot of the failures we saw with `Message-ID` lookups. And Syncs are faster now since there's no double scanning. - Added `QRESYNC` support for servers that provide it, giving us faster incremental sync. - Refactored code to better split concerns. ### Compatibility Test | Server | Status | |---------------------|------------| | Gmail | ✅ Tested | | Dovecot | ✅ Tested | | Stalwart | ✅ Tested | | Titan Email | ✅ Tested | | FastMail | ✅ Tested | image --- .../imap/messaging-imap-driver.module.ts | 6 +- .../services/imap-fetch-by-batch.service.ts | 36 ++--- .../services/imap-get-message-list.service.ts | 115 +++++++------- .../services/imap-get-messages.service.ts | 79 +++++++--- .../services/imap-incremental-sync.service.ts | 111 ++++++++++++++ .../services/imap-message-fetcher.service.ts | 114 ++++++++++++++ .../services/imap-message-locator.service.ts | 134 ---------------- .../imap-message-processor.service.ts | 144 ++++++------------ .../__tests__/parse-message-id.util.spec.ts | 52 +++++++ .../imap/utils/can-use-qresync.util.ts | 21 +++ .../imap/utils/create-sync-cursor.util.ts | 22 +++ .../imap/utils/extract-mailbox-state.util.ts | 25 +++ .../imap/utils/parse-message-id.util.ts | 22 +++ .../imap/utils/parse-sync-cursor.util.ts | 19 +++ 14 files changed, 570 insertions(+), 330 deletions(-) create mode 100644 packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/services/imap-incremental-sync.service.ts create mode 100644 packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/services/imap-message-fetcher.service.ts delete mode 100644 packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/services/imap-message-locator.service.ts create mode 100644 packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/utils/__tests__/parse-message-id.util.spec.ts create mode 100644 packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/utils/can-use-qresync.util.ts create mode 100644 packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/utils/create-sync-cursor.util.ts create mode 100644 packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/utils/extract-mailbox-state.util.ts create mode 100644 packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/utils/parse-message-id.util.ts create mode 100644 packages/twenty-server/src/modules/messaging/message-import-manager/drivers/imap/utils/parse-sync-cursor.util.ts 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; + } +};