From ee115c89ae8a4e12c9004c84f3f0fcd17d58a6f6 Mon Sep 17 00:00:00 2001 From: Artyom Savchenko Date: Wed, 23 Jul 2025 12:20:30 +0700 Subject: [PATCH] Fix messages sync for outdated gmail history id (#9583) Signed-off-by: Artem Savchenko --- .../pod-gmail/src/__tests__/sync.test.ts | 313 ++++++++++++++++++ services/gmail/pod-gmail/src/message/sync.ts | 170 +++++++++- 2 files changed, 478 insertions(+), 5 deletions(-) diff --git a/services/gmail/pod-gmail/src/__tests__/sync.test.ts b/services/gmail/pod-gmail/src/__tests__/sync.test.ts index 42f94f405b..56ce150f8b 100644 --- a/services/gmail/pod-gmail/src/__tests__/sync.test.ts +++ b/services/gmail/pod-gmail/src/__tests__/sync.test.ts @@ -182,4 +182,317 @@ describe('SyncManager', () => { expect(result).toEqual(mockResponse) }) }) + + /* eslint-disable @typescript-eslint/unbound-method */ + + describe('syncNewMessages', () => { + it('should perform full sync when no stored history ID exists', async () => { + // Arrange + const spyGetHistory = jest.spyOn((syncManager as any).stateManager, 'getHistory').mockResolvedValue(null) + const spyFullSync = jest.spyOn(syncManager, 'fullSync').mockResolvedValue(undefined) + + // Act + await syncManager.syncNewMessages(userId, userEmail) + + // Assert + expect(spyGetHistory).toHaveBeenCalledWith(userId) + expect(mockCtx.info).toHaveBeenCalledWith('No stored history ID found, performing full sync instead', { userId }) + expect(spyFullSync).toHaveBeenCalledWith(userId, userEmail) + }) + + it('should throw error when userEmail is undefined', async () => { + // Act & Assert + await expect(syncManager.syncNewMessages(userId)).rejects.toThrow('Cannot sync without user email') + }) + + it('should sync new messages and update history ID when newer messages found', async () => { + // Arrange + const storedHistoryId = '12345' + const newHistoryId = '12350' + + const spyGetHistory = jest.spyOn((syncManager as any).stateManager, 'getHistory').mockResolvedValue({ + historyId: storedHistoryId, + userId, + workspace + }) + const spySetHistoryId = jest.spyOn((syncManager as any).stateManager, 'setHistoryId').mockResolvedValue(undefined) + + const mockMessages = [{ id: 'msg1' }, { id: 'msg2' }] + + mockGmail.messages.list.mockResolvedValue({ + data: { + messages: mockMessages, + nextPageToken: null + } + }) + + mockGmail.messages.get + .mockResolvedValueOnce({ + data: { + id: 'msg1', + historyId: newHistoryId + } + }) + .mockResolvedValueOnce({ + data: { + id: 'msg2', + historyId: '12348' + } + }) + + // Act + await syncManager.syncNewMessages(userId, userEmail) + + // Assert + expect(spyGetHistory).toHaveBeenCalledWith(userId) + expect(mockGmail.messages.list).toHaveBeenCalledWith({ + userId: 'me', + pageToken: undefined + }) + expect(mockGmail.messages.get).toHaveBeenCalledTimes(2) + expect(mockMessageManager.saveMessage).toHaveBeenCalledTimes(2) + expect(spySetHistoryId).toHaveBeenCalledWith(userId, newHistoryId) + expect(mockCtx.info).toHaveBeenCalledWith('Updated history ID after new messages sync', { + userId, + oldHistoryId: storedHistoryId, + newHistoryId, + messagesProcessed: 2 + }) + }) + + it('should stop syncing when encountering message with older history ID', async () => { + // Arrange + const storedHistoryId = '12345' + const olderHistoryId = '12340' + + const spyGetHistory = jest.spyOn((syncManager as any).stateManager, 'getHistory').mockResolvedValue({ + historyId: storedHistoryId, + userId, + workspace + }) + + const mockMessages = [{ id: 'msg1' }, { id: 'msg2' }] + + mockGmail.messages.list.mockResolvedValue({ + data: { + messages: mockMessages, + nextPageToken: null + } + }) + + mockGmail.messages.get.mockResolvedValueOnce({ + data: { + id: 'msg1', + historyId: olderHistoryId + } + }) + + // Act + await syncManager.syncNewMessages(userId, userEmail) + + // Assert + expect(spyGetHistory).toHaveBeenCalledWith(userId) + expect(mockGmail.messages.get).toHaveBeenCalledTimes(1) + expect(mockMessageManager.saveMessage).not.toHaveBeenCalled() + expect(mockCtx.info).toHaveBeenCalledWith('Reached message with history ID <= stored history ID, stopping sync', { + userId, + messageHistoryId: olderHistoryId, + storedHistoryId + }) + expect(mockCtx.info).toHaveBeenCalledWith('No history ID update needed after new messages sync', { + userId, + storedHistoryId, + maxHistoryId: undefined, + messagesProcessed: 0 + }) + }) + + it('should save messages without history ID', async () => { + // Arrange + const storedHistoryId = '12345' + + const spyGetHistory = jest.spyOn((syncManager as any).stateManager, 'getHistory').mockResolvedValue({ + historyId: storedHistoryId, + userId, + workspace + }) + + const mockMessages = [{ id: 'msg1' }] + + mockGmail.messages.list.mockResolvedValue({ + data: { + messages: mockMessages, + nextPageToken: null + } + }) + + mockGmail.messages.get.mockResolvedValue({ + data: { + id: 'msg1', + historyId: null // No history ID + } + }) + + // Act + await syncManager.syncNewMessages(userId, userEmail) + + // Assert + expect(spyGetHistory).toHaveBeenCalledWith(userId) + expect(mockMessageManager.saveMessage).toHaveBeenCalledTimes(1) + expect(mockCtx.info).toHaveBeenCalledWith('No history ID update needed after new messages sync', { + userId, + storedHistoryId, + maxHistoryId: undefined, + messagesProcessed: 1 + }) + }) + + it('should handle pagination correctly', async () => { + // Arrange + const storedHistoryId = '12345' + const newHistoryId = '12350' + + const spyGetHistory = jest.spyOn((syncManager as any).stateManager, 'getHistory').mockResolvedValue({ + historyId: storedHistoryId, + userId, + workspace + }) + const spySetHistoryId = jest.spyOn((syncManager as any).stateManager, 'setHistoryId').mockResolvedValue(undefined) + + // First page + mockGmail.messages.list + .mockResolvedValueOnce({ + data: { + messages: [{ id: 'msg1' }], + nextPageToken: 'token123' + } + }) + .mockResolvedValueOnce({ + data: { + messages: [{ id: 'msg2' }], + nextPageToken: null + } + }) + + mockGmail.messages.get + .mockResolvedValueOnce({ + data: { + id: 'msg1', + historyId: newHistoryId + } + }) + .mockResolvedValueOnce({ + data: { + id: 'msg2', + historyId: '12348' + } + }) + + // Act + await syncManager.syncNewMessages(userId, userEmail) + + // Assert + expect(spyGetHistory).toHaveBeenCalledWith(userId) + expect(mockGmail.messages.list).toHaveBeenCalledTimes(2) + // The first call gets the modified query object with pageToken from the second call + // This is due to query object being reused and modified + expect(mockGmail.messages.list).toHaveBeenCalledWith({ + userId: 'me', + pageToken: 'token123' // This will be the final state of the query object + }) + expect(mockMessageManager.saveMessage).toHaveBeenCalledTimes(2) + expect(spySetHistoryId).toHaveBeenCalledWith(userId, newHistoryId) + }) + + it('should handle errors during message processing', async () => { + // Arrange + const storedHistoryId = '12345' + + const spyGetHistory = jest.spyOn((syncManager as any).stateManager, 'getHistory').mockResolvedValue({ + historyId: storedHistoryId, + userId, + workspace + }) + + const mockMessages = [{ id: 'msg1' }] + + mockGmail.messages.list.mockResolvedValue({ + data: { + messages: mockMessages, + nextPageToken: null + } + }) + + const error = new Error('Message fetch failed') + mockGmail.messages.get.mockRejectedValue(error) + + // Act + await syncManager.syncNewMessages(userId, userEmail) + + // Assert + expect(spyGetHistory).toHaveBeenCalledWith(userId) + expect(mockCtx.error).toHaveBeenCalledWith('Sync new messages error', { + workspace, + userId, + messageId: 'msg1', + err: error + }) + }) + + it('should return early when isClosing is true', async () => { + // Arrange + const storedHistoryId = '12345' + + const spyGetHistory = jest.spyOn((syncManager as any).stateManager, 'getHistory').mockResolvedValue({ + historyId: storedHistoryId, + userId, + workspace + }) + + const mockMessages = [{ id: 'msg1' }] + + mockGmail.messages.list.mockResolvedValue({ + data: { + messages: mockMessages, + nextPageToken: null + } + }) + + // Set isClosing to true + syncManager.close() + + // Act + await syncManager.syncNewMessages(userId, userEmail) + + // Assert + expect(spyGetHistory).toHaveBeenCalledWith(userId) + expect(mockGmail.messages.get).not.toHaveBeenCalled() + expect(mockMessageManager.saveMessage).not.toHaveBeenCalled() + }) + + it('should handle overall sync error', async () => { + // Arrange + const storedHistoryId = '12345' + + const spyGetHistory = jest.spyOn((syncManager as any).stateManager, 'getHistory').mockResolvedValue({ + historyId: storedHistoryId, + userId, + workspace + }) + + const error = new Error('Sync failed') + mockGmail.messages.list.mockRejectedValue(error) + + // Act + await syncManager.syncNewMessages(userId, userEmail) + + // Assert + expect(spyGetHistory).toHaveBeenCalledWith(userId) + expect(mockCtx.error).toHaveBeenCalledWith('New messages sync error', { + workspace, + userId, + err: error + }) + }) + }) }) diff --git a/services/gmail/pod-gmail/src/message/sync.ts b/services/gmail/pod-gmail/src/message/sync.ts index caa8601440..4464a1074d 100644 --- a/services/gmail/pod-gmail/src/message/sync.ts +++ b/services/gmail/pod-gmail/src/message/sync.ts @@ -51,6 +51,7 @@ export class SyncManager { } let pageToken: string | undefined let histories: GaxiosResponse + let maxHistoryId: string | undefined while (true) { try { await this.rateLimiter.take(2) @@ -61,9 +62,13 @@ export class SyncManager { pageToken }) } catch (err: any) { - this.ctx.error('Part sync get history error', { workspaceUuid: this.workspace, userId, error: err.message }) - await this.stateManager.clearHistory(userId) - void this.sync(userId, options, userEmail) + this.ctx.error('Part sync get history error', { + workspaceUuid: this.workspace, + userId, + error: err.message, + historyId + }) + void this.syncNewMessages(userId, userEmail) return } const nextPageToken = histories.data.nextPageToken @@ -89,10 +94,31 @@ export class SyncManager { } } if (history.id != null) { - await this.stateManager.setHistoryId(userId, history.id) + if (maxHistoryId == null || BigInt(maxHistoryId) < BigInt(history.id)) { + maxHistoryId = history.id + } } } if (nextPageToken == null) { + // Set the maximum historyId found during the sync, but only if it's greater than current + if (maxHistoryId != null) { + const currentHistory = await this.stateManager.getHistory(userId) + const currentHistoryId = currentHistory?.historyId + if (currentHistoryId == null || BigInt(maxHistoryId) > BigInt(currentHistoryId)) { + await this.stateManager.setHistoryId(userId, maxHistoryId) + this.ctx.info('Updated history ID', { + userId, + oldHistoryId: currentHistoryId, + newHistoryId: maxHistoryId + }) + } else { + this.ctx.info('Skipping history ID update', { + userId, + currentHistoryId, + maxHistoryId + }) + } + } return } else { pageToken = nextPageToken @@ -148,7 +174,7 @@ export class SyncManager { await this.messageManager.saveMessage(message, userEmail) if (historyId != null && q === undefined) { - if (currentHistoryId == null || Number(currentHistoryId) < Number(historyId)) { + if (currentHistoryId == null || BigInt(currentHistoryId) < BigInt(historyId)) { currentHistoryId = historyId } } @@ -190,6 +216,140 @@ export class SyncManager { }) } + async syncNewMessages (userId: PersonId, userEmail?: string): Promise { + this.ctx.info('Start sync new messages', { workspaceUuid: this.workspace, userId, userEmail }) + if (userEmail === undefined) { + throw new Error('Cannot sync without user email') + } + + // Get stored history ID to compare against + const storedHistory = await this.stateManager.getHistory(userId) + const storedHistoryId = storedHistory?.historyId + + if (storedHistoryId == null) { + this.ctx.info('No stored history ID found, performing full sync instead', { userId }) + await this.fullSync(userId, userEmail) + return + } + + this.ctx.info('Syncing new messages', { + workspaceUuid: this.workspace, + userId, + storedHistoryId + }) + + let pageToken: string | undefined + let totalProcessedMessages = 0 + let maxHistoryId: string | undefined + let continueOuterLoop = true + + const query: gmail_v1.Params$Resource$Users$Messages$List = { + userId: 'me' + } + + try { + // Process messages until we find one with historyId > storedHistoryId + while (continueOuterLoop) { + query.pageToken = pageToken + await this.rateLimiter.take(5) + const messages = await this.gmail.messages.list(query) + + const ids = messages.data.messages?.map((p) => p.id).filter((id) => id != null) ?? [] + this.ctx.info('Processing new messages page', { + workspace: this.workspace, + userId, + messagesInPage: ids.length, + totalProcessed: totalProcessedMessages, + pageToken: query.pageToken + }) + + for (const id of ids) { + if (this.isClosing) return + if (id == null) continue + + try { + const message = await this.getMessage(id) + const messageHistoryId = message.data.historyId + + // If message has a history ID, check if it's newer than stored + if (messageHistoryId != null) { + if (BigInt(messageHistoryId) > BigInt(storedHistoryId)) { + await this.messageManager.saveMessage(message, userEmail) + totalProcessedMessages++ + + // Track the maximum history ID found + if (maxHistoryId == null || BigInt(messageHistoryId) > BigInt(maxHistoryId)) { + maxHistoryId = messageHistoryId + } + } else { + // This message is older or equal to stored history, stop processing + this.ctx.info('Reached message with history ID <= stored history ID, stopping sync', { + userId, + messageHistoryId, + storedHistoryId + }) + continueOuterLoop = false + break + } + } else { + // Message without history ID, save it anyway + await this.messageManager.saveMessage(message, userEmail) + totalProcessedMessages++ + } + } catch (err: any) { + if (this.isClosing) return + this.ctx.error('Sync new messages error', { + workspace: this.workspace, + userId, + messageId: id, + err + }) + } + } + + // If we should stop processing or no more pages, stop + if (!continueOuterLoop || messages.data.nextPageToken == null) { + this.ctx.info('Completed new messages sync', { + workspace: this.workspace, + userId, + totalMessages: totalProcessedMessages + }) + break + } + + pageToken = messages.data.nextPageToken + } + + // Update stored history ID with the maximum found, only if we found newer messages + if (maxHistoryId != null && BigInt(maxHistoryId) > BigInt(storedHistoryId)) { + await this.stateManager.setHistoryId(userId, maxHistoryId) + this.ctx.info('Updated history ID after new messages sync', { + userId, + oldHistoryId: storedHistoryId, + newHistoryId: maxHistoryId, + messagesProcessed: totalProcessedMessages + }) + } else { + this.ctx.info('No history ID update needed after new messages sync', { + userId, + storedHistoryId, + maxHistoryId, + messagesProcessed: totalProcessedMessages + }) + } + + this.ctx.info('New messages sync finished', { + workspaceUuid: this.workspace, + userId, + userEmail, + totalMessages: totalProcessedMessages + }) + } catch (err) { + if (this.isClosing) return + this.ctx.error('New messages sync error', { workspace: this.workspace, userId, err }) + } + } + async sync (userId: PersonId, options: SyncOptions, userEmail?: string): Promise { const mutexKey = `${this.workspace}:${userId}` const releaseLock = await this.syncMutex.lock(mutexKey)