From 39017f43b833ced2af721a3d5ab46023beb64ebf Mon Sep 17 00:00:00 2001 From: neo773 <62795688+neo773@users.noreply.github.com> Date: Mon, 10 Nov 2025 22:48:17 +0530 Subject: [PATCH] handle invalid/expired syncCursor (#15715) --- .../gmail-get-history.service.spec.ts | 149 ++++++++++++++++++ ...rse-gmail-message-list-fetch-error.spec.ts | 6 +- ...rse-gmail-message-list-fetch-error.util.ts | 2 +- .../parse-microsoft-messages-import.util.ts | 8 + .../messaging-import-manager.module.ts | 2 + .../messaging-clear-cursors.module.ts | 9 ++ .../messaging-clear-cursors.service.spec.ts | 62 ++++++++ .../messaging-clear-cursors.service.ts | 48 ++++++ ...saging-import-exception-handler.service.ts | 22 ++- 9 files changed, 304 insertions(+), 4 deletions(-) create mode 100644 packages/twenty-server/src/modules/messaging/message-import-manager/drivers/gmail/services/gmail-get-history.service.spec.ts create mode 100644 packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-clear-cursors.module.ts create mode 100644 packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-clear-cursors.service.spec.ts create mode 100644 packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-clear-cursors.service.ts diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/gmail/services/gmail-get-history.service.spec.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/gmail/services/gmail-get-history.service.spec.ts new file mode 100644 index 0000000000..ff0f093fb3 --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/gmail/services/gmail-get-history.service.spec.ts @@ -0,0 +1,149 @@ +import { Test, type TestingModule } from '@nestjs/testing'; + +import { type gmail_v1 } from 'googleapis'; + +import { MessageImportDriverExceptionCode } from 'src/modules/messaging/message-import-manager/drivers/exceptions/message-import-driver.exception'; +import { GmailGetHistoryService } from 'src/modules/messaging/message-import-manager/drivers/gmail/services/gmail-get-history.service'; +import { GmailMessageListFetchErrorHandler } from 'src/modules/messaging/message-import-manager/drivers/gmail/services/gmail-message-list-fetch-error-handler.service'; + +describe('GmailGetHistoryService', () => { + let service: GmailGetHistoryService; + let errorHandler: GmailMessageListFetchErrorHandler; + + beforeEach(async () => { + const module: TestingModule = await Test.createTestingModule({ + providers: [ + GmailGetHistoryService, + { + provide: GmailMessageListFetchErrorHandler, + useValue: { handleError: jest.fn() }, + }, + ], + }).compile(); + + service = module.get(GmailGetHistoryService); + errorHandler = module.get(GmailMessageListFetchErrorHandler); + }); + + describe('getHistory', () => { + it('should fetch history', async () => { + const mockClient = { + users: { + history: { + list: jest.fn().mockResolvedValue({ + data: { + history: [{ messagesAdded: [{ message: { id: 'msg1' } }] }], + historyId: '200', + }, + }), + }, + }, + } as unknown as gmail_v1.Gmail; + + const result = await service.getHistory(mockClient, '100'); + + expect(result.historyId).toBe('200'); + expect(result.history).toHaveLength(1); + }); + + it('should handle pagination', async () => { + const mockClient = { + users: { + history: { + list: jest + .fn() + .mockResolvedValueOnce({ + data: { + history: [{ id: '1' }], + nextPageToken: 'token', + historyId: '200', + }, + }) + .mockResolvedValueOnce({ + data: { history: [{ id: '2' }], historyId: '201' }, + }), + }, + }, + } as unknown as gmail_v1.Gmail; + + const result = await service.getHistory(mockClient, '100'); + + expect(result.history).toHaveLength(2); + expect(mockClient.users.history.list).toHaveBeenCalledTimes(2); + }); + + it('should throw SYNC_CURSOR_ERROR when historyId is expired', async () => { + const error = { code: 404, message: 'Not found' }; + const mockClient = { + users: { + history: { + list: jest.fn().mockRejectedValue(error), + }, + }, + } as unknown as gmail_v1.Gmail; + + (errorHandler.handleError as jest.Mock).mockImplementation(() => { + throw { + code: MessageImportDriverExceptionCode.SYNC_CURSOR_ERROR, + message: 'Sync cursor error', + }; + }); + + await expect( + service.getHistory(mockClient, 'expired-id'), + ).rejects.toMatchObject({ + code: MessageImportDriverExceptionCode.SYNC_CURSOR_ERROR, + }); + + expect(errorHandler.handleError).toHaveBeenCalledWith(error); + }); + + it('should delegate other errors to error handler', async () => { + const error = { code: 500, message: 'Server error' }; + const mockClient = { + users: { + history: { + list: jest.fn().mockRejectedValue(error), + }, + }, + } as unknown as gmail_v1.Gmail; + + await service.getHistory(mockClient, '100'); + + expect(errorHandler.handleError).toHaveBeenCalledWith(error); + }); + }); + + describe('getMessageIdsFromHistory', () => { + it('should extract message IDs', async () => { + const history: gmail_v1.Schema$History[] = [ + { + messagesAdded: [{ message: { id: 'add1' } }], + messagesDeleted: [{ message: { id: 'del1' } }], + }, + ]; + + const result = await service.getMessageIdsFromHistory(history); + + expect(result.messagesAdded).toEqual(['add1']); + expect(result.messagesDeleted).toEqual(['del1']); + }); + + it('should deduplicate messages that appear in both lists', async () => { + const history: gmail_v1.Schema$History[] = [ + { + messagesAdded: [ + { message: { id: 'msg1' } }, + { message: { id: 'msg2' } }, + ], + messagesDeleted: [{ message: { id: 'msg1' } }], + }, + ]; + + const result = await service.getMessageIdsFromHistory(history); + + expect(result.messagesAdded).toEqual(['msg2']); + expect(result.messagesDeleted).toEqual([]); + }); + }); +}); diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/gmail/utils/__tests__/parse-gmail-message-list-fetch-error.spec.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/gmail/utils/__tests__/parse-gmail-message-list-fetch-error.spec.ts index 1f41d47d09..07e78c864c 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/gmail/utils/__tests__/parse-gmail-message-list-fetch-error.spec.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/gmail/utils/__tests__/parse-gmail-message-list-fetch-error.spec.ts @@ -84,12 +84,14 @@ describe('parseGmailMessageListFetchError', () => { ); }); - it('should handle 404 Not Found', () => { + it('should handle 404 as sync cursor error', () => { const error = gmailApiErrorMocks.getError(404); const exception = parseGmailMessageListFetchError(error.error); expect(exception).toBeInstanceOf(MessageImportDriverException); - expect(exception.code).toBe(MessageImportDriverExceptionCode.NOT_FOUND); + expect(exception.code).toBe( + MessageImportDriverExceptionCode.SYNC_CURSOR_ERROR, + ); }); it('should handle 410 Gone', () => { diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/gmail/utils/parse-gmail-message-list-fetch-error.util.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/gmail/utils/parse-gmail-message-list-fetch-error.util.ts index 9fc3cb541c..c0b8ec1390 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/gmail/utils/parse-gmail-message-list-fetch-error.util.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/gmail/utils/parse-gmail-message-list-fetch-error.util.ts @@ -52,7 +52,7 @@ export const parseGmailMessageListFetchError = ( case 404: return new MessageImportDriverException( message, - MessageImportDriverExceptionCode.NOT_FOUND, + MessageImportDriverExceptionCode.SYNC_CURSOR_ERROR, { cause: options?.cause }, ); diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/microsoft/utils/parse-microsoft-messages-import.util.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/microsoft/utils/parse-microsoft-messages-import.util.ts index 405e1957d4..5c8160dba1 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/microsoft/utils/parse-microsoft-messages-import.util.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/drivers/microsoft/utils/parse-microsoft-messages-import.util.ts @@ -59,6 +59,14 @@ export const parseMicrosoftMessagesImportError = ( ); } + if (error.statusCode === 410) { + return new MessageImportDriverException( + `Sync cursor error: ${error.message}`, + MessageImportDriverExceptionCode.SYNC_CURSOR_ERROR, + { cause: options?.cause }, + ); + } + return new MessageImportDriverException( `Microsoft Graph API unknown error: ${error} with status code ${error.statusCode}`, MessageImportDriverExceptionCode.UNKNOWN, diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/messaging-import-manager.module.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/messaging-import-manager.module.ts index 885afbf4ae..5736b5c404 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/messaging-import-manager.module.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/messaging-import-manager.module.ts @@ -35,6 +35,7 @@ import { MessagingOngoingStaleJob } from 'src/modules/messaging/message-import-m import { MessagingRelaunchFailedMessageChannelJob } from 'src/modules/messaging/message-import-manager/jobs/messaging-relaunch-failed-message-channel.job'; import { MessagingMessageImportManagerMessageChannelListener } from 'src/modules/messaging/message-import-manager/listeners/messaging-import-manager-message-channel.listener'; import { MessagingAccountAuthenticationService } from 'src/modules/messaging/message-import-manager/services/messaging-account-authentication.service'; +import { MessagingClearCursorsModule } from 'src/modules/messaging/message-import-manager/services/messaging-clear-cursors.module'; import { MessagingCursorService } from 'src/modules/messaging/message-import-manager/services/messaging-cursor.service'; import { MessagingGetMessageListService } from 'src/modules/messaging/message-import-manager/services/messaging-get-message-list.service'; import { MessagingGetMessagesService } from 'src/modules/messaging/message-import-manager/services/messaging-get-messages.service'; @@ -65,6 +66,7 @@ import { MessagingMonitoringModule } from 'src/modules/messaging/monitoring/mess FeatureFlagModule, MessageParticipantManagerModule, MessagingFolderSyncManagerModule, + MessagingClearCursorsModule, MessagingMonitoringModule, MessagingMessageCleanerModule, WorkspaceEventEmitterModule, diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-clear-cursors.module.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-clear-cursors.module.ts new file mode 100644 index 0000000000..72927d99f2 --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-clear-cursors.module.ts @@ -0,0 +1,9 @@ +import { Module } from '@nestjs/common'; + +import { MessagingClearCursorsService } from 'src/modules/messaging/message-import-manager/services/messaging-clear-cursors.service'; + +@Module({ + providers: [MessagingClearCursorsService], + exports: [MessagingClearCursorsService], +}) +export class MessagingClearCursorsModule {} diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-clear-cursors.service.spec.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-clear-cursors.service.spec.ts new file mode 100644 index 0000000000..c5dda9cfa3 --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-clear-cursors.service.spec.ts @@ -0,0 +1,62 @@ +import { Test, type TestingModule } from '@nestjs/testing'; + +import { TwentyORMManager } from 'src/engine/twenty-orm/twenty-orm.manager'; +import { MessagingClearCursorsService } from 'src/modules/messaging/message-import-manager/services/messaging-clear-cursors.service'; + +describe('MessagingClearCursorsService', () => { + let service: MessagingClearCursorsService; + + const mockMessageChannelRepository = { + update: jest.fn(), + }; + + const mockMessageFolderRepository = { + update: jest.fn(), + }; + + beforeEach(async () => { + const module: TestingModule = await Test.createTestingModule({ + providers: [ + MessagingClearCursorsService, + { + provide: TwentyORMManager, + useValue: { + getRepository: jest.fn((entityName) => { + if (entityName === 'messageChannel') { + return mockMessageChannelRepository; + } + if (entityName === 'messageFolder') { + return mockMessageFolderRepository; + } + }), + }, + }, + ], + }).compile(); + + service = module.get(MessagingClearCursorsService); + }); + + afterEach(() => { + jest.clearAllMocks(); + }); + + describe('clearAllMessageChannelCursors', () => { + const messageChannelId = 'test-channel-id'; + + it('should clear message channel and folder cursors', async () => { + await service.clearAllMessageChannelCursors(messageChannelId); + + expect(mockMessageChannelRepository.update).toHaveBeenCalledWith( + { id: messageChannelId }, + { syncCursor: '' }, + undefined, + ); + expect(mockMessageFolderRepository.update).toHaveBeenCalledWith( + { messageChannelId }, + { syncCursor: '' }, + undefined, + ); + }); + }); +}); diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-clear-cursors.service.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-clear-cursors.service.ts new file mode 100644 index 0000000000..feabc194b7 --- /dev/null +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-clear-cursors.service.ts @@ -0,0 +1,48 @@ +import { Injectable, Logger } from '@nestjs/common'; + +import { type WorkspaceEntityManager } from 'src/engine/twenty-orm/entity-manager/workspace-entity-manager'; +import { TwentyORMManager } from 'src/engine/twenty-orm/twenty-orm.manager'; +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'; + +@Injectable() +export class MessagingClearCursorsService { + private readonly logger = new Logger(MessagingClearCursorsService.name); + + constructor(private readonly twentyORMManager: TwentyORMManager) {} + + async clearAllMessageChannelCursors( + messageChannelId: string, + transactionManager?: WorkspaceEntityManager, + ): Promise { + this.logger.log( + `MessageChannelId: ${messageChannelId} - Clearing all sync cursors`, + ); + + const messageChannelRepository = + await this.twentyORMManager.getRepository( + 'messageChannel', + ); + + const messageFolderRepository = + await this.twentyORMManager.getRepository( + 'messageFolder', + ); + + await messageChannelRepository.update( + { id: messageChannelId }, + { syncCursor: '' }, + transactionManager, + ); + + await messageFolderRepository.update( + { messageChannelId }, + { syncCursor: '' }, + transactionManager, + ); + + this.logger.log( + `MessageChannelId: ${messageChannelId} - Cleared all sync cursors`, + ); + } +} diff --git a/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-import-exception-handler.service.ts b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-import-exception-handler.service.ts index 02013355ca..4bb42c6ff6 100644 --- a/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-import-exception-handler.service.ts +++ b/packages/twenty-server/src/modules/messaging/message-import-manager/services/messaging-import-exception-handler.service.ts @@ -17,6 +17,7 @@ import { MessageImportDriverExceptionCode, } from 'src/modules/messaging/message-import-manager/drivers/exceptions/message-import-driver.exception'; import { MessageNetworkExceptionCode } from 'src/modules/messaging/message-import-manager/drivers/exceptions/message-network.exception'; +import { MessagingClearCursorsService } from 'src/modules/messaging/message-import-manager/services/messaging-clear-cursors.service'; export enum MessageImportSyncStep { MESSAGE_LIST_FETCH = 'MESSAGE_LIST_FETCH', @@ -29,6 +30,7 @@ export class MessageImportExceptionHandlerService { constructor( private readonly twentyORMManager: TwentyORMManager, private readonly messageChannelSyncStatusService: MessageChannelSyncStatusService, + private readonly messagingClearCursorsService: MessagingClearCursorsService, private readonly exceptionHandlerService: ExceptionHandlerService, ) {} @@ -81,7 +83,10 @@ export class MessageImportExceptionHandlerService { ); break; case MessageImportDriverExceptionCode.SYNC_CURSOR_ERROR: - await this.handlePermanentException(messageChannel, workspaceId); + await this.handleSyncCursorErrorException( + messageChannel, + workspaceId, + ); break; case MessageImportDriverExceptionCode.UNKNOWN: case MessageImportDriverExceptionCode.UNKNOWN_NETWORK_ERROR: @@ -98,6 +103,21 @@ export class MessageImportExceptionHandlerService { } } + private async handleSyncCursorErrorException( + messageChannel: Pick, + workspaceId: string, + ): Promise { + await this.messageChannelSyncStatusService.markAsFailed( + [messageChannel.id], + workspaceId, + MessageChannelSyncStatus.FAILED_UNKNOWN, + ); + + await this.messagingClearCursorsService.clearAllMessageChannelCursors( + messageChannel.id, + ); + } + private async handleTemporaryException( syncStep: MessageImportSyncStep, messageChannel: Pick<