messaging: gmail folder backfill (#21753)

demo


https://github.com/user-attachments/assets/a157cee1-a8fa-4050-af1b-c31a83fb75da

/closes #17095


<!-- This is an auto-generated description by cubic. -->
<a
href="https://cubic.dev/pr/twentyhq/twenty/pull/21753?utm_source=github"
target="_blank" rel="noopener noreferrer"
data-no-image-dialog="true"><picture><source
media="(prefers-color-scheme: dark)"
srcset="https://www.cubic.dev/buttons/review-in-cubic-dark.svg"><source
media="(prefers-color-scheme: light)"
srcset="https://www.cubic.dev/buttons/review-in-cubic-light.svg"><img
alt="Review in cubic"
src="https://www.cubic.dev/buttons/review-in-cubic-dark.svg"></picture></a>
<!-- End of auto-generated description by cubic. -->
This commit is contained in:
neo773
2026-06-19 05:29:35 +05:30
committed by GitHub
parent d88eb6c16b
commit 616d58bc7e
21 changed files with 223 additions and 56 deletions
@@ -0,0 +1,25 @@
import { QueryRunner } from 'typeorm';
import { RegisteredInstanceCommand } from 'src/engine/core-modules/upgrade/decorators/registered-instance-command.decorator';
import { FastInstanceCommand } from 'src/engine/core-modules/upgrade/interfaces/fast-instance-command.interface';
@RegisteredInstanceCommand('2.15.0', 1781714499016)
export class AddFolderImportToMessageFolderPendingSyncActionFastInstanceCommand implements FastInstanceCommand {
public async up(queryRunner: QueryRunner): Promise<void> {
await queryRunner.query('ALTER TYPE "core"."messageFolder_pendingsyncaction_enum" RENAME TO "messageFolder_pendingsyncaction_enum_old"');
await queryRunner.query('CREATE TYPE "core"."messageFolder_pendingsyncaction_enum" AS ENUM(\'FOLDER_DELETION\', \'FOLDER_IMPORT\', \'NONE\')');
await queryRunner.query('ALTER TABLE "core"."messageFolder" ALTER COLUMN "pendingSyncAction" DROP DEFAULT');
await queryRunner.query('ALTER TABLE "core"."messageFolder" ALTER COLUMN "pendingSyncAction" TYPE "core"."messageFolder_pendingsyncaction_enum" USING "pendingSyncAction"::"text"::"core"."messageFolder_pendingsyncaction_enum"');
await queryRunner.query('ALTER TABLE "core"."messageFolder" ALTER COLUMN "pendingSyncAction" SET DEFAULT \'NONE\'');
await queryRunner.query('DROP TYPE "core"."messageFolder_pendingsyncaction_enum_old"');
}
public async down(queryRunner: QueryRunner): Promise<void> {
await queryRunner.query('CREATE TYPE "core"."messageFolder_pendingsyncaction_enum_old" AS ENUM(\'FOLDER_DELETION\', \'NONE\')');
await queryRunner.query('ALTER TABLE "core"."messageFolder" ALTER COLUMN "pendingSyncAction" DROP DEFAULT');
await queryRunner.query('ALTER TABLE "core"."messageFolder" ALTER COLUMN "pendingSyncAction" TYPE "core"."messageFolder_pendingsyncaction_enum_old" USING (CASE WHEN "pendingSyncAction"::"text" = \'FOLDER_IMPORT\' THEN \'NONE\' ELSE "pendingSyncAction"::"text" END)::"core"."messageFolder_pendingsyncaction_enum_old"');
await queryRunner.query('ALTER TABLE "core"."messageFolder" ALTER COLUMN "pendingSyncAction" SET DEFAULT \'NONE\'');
await queryRunner.query('DROP TYPE "core"."messageFolder_pendingsyncaction_enum"');
await queryRunner.query('ALTER TYPE "core"."messageFolder_pendingsyncaction_enum_old" RENAME TO "messageFolder_pendingsyncaction_enum"');
}
}
@@ -73,6 +73,7 @@ import { AddLogicFunctionExecutionModeFastInstanceCommand } from 'src/database/c
import { EncryptNonSecretApplicationVariableSlowInstanceCommand } from 'src/database/commands/upgrade-version-command/2-9/2-9-instance-command-slow-1798400000000-encrypt-non-secret-application-variable';
import { MigrateAiModelPreferencesSlowInstanceCommand } from 'src/database/commands/upgrade-version-command/2-9/2-9-instance-command-slow-1799000010000-migrate-ai-model-preferences';
import { AddHasPaymentMethodToBillingCustomerFastInstanceCommand } from 'src/database/commands/upgrade-version-command/2-15/2-15-instance-command-fast-1781280240009-add-has-payment-method-to-billing-customer';
import { AddFolderImportToMessageFolderPendingSyncActionFastInstanceCommand } from './2-15/2-15-instance-command-fast-1781714499016-add-folder-import-to-message-folder-pending-sync-action';
export const INSTANCE_COMMANDS = [
AddViewFieldGroupIdIndexOnViewFieldFastInstanceCommand,
@@ -148,4 +149,5 @@ export const INSTANCE_COMMANDS = [
SetTableWidgetViewsVisibilityToWorkspaceSlowInstanceCommand,
AddIsSystemSideEffectFastInstanceCommand,
BackfillConnectionSecuritySlowInstanceCommand,
AddFolderImportToMessageFolderPendingSyncActionFastInstanceCommand,
];
@@ -1,22 +1,15 @@
import { Field, InputType } from '@nestjs/graphql';
import { Type } from 'class-transformer';
import {
IsBoolean,
IsNotEmpty,
IsOptional,
IsUUID,
ValidateNested,
} from 'class-validator';
import { IsBoolean, IsNotEmpty, IsUUID, ValidateNested } from 'class-validator';
import { UUIDScalarType } from 'src/engine/api/graphql/workspace-schema-builder/graphql-types/scalars';
@InputType()
export class UpdateMessageFolderInputUpdates {
@IsOptional()
@IsBoolean()
@Field({ nullable: true })
isSynced?: boolean;
@Field()
isSynced: boolean;
}
@InputType()
@@ -3,7 +3,10 @@ import { InjectRepository } from '@nestjs/typeorm';
import { In, Repository } from 'typeorm';
import { MessageFolderPendingSyncAction } from 'twenty-shared/types';
import {
ConnectedAccountProvider,
MessageFolderPendingSyncAction,
} from 'twenty-shared/types';
import { ConnectedAccountMetadataService } from 'src/engine/metadata-modules/connected-account/connected-account-metadata.service';
import { MessageFolderDTO } from 'src/engine/metadata-modules/message-folder/dtos/message-folder.dto';
@@ -183,7 +186,7 @@ export class MessageFolderMetadataService {
return this.repository.findOneOrFail({ where: { id, workspaceId } });
}
async updateMany({
async setSyncStatus({
ids,
workspaceId,
data,
@@ -192,10 +195,53 @@ export class MessageFolderMetadataService {
workspaceId: string;
data: Partial<MessageFolderEntity>;
}): Promise<MessageFolderDTO[]> {
await this.repository.update(
{ id: In(ids), workspaceId },
data as Record<string, unknown>,
);
await this.repository.manager.transaction(async (manager) => {
if (!data.isSynced) {
await manager.update(
MessageFolderEntity,
{ id: In(ids), workspaceId },
{ isSynced: false },
);
await manager.update(
MessageFolderEntity,
{
id: In(ids),
workspaceId,
pendingSyncAction: MessageFolderPendingSyncAction.FOLDER_IMPORT,
},
{ pendingSyncAction: MessageFolderPendingSyncAction.NONE },
);
return;
}
const folderIdsToBackfill = (
await manager.find(MessageFolderEntity, {
where: {
id: In(ids),
workspaceId,
isSynced: false,
messageChannel: {
connectedAccount: { provider: ConnectedAccountProvider.GOOGLE },
},
},
})
).map((folder) => folder.id);
await manager.update(
MessageFolderEntity,
{ id: In(ids), workspaceId },
{ isSynced: true },
);
if (folderIdsToBackfill.length > 0) {
await manager.update(
MessageFolderEntity,
{ id: In(folderIdsToBackfill), workspaceId },
{ pendingSyncAction: MessageFolderPendingSyncAction.FOLDER_IMPORT },
);
}
});
return this.repository.find({ where: { id: In(ids), workspaceId } });
}
@@ -86,7 +86,7 @@ export class MessageFolderResolver {
),
);
return this.messageFolderMetadataService.updateMany({
return this.messageFolderMetadataService.setSyncStatus({
ids: input.ids,
workspaceId: workspace.id,
data: input.update,
@@ -28,7 +28,7 @@ export class GmailGetMessageListService {
private readonly gmailMessageListFetchErrorHandler: GmailMessageListFetchErrorHandler,
) {}
private async getMessageListWithoutCursor(
async getMessageListWithoutCursor(
connectedAccount: Pick<
ConnectedAccountEntity,
'provider' | 'id' | 'handle'
@@ -18,7 +18,6 @@ import { RefreshTokensManagerModule } from 'src/modules/connected-account/refres
import { MessagingCommonModule } from 'src/modules/messaging/common/messaging-common.module';
import { MessagingMessageCleanerModule } from 'src/modules/messaging/message-cleaner/messaging-message-cleaner.module';
import { MessagingFolderSyncManagerModule } from 'src/modules/messaging/message-folder-manager/messaging-folder-sync-manager.module';
import { MessagingSingleMessageImportCommand } from 'src/modules/messaging/message-import-manager/commands/messaging-single-message-import.command';
import { MessagingTriggerMessageListFetchCommand } from 'src/modules/messaging/message-import-manager/commands/messaging-trigger-message-list-fetch.command';
import { MessagingMessageListFetchCronCommand } from 'src/modules/messaging/message-import-manager/crons/commands/messaging-message-list-fetch.cron.command';
import { MessagingMessagesImportCronCommand } from 'src/modules/messaging/message-import-manager/crons/commands/messaging-messages-import.cron.command';
@@ -34,7 +33,6 @@ import { MessagingInboundEmailDriverModule } from 'src/modules/messaging/message
import { InboundEmailImportService } from 'src/modules/messaging/message-import-manager/drivers/inbound-email/services/inbound-email-import.service';
import { MessagingMicrosoftDriverModule } from 'src/modules/messaging/message-import-manager/drivers/microsoft/messaging-microsoft-driver.module';
import { MessagingSmtpDriverModule } from 'src/modules/messaging/message-import-manager/drivers/smtp/messaging-smtp-driver.module';
import { MessagingAddSingleMessageToCacheForImportJob } from 'src/modules/messaging/message-import-manager/jobs/messaging-add-single-message-to-cache-for-import.job';
import { MessagingCleanCacheJob } from 'src/modules/messaging/message-import-manager/jobs/messaging-clean-cache';
import { MessagingInboundEmailImportJob } from 'src/modules/messaging/message-import-manager/jobs/messaging-inbound-email-import.job';
import { MessagingMessageListFetchJob } from 'src/modules/messaging/message-import-manager/jobs/messaging-message-list-fetch.job';
@@ -44,6 +42,7 @@ import { MessagingRelaunchFailedMessageChannelJob } from 'src/modules/messaging/
import { MessagingCursorService } from 'src/modules/messaging/message-import-manager/services/messaging-cursor.service';
import { MessagingDeleteFolderMessagesService } from 'src/modules/messaging/message-import-manager/services/messaging-delete-folder-messages.service';
import { MessagingDeleteGroupEmailMessagesService } from 'src/modules/messaging/message-import-manager/services/messaging-delete-group-email-messages.service';
import { MessagingImportFolderMessagesService } from 'src/modules/messaging/message-import-manager/services/messaging-import-folder-messages.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';
import { MessageImportExceptionHandlerService } from 'src/modules/messaging/message-import-manager/services/messaging-import-exception-handler.service';
@@ -92,7 +91,6 @@ import { MessagingMonitoringModule } from 'src/modules/messaging/monitoring/mess
MessagingMessagesImportCronCommand,
MessagingOngoingStaleCronCommand,
MessagingRelaunchFailedMessageChannelsCronCommand,
MessagingSingleMessageImportCommand,
MessagingTriggerMessageListFetchCommand,
MessagingMessageListFetchJob,
MessagingMessagesImportJob,
@@ -102,7 +100,6 @@ import { MessagingMonitoringModule } from 'src/modules/messaging/monitoring/mess
MessagingMessagesImportCronJob,
MessagingOngoingStaleCronJob,
MessagingRelaunchFailedMessageChannelsCronJob,
MessagingAddSingleMessageToCacheForImportJob,
MessagingCleanCacheJob,
MessagingInboundEmailImportJob,
MessagingMessageService,
@@ -117,6 +114,7 @@ import { MessagingMonitoringModule } from 'src/modules/messaging/monitoring/mess
MessagingProcessFolderActionsService,
MessagingProcessGroupEmailActionsService,
MessagingDeleteFolderMessagesService,
MessagingImportFolderMessagesService,
MessagingDeleteGroupEmailMessagesService,
InboundEmailImportService,
],
@@ -0,0 +1,49 @@
import { Injectable } from '@nestjs/common';
import {
ConnectedAccountProvider,
MessageFolderImportPolicy,
} from 'twenty-shared/types';
import { type MessageChannelEntity } from 'src/engine/metadata-modules/message-channel/entities/message-channel.entity';
import { type MessageFolderEntity } from 'src/engine/metadata-modules/message-folder/entities/message-folder.entity';
import { GmailGetMessageListService } from 'src/modules/messaging/message-import-manager/drivers/gmail/services/gmail-get-message-list.service';
@Injectable()
export class MessagingImportFolderMessagesService {
constructor(
private readonly gmailGetMessageListService: GmailGetMessageListService,
) {}
async getFolderMessageIdsToImport(
messageChannel: MessageChannelEntity,
messageFolder: MessageFolderEntity,
): Promise<string[]> {
switch (messageChannel.connectedAccount.provider) {
case ConnectedAccountProvider.GOOGLE: {
const foldersScopedToImportedFolder = messageChannel.messageFolders.map(
(folder) => ({
name: folder.name,
externalId: folder.externalId,
parentFolderId: folder.parentFolderId,
isSynced: folder.id === messageFolder.id,
}),
);
const [messageList] =
await this.gmailGetMessageListService.getMessageListWithoutCursor(
messageChannel.connectedAccount,
foldersScopedToImportedFolder,
{
messageFolderImportPolicy:
MessageFolderImportPolicy.SELECTED_FOLDERS,
},
);
return messageList?.messageExternalIds ?? [];
}
default:
return [];
}
}
}
@@ -27,7 +27,10 @@ import {
MessageImportSyncStep,
} from 'src/modules/messaging/message-import-manager/services/messaging-import-exception-handler.service';
import { MessagingMessagesImportService } from 'src/modules/messaging/message-import-manager/services/messaging-messages-import.service';
import { MessagingProcessFolderActionsService } from 'src/modules/messaging/message-import-manager/services/messaging-process-folder-actions.service';
import {
MessagingProcessFolderActionsService,
type ProcessFolderActionsResult,
} from 'src/modules/messaging/message-import-manager/services/messaging-process-folder-actions.service';
import { MessagingProcessGroupEmailActionsService } from 'src/modules/messaging/message-import-manager/services/messaging-process-group-email-actions.service';
import { MessageChannelEntity } from 'src/engine/metadata-modules/message-channel/entities/message-channel.entity';
@@ -68,7 +71,7 @@ export class MessagingMessageListFetchService {
workspaceId,
);
const pendingFolderActionsProcessed =
const processedfolderActionsResult =
await this.processPendingFolderActions(messageChannel, workspaceId);
await this.messageChannelSyncStatusService.markAsMessagesListFetchOngoing(
@@ -81,7 +84,8 @@ export class MessagingMessageListFetchService {
);
const freshMessageChannel =
pendingGroupEmailActionsProcessed || pendingFolderActionsProcessed
pendingGroupEmailActionsProcessed ||
isDefined(processedfolderActionsResult)
? await this.messageChannelRepository.findOne({
where: {
id: messageChannel.id,
@@ -120,9 +124,12 @@ export class MessagingMessageListFetchService {
`messages-to-import:${workspaceId}:${freshMessageChannel.id}`,
);
const messageExternalIds = messageLists.flatMap(
(messageList) => messageList.messageExternalIds,
);
const messageExternalIds = [
...messageLists.flatMap(
(messageList) => messageList.messageExternalIds,
),
...(processedfolderActionsResult?.messageExternalIdsToImport ?? []),
];
const messageExternalIdsToDelete = messageLists.flatMap(
(messageList) => messageList.messageExternalIdsToDelete,
@@ -312,7 +319,7 @@ export class MessagingMessageListFetchService {
private async processPendingFolderActions(
messageChannel: MessageChannelEntity,
workspaceId: string,
): Promise<boolean> {
): Promise<ProcessFolderActionsResult | null> {
const foldersWithPendingActions = messageChannel.messageFolders.filter(
(folder) =>
isDefined(folder.pendingSyncAction) &&
@@ -320,20 +327,18 @@ export class MessagingMessageListFetchService {
);
if (foldersWithPendingActions.length === 0) {
return false;
return null;
}
this.logger.log(
`messageChannelId: ${messageChannel.id} Processing pending folder actions before message list fetch`,
);
await this.messagingProcessFolderActionsService.processFolderActions(
return this.messagingProcessFolderActionsService.processFolderActions(
messageChannel,
foldersWithPendingActions,
workspaceId,
);
return true;
}
private async computeFullSyncMessageChannelMessageAssociationsToDelete(
@@ -10,6 +10,11 @@ import { GlobalWorkspaceOrmManager } from 'src/engine/twenty-orm/global-workspac
import { buildSystemAuthContext } from 'src/engine/twenty-orm/utils/build-system-auth-context.util';
import { MessageFolderPendingSyncAction } from 'twenty-shared/types';
import { MessagingDeleteFolderMessagesService } from 'src/modules/messaging/message-import-manager/services/messaging-delete-folder-messages.service';
import { MessagingImportFolderMessagesService } from 'src/modules/messaging/message-import-manager/services/messaging-import-folder-messages.service';
export type ProcessFolderActionsResult = {
messageExternalIdsToImport: string[];
};
@Injectable()
export class MessagingProcessFolderActionsService {
@@ -22,13 +27,14 @@ export class MessagingProcessFolderActionsService {
@InjectRepository(MessageFolderEntity)
private readonly messageFolderRepository: Repository<MessageFolderEntity>,
private readonly messagingDeleteFolderMessagesService: MessagingDeleteFolderMessagesService,
private readonly messagingImportFolderMessagesService: MessagingImportFolderMessagesService,
) {}
async processFolderActions(
messageChannel: MessageChannelEntity,
messageFolders: MessageFolderEntity[],
workspaceId: string,
): Promise<void> {
): Promise<ProcessFolderActionsResult> {
const foldersWithPendingActions = messageFolders.filter(
(folder) =>
isDefined(folder.pendingSyncAction) &&
@@ -36,13 +42,14 @@ export class MessagingProcessFolderActionsService {
);
if (foldersWithPendingActions.length === 0) {
return;
return { messageExternalIdsToImport: [] };
}
this.logger.log(
`WorkspaceId: ${workspaceId}, MessageChannelId: ${messageChannel.id} - Processing ${foldersWithPendingActions.length} folders with pending actions`,
);
const messageExternalIdsToImport: string[] = [];
const folderIdsToDelete: string[] = [];
const processedFolderIds: string[] = [];
const failedFolderIds: Array<{ folderId: string; error: Error }> = [];
@@ -53,21 +60,37 @@ export class MessagingProcessFolderActionsService {
`WorkspaceId: ${workspaceId}, MessageChannelId: ${messageChannel.id}, FolderId: ${folder.id} - Processing folder action: ${folder.pendingSyncAction}`,
);
if (
folder.pendingSyncAction ===
MessageFolderPendingSyncAction.FOLDER_DELETION
) {
await this.messagingDeleteFolderMessagesService.deleteFolderMessages(
workspaceId,
messageChannel,
folder,
);
switch (folder.pendingSyncAction) {
case MessageFolderPendingSyncAction.FOLDER_DELETION: {
await this.messagingDeleteFolderMessagesService.deleteFolderMessages(
workspaceId,
messageChannel,
folder,
);
folderIdsToDelete.push(folder.id);
folderIdsToDelete.push(folder.id);
this.logger.debug(
`WorkspaceId: ${workspaceId}, MessageChannelId: ${messageChannel.id}, FolderId: ${folder.id} - Completed FOLDER_DELETION action`,
);
this.logger.debug(
`WorkspaceId: ${workspaceId}, MessageChannelId: ${messageChannel.id}, FolderId: ${folder.id} - Completed FOLDER_DELETION action`,
);
break;
}
case MessageFolderPendingSyncAction.FOLDER_IMPORT: {
const folderMessageExternalIdsToImport =
await this.messagingImportFolderMessagesService.getFolderMessageIdsToImport(
messageChannel,
folder,
);
messageExternalIdsToImport.push(
...folderMessageExternalIdsToImport,
);
this.logger.debug(
`WorkspaceId: ${workspaceId}, MessageChannelId: ${messageChannel.id}, FolderId: ${folder.id} - Completed FOLDER_IMPORT action`,
);
break;
}
}
processedFolderIds.push(folder.id);
@@ -117,5 +140,9 @@ export class MessagingProcessFolderActionsService {
{ lite: true },
);
}
return {
messageExternalIdsToImport: [...new Set(messageExternalIdsToImport)],
};
}
}