handle invalid/expired syncCursor (#15715)

This commit is contained in:
neo773
2025-11-10 22:48:17 +05:30
committed by GitHub
parent 22a81f705b
commit 39017f43b8
9 changed files with 304 additions and 4 deletions
@@ -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([]);
});
});
});
@@ -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', () => {
@@ -52,7 +52,7 @@ export const parseGmailMessageListFetchError = (
case 404:
return new MessageImportDriverException(
message,
MessageImportDriverExceptionCode.NOT_FOUND,
MessageImportDriverExceptionCode.SYNC_CURSOR_ERROR,
{ cause: options?.cause },
);
@@ -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,
@@ -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,
@@ -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 {}
@@ -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,
);
});
});
});
@@ -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<void> {
this.logger.log(
`MessageChannelId: ${messageChannelId} - Clearing all sync cursors`,
);
const messageChannelRepository =
await this.twentyORMManager.getRepository<MessageChannelWorkspaceEntity>(
'messageChannel',
);
const messageFolderRepository =
await this.twentyORMManager.getRepository<MessageFolderWorkspaceEntity>(
'messageFolder',
);
await messageChannelRepository.update(
{ id: messageChannelId },
{ syncCursor: '' },
transactionManager,
);
await messageFolderRepository.update(
{ messageChannelId },
{ syncCursor: '' },
transactionManager,
);
this.logger.log(
`MessageChannelId: ${messageChannelId} - Cleared all sync cursors`,
);
}
}
@@ -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<MessageChannelWorkspaceEntity, 'id'>,
workspaceId: string,
): Promise<void> {
await this.messageChannelSyncStatusService.markAsFailed(
[messageChannel.id],
workspaceId,
MessageChannelSyncStatus.FAILED_UNKNOWN,
);
await this.messagingClearCursorsService.clearAllMessageChannelCursors(
messageChannel.id,
);
}
private async handleTemporaryException(
syncStep: MessageImportSyncStep,
messageChannel: Pick<