From abe774da150a98d624c66d09355939e47f58b301 Mon Sep 17 00:00:00 2001 From: neo773 <62795688+neo773@users.noreply.github.com> Date: Fri, 19 Dec 2025 21:34:48 +0530 Subject: [PATCH] Message folders optimization (#16479) - Batch Gmail API calls using `googleapis-batcher` for folder processing - Add concurrency limit for Microsoft Graph folder processing - Skip IMAP folder sync when no new messages (checks UIDVALIDITY/MODSEQ) - Refactored `syncMessageFolders` to return folder state directly, avoiding extra DB round-trips - Refactored `processPendingFolderActions` to reuse state instead of querying DB again - Add unique index on message folders entity --------- Co-authored-by: Charles Bochet --- packages/twenty-server/package.json | 1 + .../create-message-channel.service.ts | 41 +- .../message-folder.workspace-entity.ts | 5 + .../services/gmail-get-all-folders.service.ts | 6 +- .../services/imap-get-all-folders.service.ts | 10 +- .../microsoft-get-all-folders.service.ts | 6 +- .../message-folder-driver.interface.ts | 20 +- .../sync-message-folders.service.spec.ts | 543 ++++++++++++++++++ .../services/sync-message-folders.service.ts | 330 ++++------- .../compute-folder-ids-to-delete.util.spec.ts | 78 +++ .../compute-folders-to-create.util.spec.ts | 69 +++ .../compute-folders-to-update.util.spec.ts | 127 ++++ .../compute-updated-folders.util.spec.ts | 90 +++ .../compute-folder-ids-to-delete.util.ts | 22 + .../utils/compute-folders-to-create.util.ts | 39 ++ .../utils/compute-folders-to-update.util.ts | 58 ++ .../utils/compute-updated-folders.util.ts | 38 ++ .../gmail-get-message-list.service.ts | 56 +- .../services/imap-get-message-list.service.ts | 93 ++- ...osoft-get-message-list.service.dev.spec.ts | 16 +- .../microsoft-get-message-list.service.ts | 36 +- ...ssaging-message-list-fetch.service.spec.ts | 21 +- .../messaging-get-message-list.service.ts | 8 +- .../messaging-message-list-fetch.service.ts | 55 +- .../types/get-message-lists-args.type.ts | 8 +- yarn.lock | 61 +- 26 files changed, 1479 insertions(+), 358 deletions(-) create mode 100644 packages/twenty-server/src/modules/messaging/message-folder-manager/services/sync-message-folders.service.spec.ts create mode 100644 packages/twenty-server/src/modules/messaging/message-folder-manager/utils/__tests__/compute-folder-ids-to-delete.util.spec.ts create mode 100644 packages/twenty-server/src/modules/messaging/message-folder-manager/utils/__tests__/compute-folders-to-create.util.spec.ts create mode 100644 packages/twenty-server/src/modules/messaging/message-folder-manager/utils/__tests__/compute-folders-to-update.util.spec.ts create mode 100644 packages/twenty-server/src/modules/messaging/message-folder-manager/utils/__tests__/compute-updated-folders.util.spec.ts create mode 100644 packages/twenty-server/src/modules/messaging/message-folder-manager/utils/compute-folder-ids-to-delete.util.ts create mode 100644 packages/twenty-server/src/modules/messaging/message-folder-manager/utils/compute-folders-to-create.util.ts create mode 100644 packages/twenty-server/src/modules/messaging/message-folder-manager/utils/compute-folders-to-update.util.ts create mode 100644 packages/twenty-server/src/modules/messaging/message-folder-manager/utils/compute-updated-folders.util.ts diff --git a/packages/twenty-server/package.json b/packages/twenty-server/package.json index 4a5d59e6de..ae628a2edf 100644 --- a/packages/twenty-server/package.json +++ b/packages/twenty-server/package.json @@ -37,6 +37,7 @@ "@graphql-tools/schema": "10.0.4", "@graphql-tools/utils": "9.2.1", "@graphql-yoga/nestjs": "patch:@graphql-yoga/nestjs@2.1.0#./patches/@graphql-yoga+nestjs+2.1.0.patch", + "@jrmdayn/googleapis-batcher": "^0.10.1", "@lingui/conf": "5.1.2", "@lingui/core": "^5.1.2", "@lingui/format-po": "5.1.2", diff --git a/packages/twenty-server/src/engine/core-modules/auth/services/create-message-channel.service.ts b/packages/twenty-server/src/engine/core-modules/auth/services/create-message-channel.service.ts index fba5a1fb97..f49cdbe4cc 100644 --- a/packages/twenty-server/src/engine/core-modules/auth/services/create-message-channel.service.ts +++ b/packages/twenty-server/src/engine/core-modules/auth/services/create-message-channel.service.ts @@ -6,7 +6,6 @@ import { v4 } from 'uuid'; import { type WorkspaceEntityManager } from 'src/engine/twenty-orm/entity-manager/workspace-entity-manager'; import { GlobalWorkspaceOrmManager } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager'; import { buildSystemAuthContext } from 'src/engine/twenty-orm/utils/build-system-auth-context.util'; -import { type ConnectedAccountWorkspaceEntity } from 'src/modules/connected-account/standard-objects/connected-account.workspace-entity'; import { MessageChannelPendingGroupEmailsAction, MessageChannelSyncStage, @@ -48,25 +47,17 @@ export class CreateMessageChannelService { return this.globalWorkspaceOrmManager.executeInWorkspaceContext( authContext, async () => { - const connectedAccountRepository = - await this.globalWorkspaceOrmManager.getRepository( - workspaceId, - 'connectedAccount', - ); - - const connectedAccount = await connectedAccountRepository.findOne({ - where: { id: connectedAccountId }, - }); - const messageChannelRepository = await this.globalWorkspaceOrmManager.getRepository( workspaceId, 'messageChannel', ); - const newMessageChannel = await messageChannelRepository.save( + const newMessageChannelId = v4(); + + await messageChannelRepository.insert( { - id: v4(), + id: newMessageChannelId, connectedAccountId, type: MessageChannelType.EMAIL, handle, @@ -77,22 +68,24 @@ export class CreateMessageChannelService { pendingGroupEmailsAction: MessageChannelPendingGroupEmailsAction.NONE, }, - {}, manager, ); - if (isDefined(connectedAccount)) { - await this.syncMessageFoldersService.syncMessageFolders({ - workspaceId, - messageChannel: { - ...newMessageChannel, - connectedAccount, - }, - manager, - }); + const createdMessageChannel = await messageChannelRepository.findOne({ + where: { id: newMessageChannelId }, + relations: ['connectedAccount', 'messageFolders'], + }); + + if (!isDefined(createdMessageChannel)) { + throw new Error('Message channel not found'); } - return newMessageChannel.id; + await this.syncMessageFoldersService.syncMessageFolders({ + messageChannel: createdMessageChannel, + workspaceId, + }); + + return newMessageChannelId; }, ); } diff --git a/packages/twenty-server/src/modules/messaging/common/standard-objects/message-folder.workspace-entity.ts b/packages/twenty-server/src/modules/messaging/common/standard-objects/message-folder.workspace-entity.ts index 3cc848eed1..b031428e38 100644 --- a/packages/twenty-server/src/modules/messaging/common/standard-objects/message-folder.workspace-entity.ts +++ b/packages/twenty-server/src/modules/messaging/common/standard-objects/message-folder.workspace-entity.ts @@ -10,6 +10,7 @@ import { RelationType } from 'src/engine/metadata-modules/field-metadata/interfa import { BaseWorkspaceEntity } from 'src/engine/twenty-orm/base.workspace-entity'; import { WorkspaceEntity } from 'src/engine/twenty-orm/decorators/workspace-entity.decorator'; import { WorkspaceField } from 'src/engine/twenty-orm/decorators/workspace-field.decorator'; +import { WorkspaceIndex } from 'src/engine/twenty-orm/decorators/workspace-index.decorator'; import { WorkspaceIsNotAuditLogged } from 'src/engine/twenty-orm/decorators/workspace-is-not-audit-logged.decorator'; import { WorkspaceIsNullable } from 'src/engine/twenty-orm/decorators/workspace-is-nullable.decorator'; import { WorkspaceIsSystem } from 'src/engine/twenty-orm/decorators/workspace-is-system.decorator'; @@ -39,6 +40,10 @@ registerEnumType(MessageFolderPendingSyncAction, { }) @WorkspaceIsNotAuditLogged() @WorkspaceIsSystem() +@WorkspaceIndex(['messageChannelId', 'externalId'], { + isUnique: true, + indexWhereClause: '"deletedAt" IS NULL', +}) export class MessageFolderWorkspaceEntity extends BaseWorkspaceEntity { @WorkspaceField({ standardId: MESSAGE_FOLDER_STANDARD_FIELD_IDS.name, diff --git a/packages/twenty-server/src/modules/messaging/message-folder-manager/drivers/gmail/services/gmail-get-all-folders.service.ts b/packages/twenty-server/src/modules/messaging/message-folder-manager/drivers/gmail/services/gmail-get-all-folders.service.ts index 9d7d1745f5..3319e0349d 100644 --- a/packages/twenty-server/src/modules/messaging/message-folder-manager/drivers/gmail/services/gmail-get-all-folders.service.ts +++ b/packages/twenty-server/src/modules/messaging/message-folder-manager/drivers/gmail/services/gmail-get-all-folders.service.ts @@ -3,7 +3,7 @@ import { Injectable, Logger } from '@nestjs/common'; import { google } from 'googleapis'; import { - MessageFolder, + DiscoveredMessageFolder, MessageFolderDriver, } from 'src/modules/messaging/message-folder-manager/interfaces/message-folder-driver.interface'; @@ -34,7 +34,7 @@ export class GmailGetAllFoldersService implements MessageFolderDriver { MessageChannelWorkspaceEntity, 'messageFolderImportPolicy' >, - ): Promise { + ): Promise { try { const oAuth2Client = await this.oAuth2ClientManagerService.getGoogleOAuth2Client( @@ -60,7 +60,7 @@ export class GmailGetAllFoldersService implements MessageFolderDriver { const labels = response.data.labels || []; - const folders: MessageFolder[] = []; + const folders: DiscoveredMessageFolder[] = []; const labelNameToIdMap = new Map(); diff --git a/packages/twenty-server/src/modules/messaging/message-folder-manager/drivers/imap/services/imap-get-all-folders.service.ts b/packages/twenty-server/src/modules/messaging/message-folder-manager/drivers/imap/services/imap-get-all-folders.service.ts index d781b38e69..b8eb1ecd83 100644 --- a/packages/twenty-server/src/modules/messaging/message-folder-manager/drivers/imap/services/imap-get-all-folders.service.ts +++ b/packages/twenty-server/src/modules/messaging/message-folder-manager/drivers/imap/services/imap-get-all-folders.service.ts @@ -4,7 +4,7 @@ import { ImapFlow, type ListResponse } from 'imapflow'; import { isDefined } from 'twenty-shared/utils'; import { - MessageFolder, + DiscoveredMessageFolder, MessageFolderDriver, } from 'src/modules/messaging/message-folder-manager/interfaces/message-folder-driver.interface'; @@ -34,7 +34,7 @@ export class ImapGetAllFoldersService implements MessageFolderDriver { MessageChannelWorkspaceEntity, 'messageFolderImportPolicy' >, - ): Promise { + ): Promise { try { const client = await this.imapClientProvider.getClient(connectedAccount); @@ -66,8 +66,8 @@ export class ImapGetAllFoldersService implements MessageFolderDriver { MessageChannelWorkspaceEntity, 'messageFolderImportPolicy' >, - ): Promise { - const folders: MessageFolder[] = []; + ): Promise { + const folders: DiscoveredMessageFolder[] = []; const pathToExternalIdMap = new Map(); const sentFolder = await this.imapFindSentFolderService.findSentFolder(client); @@ -135,7 +135,7 @@ export class ImapGetAllFoldersService implements MessageFolderDriver { private isValidMailbox( mailbox: ListResponse, - existingFolders: MessageFolder[], + existingFolders: DiscoveredMessageFolder[], ): boolean { if (mailbox.flags?.has('\\Noselect')) { return false; diff --git a/packages/twenty-server/src/modules/messaging/message-folder-manager/drivers/microsoft/services/microsoft-get-all-folders.service.ts b/packages/twenty-server/src/modules/messaging/message-folder-manager/drivers/microsoft/services/microsoft-get-all-folders.service.ts index 9d3da2d4e0..d155b7618d 100644 --- a/packages/twenty-server/src/modules/messaging/message-folder-manager/drivers/microsoft/services/microsoft-get-all-folders.service.ts +++ b/packages/twenty-server/src/modules/messaging/message-folder-manager/drivers/microsoft/services/microsoft-get-all-folders.service.ts @@ -3,7 +3,7 @@ import { Injectable, Logger } from '@nestjs/common'; import { isDefined } from 'twenty-shared/utils'; import { - MessageFolder, + DiscoveredMessageFolder, MessageFolderDriver, } from 'src/modules/messaging/message-folder-manager/interfaces/message-folder-driver.interface'; @@ -44,7 +44,7 @@ export class MicrosoftGetAllFoldersService implements MessageFolderDriver { MessageChannelWorkspaceEntity, 'messageFolderImportPolicy' >, - ): Promise { + ): Promise { try { const microsoftClient = await this.oAuth2ClientManagerService.getMicrosoftOAuth2Client( @@ -67,7 +67,7 @@ export class MicrosoftGetAllFoldersService implements MessageFolderDriver { const folders = (response.value as MicrosoftGraphFolder[]) || []; const rootFolderId = this.getRootFolderId(folders); - const folderInfos: MessageFolder[] = []; + const folderInfos: DiscoveredMessageFolder[] = []; for (const folder of folders) { if (!folder.displayName) { diff --git a/packages/twenty-server/src/modules/messaging/message-folder-manager/interfaces/message-folder-driver.interface.ts b/packages/twenty-server/src/modules/messaging/message-folder-manager/interfaces/message-folder-driver.interface.ts index 783a77fa07..706d601426 100644 --- a/packages/twenty-server/src/modules/messaging/message-folder-manager/interfaces/message-folder-driver.interface.ts +++ b/packages/twenty-server/src/modules/messaging/message-folder-manager/interfaces/message-folder-driver.interface.ts @@ -2,12 +2,24 @@ import { type ConnectedAccountWorkspaceEntity } from 'src/modules/connected-acco import { type MessageChannelWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-channel.workspace-entity'; import { type MessageFolderWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; -export type MessageFolder = Pick< +export type DiscoveredMessageFolder = Pick< MessageFolderWorkspaceEntity, 'name' | 'isSynced' | 'isSentFolder' | 'externalId' | 'parentFolderId' >; -export interface MessageFolderDriver { +export type MessageFolder = Pick< + MessageFolderWorkspaceEntity, + | 'name' + | 'isSynced' + | 'isSentFolder' + | 'externalId' + | 'parentFolderId' + | 'id' + | 'syncCursor' + | 'pendingSyncAction' +>; + +export type MessageFolderDriver = { getAllMessageFolders( connectedAccount: Pick< ConnectedAccountWorkspaceEntity, @@ -22,5 +34,5 @@ export interface MessageFolderDriver { MessageChannelWorkspaceEntity, 'messageFolderImportPolicy' >, - ): Promise; -} + ): Promise; +}; diff --git a/packages/twenty-server/src/modules/messaging/message-folder-manager/services/sync-message-folders.service.spec.ts b/packages/twenty-server/src/modules/messaging/message-folder-manager/services/sync-message-folders.service.spec.ts new file mode 100644 index 0000000000..1d3ae5cf13 --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-folder-manager/services/sync-message-folders.service.spec.ts @@ -0,0 +1,543 @@ +import { Test, type TestingModule } from '@nestjs/testing'; + +import { ConnectedAccountProvider } from 'twenty-shared/types'; + +import { type DiscoveredMessageFolder } from 'src/modules/messaging/message-folder-manager/interfaces/message-folder-driver.interface'; + +import { GlobalWorkspaceOrmManager } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager'; +import { + MessageChannelContactAutoCreationPolicy, + MessageChannelType, + MessageChannelVisibility, + MessageFolderImportPolicy, +} from 'src/modules/messaging/common/standard-objects/message-channel.workspace-entity'; +import { + MessageFolderPendingSyncAction, + type MessageFolderWorkspaceEntity, +} from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; +import { GmailGetAllFoldersService } from 'src/modules/messaging/message-folder-manager/drivers/gmail/services/gmail-get-all-folders.service'; +import { ImapGetAllFoldersService } from 'src/modules/messaging/message-folder-manager/drivers/imap/services/imap-get-all-folders.service'; +import { MicrosoftGetAllFoldersService } from 'src/modules/messaging/message-folder-manager/drivers/microsoft/services/microsoft-get-all-folders.service'; +import { SyncMessageFoldersService } from 'src/modules/messaging/message-folder-manager/services/sync-message-folders.service'; + +type SyncedMessageFolder = Pick< + MessageFolderWorkspaceEntity, + | 'id' + | 'name' + | 'isSynced' + | 'isSentFolder' + | 'externalId' + | 'syncCursor' + | 'parentFolderId' + | 'pendingSyncAction' +>; + +const createMockMessageChannel = ( + overrides: { + provider?: ConnectedAccountProvider; + messageFolders?: SyncedMessageFolder[]; + } = {}, +) => ({ + id: 'channel-123', + handle: 'test@gmail.com', + type: MessageChannelType.EMAIL, + messageFolderImportPolicy: MessageFolderImportPolicy.ALL_FOLDERS, + connectedAccount: { + id: 'account-456', + handle: 'test@gmail.com', + provider: overrides.provider ?? ConnectedAccountProvider.GOOGLE, + accessToken: 'mock-access-token', + refreshToken: 'mock-refresh-token', + connectionParameters: {}, + }, + messageFolders: overrides.messageFolders ?? [], + visibility: MessageChannelVisibility.SHARE_EVERYTHING, + isContactAutoCreationEnabled: false, + contactAutoCreationPolicy: MessageChannelContactAutoCreationPolicy.NONE, + excludeNonProfessionalEmails: false, + excludeGroupEmails: false, +}); + +const createMockDiscoveredFolder = ( + overrides: Partial = {}, +): DiscoveredMessageFolder => ({ + externalId: `external-${Math.random().toString(36).substring(7)}`, + name: 'Test Folder', + isSynced: false, + isSentFolder: false, + parentFolderId: null, + ...overrides, +}); + +const createMockExistingFolder = ( + overrides: Partial = {}, +): SyncedMessageFolder => ({ + id: `folder-${Math.random().toString(36).substring(7)}`, + externalId: `external-${Math.random().toString(36).substring(7)}`, + name: 'Existing Folder', + isSynced: true, + isSentFolder: false, + syncCursor: null, + parentFolderId: null, + pendingSyncAction: MessageFolderPendingSyncAction.NONE, + ...overrides, +}); + +describe('SyncMessageFoldersService', () => { + let service: SyncMessageFoldersService; + let gmailGetAllFoldersService: jest.Mocked; + + let mockRepository: { + delete: jest.Mock; + update: jest.Mock; + updateMany: jest.Mock; + save: jest.Mock; + }; + let mockTransactionManager: object; + + beforeEach(async () => { + mockRepository = { + delete: jest.fn(), + update: jest.fn(), + updateMany: jest.fn(), + save: jest.fn().mockImplementation((folders) => + folders.map((folder: Partial) => ({ + ...folder, + id: `new-folder-${Math.random().toString(36).substring(7)}`, + isSynced: false, + syncCursor: null, + })), + ), + }; + + mockTransactionManager = {}; + + const mockDataSource = { + transaction: jest + .fn() + .mockImplementation((callback) => callback(mockTransactionManager)), + }; + + const module: TestingModule = await Test.createTestingModule({ + providers: [ + SyncMessageFoldersService, + { + provide: GlobalWorkspaceOrmManager, + useValue: { + executeInWorkspaceContext: jest + .fn() + .mockImplementation((_, callback) => callback()), + getRepository: jest.fn().mockResolvedValue(mockRepository), + getDataSourceForWorkspace: jest + .fn() + .mockResolvedValue(mockDataSource), + getGlobalWorkspaceDataSource: jest + .fn() + .mockResolvedValue(mockDataSource), + }, + }, + { + provide: GmailGetAllFoldersService, + useValue: { + getAllMessageFolders: jest.fn(), + }, + }, + { + provide: MicrosoftGetAllFoldersService, + useValue: { + getAllMessageFolders: jest.fn(), + }, + }, + { + provide: ImapGetAllFoldersService, + useValue: { + getAllMessageFolders: jest.fn(), + }, + }, + ], + }).compile(); + + service = module.get(SyncMessageFoldersService); + gmailGetAllFoldersService = module.get(GmailGetAllFoldersService); + }); + + describe('syncMessageFolders', () => { + const workspaceId = 'workspace-789'; + + describe('folder creation scenarios', () => { + it('should create new folders when none exist locally', async () => { + const discoveredFolders = [ + createMockDiscoveredFolder({ + externalId: 'inbox-ext', + name: 'INBOX', + }), + createMockDiscoveredFolder({ + externalId: 'sent-ext', + name: 'Sent', + isSentFolder: true, + }), + ]; + const messageChannel = createMockMessageChannel({ + messageFolders: [], + }); + + gmailGetAllFoldersService.getAllMessageFolders.mockResolvedValue( + discoveredFolders, + ); + + const result = await service.syncMessageFolders({ + messageChannel, + workspaceId, + }); + + expect(mockRepository.save).toHaveBeenCalledWith( + expect.arrayContaining([ + expect.objectContaining({ + name: 'INBOX', + externalId: 'inbox-ext', + messageChannelId: 'channel-123', + isSentFolder: false, + }), + expect.objectContaining({ + name: 'Sent', + externalId: 'sent-ext', + messageChannelId: 'channel-123', + isSentFolder: true, + }), + ]), + {}, + mockTransactionManager, + ); + expect(result).toHaveLength(2); + }); + + it('should handle nested folder creation with parent references', async () => { + const discoveredFolders = [ + createMockDiscoveredFolder({ + externalId: 'parent-ext', + name: 'Work', + parentFolderId: null, + }), + createMockDiscoveredFolder({ + externalId: 'child-ext', + name: 'Projects', + parentFolderId: 'parent-folder-id', + }), + ]; + const messageChannel = createMockMessageChannel({ + messageFolders: [], + }); + + gmailGetAllFoldersService.getAllMessageFolders.mockResolvedValue( + discoveredFolders, + ); + + await service.syncMessageFolders({ + messageChannel, + workspaceId, + }); + + expect(mockRepository.save).toHaveBeenCalledWith( + expect.arrayContaining([ + expect.objectContaining({ + name: 'Projects', + parentFolderId: 'parent-folder-id', + }), + ]), + {}, + mockTransactionManager, + ); + }); + }); + + describe('folder update scenarios', () => { + it('should update folder when name changes', async () => { + const existingFolder = createMockExistingFolder({ + id: 'folder-1', + externalId: 'inbox-ext', + name: 'INBOX', + }); + const discoveredFolders = [ + createMockDiscoveredFolder({ + externalId: 'inbox-ext', + name: 'Primary Inbox', + }), + ]; + const messageChannel = createMockMessageChannel({ + messageFolders: [existingFolder], + }); + + gmailGetAllFoldersService.getAllMessageFolders.mockResolvedValue( + discoveredFolders, + ); + + const result = await service.syncMessageFolders({ + messageChannel, + workspaceId, + }); + + expect(mockRepository.updateMany).toHaveBeenCalledWith( + expect.arrayContaining([ + expect.objectContaining({ + criteria: 'folder-1', + partialEntity: expect.objectContaining({ name: 'Primary Inbox' }), + }), + ]), + mockTransactionManager, + ); + expect(result).toContainEqual( + expect.objectContaining({ + id: 'folder-1', + name: 'Primary Inbox', + }), + ); + }); + + it('should update folder when parent folder changes', async () => { + const existingFolder = createMockExistingFolder({ + id: 'folder-1', + externalId: 'child-ext', + name: 'Projects', + parentFolderId: 'old-parent-id', + }); + const discoveredFolders = [ + createMockDiscoveredFolder({ + externalId: 'child-ext', + name: 'Projects', + parentFolderId: 'new-parent-id', + }), + ]; + const messageChannel = createMockMessageChannel({ + messageFolders: [existingFolder], + }); + + gmailGetAllFoldersService.getAllMessageFolders.mockResolvedValue( + discoveredFolders, + ); + + await service.syncMessageFolders({ + messageChannel, + workspaceId, + }); + + expect(mockRepository.updateMany).toHaveBeenCalledWith( + expect.arrayContaining([ + expect.objectContaining({ + criteria: 'folder-1', + partialEntity: expect.objectContaining({ + parentFolderId: 'new-parent-id', + }), + }), + ]), + mockTransactionManager, + ); + }); + + it('should not update folder when nothing has changed', async () => { + const existingFolder = createMockExistingFolder({ + id: 'folder-1', + externalId: 'inbox-ext', + name: 'INBOX', + isSentFolder: false, + parentFolderId: null, + }); + const discoveredFolders = [ + createMockDiscoveredFolder({ + externalId: 'inbox-ext', + name: 'INBOX', + isSentFolder: false, + parentFolderId: null, + }), + ]; + const messageChannel = createMockMessageChannel({ + messageFolders: [existingFolder], + }); + + gmailGetAllFoldersService.getAllMessageFolders.mockResolvedValue( + discoveredFolders, + ); + + await service.syncMessageFolders({ + messageChannel, + workspaceId, + }); + + expect(mockRepository.updateMany).not.toHaveBeenCalled(); + }); + }); + + describe('folder deletion scenarios', () => { + it('should delete folders that no longer exist remotely', async () => { + const existingFolders = [ + createMockExistingFolder({ + id: 'folder-1', + externalId: 'inbox-ext', + name: 'INBOX', + }), + createMockExistingFolder({ + id: 'folder-2', + externalId: 'deleted-ext', + name: 'Old Folder', + }), + ]; + const discoveredFolders = [ + createMockDiscoveredFolder({ + externalId: 'inbox-ext', + name: 'INBOX', + }), + ]; + const messageChannel = createMockMessageChannel({ + messageFolders: existingFolders, + }); + + gmailGetAllFoldersService.getAllMessageFolders.mockResolvedValue( + discoveredFolders, + ); + + const result = await service.syncMessageFolders({ + messageChannel: messageChannel, + workspaceId, + }); + + expect(mockRepository.updateMany).toHaveBeenCalledWith( + expect.arrayContaining([ + expect.objectContaining({ + criteria: 'folder-2', + partialEntity: expect.objectContaining({ + pendingSyncAction: 'FOLDER_DELETION', + }), + }), + ]), + mockTransactionManager, + ); + expect(result).toContainEqual( + expect.objectContaining({ + id: 'folder-2', + pendingSyncAction: MessageFolderPendingSyncAction.FOLDER_DELETION, + }), + ); + expect(result).toHaveLength(2); + }); + }); + + describe('complex sync scenarios', () => { + it('should handle simultaneous create, update, and delete operations', async () => { + const existingFolders = [ + createMockExistingFolder({ + id: 'folder-to-update', + externalId: 'update-ext', + name: 'Old Name', + }), + createMockExistingFolder({ + id: 'folder-to-delete', + externalId: 'delete-ext', + name: 'To Delete', + }), + createMockExistingFolder({ + id: 'folder-unchanged', + externalId: 'unchanged-ext', + name: 'Unchanged', + }), + ]; + const discoveredFolders = [ + createMockDiscoveredFolder({ + externalId: 'update-ext', + name: 'New Name', + }), + createMockDiscoveredFolder({ + externalId: 'unchanged-ext', + name: 'Unchanged', + }), + createMockDiscoveredFolder({ + externalId: 'new-ext', + name: 'New Folder', + }), + ]; + const messageChannel = createMockMessageChannel({ + messageFolders: existingFolders, + }); + + gmailGetAllFoldersService.getAllMessageFolders.mockResolvedValue( + discoveredFolders, + ); + + const result = await service.syncMessageFolders({ + messageChannel, + workspaceId, + }); + + expect(mockRepository.updateMany).toHaveBeenCalledWith( + expect.arrayContaining([ + expect.objectContaining({ + criteria: 'folder-to-delete', + partialEntity: expect.objectContaining({ + pendingSyncAction: 'FOLDER_DELETION', + }), + }), + ]), + mockTransactionManager, + ); + expect(mockRepository.updateMany).toHaveBeenCalledWith( + expect.arrayContaining([ + expect.objectContaining({ + criteria: 'folder-to-update', + partialEntity: expect.objectContaining({ name: 'New Name' }), + }), + ]), + mockTransactionManager, + ); + expect(mockRepository.save).toHaveBeenCalledWith( + expect.arrayContaining([ + expect.objectContaining({ externalId: 'new-ext' }), + ]), + {}, + mockTransactionManager, + ); + expect(result).toHaveLength(4); + expect(result).toContainEqual( + expect.objectContaining({ + id: 'folder-to-delete', + pendingSyncAction: MessageFolderPendingSyncAction.FOLDER_DELETION, + }), + ); + }); + + it('should preserve syncCursor and isSynced for unchanged folders', async () => { + const existingFolder = createMockExistingFolder({ + id: 'folder-1', + externalId: 'inbox-ext', + name: 'INBOX', + isSynced: true, + syncCursor: 'cursor-abc123', + }); + const discoveredFolders = [ + createMockDiscoveredFolder({ + externalId: 'inbox-ext', + name: 'INBOX', + }), + ]; + const messageChannel = createMockMessageChannel({ + messageFolders: [existingFolder], + }); + + gmailGetAllFoldersService.getAllMessageFolders.mockResolvedValue( + discoveredFolders, + ); + + const result = await service.syncMessageFolders({ + messageChannel, + workspaceId, + }); + + expect(result).toContainEqual( + expect.objectContaining({ + id: 'folder-1', + isSynced: true, + syncCursor: 'cursor-abc123', + pendingSyncAction: MessageFolderPendingSyncAction.NONE, + }), + ); + }); + }); + }); +}); diff --git a/packages/twenty-server/src/modules/messaging/message-folder-manager/services/sync-message-folders.service.ts b/packages/twenty-server/src/modules/messaging/message-folder-manager/services/sync-message-folders.service.ts index 11e632a467..9eac1f436a 100644 --- a/packages/twenty-server/src/modules/messaging/message-folder-manager/services/sync-message-folders.service.ts +++ b/packages/twenty-server/src/modules/messaging/message-folder-manager/services/sync-message-folders.service.ts @@ -1,54 +1,28 @@ import { Injectable } from '@nestjs/common'; -import { isNonEmptyString } from '@sniptt/guards'; -import deepEqual from 'deep-equal'; import { ConnectedAccountProvider } from 'twenty-shared/types'; -import { isDefined } from 'twenty-shared/utils'; -import { v4 } from 'uuid'; -import { MessageFolder } from 'src/modules/messaging/message-folder-manager/interfaces/message-folder-driver.interface'; +import { + DiscoveredMessageFolder, + MessageFolder, +} from 'src/modules/messaging/message-folder-manager/interfaces/message-folder-driver.interface'; -import { type WorkspaceEntityManager } from 'src/engine/twenty-orm/entity-manager/workspace-entity-manager'; +import { WorkspaceEntityManager } from 'src/engine/twenty-orm/entity-manager/workspace-entity-manager'; import { GlobalWorkspaceOrmManager } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager'; -import { WorkspaceRepository } from 'src/engine/twenty-orm/repository/workspace.repository'; import { buildSystemAuthContext } from 'src/engine/twenty-orm/utils/build-system-auth-context.util'; +import { ConnectedAccountWorkspaceEntity } from 'src/modules/connected-account/standard-objects/connected-account.workspace-entity'; import { type MessageChannelWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-channel.workspace-entity'; import { MessageFolderPendingSyncAction, - type MessageFolderWorkspaceEntity, + MessageFolderWorkspaceEntity, } from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; import { GmailGetAllFoldersService } from 'src/modules/messaging/message-folder-manager/drivers/gmail/services/gmail-get-all-folders.service'; import { ImapGetAllFoldersService } from 'src/modules/messaging/message-folder-manager/drivers/imap/services/imap-get-all-folders.service'; import { MicrosoftGetAllFoldersService } from 'src/modules/messaging/message-folder-manager/drivers/microsoft/services/microsoft-get-all-folders.service'; -import { MessageFolderName } from 'src/modules/messaging/message-import-manager/drivers/microsoft/types/folders'; - -type SyncMessageFoldersInput = { - workspaceId: string; - messageChannel: Pick< - MessageChannelWorkspaceEntity, - 'messageFolderImportPolicy' | 'connectedAccount' | 'id' - >; - manager: WorkspaceEntityManager; -}; - -type MessageFolderToInsert = Pick< - MessageFolderWorkspaceEntity, - | 'id' - | 'messageChannelId' - | 'name' - | 'syncCursor' - | 'isSynced' - | 'isSentFolder' - | 'externalId' - | 'parentFolderId' ->; - -type MessageFolderToUpdate = Partial< - Pick< - MessageFolderWorkspaceEntity, - 'name' | 'externalId' | 'isSentFolder' | 'parentFolderId' - > ->; +import { computeFolderIdsToDelete } from 'src/modules/messaging/message-folder-manager/utils/compute-folder-ids-to-delete.util'; +import { computeFoldersToCreate } from 'src/modules/messaging/message-folder-manager/utils/compute-folders-to-create.util'; +import { computeFoldersToUpdate } from 'src/modules/messaging/message-folder-manager/utils/compute-folders-to-update.util'; +import { computeUpdatedFolders } from 'src/modules/messaging/message-folder-manager/utils/compute-updated-folders.util'; @Injectable() export class SyncMessageFoldersService { @@ -59,160 +33,71 @@ export class SyncMessageFoldersService { private readonly imapGetAllFoldersService: ImapGetAllFoldersService, ) {} - async syncMessageFolders(input: SyncMessageFoldersInput): Promise { - const { workspaceId, messageChannel, manager } = input; - - const authContext = buildSystemAuthContext(workspaceId); - - await this.globalWorkspaceOrmManager.executeInWorkspaceContext( - authContext, - async () => { - const folders = await this.discoverAllFolders( - messageChannel.connectedAccount, - messageChannel, - ); - - await this.upsertDiscoveredFolders({ - workspaceId, - messageChannelId: messageChannel.id, - folders, - manager, - }); - }, - ); - } - - private async upsertDiscoveredFolders({ + async syncMessageFolders({ + messageChannel, workspaceId, - messageChannelId, - folders, - manager, }: { + messageChannel: Pick< + MessageChannelWorkspaceEntity, + 'id' | 'messageFolderImportPolicy' + > & { + connectedAccount: Pick< + ConnectedAccountWorkspaceEntity, + | 'provider' + | 'accessToken' + | 'refreshToken' + | 'id' + | 'handle' + | 'connectionParameters' + >; + messageFolders: MessageFolder[]; + }; workspaceId: string; - messageChannelId: string; - folders: MessageFolder[]; - manager: WorkspaceEntityManager; - }): Promise { - const messageFolderRepository = - await this.globalWorkspaceOrmManager.getRepository( - workspaceId, - 'messageFolder', - ); - - const existingFolderMap = await this.buildExistingFolderMap({ - messageChannelId, - messageFolderRepository, - }); - - const inserts: MessageFolderToInsert[] = []; - const updates: [string, MessageFolderToUpdate][] = []; - const foldersToMarkForDeletion: string[] = []; - - const discoveredExternalIds = new Set( - folders - .filter((folder) => folder.externalId) - .map((folder) => folder.externalId!), + }): Promise { + const discoveredFolders = await this.discoverAllFolders( + messageChannel.connectedAccount, + messageChannel, ); - for (const existingFolder of existingFolderMap.values()) { - if ( - existingFolder.externalId && - !discoveredExternalIds.has(existingFolder.externalId) - ) { - foldersToMarkForDeletion.push(existingFolder.id); - } - } + const { messageFolders: existingFolders, id: messageChannelId } = + messageChannel; - for (const folder of folders) { - const existingFolder = this.findExistingFolderInMap( - existingFolderMap, - folder, - ); - - if (existingFolder) { - const folderSyncData = { - name: folder.name, - externalId: folder.externalId, - isSentFolder: folder.isSentFolder, - parentFolderId: isNonEmptyString(folder.parentFolderId) - ? folder.parentFolderId - : null, - }; - - const existingFolderData = { - name: existingFolder.name, - externalId: existingFolder.externalId, - isSentFolder: existingFolder.isSentFolder, - parentFolderId: isNonEmptyString(existingFolder.parentFolderId) - ? existingFolder.parentFolderId - : null, - }; - - if (!deepEqual(folderSyncData, existingFolderData)) { - updates.push([existingFolder.id, folderSyncData]); - } - continue; - } - - inserts.push({ - id: v4(), - messageChannelId, - name: folder.name, - syncCursor: '', - isSynced: folder.isSynced, - isSentFolder: folder.isSentFolder, - externalId: folder.externalId, - parentFolderId: folder.parentFolderId, - }); - } - - if (inserts.length > 0) { - await messageFolderRepository.insert(inserts, manager); - } - - if (updates.length > 0) { - await messageFolderRepository.updateMany( - updates.map(([id, data]) => ({ - criteria: id, - partialEntity: data, - })), - manager, - ); - } - - if (foldersToMarkForDeletion.length > 0) { - await messageFolderRepository.updateMany( - foldersToMarkForDeletion.map((id) => ({ - criteria: id, - partialEntity: { - pendingSyncAction: MessageFolderPendingSyncAction.FOLDER_DELETION, - }, - })), - manager, - ); - } + return this.syncFolderChanges( + discoveredFolders, + existingFolders, + messageChannelId, + workspaceId, + ); } async discoverAllFolders( - connectedAccount: MessageChannelWorkspaceEntity['connectedAccount'], + connectedAccount: Pick< + ConnectedAccountWorkspaceEntity, + | 'accessToken' + | 'refreshToken' + | 'id' + | 'handle' + | 'provider' + | 'connectionParameters' + >, messageChannel: Pick< MessageChannelWorkspaceEntity, 'messageFolderImportPolicy' >, - ): Promise { + ): Promise { switch (connectedAccount.provider) { case ConnectedAccountProvider.GOOGLE: - return await this.gmailGetAllFoldersService.getAllMessageFolders( + return this.gmailGetAllFoldersService.getAllMessageFolders( connectedAccount, messageChannel, ); case ConnectedAccountProvider.MICROSOFT: - return await this.microsoftGetAllFoldersService.getAllMessageFolders( + return this.microsoftGetAllFoldersService.getAllMessageFolders( connectedAccount, messageChannel, ); case ConnectedAccountProvider.IMAP_SMTP_CALDAV: - return await this.imapGetAllFoldersService.getAllMessageFolders( + return this.imapGetAllFoldersService.getAllMessageFolders( connectedAccount, messageChannel, ); @@ -223,55 +108,86 @@ export class SyncMessageFoldersService { } } - private async buildExistingFolderMap({ - messageChannelId, - messageFolderRepository, - }: { - messageChannelId: string; - messageFolderRepository: WorkspaceRepository; - }): Promise> { - const existingFolders = await messageFolderRepository.find({ - where: { messageChannelId }, + private async syncFolderChanges( + discoveredFolders: DiscoveredMessageFolder[], + existingFolders: MessageFolder[], + messageChannelId: string, + workspaceId: string, + ): Promise { + const foldersToCreate = computeFoldersToCreate({ + discoveredFolders, + existingFolders, + messageChannelId, }); - const existingFolderMap = new Map(); + const foldersToUpdate = computeFoldersToUpdate({ + discoveredFolders, + existingFolders, + }); - for (const existingFolder of existingFolders) { - if (isDefined(existingFolder.externalId)) { - existingFolderMap.set(existingFolder.externalId, existingFolder); - } - existingFolderMap.set(existingFolder.name ?? '', existingFolder); - } + const folderIdsToDelete = computeFolderIdsToDelete({ + discoveredFolders, + existingFolders, + }); - return existingFolderMap; - } + const authContext = buildSystemAuthContext(workspaceId); - private findExistingFolderInMap( - existingFolderMap: Map, - folder: MessageFolder, - ): MessageFolderWorkspaceEntity | undefined { - if (isDefined(folder.externalId)) { - const existingFolder = existingFolderMap.get(folder.externalId); + return this.globalWorkspaceOrmManager.executeInWorkspaceContext( + authContext, + async () => { + const messageFolderRepository = + await this.globalWorkspaceOrmManager.getRepository( + workspaceId, + 'messageFolder', + ); - if (existingFolder) { - return existingFolder; - } - } + const workspaceDataSource = + await this.globalWorkspaceOrmManager.getGlobalWorkspaceDataSource(); - const legacyFolderName = this.getLegacyFolderName(folder); + return workspaceDataSource.transaction( + async (transactionManager: WorkspaceEntityManager) => { + if (folderIdsToDelete.length > 0) { + await messageFolderRepository.updateMany( + folderIdsToDelete.map((id) => ({ + criteria: id, + partialEntity: { + pendingSyncAction: + MessageFolderPendingSyncAction.FOLDER_DELETION, + }, + })), + transactionManager, + ); + } - return existingFolderMap.get(legacyFolderName); - } + if (foldersToUpdate.size > 0) { + await messageFolderRepository.updateMany( + Array.from(foldersToUpdate.entries()).map(([id, data]) => ({ + criteria: id, + partialEntity: data, + })), + transactionManager, + ); + } - private getLegacyFolderName(folder: MessageFolder): string { - if (folder.isSynced && !folder.isSentFolder) { - return MessageFolderName.INBOX; - } + const createdFolders = + foldersToCreate.length > 0 + ? await messageFolderRepository.save( + foldersToCreate, + {}, + transactionManager, + ) + : []; - if (folder.isSynced && folder.isSentFolder) { - return MessageFolderName.SENT_ITEMS; - } + const updatedExistingFolders = computeUpdatedFolders({ + existingFolders, + foldersToUpdate, + folderIdsToDelete, + }); - return folder.name ?? ''; + return [...updatedExistingFolders, ...createdFolders]; + }, + ); + }, + ); } } diff --git a/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/__tests__/compute-folder-ids-to-delete.util.spec.ts b/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/__tests__/compute-folder-ids-to-delete.util.spec.ts new file mode 100644 index 0000000000..be882130fa --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/__tests__/compute-folder-ids-to-delete.util.spec.ts @@ -0,0 +1,78 @@ +import { MessageFolderPendingSyncAction } from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; +import { computeFolderIdsToDelete } from 'src/modules/messaging/message-folder-manager/utils/compute-folder-ids-to-delete.util'; + +describe('computeFolderIdsToDelete', () => { + it('should mark folders deleted from provider for deletion', () => { + const discoveredFolders = [ + { + name: 'Inbox', + externalId: 'INBOX', + isSynced: true, + isSentFolder: false, + parentFolderId: null, + }, + ]; + + const existingFolders = [ + { + id: 'inbox-id', + name: 'Inbox', + externalId: 'INBOX', + isSynced: true, + isSentFolder: false, + parentFolderId: null, + syncCursor: 'cursor', + pendingSyncAction: MessageFolderPendingSyncAction.NONE, + }, + { + id: 'deleted-label-id', + name: 'Old Label', + externalId: 'deleted-label', + isSynced: true, + isSentFolder: false, + parentFolderId: null, + syncCursor: 'cursor', + pendingSyncAction: MessageFolderPendingSyncAction.NONE, + }, + ]; + + const result = computeFolderIdsToDelete({ + discoveredFolders, + existingFolders, + }); + + expect(result).toEqual(['deleted-label-id']); + }); + + it('should return empty when all local folders still exist in provider', () => { + const discoveredFolders = [ + { + name: 'Inbox', + externalId: 'INBOX', + isSynced: true, + isSentFolder: false, + parentFolderId: null, + }, + ]; + + const existingFolders = [ + { + id: 'inbox-id', + name: 'Inbox', + externalId: 'INBOX', + isSynced: true, + isSentFolder: false, + parentFolderId: null, + syncCursor: 'cursor', + pendingSyncAction: MessageFolderPendingSyncAction.NONE, + }, + ]; + + const result = computeFolderIdsToDelete({ + discoveredFolders, + existingFolders, + }); + + expect(result).toEqual([]); + }); +}); diff --git a/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/__tests__/compute-folders-to-create.util.spec.ts b/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/__tests__/compute-folders-to-create.util.spec.ts new file mode 100644 index 0000000000..9bc9f33f53 --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/__tests__/compute-folders-to-create.util.spec.ts @@ -0,0 +1,69 @@ +import { MessageFolderPendingSyncAction } from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; +import { computeFoldersToCreate } from 'src/modules/messaging/message-folder-manager/utils/compute-folders-to-create.util'; + +describe('computeFoldersToCreate', () => { + const messageChannelId = 'channel-123'; + + it('should create folders that exist in provider but not locally', () => { + const discoveredFolders = [ + { + name: 'Inbox', + externalId: 'INBOX', + isSynced: true, + isSentFolder: false, + parentFolderId: null, + }, + { + name: 'Sent', + externalId: 'SENT', + isSynced: true, + isSentFolder: true, + parentFolderId: null, + }, + ]; + + const existingFolders = [ + { + id: 'existing-id', + name: 'Inbox', + externalId: 'INBOX', + isSynced: true, + isSentFolder: false, + parentFolderId: null, + syncCursor: 'cursor', + pendingSyncAction: MessageFolderPendingSyncAction.NONE, + }, + ]; + + const result = computeFoldersToCreate({ + discoveredFolders, + existingFolders, + messageChannelId, + }); + + expect(result).toHaveLength(1); + expect(result[0].externalId).toBe('SENT'); + expect(result[0].messageChannelId).toBe(messageChannelId); + expect(result[0].syncCursor).toBeNull(); + }); + + it('should normalize empty parentFolderId to null', () => { + const discoveredFolders = [ + { + name: 'Subfolder', + externalId: 'sub-1', + isSynced: true, + isSentFolder: false, + parentFolderId: '', + }, + ]; + + const result = computeFoldersToCreate({ + discoveredFolders, + existingFolders: [], + messageChannelId, + }); + + expect(result[0].parentFolderId).toBeNull(); + }); +}); diff --git a/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/__tests__/compute-folders-to-update.util.spec.ts b/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/__tests__/compute-folders-to-update.util.spec.ts new file mode 100644 index 0000000000..304b8c2ec7 --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/__tests__/compute-folders-to-update.util.spec.ts @@ -0,0 +1,127 @@ +import { MessageFolderPendingSyncAction } from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; +import { computeFoldersToUpdate } from 'src/modules/messaging/message-folder-manager/utils/compute-folders-to-update.util'; + +describe('computeFoldersToUpdate', () => { + it('should detect folder rename from provider', () => { + const discoveredFolders = [ + { + name: 'Work Emails', + externalId: 'label-123', + isSynced: true, + isSentFolder: false, + parentFolderId: null, + }, + ]; + + const existingFolders = [ + { + id: 'folder-id', + name: 'Old Label Name', + externalId: 'label-123', + isSynced: true, + isSentFolder: false, + parentFolderId: null, + syncCursor: 'cursor', + pendingSyncAction: MessageFolderPendingSyncAction.NONE, + }, + ]; + + const result = computeFoldersToUpdate({ + discoveredFolders, + existingFolders, + }); + + expect(result.get('folder-id')?.name).toBe('Work Emails'); + }); + + it('should detect folder moved to different parent', () => { + const discoveredFolders = [ + { + name: 'Subfolder', + externalId: 'sub-1', + isSynced: true, + isSentFolder: false, + parentFolderId: 'new-parent-id', + }, + ]; + + const existingFolders = [ + { + id: 'folder-id', + name: 'Subfolder', + externalId: 'sub-1', + isSynced: true, + isSentFolder: false, + parentFolderId: 'old-parent-id', + syncCursor: 'cursor', + pendingSyncAction: MessageFolderPendingSyncAction.NONE, + }, + ]; + + const result = computeFoldersToUpdate({ + discoveredFolders, + existingFolders, + }); + + expect(result.get('folder-id')?.parentFolderId).toBe('new-parent-id'); + }); + + it('should not flag unchanged folders for update', () => { + const folder = { + name: 'Inbox', + externalId: 'INBOX', + isSynced: true, + isSentFolder: false, + parentFolderId: null, + }; + + const discoveredFolders = [folder]; + const existingFolders = [ + { + ...folder, + id: 'folder-id', + syncCursor: 'cursor', + pendingSyncAction: MessageFolderPendingSyncAction.NONE, + }, + ]; + + const result = computeFoldersToUpdate({ + discoveredFolders, + existingFolders, + }); + + expect(result.size).toBe(0); + }); + + it('should treat empty string parentFolderId same as null', () => { + const discoveredFolders = [ + { + name: 'Folder', + externalId: 'ext-1', + isSynced: true, + isSentFolder: false, + parentFolderId: '', + }, + ]; + + const existingFolders = [ + { + id: 'folder-id', + name: 'Folder', + externalId: 'ext-1', + isSynced: true, + isSentFolder: false, + parentFolderId: null, + syncCursor: 'cursor', + pendingSyncAction: MessageFolderPendingSyncAction.NONE, + }, + ]; + + const result = computeFoldersToUpdate({ + discoveredFolders, + existingFolders, + }); + + expect(result.size).toBe(0); + }); +}); diff --git a/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/__tests__/compute-updated-folders.util.spec.ts b/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/__tests__/compute-updated-folders.util.spec.ts new file mode 100644 index 0000000000..d7c8813367 --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/__tests__/compute-updated-folders.util.spec.ts @@ -0,0 +1,90 @@ +import { MessageFolderPendingSyncAction } from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; +import { computeUpdatedFolders } from 'src/modules/messaging/message-folder-manager/utils/compute-updated-folders.util'; + +describe('computeUpdatedFolders', () => { + it('should apply updates and set correct pendingSyncAction', () => { + const existingFolders = [ + { + id: 'folder-1', + name: 'Old Name', + externalId: 'ext-1', + isSynced: true, + isSentFolder: false, + parentFolderId: null, + syncCursor: 'cursor', + pendingSyncAction: MessageFolderPendingSyncAction.NONE, + }, + { + id: 'folder-2', + name: 'To Delete', + externalId: 'ext-2', + isSynced: true, + isSentFolder: false, + parentFolderId: null, + syncCursor: 'cursor', + pendingSyncAction: MessageFolderPendingSyncAction.NONE, + }, + { + id: 'folder-3', + name: 'Unchanged', + externalId: 'ext-3', + isSynced: true, + isSentFolder: false, + parentFolderId: null, + syncCursor: 'cursor', + pendingSyncAction: MessageFolderPendingSyncAction.NONE, + }, + ]; + + const foldersToUpdate = new Map([['folder-1', { name: 'New Name' }]]); + const folderIdsToDelete = ['folder-2']; + + const result = computeUpdatedFolders({ + existingFolders, + foldersToUpdate, + folderIdsToDelete, + }); + + expect(result[0].name).toBe('New Name'); + expect(result[0].pendingSyncAction).toBe( + MessageFolderPendingSyncAction.NONE, + ); + + expect(result[1].pendingSyncAction).toBe( + MessageFolderPendingSyncAction.FOLDER_DELETION, + ); + + expect(result[2].name).toBe('Unchanged'); + expect(result[2].pendingSyncAction).toBe( + MessageFolderPendingSyncAction.NONE, + ); + }); + + it('should preserve existing properties when applying partial updates', () => { + const existingFolders = [ + { + id: 'folder-1', + name: 'Original', + externalId: 'ext-1', + isSynced: true, + isSentFolder: true, + parentFolderId: 'parent-123', + syncCursor: 'cursor-abc', + pendingSyncAction: MessageFolderPendingSyncAction.NONE, + }, + ]; + + const foldersToUpdate = new Map([['folder-1', { name: 'Updated' }]]); + + const result = computeUpdatedFolders({ + existingFolders, + foldersToUpdate, + folderIdsToDelete: [], + }); + + expect(result[0].name).toBe('Updated'); + expect(result[0].isSentFolder).toBe(true); + expect(result[0].parentFolderId).toBe('parent-123'); + expect(result[0].syncCursor).toBe('cursor-abc'); + }); +}); diff --git a/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/compute-folder-ids-to-delete.util.ts b/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/compute-folder-ids-to-delete.util.ts new file mode 100644 index 0000000000..fa7eda77bb --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/compute-folder-ids-to-delete.util.ts @@ -0,0 +1,22 @@ +import { + type DiscoveredMessageFolder, + type MessageFolder, +} from 'src/modules/messaging/message-folder-manager/interfaces/message-folder-driver.interface'; + +export const computeFolderIdsToDelete = ({ + discoveredFolders, + existingFolders, +}: { + discoveredFolders: DiscoveredMessageFolder[]; + existingFolders: MessageFolder[]; +}): string[] => { + const discoveredExternalIds = new Set( + discoveredFolders.map((discoveredFolder) => discoveredFolder.externalId), + ); + + return existingFolders + .filter( + (existingFolder) => !discoveredExternalIds.has(existingFolder.externalId), + ) + .map((existingFolder) => existingFolder.id); +}; diff --git a/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/compute-folders-to-create.util.ts b/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/compute-folders-to-create.util.ts new file mode 100644 index 0000000000..2e45c4e1ac --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/compute-folders-to-create.util.ts @@ -0,0 +1,39 @@ +import { isNonEmptyString } from '@sniptt/guards'; + +import { + type DiscoveredMessageFolder, + type MessageFolder, +} from 'src/modules/messaging/message-folder-manager/interfaces/message-folder-driver.interface'; + +import { type MessageFolderWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; + +export const computeFoldersToCreate = ({ + discoveredFolders, + existingFolders, + messageChannelId, +}: { + discoveredFolders: DiscoveredMessageFolder[]; + existingFolders: MessageFolder[]; + messageChannelId: string; +}): Partial[] => { + const existingFoldersByExternalId = new Map( + existingFolders.map((folder) => [folder.externalId, folder]), + ); + + return discoveredFolders + .filter( + (discoveredFolder) => + !existingFoldersByExternalId.has(discoveredFolder.externalId), + ) + .map((discoveredFolder) => ({ + name: discoveredFolder.name, + externalId: discoveredFolder.externalId, + messageChannelId, + isSentFolder: discoveredFolder.isSentFolder, + isSynced: discoveredFolder.isSynced, + syncCursor: null, + parentFolderId: isNonEmptyString(discoveredFolder.parentFolderId) + ? discoveredFolder.parentFolderId + : null, + })); +}; diff --git a/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/compute-folders-to-update.util.ts b/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/compute-folders-to-update.util.ts new file mode 100644 index 0000000000..af16be4021 --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/compute-folders-to-update.util.ts @@ -0,0 +1,58 @@ +import { isNonEmptyString } from '@sniptt/guards'; +import deepEqual from 'deep-equal'; + +import { + type DiscoveredMessageFolder, + type MessageFolder, +} from 'src/modules/messaging/message-folder-manager/interfaces/message-folder-driver.interface'; + +import { type MessageFolderWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; + +export const computeFoldersToUpdate = ({ + discoveredFolders, + existingFolders, +}: { + discoveredFolders: DiscoveredMessageFolder[]; + existingFolders: MessageFolder[]; +}): Map> => { + const existingFoldersByExternalId = new Map( + existingFolders.map((folder) => [folder.externalId, folder]), + ); + + const foldersToUpdate = new Map< + string, + Partial + >(); + + for (const discoveredFolder of discoveredFolders) { + const existingFolder = existingFoldersByExternalId.get( + discoveredFolder.externalId, + ); + + if (!existingFolder) { + continue; + } + + const discoveredFolderData = { + name: discoveredFolder.name, + isSentFolder: discoveredFolder.isSentFolder, + parentFolderId: isNonEmptyString(discoveredFolder.parentFolderId) + ? discoveredFolder.parentFolderId + : null, + }; + + const existingFolderData = { + name: existingFolder.name, + isSentFolder: existingFolder.isSentFolder, + parentFolderId: isNonEmptyString(existingFolder.parentFolderId) + ? existingFolder.parentFolderId + : null, + }; + + if (!deepEqual(discoveredFolderData, existingFolderData)) { + foldersToUpdate.set(existingFolder.id, discoveredFolderData); + } + } + + return foldersToUpdate; +}; diff --git a/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/compute-updated-folders.util.ts b/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/compute-updated-folders.util.ts new file mode 100644 index 0000000000..7387a4631d --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-folder-manager/utils/compute-updated-folders.util.ts @@ -0,0 +1,38 @@ +import { type MessageFolder } from 'src/modules/messaging/message-folder-manager/interfaces/message-folder-driver.interface'; + +import { + MessageFolderPendingSyncAction, + type MessageFolderWorkspaceEntity, +} from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; + +export const computeUpdatedFolders = ({ + existingFolders, + foldersToUpdate, + folderIdsToDelete, +}: { + existingFolders: MessageFolder[]; + foldersToUpdate: Map>; + folderIdsToDelete: string[]; +}): MessageFolder[] => { + return existingFolders.map((existingFolder) => { + const update = foldersToUpdate.get(existingFolder.id); + const isMarkedForDeletion = folderIdsToDelete.includes(existingFolder.id); + + const pendingSyncAction = isMarkedForDeletion + ? MessageFolderPendingSyncAction.FOLDER_DELETION + : MessageFolderPendingSyncAction.NONE; + + if (update) { + return { + ...existingFolder, + ...update, + pendingSyncAction, + }; + } + + return { + ...existingFolder, + pendingSyncAction, + }; + }); +}; diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/gmail/services/gmail-get-message-list.service.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/gmail/services/gmail-get-message-list.service.ts index badbc3bdab..833ed0b190 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/gmail/services/gmail-get-message-list.service.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/gmail/services/gmail-get-message-list.service.ts @@ -1,7 +1,8 @@ import { Injectable, Logger } from '@nestjs/common'; +import { batchFetchImplementation } from '@jrmdayn/googleapis-batcher'; import { isNonEmptyString } from '@sniptt/guards'; -import { google, type gmail_v1 as gmailV1 } from 'googleapis'; +import { google } from 'googleapis'; import { isDefined } from 'twenty-shared/utils'; import { OAuth2ClientManagerService } from 'src/modules/connected-account/oauth2-client-manager/services/oauth2-client-manager.service'; @@ -19,6 +20,8 @@ import { type GetMessageListsArgs } from 'src/modules/messaging/message-import-m import { type GetMessageListsResponse } from 'src/modules/messaging/message-import-manager/types/get-message-lists-response.type'; import { assertNotNull } from 'src/utils/assert'; +const GMAIL_BATCH_REQUEST_MAX_SIZE = 50; + @Injectable() export class GmailGetMessageListService { private readonly logger = new Logger(GmailGetMessageListService.name); @@ -165,7 +168,7 @@ export class GmailGetMessageListService { await this.gmailGetHistoryService.getMessageIdsFromHistory(history); const messageIdsToFilter = await this.getEmailIdsFromExcludedFolders( - gmailClient, + connectedAccount, messageChannel.syncCursor, messageFolders, ); @@ -193,33 +196,54 @@ export class GmailGetMessageListService { } private async getEmailIdsFromExcludedFolders( - gmailClient: gmailV1.Gmail, + connectedAccount: Pick< + ConnectedAccountWorkspaceEntity, + 'provider' | 'accessToken' | 'refreshToken' | 'id' | 'handle' + >, lastSyncHistoryId: string, messageFolders: Pick< MessageFolderWorkspaceEntity, 'name' | 'externalId' | 'isSynced' >[], ): Promise { - const emailIds: string[] = []; - const toBeExcludedFolders = messageFolders.filter( - (folder) => !folder.isSynced, + (folder) => !folder.isSynced && isDefined(folder.externalId), ); - for (const folder of toBeExcludedFolders) { - if (!isDefined(folder.externalId)) { - continue; - } + if (toBeExcludedFolders.length === 0) { + return []; + } - const { history } = await this.gmailGetHistoryService.getHistory( - gmailClient, - lastSyncHistoryId, - ['messageAdded'], - folder.externalId, + const oAuth2Client = + await this.oAuth2ClientManagerService.getGoogleOAuth2Client( + connectedAccount, ); + const batchedFetchImplementation = batchFetchImplementation({ + maxBatchSize: GMAIL_BATCH_REQUEST_MAX_SIZE, + }); + const batchedGmailClient = google.gmail({ + version: 'v1', + auth: oAuth2Client, + fetchImplementation: batchedFetchImplementation, + }); + + const historyPromises = toBeExcludedFolders.map((folder) => + this.gmailGetHistoryService.getHistory( + batchedGmailClient, + lastSyncHistoryId, + ['messageAdded'], + folder.externalId!, + ), + ); + + const historyResults = await Promise.all(historyPromises); + + const emailIds: string[] = []; + + for (const { history } of historyResults) { const emailIdsFromCategory = history - .map((history) => history.messagesAdded) + .map((historyItem) => historyItem.messagesAdded) .flat() .map((message) => message?.message?.id) .filter((id) => id) 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 6c1c9e33d3..ea890b285a 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,8 +1,10 @@ import { Injectable, Logger } from '@nestjs/common'; import { type ImapFlow } from 'imapflow'; +import { isDefined } from 'twenty-shared/utils'; + +import { MessageFolder } from 'src/modules/messaging/message-folder-manager/interfaces/message-folder-driver.interface'; -import { type MessageFolderWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; import { MessageImportDriverException, MessageImportDriverExceptionCode, @@ -12,6 +14,7 @@ import { ImapMessageListFetchErrorHandler } from 'src/modules/messaging/message- import { ImapSyncService } from 'src/modules/messaging/message-import-manager/drivers/imap/services/imap-sync.service'; 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 { parseMessageId } from 'src/modules/messaging/message-import-manager/drivers/imap/utils/parse-message-id.util'; import { 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 { @@ -19,11 +22,6 @@ import { type GetOneMessageListResponse, } from 'src/modules/messaging/message-import-manager/types/get-message-lists-response.type'; -type MessageFolder = Pick< - MessageFolderWorkspaceEntity, - 'id' | 'name' | 'syncCursor' | 'externalId' ->; - @Injectable() export class ImapGetMessageListService { private readonly logger = new Logger(ImapGetMessageListService.name); @@ -65,17 +63,31 @@ export class ImapGetMessageListService { client: ImapFlow, folder: MessageFolder, ): Promise { - const folderPath = folder.externalId?.split(':')[0]; + const folderPath = parseMessageId(folder.externalId ?? '')?.folder; - if (!folderPath) { + if (!isDefined(folderPath)) { throw new MessageImportDriverException( `Folder ${folder.name} has no path`, MessageImportDriverExceptionCode.NOT_FOUND, ); } + if (await this.canSkipFolderSync(client, folder)) { + this.logger.log(`Skipping folder ${folder.name}: no new messages`); + + return { + messageExternalIds: [], + messageExternalIdsToDelete: [], + nextSyncCursor: folder.syncCursor ?? '', + previousSyncCursor: folder.syncCursor, + folderId: folder.id, + }; + } + this.logger.log(`Processing folder: ${folder.name}`); + const previousCursor = parseSyncCursor(folder.syncCursor); + const lock = await client.getMailboxLock(folderPath); try { @@ -89,7 +101,6 @@ export class ImapGetMessageListService { } const mailboxState = extractMailboxState(mailbox); - const previousCursor = parseSyncCursor(folder.syncCursor); const { messageUids } = await this.imapSyncService.syncFolder( client, @@ -125,4 +136,68 @@ export class ImapGetMessageListService { lock.release(); } } + + private async canSkipFolderSync( + client: ImapFlow, + folder: MessageFolder, + ): Promise { + const folderPath = parseMessageId(folder.externalId ?? '')?.folder; + const previousCursor = parseSyncCursor(folder.syncCursor); + + if (!isDefined(folderPath) || !isDefined(previousCursor)) { + return false; + } + + try { + const supportsCondstore = client.capabilities.has('CONDSTORE'); + + const status = await client.status(folderPath, { + uidNext: true, + uidValidity: true, + ...(supportsCondstore && { highestModseq: true }), + }); + + if (!isDefined(status.uidValidity)) { + this.logger.debug( + `Folder ${folderPath}: Server missing UIDVALIDITY. Sync required.`, + ); + + return false; + } + + const uidNext = Number(status.uidNext ?? 1); + const uidValidity = Number(status.uidValidity); + + if (previousCursor.uidValidity !== uidValidity) { + this.logger.debug( + `Folder ${folderPath}: UIDVALIDITY changed (${previousCursor.uidValidity} → ${uidValidity}). Full sync required.`, + ); + + return false; + } + + const hasModSeqChanged = + isDefined(previousCursor.modSeq) && + isDefined(status.highestModseq) && + previousCursor.modSeq !== status.highestModseq.toString(); + + if (hasModSeqChanged) { + this.logger.debug( + `Folder ${folderPath}: MODSEQ changed (${previousCursor.modSeq} → ${status.highestModseq}). Sync required.`, + ); + + return false; + } + + const maxUid = Math.max(0, uidNext - 1); + + return previousCursor.highestUid >= maxUid; + } catch (error) { + this.logger.warn( + `Failed to get status for folder ${folderPath}: ${error.message}`, + ); + + return false; + } + } } diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/microsoft/services/microsoft-get-message-list.service.dev.spec.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/microsoft/services/microsoft-get-message-list.service.dev.spec.ts index d2a72c7618..2fc97bf457 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/microsoft/services/microsoft-get-message-list.service.dev.spec.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/microsoft/services/microsoft-get-message-list.service.dev.spec.ts @@ -8,7 +8,10 @@ import { MicrosoftOAuth2ClientManagerService } from 'src/modules/connected-accou import { OAuth2ClientManagerService } from 'src/modules/connected-account/oauth2-client-manager/services/oauth2-client-manager.service'; import { type ConnectedAccountWorkspaceEntity } from 'src/modules/connected-account/standard-objects/connected-account.workspace-entity'; import { MessageChannelWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-channel.workspace-entity'; -import { MessageFolderWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; +import { + MessageFolderPendingSyncAction, + MessageFolderWorkspaceEntity, +} from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; import { microsoftGraphWithMessagesDeltaLink } from 'src/modules/messaging/message-import-manager/drivers/microsoft/mocks/microsoft-api-examples'; import { MessageFolderName } from 'src/modules/messaging/message-import-manager/drivers/microsoft/types/folders'; @@ -79,6 +82,8 @@ xdescribe('Microsoft dev tests : get message list service', () => { isSynced: false, isSentFolder: false, externalId: null, + parentFolderId: null, + pendingSyncAction: MessageFolderPendingSyncAction.NONE, }, ], }); @@ -108,6 +113,8 @@ xdescribe('Microsoft dev tests : get message list service', () => { isSynced: false, isSentFolder: false, externalId: null, + parentFolderId: null, + pendingSyncAction: MessageFolderPendingSyncAction.NONE, }, ], }), @@ -127,6 +134,8 @@ xdescribe('Microsoft dev tests : get message list service', () => { isSynced: false, isSentFolder: false, externalId: null, + parentFolderId: null, + pendingSyncAction: MessageFolderPendingSyncAction.NONE, }, ], }); @@ -150,6 +159,8 @@ xdescribe('Microsoft dev tests : get message list service', () => { isSynced: false, isSentFolder: false, externalId: null, + parentFolderId: null, + pendingSyncAction: MessageFolderPendingSyncAction.NONE, }, ], }), @@ -168,6 +179,7 @@ xdescribe('Microsoft dev tests : get message list service for folders', () => { inboxFolder.name = MessageFolderName.INBOX; inboxFolder.syncCursor = 'inbox-sync-cursor'; inboxFolder.messageChannelId = 'message-channel-1'; + inboxFolder.parentFolderId = null; const sentFolder = new MessageFolderWorkspaceEntity(); @@ -175,6 +187,7 @@ xdescribe('Microsoft dev tests : get message list service for folders', () => { sentFolder.name = MessageFolderName.SENT_ITEMS; sentFolder.syncCursor = 'sent-sync-cursor'; sentFolder.messageChannelId = 'message-channel-1'; + sentFolder.parentFolderId = null; const otherFolder = new MessageFolderWorkspaceEntity(); @@ -182,6 +195,7 @@ xdescribe('Microsoft dev tests : get message list service for folders', () => { otherFolder.name = 'other'; otherFolder.syncCursor = 'other-sync-cursor'; otherFolder.messageChannelId = 'message-channel-2'; + otherFolder.parentFolderId = null; const messageChannelNoFolders = new MessageChannelWorkspaceEntity(); diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/microsoft/services/microsoft-get-message-list.service.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/microsoft/services/microsoft-get-message-list.service.ts index 7c27d67297..82f7332a4a 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/microsoft/services/microsoft-get-message-list.service.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/microsoft/services/microsoft-get-message-list.service.ts @@ -6,7 +6,7 @@ import { type PageIteratorCallback, } from '@microsoft/microsoft-graph-client'; import { isNonEmptyString } from '@sniptt/guards'; -import { isDefined } from 'twenty-shared/utils'; +import pLimit from 'p-limit'; import { OAuth2ClientManagerService } from 'src/modules/connected-account/oauth2-client-manager/services/oauth2-client-manager.service'; import { type ConnectedAccountWorkspaceEntity } from 'src/modules/connected-account/standard-objects/connected-account.workspace-entity'; @@ -25,6 +25,9 @@ import { // Microsoft API limit is 999 messages per request on this endpoint const MESSAGING_MICROSOFT_USERS_MESSAGES_LIST_MAX_RESULT = 999; +/* reference: https://learn.microsoft.com/en-us/graph/throttling-limits#limits-per-mailbox */ +const FOLDER_PROCESSING_CONCURRENCY = 4; + @Injectable() export class MicrosoftGetMessageListService { private readonly logger = new Logger(MicrosoftGetMessageListService.name); @@ -38,8 +41,6 @@ export class MicrosoftGetMessageListService { connectedAccount, messageFolders, }: GetMessageListsArgs): Promise { - const result: GetMessageListsResponse = []; - if (messageFolders.length === 0) { throw new MessageImportDriverException( `Message channel ${messageChannel.id} has no message folders`, @@ -47,16 +48,22 @@ export class MicrosoftGetMessageListService { ); } - for (const folder of messageFolders) { - const response = await this.getMessageList(connectedAccount, folder); + const limit = pLimit(FOLDER_PROCESSING_CONCURRENCY); - result.push({ - ...response, - folderId: folder.id, - }); - } + const results = await Promise.all( + messageFolders.map((folder) => + limit(async () => { + const response = await this.getMessageList(connectedAccount, folder); - return result; + return { + ...response, + folderId: folder.id, + }; + }), + ), + ); + + return results; } public async getMessageList( @@ -116,13 +123,6 @@ export class MicrosoftGetMessageListService { this.microsoftMessageListFetchErrorHandler.handleError(error); }); - if (!isDefined(messageFolder.syncCursor)) { - throw new MessageImportDriverException( - 'Message folder sync cursor is required', - MessageImportDriverExceptionCode.SYNC_CURSOR_ERROR, - ); - } - return { messageExternalIds, messageExternalIdsToDelete, diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/services/__tests__/messaging-message-list-fetch.service.spec.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/services/__tests__/messaging-message-list-fetch.service.spec.ts index ca4b358601..73b9c81ade 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/services/__tests__/messaging-message-list-fetch.service.spec.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/services/__tests__/messaging-message-list-fetch.service.spec.ts @@ -7,7 +7,10 @@ import { CacheStorageNamespace } from 'src/engine/core-modules/cache-storage/typ import { GlobalWorkspaceOrmManager } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager'; import { MessageChannelSyncStatusService } from 'src/modules/messaging/common/services/message-channel-sync-status.service'; import { type MessageChannelWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-channel.workspace-entity'; -import { type MessageFolderWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; +import { + MessageFolderPendingSyncAction, + type MessageFolderWorkspaceEntity, +} from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; import { MessagingMessageCleanerService } from 'src/modules/messaging/message-cleaner/services/messaging-message-cleaner.service'; import { SyncMessageFoldersService } from 'src/modules/messaging/message-folder-manager/services/sync-message-folders.service'; import { MessagingAccountAuthenticationService } from 'src/modules/messaging/message-import-manager/services/messaging-account-authentication.service'; @@ -64,7 +67,8 @@ describe('MessagingMessageListFetchService', () => { handleAliases: '', }, syncCursor: 'google-sync-cursor', - } as MessageChannelWorkspaceEntity; + messageFolders: [], + } as unknown as MessageChannelWorkspaceEntity; }); beforeEach(async () => { @@ -248,7 +252,16 @@ describe('MessagingMessageListFetchService', () => { { provide: SyncMessageFoldersService, useValue: { - syncMessageFolders: jest.fn().mockResolvedValue(undefined), + syncMessageFolders: jest.fn().mockResolvedValue([ + { + id: 'inbox-folder-id', + name: 'inbox', + syncCursor: 'inbox-sync-cursor', + messageChannelId: 'microsoft-message-channel-id', + isSynced: true, + pendingSyncAction: MessageFolderPendingSyncAction.NONE, + }, + ]), }, }, { @@ -323,6 +336,7 @@ describe('MessagingMessageListFetchService', () => { syncCursor: 'inbox-sync-cursor', messageChannelId: 'microsoft-message-channel-id', isSynced: true, + pendingSyncAction: MessageFolderPendingSyncAction.NONE, }, ], ); @@ -384,6 +398,7 @@ describe('MessagingMessageListFetchService', () => { syncCursor: 'inbox-sync-cursor', messageChannelId: 'microsoft-message-channel-id', isSynced: true, + pendingSyncAction: MessageFolderPendingSyncAction.NONE, }, ], ); diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-get-message-list.service.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-get-message-list.service.ts index d4a574b5cf..512ee6f003 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-get-message-list.service.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-get-message-list.service.ts @@ -2,8 +2,9 @@ import { Injectable } from '@nestjs/common'; import { ConnectedAccountProvider } from 'twenty-shared/types'; +import { MessageFolder } from 'src/modules/messaging/message-folder-manager/interfaces/message-folder-driver.interface'; + import { type MessageChannelWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-channel.workspace-entity'; -import { MessageFolderWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; import { MessageImportDriverException, MessageImportDriverExceptionCode, @@ -13,11 +14,6 @@ import { ImapGetMessageListService } from 'src/modules/messaging/message-import- import { MicrosoftGetMessageListService } from 'src/modules/messaging/message-import-manager/drivers/microsoft/services/microsoft-get-message-list.service'; import { type GetMessageListsResponse } from 'src/modules/messaging/message-import-manager/types/get-message-lists-response.type'; -type MessageFolder = Pick< - MessageFolderWorkspaceEntity, - 'name' | 'isSynced' | 'isSentFolder' | 'externalId' | 'syncCursor' | 'id' ->; - @Injectable() export class MessagingGetMessageListService { constructor( 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 25cf9ae976..28b511f263 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 @@ -8,7 +8,6 @@ import { In, MoreThanOrEqual } from 'typeorm'; import { InjectCacheStorage } from 'src/engine/core-modules/cache-storage/decorators/cache-storage.decorator'; import { CacheStorageService } from 'src/engine/core-modules/cache-storage/services/cache-storage.service'; import { CacheStorageNamespace } from 'src/engine/core-modules/cache-storage/types/cache-storage-namespace.enum'; -import { type WorkspaceEntityManager } from 'src/engine/twenty-orm/entity-manager/workspace-entity-manager'; import { GlobalWorkspaceOrmManager } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager'; import { buildSystemAuthContext } from 'src/engine/twenty-orm/utils/build-system-auth-context.util'; import { MessageChannelSyncStatusService } from 'src/modules/messaging/common/services/message-channel-sync-status.service'; @@ -19,10 +18,7 @@ import { MessageChannelWorkspaceEntity, MessageFolderImportPolicy, } from 'src/modules/messaging/common/standard-objects/message-channel.workspace-entity'; -import { - MessageFolderPendingSyncAction, - MessageFolderWorkspaceEntity, -} from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; +import { MessageFolderPendingSyncAction } from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; import { MessagingMessageCleanerService } from 'src/modules/messaging/message-cleaner/services/messaging-message-cleaner.service'; import { SyncMessageFoldersService } from 'src/modules/messaging/message-folder-manager/services/sync-message-folders.service'; import { MessagingAccountAuthenticationService } from 'src/modules/messaging/message-import-manager/services/messaging-account-authentication.service'; @@ -127,33 +123,21 @@ export class MessagingMessageListFetchService { }, }; - const datasource = - await this.globalWorkspaceOrmManager.getGlobalWorkspaceDataSource(); - - await this.syncMessageFoldersService.syncMessageFolders({ - workspaceId, - messageChannel: messageChannelWithFreshTokens, - manager: datasource.manager as WorkspaceEntityManager, - }); - - const messageFolderRepository = - await this.globalWorkspaceOrmManager.getRepository( + const messageFolders = + await this.syncMessageFoldersService.syncMessageFolders({ + messageChannel: messageChannelWithFreshTokens, workspaceId, - 'messageFolder', - ); + }); - const messageFolders = await messageFolderRepository.find({ - where: { - messageChannelId: freshMessageChannel.id, - pendingSyncAction: MessageFolderPendingSyncAction.NONE, - }, - }); - - const messageFoldersToSync = + const messageFoldersToSync = ( messageChannelWithFreshTokens.messageFolderImportPolicy === MessageFolderImportPolicy.ALL_FOLDERS ? messageFolders - : messageFolders.filter((folder) => folder.isSynced); + : messageFolders.filter((folder) => folder.isSynced) + ).filter( + (folder) => + folder.pendingSyncAction === MessageFolderPendingSyncAction.NONE, + ); const messageLists = await this.messagingGetMessageListService.getMessageLists( @@ -360,18 +344,11 @@ export class MessagingMessageListFetchService { messageChannel: MessageChannelWorkspaceEntity, workspaceId: string, ): Promise { - const messageFolderRepository = - await this.globalWorkspaceOrmManager.getRepository( - workspaceId, - 'messageFolder', - ); - - const foldersWithPendingActions = await messageFolderRepository.find({ - where: { - messageChannelId: messageChannel.id, - pendingSyncAction: MessageFolderPendingSyncAction.FOLDER_DELETION, - }, - }); + const foldersWithPendingActions = messageChannel.messageFolders.filter( + (folder) => + isDefined(folder.pendingSyncAction) && + folder.pendingSyncAction !== MessageFolderPendingSyncAction.NONE, + ); if (foldersWithPendingActions.length === 0) { return false; diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/types/get-message-lists-args.type.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/types/get-message-lists-args.type.ts index 078b71a66a..20fcef8144 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/types/get-message-lists-args.type.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/types/get-message-lists-args.type.ts @@ -1,6 +1,7 @@ +import { type MessageFolder } from 'src/modules/messaging/message-folder-manager/interfaces/message-folder-driver.interface'; + import { type ConnectedAccountWorkspaceEntity } from 'src/modules/connected-account/standard-objects/connected-account.workspace-entity'; import { type MessageChannelWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-channel.workspace-entity'; -import { type MessageFolderWorkspaceEntity } from 'src/modules/messaging/common/standard-objects/message-folder.workspace-entity'; export type GetMessageListsArgs = { messageChannel: Pick; @@ -13,8 +14,5 @@ export type GetMessageListsArgs = { | 'handle' | 'connectionParameters' >; - messageFolders: Pick< - MessageFolderWorkspaceEntity, - 'name' | 'syncCursor' | 'id' | 'isSynced' | 'isSentFolder' | 'externalId' - >[]; + messageFolders: MessageFolder[]; }; diff --git a/yarn.lock b/yarn.lock index abb04ea719..e956802407 100644 --- a/yarn.lock +++ b/yarn.lock @@ -9726,6 +9726,22 @@ __metadata: languageName: node linkType: hard +"@jrmdayn/googleapis-batcher@npm:^0.10.1": + version: 0.10.1 + resolution: "@jrmdayn/googleapis-batcher@npm:0.10.1" + dependencies: + dataloader: "npm:^2.2.1" + debug: "npm:^4.3.4" + handlebars: "npm:^4.7.7" + next-line: "npm:^1.1.0" + node-fetch: "npm:2" + peerDependencies: + gaxios: ">= 5" + googleapis: ">= 109" + checksum: 10c0/ffc61154fa03227c3c865cc44f33d545a50b66804c87e3d93b28afe9d297fd6ac91bd55205d0d414aba6f907840ca2f22ec5531617fd09c184811e6eba388ab1 + languageName: node + linkType: hard + "@js-sdsl/ordered-map@npm:^4.4.2": version: 4.4.2 resolution: "@js-sdsl/ordered-map@npm:4.4.2" @@ -32098,6 +32114,13 @@ __metadata: languageName: node linkType: hard +"dataloader@npm:^2.2.1": + version: 2.2.3 + resolution: "dataloader@npm:2.2.3" + checksum: 10c0/9b9a056fbc863ca86da87d59e053e871e263b4966aa4d55e40d61a65e96815fae5530ca220629064ca5f8e3000c0c4ec93292e170c38ff393fb34256b4d7c1aa + languageName: node + linkType: hard + "date-fns-tz@npm:^2.0.0": version: 2.0.1 resolution: "date-fns-tz@npm:2.0.1" @@ -37595,7 +37618,7 @@ __metadata: languageName: node linkType: hard -"handlebars@npm:*, handlebars@npm:^4.7.8": +"handlebars@npm:*, handlebars@npm:^4.7.7, handlebars@npm:^4.7.8": version: 4.7.8 resolution: "handlebars@npm:4.7.8" dependencies: @@ -45952,6 +45975,13 @@ __metadata: languageName: node linkType: hard +"next-line@npm:^1.1.0": + version: 1.1.0 + resolution: "next-line@npm:1.1.0" + checksum: 10c0/c17aed62e7af05e612cb852d601298f8bef47822231f06fb68c4a8ae7b272c2d49923eefcff56082f173da36a1b1cbd00dad11b05559957d268c38daab2e721b + languageName: node + linkType: hard + "next-mdx-remote-client@npm:^1.0.3": version: 1.1.4 resolution: "next-mdx-remote-client@npm:1.1.4" @@ -46255,6 +46285,20 @@ __metadata: languageName: node linkType: hard +"node-fetch@npm:2, node-fetch@npm:2.7.0, node-fetch@npm:^2.6.0, node-fetch@npm:^2.6.1, node-fetch@npm:^2.6.12, node-fetch@npm:^2.6.7, node-fetch@npm:^2.6.9, node-fetch@npm:^2.7.0": + version: 2.7.0 + resolution: "node-fetch@npm:2.7.0" + dependencies: + whatwg-url: "npm:^5.0.0" + peerDependencies: + encoding: ^0.1.0 + peerDependenciesMeta: + encoding: + optional: true + checksum: 10c0/b55786b6028208e6fbe594ccccc213cab67a72899c9234eb59dba51062a299ea853210fcf526998eaa2867b0963ad72338824450905679ff0fa304b8c5093ae8 + languageName: node + linkType: hard + "node-fetch@npm:2.6.7": version: 2.6.7 resolution: "node-fetch@npm:2.6.7" @@ -46269,20 +46313,6 @@ __metadata: languageName: node linkType: hard -"node-fetch@npm:2.7.0, node-fetch@npm:^2.6.0, node-fetch@npm:^2.6.1, node-fetch@npm:^2.6.12, node-fetch@npm:^2.6.7, node-fetch@npm:^2.6.9, node-fetch@npm:^2.7.0": - version: 2.7.0 - resolution: "node-fetch@npm:2.7.0" - dependencies: - whatwg-url: "npm:^5.0.0" - peerDependencies: - encoding: ^0.1.0 - peerDependenciesMeta: - encoding: - optional: true - checksum: 10c0/b55786b6028208e6fbe594ccccc213cab67a72899c9234eb59dba51062a299ea853210fcf526998eaa2867b0963ad72338824450905679ff0fa304b8c5093ae8 - languageName: node - linkType: hard - "node-forge@npm:^1.3.1": version: 1.3.2 resolution: "node-forge@npm:1.3.2" @@ -56494,6 +56524,7 @@ __metadata: "@graphql-tools/schema": "npm:10.0.4" "@graphql-tools/utils": "npm:9.2.1" "@graphql-yoga/nestjs": "patch:@graphql-yoga/nestjs@2.1.0#./patches/@graphql-yoga+nestjs+2.1.0.patch" + "@jrmdayn/googleapis-batcher": "npm:^0.10.1" "@lingui/cli": "npm:^5.1.2" "@lingui/conf": "npm:5.1.2" "@lingui/core": "npm:^5.1.2"