Improve messaging batch (#13747)

When users have more than 100k messages, this part is also starting to
slow down
This commit is contained in:
Charles Bochet
2025-08-07 22:45:19 +02:00
committed by GitHub
parent b44b46109c
commit 5c7829c19a
2 changed files with 170 additions and 52 deletions
@@ -4,13 +4,11 @@ import { Process } from 'src/engine/core-modules/message-queue/decorators/proces
import { Processor } from 'src/engine/core-modules/message-queue/decorators/processor.decorator';
import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants';
import { TwentyORMManager } from 'src/engine/twenty-orm/twenty-orm.manager';
import { ConnectedAccountRefreshTokensService } from 'src/modules/connected-account/refresh-tokens-manager/services/connected-account-refresh-tokens.service';
import { isThrottled } from 'src/modules/connected-account/utils/is-throttled';
import {
MessageChannelSyncStage,
MessageChannelWorkspaceEntity,
} from 'src/modules/messaging/common/standard-objects/message-channel.workspace-entity';
import { MessagingAccountAuthenticationService } from 'src/modules/messaging/message-import-manager/services/messaging-account-authentication.service';
import {
MessageImportExceptionHandlerService,
MessageImportSyncStep,
@@ -32,9 +30,7 @@ export class MessagingMessageListFetchJob {
private readonly messagingMessageListFetchService: MessagingMessageListFetchService,
private readonly messagingMonitoringService: MessagingMonitoringService,
private readonly twentyORMManager: TwentyORMManager,
private readonly connectedAccountRefreshTokensService: ConnectedAccountRefreshTokensService,
private readonly messageImportErrorHandlerService: MessageImportExceptionHandlerService,
private readonly messagingAccountAuthenticationService: MessagingAccountAuthenticationService,
) {}
@Process(MessagingMessageListFetchJob.name)
@@ -1,6 +1,9 @@
import { Injectable } from '@nestjs/common';
import { Injectable, Logger } from '@nestjs/common';
import { In } from 'typeorm';
import { isNonEmptyString } from '@sniptt/guards';
import chunk from 'lodash.chunk';
import { isDefined } from 'twenty-shared/utils';
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';
@@ -26,6 +29,7 @@ const ONE_WEEK_IN_MILLISECONDS = 7 * 24 * 60 * 60 * 1000;
@Injectable()
export class MessagingMessageListFetchService {
private readonly logger = new Logger(MessagingMessageListFetchService.name);
constructor(
@InjectCacheStorage(CacheStorageNamespace.ModuleMessaging)
private readonly cacheStorage: CacheStorageService,
@@ -48,6 +52,10 @@ export class MessagingMessageListFetchService {
[messageChannel.id],
);
this.logger.log(
`messageChannelId: ${messageChannel.id} Processing message list fetch`,
);
const { accessToken, refreshToken } =
await this.messagingAccountAuthenticationService.validateAndRefreshConnectedAccountAuthentication(
{
@@ -75,33 +83,40 @@ export class MessagingMessageListFetchService {
`messages-to-import:${workspaceId}:${messageChannel.id}`,
);
const totalMessageCount = messageLists.reduce(
(acc, messageList) => acc + messageList.messageExternalIds.length,
0,
const messageExternalIds = messageLists.flatMap(
(messageList) => messageList.messageExternalIds,
);
for (const messageList of messageLists) {
if (messageList.messageExternalIds.length === 0) {
continue;
}
const messageExternalIdsToDelete = messageLists.flatMap(
(messageList) => messageList.messageExternalIdsToDelete,
);
const {
messageExternalIds,
nextSyncCursor,
folderId,
messageExternalIdsToDelete,
previousSyncCursor,
} = messageList;
const isFullSync = messageLists.every(
(messageList) => !isNonEmptyString(messageList.previousSyncCursor),
);
const messageChannelMessageAssociationRepository =
await this.twentyORMManager.getRepository<MessageChannelMessageAssociationWorkspaceEntity>(
'messageChannelMessageAssociation',
);
let totalMessagesToImportCount = 0;
this.logger.log(
`messageChannelId: ${messageChannel.id} Is full sync: ${isFullSync} and message lists length: ${messageLists.length}`,
);
const messageChannelMessageAssociationRepository =
await this.twentyORMManager.getRepository<MessageChannelMessageAssociationWorkspaceEntity>(
'messageChannelMessageAssociation',
);
const messageExternalIdsChunks = chunk(messageExternalIds, 200);
for (const [
index,
messageExternalIdsChunk,
] of messageExternalIdsChunks.entries()) {
const existingMessageChannelMessageAssociations =
await messageChannelMessageAssociationRepository.find({
where: {
messageChannelId: messageChannel.id,
messageExternalId: In(messageExternalIdsChunk),
},
});
@@ -111,46 +126,30 @@ export class MessagingMessageListFetchService {
messageChannelMessageAssociation.messageExternalId,
);
const messageExternalIdsToImport = messageExternalIds.filter(
const messageExternalIdsToImport = messageExternalIdsChunk.filter(
(messageExternalId) =>
!existingMessageChannelMessageAssociationsExternalIds.includes(
messageExternalId,
),
);
const isFullSync = !previousSyncCursor;
const additionalMessageExternalIdsToDelete = isFullSync
? existingMessageChannelMessageAssociationsExternalIds.filter(
(existingMessageCMAExternalId) =>
existingMessageCMAExternalId &&
!messageExternalIds.includes(existingMessageCMAExternalId),
)
: [];
const allMessageExternalIdsToDelete = [
...messageExternalIdsToDelete,
...additionalMessageExternalIdsToDelete,
];
if (allMessageExternalIdsToDelete.length) {
await messageChannelMessageAssociationRepository.delete({
messageChannelId: messageChannelWithFreshTokens.id,
messageExternalId: In(allMessageExternalIdsToDelete),
});
await this.messagingMessageCleanerService.cleanWorkspaceThreads(
workspaceId,
);
}
if (messageExternalIdsToImport.length) {
this.logger.log(
`messageChannelId: ${messageChannel.id} Adding ${messageExternalIdsToImport.length} message external ids to import in batch ${index + 1}`,
);
totalMessagesToImportCount += messageExternalIdsToImport.length;
await this.cacheStorage.setAdd(
`messages-to-import:${workspaceId}:${messageChannelWithFreshTokens.id}`,
messageExternalIdsToImport,
ONE_WEEK_IN_MILLISECONDS,
);
}
}
for (const messageList of messageLists) {
const { nextSyncCursor, folderId } = messageList;
await this.messagingCursorService.updateCursor(
messageChannelWithFreshTokens,
@@ -159,12 +158,135 @@ export class MessagingMessageListFetchService {
);
}
if (totalMessageCount === 0) {
const fullSyncMessageChannelMessageAssociationsToDelete = [];
if (isFullSync) {
const firstMessageChannelMessageAssociation =
await messageChannelMessageAssociationRepository.findOne({
where: {
messageChannelId: messageChannelWithFreshTokens.id,
},
order: {
id: 'ASC',
},
});
if (!isDefined(firstMessageChannelMessageAssociation)) {
this.logger.log(
`messageChannelId: ${messageChannel.id} Full sync: No message channel message associations found`,
);
return;
}
this.logger.log(
`messageChannelId: ${messageChannel.id} Full sync: First message channel message association id: ${firstMessageChannelMessageAssociation.id}`,
);
let nextFirstBatchMessageChannelMessageAssociationId:
| string
| undefined = firstMessageChannelMessageAssociation.id;
let batchIndex = 0;
while (isDefined(nextFirstBatchMessageChannelMessageAssociationId)) {
const existingMessageChannelMessageAssociations =
await messageChannelMessageAssociationRepository.find({
where: {
messageChannelId: messageChannelWithFreshTokens.id,
id: MoreThanOrEqual(
nextFirstBatchMessageChannelMessageAssociationId,
),
},
order: {
id: 'ASC',
},
take: 200,
});
if (existingMessageChannelMessageAssociations.length < 200) {
nextFirstBatchMessageChannelMessageAssociationId = undefined;
break;
}
nextFirstBatchMessageChannelMessageAssociationId =
existingMessageChannelMessageAssociations[
existingMessageChannelMessageAssociations.length - 1
].id;
batchIndex++;
const messageChannelMessageAssociationsToDelete =
existingMessageChannelMessageAssociations.filter(
(existingMessageChannelMessageAssociation) =>
isDefined(
existingMessageChannelMessageAssociation.messageExternalId,
) &&
!messageExternalIds.includes(
existingMessageChannelMessageAssociation.messageExternalId,
),
);
this.logger.log(
`messageChannelId: ${messageChannel.id} Full sync: Message channel message associations to delete in batch ${batchIndex}: ${messageChannelMessageAssociationsToDelete.length}`,
);
fullSyncMessageChannelMessageAssociationsToDelete.push(
...messageChannelMessageAssociationsToDelete,
);
}
}
const allMessageExternalIdsToDelete = [
...messageExternalIdsToDelete,
...fullSyncMessageChannelMessageAssociationsToDelete.map(
(messageChannelMessageAssociation) =>
messageChannelMessageAssociation.messageExternalId,
),
];
if (allMessageExternalIdsToDelete.length) {
this.logger.log(
`messageChannelId: ${messageChannel.id} Deleting ${allMessageExternalIdsToDelete.length} message channel message associations`,
);
const toDeleteChunks = chunk(allMessageExternalIdsToDelete, 200);
for (const [index, toDeleteChunk] of toDeleteChunks.entries()) {
await messageChannelMessageAssociationRepository.delete({
messageChannelId: messageChannelWithFreshTokens.id,
messageExternalId: In(toDeleteChunk),
});
this.logger.log(
`messageChannelId: ${messageChannel.id} Deleted ${toDeleteChunk.length} message channel message associations in batch ${index + 1}`,
);
}
}
this.logger.log(
`messageChannelId: ${messageChannel.id} launching workspace thread cleanup`,
);
await this.messagingMessageCleanerService.cleanWorkspaceThreads(
workspaceId,
);
this.logger.log(
`messageChannelId: ${messageChannel.id} Total messages to import count: ${totalMessagesToImportCount}`,
);
if (totalMessagesToImportCount === 0) {
await this.messageChannelSyncStatusService.markAsCompletedAndScheduleMessageListFetch(
[messageChannelWithFreshTokens.id],
);
return;
}
this.logger.log(
`messageChannelId: ${messageChannel.id} Scheduling direct messages import`,
);
await this.messageChannelSyncStatusService.scheduleMessagesImport([
messageChannelWithFreshTokens.id,
]);