Other messaging improvements (#13750)

This commit is contained in:
Charles Bochet
2025-08-08 00:19:22 +02:00
committed by GitHub
parent 5c7829c19a
commit 525c8bfe3c
4 changed files with 118 additions and 93 deletions
@@ -1,4 +1,4 @@
import { Injectable } from '@nestjs/common';
import { Injectable, Logger } from '@nestjs/common';
import { IsNull } from 'typeorm';
@@ -10,6 +10,7 @@ import { deleteUsingPagination } from 'src/modules/messaging/message-cleaner/uti
@Injectable()
export class MessagingMessageCleanerService {
private readonly logger = new Logger(MessagingMessageCleanerService.name);
constructor(private readonly twentyORMManager: TwentyORMManager) {}
public async cleanWorkspaceThreads(workspaceId: string) {
@@ -57,6 +58,9 @@ export class MessagingMessageCleanerService {
workspaceId: string,
transactionManager?: WorkspaceEntityManager,
) => {
this.logger.log(
`WorkspaceId: ${workspaceId} Deleting ${ids.length} messages from message cleaner`,
);
await messageRepository.delete(ids, transactionManager);
},
transactionManager,
@@ -91,14 +91,15 @@ export class MessagingMessageListFetchService {
(messageList) => messageList.messageExternalIdsToDelete,
);
const isFullSync = messageLists.every(
(messageList) => !isNonEmptyString(messageList.previousSyncCursor),
);
const isFullSync =
messageLists.every(
(messageList) => !isNonEmptyString(messageList.previousSyncCursor),
) && isNonEmptyString(messageChannel.syncCursor);
let totalMessagesToImportCount = 0;
this.logger.log(
`messageChannelId: ${messageChannel.id} Is full sync: ${isFullSync} and message lists length: ${messageLists.length}`,
`messageChannelId: ${messageChannel.id} Is full sync: ${isFullSync} and toImportCount: ${messageExternalIds.length}, toDeleteCount: ${messageExternalIdsToDelete.length}`,
);
const messageChannelMessageAssociationRepository =
@@ -158,83 +159,12 @@ export class MessagingMessageListFetchService {
);
}
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 fullSyncMessageChannelMessageAssociationsToDelete = isFullSync
? await this.computeFullSyncMessageChannelMessageAssociationsToDelete(
messageChannel,
messageExternalIds,
)
: [];
const allMessageExternalIdsToDelete = [
...messageExternalIdsToDelete,
@@ -308,4 +238,91 @@ export class MessagingMessageListFetchService {
);
}
}
private async computeFullSyncMessageChannelMessageAssociationsToDelete(
messageChannel: Pick<MessageChannelWorkspaceEntity, 'id'>,
messageExternalIds: string[],
) {
const messageChannelMessageAssociationRepository =
await this.twentyORMManager.getRepository<MessageChannelMessageAssociationWorkspaceEntity>(
'messageChannelMessageAssociation',
);
const fullSyncMessageChannelMessageAssociationsToDelete = [];
const firstMessageChannelMessageAssociation =
await messageChannelMessageAssociationRepository.findOne({
where: {
messageChannelId: messageChannel.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: messageChannel.id,
id: MoreThanOrEqual(
nextFirstBatchMessageChannelMessageAssociationId,
),
},
order: {
id: 'ASC',
},
take: 200,
});
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,
);
if (existingMessageChannelMessageAssociations.length < 200) {
nextFirstBatchMessageChannelMessageAssociationId = undefined;
break;
}
nextFirstBatchMessageChannelMessageAssociationId =
existingMessageChannelMessageAssociations[
existingMessageChannelMessageAssociations.length - 1
].id;
batchIndex++;
}
return fullSyncMessageChannelMessageAssociationsToDelete;
}
}
@@ -123,12 +123,14 @@ export class MessagingMessagesImportService {
blocklist.map((blocklistItem) => blocklistItem.handle),
);
await this.saveMessagesAndEnqueueContactCreationService.saveMessagesAndEnqueueContactCreation(
messagesToSave,
messageChannel,
connectedAccountWithFreshTokens,
workspaceId,
);
if (messagesToSave.length > 0) {
await this.saveMessagesAndEnqueueContactCreationService.saveMessagesAndEnqueueContactCreation(
messagesToSave,
messageChannel,
connectedAccountWithFreshTokens,
workspaceId,
);
}
if (
messageIdsToFetch.length < MESSAGING_GMAIL_USERS_MESSAGES_GET_BATCH_SIZE
@@ -10,17 +10,19 @@ export const filterEmails = (
messages: MessageWithParticipants[],
blocklist: string[],
) => {
const messagesWithoutIcsAttachments = filterOutIcsAttachments(messages);
const messagesWithoutBlocklisted = filterOutBlocklistedMessages(
[primaryHandle, ...handleAliases],
filterOutIcsAttachments(messages),
messagesWithoutIcsAttachments,
blocklist,
);
if (isWorkEmail(primaryHandle)) {
return filterOutInternals(primaryHandle, messagesWithoutBlocklisted);
}
const messagesWithoutInternals = isWorkEmail(primaryHandle)
? filterOutInternals(primaryHandle, messagesWithoutBlocklisted)
: messagesWithoutBlocklisted;
return messagesWithoutBlocklisted;
return messagesWithoutInternals;
};
const filterOutBlocklistedMessages = (