fix(messaging): emit channel and account deletion events from core metadata services (#21491)

/closes #21425

<!-- This is an auto-generated description by cubic. -->
<a
href="https://cubic.dev/pr/twentyhq/twenty/pull/21491?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-13 17:03:37 +05:30
committed by GitHub
parent b14195428a
commit d2fbc165b6
16 changed files with 164 additions and 75 deletions
@@ -8,6 +8,7 @@ import { CalendarChannelGraphqlApiExceptionInterceptor } from 'src/engine/metada
import { CalendarChannelResolver } from 'src/engine/metadata-modules/calendar-channel/resolvers/calendar-channel.resolver';
import { ConnectedAccountMetadataModule } from 'src/engine/metadata-modules/connected-account/connected-account-metadata.module';
import { PermissionsModule } from 'src/engine/metadata-modules/permissions/permissions.module';
import { WorkspaceEventEmitterModule } from 'src/engine/workspace-event-emitter/workspace-event-emitter.module';
@Module({
imports: [
@@ -15,6 +16,7 @@ import { PermissionsModule } from 'src/engine/metadata-modules/permissions/permi
PermissionsModule,
FeatureFlagModule,
ConnectedAccountMetadataModule,
WorkspaceEventEmitterModule,
],
providers: [
CalendarChannelMetadataService,
@@ -12,9 +12,12 @@ import {
CalendarChannelException,
CalendarChannelExceptionCode,
} from 'src/engine/metadata-modules/calendar-channel/calendar-channel.exception';
import { CALENDAR_CHANNEL_DELETED_EVENT } from 'src/engine/metadata-modules/calendar-channel/constants/calendar-channel-deleted.constant';
import { CalendarChannelDTO } from 'src/engine/metadata-modules/calendar-channel/dtos/calendar-channel.dto';
import { CalendarChannelEntity } from 'src/engine/metadata-modules/calendar-channel/entities/calendar-channel.entity';
import { type CalendarChannelDeletedEvent } from 'src/engine/metadata-modules/calendar-channel/types/calendar-channel-deleted.type';
import { ConnectedAccountMetadataService } from 'src/engine/metadata-modules/connected-account/connected-account-metadata.service';
import { WorkspaceEventEmitter } from 'src/engine/workspace-event-emitter/workspace-event-emitter';
@Injectable()
export class CalendarChannelMetadataService {
@@ -22,6 +25,7 @@ export class CalendarChannelMetadataService {
@InjectRepository(CalendarChannelEntity)
private readonly repository: Repository<CalendarChannelEntity>,
private readonly connectedAccountMetadataService: ConnectedAccountMetadataService,
private readonly workspaceEventEmitter: WorkspaceEventEmitter,
) {}
async findAll(workspaceId: string): Promise<CalendarChannelDTO[]> {
@@ -183,6 +187,12 @@ export class CalendarChannelMetadataService {
await this.repository.delete({ id, workspaceId });
this.workspaceEventEmitter.emitCustomBatchEvent<CalendarChannelDeletedEvent>(
CALENDAR_CHANNEL_DELETED_EVENT,
[{ calendarChannelId: id }],
workspaceId,
);
return calendarChannel;
}
}
@@ -0,0 +1 @@
export const CALENDAR_CHANNEL_DELETED_EVENT = 'calendarChannel_deleted';
@@ -0,0 +1,3 @@
export type CalendarChannelDeletedEvent = {
calendarChannelId: string;
};
@@ -4,13 +4,20 @@ import { InjectRepository } from '@nestjs/typeorm';
import { Repository } from 'typeorm';
import { AppOAuthRevokeService } from 'src/engine/core-modules/application/connection-provider/refresh/services/app-oauth-revoke.service';
import { CALENDAR_CHANNEL_DELETED_EVENT } from 'src/engine/metadata-modules/calendar-channel/constants/calendar-channel-deleted.constant';
import { CalendarChannelEntity } from 'src/engine/metadata-modules/calendar-channel/entities/calendar-channel.entity';
import { type CalendarChannelDeletedEvent } from 'src/engine/metadata-modules/calendar-channel/types/calendar-channel-deleted.type';
import { CONNECTED_ACCOUNT_DELETED_EVENT } from 'src/engine/metadata-modules/connected-account/constants/connected-account-deleted.constant';
import {
ConnectedAccountException,
ConnectedAccountExceptionCode,
} from 'src/engine/metadata-modules/connected-account/connected-account.exception';
import { ConnectedAccountEntity } from 'src/engine/metadata-modules/connected-account/entities/connected-account.entity';
import { type ConnectedAccountDeletedEvent } from 'src/engine/metadata-modules/connected-account/types/connected-account-deleted.type';
import { MESSAGE_CHANNEL_DELETED_EVENT } from 'src/engine/metadata-modules/message-channel/constants/message-channel-deleted.constant';
import { MessageChannelEntity } from 'src/engine/metadata-modules/message-channel/entities/message-channel.entity';
import { type MessageChannelDeletedEvent } from 'src/engine/metadata-modules/message-channel/types/message-channel-deleted.type';
import { WorkspaceEventEmitter } from 'src/engine/workspace-event-emitter/workspace-event-emitter';
@Injectable()
export class ConnectedAccountMetadataService {
@@ -24,6 +31,7 @@ export class ConnectedAccountMetadataService {
@InjectRepository(MessageChannelEntity)
private readonly messageChannelRepository: Repository<MessageChannelEntity>,
private readonly appOAuthRevokeService: AppOAuthRevokeService,
private readonly workspaceEventEmitter: WorkspaceEventEmitter,
) {}
async findByUserWorkspaceId({
@@ -164,23 +172,52 @@ export class ConnectedAccountMetadataService {
where: { id, workspaceId },
});
const [messageChannelCount, calendarChannelCount] = await Promise.all([
this.messageChannelRepository.count({
const [messageChannels, calendarChannels] = await Promise.all([
this.messageChannelRepository.find({
where: { connectedAccountId: id, workspaceId },
select: { id: true },
}),
this.calendarChannelRepository.count({
this.calendarChannelRepository.find({
where: { connectedAccountId: id, workspaceId },
select: { id: true },
}),
]);
this.logger.log(
`WorkspaceId: ${workspaceId} Deleting connected account ${id} with ${messageChannelCount} message channel(s) and ${calendarChannelCount} calendar channel(s)`,
`WorkspaceId: ${workspaceId} Deleting connected account ${id} with ${messageChannels.length} message channel(s) and ${calendarChannels.length} calendar channel(s)`,
);
await this.appOAuthRevokeService.revokeIfApp(connectedAccount);
await this.repository.delete({ id, workspaceId });
this.workspaceEventEmitter.emitCustomBatchEvent<MessageChannelDeletedEvent>(
MESSAGE_CHANNEL_DELETED_EVENT,
messageChannels.map((messageChannel) => ({
messageChannelId: messageChannel.id,
})),
workspaceId,
);
this.workspaceEventEmitter.emitCustomBatchEvent<CalendarChannelDeletedEvent>(
CALENDAR_CHANNEL_DELETED_EVENT,
calendarChannels.map((calendarChannel) => ({
calendarChannelId: calendarChannel.id,
})),
workspaceId,
);
this.workspaceEventEmitter.emitCustomBatchEvent<ConnectedAccountDeletedEvent>(
CONNECTED_ACCOUNT_DELETED_EVENT,
[
{
connectedAccountId: id,
userWorkspaceId: connectedAccount.userWorkspaceId,
},
],
workspaceId,
);
return connectedAccount;
}
}
@@ -0,0 +1 @@
export const CONNECTED_ACCOUNT_DELETED_EVENT = 'connectedAccount_deleted';
@@ -0,0 +1,4 @@
export type ConnectedAccountDeletedEvent = {
connectedAccountId: string;
userWorkspaceId: string;
};
@@ -0,0 +1 @@
export const MESSAGE_CHANNEL_DELETED_EVENT = 'messageChannel_deleted';
@@ -8,6 +8,7 @@ import { MessageChannelMetadataService } from 'src/engine/metadata-modules/messa
import { MessageChannelResolver } from 'src/engine/metadata-modules/message-channel/resolvers/message-channel.resolver';
import { MessageFolderEntity } from 'src/engine/metadata-modules/message-folder/entities/message-folder.entity';
import { PermissionsModule } from 'src/engine/metadata-modules/permissions/permissions.module';
import { WorkspaceEventEmitterModule } from 'src/engine/workspace-event-emitter/workspace-event-emitter.module';
import { MessagingImportManagerModule } from 'src/modules/messaging/message-import-manager/messaging-import-manager.module';
@Module({
@@ -16,6 +17,7 @@ import { MessagingImportManagerModule } from 'src/modules/messaging/message-impo
PermissionsModule,
ConnectedAccountMetadataModule,
MessagingImportManagerModule,
WorkspaceEventEmitterModule,
],
providers: [
MessageChannelMetadataService,
@@ -19,6 +19,7 @@ import {
import { StorageDriverType } from 'src/engine/core-modules/file-storage/interfaces/file-storage.interface';
import { TwentyConfigService } from 'src/engine/core-modules/twenty-config/twenty-config.service';
import { ConnectedAccountMetadataService } from 'src/engine/metadata-modules/connected-account/connected-account-metadata.service';
import { MESSAGE_CHANNEL_DELETED_EVENT } from 'src/engine/metadata-modules/message-channel/constants/message-channel-deleted.constant';
import { CreateEmailGroupChannelOutput } from 'src/engine/metadata-modules/message-channel/dtos/create-email-group-channel.output';
import { MessageChannelDTO } from 'src/engine/metadata-modules/message-channel/dtos/message-channel.dto';
import { MessageChannelEntity } from 'src/engine/metadata-modules/message-channel/entities/message-channel.entity';
@@ -26,6 +27,8 @@ import {
MessageChannelException,
MessageChannelExceptionCode,
} from 'src/engine/metadata-modules/message-channel/message-channel.exception';
import { type MessageChannelDeletedEvent } from 'src/engine/metadata-modules/message-channel/types/message-channel-deleted.type';
import { WorkspaceEventEmitter } from 'src/engine/workspace-event-emitter/workspace-event-emitter';
import { INBOUND_EMAIL_LOCAL_PART_PREFIX } from 'src/modules/messaging/message-import-manager/drivers/inbound-email/constants/inbound-email-local-part-prefix.constant';
import { INBOUND_EMAIL_LOCAL_PART_RANDOM_BYTES } from 'src/modules/messaging/message-import-manager/drivers/inbound-email/constants/inbound-email-local-part-random-bytes.constant';
@@ -36,6 +39,7 @@ export class MessageChannelMetadataService {
private readonly repository: Repository<MessageChannelEntity>,
private readonly connectedAccountMetadataService: ConnectedAccountMetadataService,
private readonly twentyConfigService: TwentyConfigService,
private readonly workspaceEventEmitter: WorkspaceEventEmitter,
) {}
async findAll(workspaceId: string): Promise<MessageChannelDTO[]> {
@@ -275,6 +279,12 @@ export class MessageChannelMetadataService {
await this.repository.delete({ id, workspaceId });
this.workspaceEventEmitter.emitCustomBatchEvent<MessageChannelDeletedEvent>(
MESSAGE_CHANNEL_DELETED_EVENT,
[{ messageChannelId: id }],
workspaceId,
);
return messageChannel;
}
}
@@ -0,0 +1,3 @@
export type MessageChannelDeletedEvent = {
messageChannelId: string;
};
@@ -1,18 +1,18 @@
import { Injectable } from '@nestjs/common';
import { type ObjectRecordDeleteEvent } from 'twenty-shared/database-events';
import { isDefined } from 'twenty-shared/utils';
import { OnCustomBatchEvent } from 'src/engine/api/graphql/graphql-query-runner/decorators/on-custom-batch-event.decorator';
import { InjectMessageQueue } from 'src/engine/core-modules/message-queue/decorators/message-queue.decorator';
import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants';
import { MessageQueueService } from 'src/engine/core-modules/message-queue/services/message-queue.service';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import { type CalendarChannelEntity } from 'src/engine/metadata-modules/calendar-channel/entities/calendar-channel.entity';
import { CALENDAR_CHANNEL_DELETED_EVENT } from 'src/engine/metadata-modules/calendar-channel/constants/calendar-channel-deleted.constant';
import { type CalendarChannelDeletedEvent } from 'src/engine/metadata-modules/calendar-channel/types/calendar-channel-deleted.type';
import { CustomWorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/custom-workspace-batch-event.type';
import {
CalendarChannelDeletionCleanupJob,
type CalendarChannelDeletionCleanupJobData,
} from 'src/modules/calendar/calendar-event-cleaner/jobs/calendar-channel-deletion-cleanup.job';
import { OnDatabaseBatchEvent } from 'src/engine/api/graphql/graphql-query-runner/decorators/on-database-batch-event.decorator';
import { DatabaseEventAction } from 'src/engine/api/graphql/graphql-query-runner/enums/database-event-action';
@Injectable()
export class CalendarEventCleanerCalendarChannelListener {
@@ -21,19 +21,23 @@ export class CalendarEventCleanerCalendarChannelListener {
private readonly calendarQueueService: MessageQueueService,
) {}
@OnDatabaseBatchEvent('calendarChannel', DatabaseEventAction.DESTROYED)
async handleDestroyedEvent(
payload: WorkspaceEventBatch<
ObjectRecordDeleteEvent<CalendarChannelEntity>
>,
@OnCustomBatchEvent(CALENDAR_CHANNEL_DELETED_EVENT)
async handleDeletedEvent(
batchEvent: CustomWorkspaceEventBatch<CalendarChannelDeletedEvent>,
) {
const { workspaceId } = batchEvent;
if (!isDefined(workspaceId)) {
return;
}
await Promise.all(
payload.events.map((eventPayload) =>
batchEvent.events.map((event) =>
this.calendarQueueService.add<CalendarChannelDeletionCleanupJobData>(
CalendarChannelDeletionCleanupJob.name,
{
workspaceId: payload.workspaceId,
calendarChannelId: eventPayload.recordId,
workspaceId,
calendarChannelId: event.calendarChannelId,
},
),
),
@@ -1,18 +1,18 @@
import { Injectable } from '@nestjs/common';
import { type ObjectRecordDeleteEvent } from 'twenty-shared/database-events';
import { isDefined } from 'twenty-shared/utils';
import { OnCustomBatchEvent } from 'src/engine/api/graphql/graphql-query-runner/decorators/on-custom-batch-event.decorator';
import { InjectMessageQueue } from 'src/engine/core-modules/message-queue/decorators/message-queue.decorator';
import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants';
import { MessageQueueService } from 'src/engine/core-modules/message-queue/services/message-queue.service';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import { CONNECTED_ACCOUNT_DELETED_EVENT } from 'src/engine/metadata-modules/connected-account/constants/connected-account-deleted.constant';
import { type ConnectedAccountDeletedEvent } from 'src/engine/metadata-modules/connected-account/types/connected-account-deleted.type';
import { CustomWorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/custom-workspace-batch-event.type';
import {
DeleteConnectedAccountAssociatedCalendarDataJob,
type DeleteConnectedAccountAssociatedCalendarDataJobData,
} from 'src/modules/calendar/calendar-event-cleaner/jobs/delete-connected-account-associated-calendar-data.job';
import { type ConnectedAccountEntity } from 'src/engine/metadata-modules/connected-account/entities/connected-account.entity';
import { OnDatabaseBatchEvent } from 'src/engine/api/graphql/graphql-query-runner/decorators/on-database-batch-event.decorator';
import { DatabaseEventAction } from 'src/engine/api/graphql/graphql-query-runner/enums/database-event-action';
@Injectable()
export class CalendarEventCleanerConnectedAccountListener {
@@ -21,19 +21,23 @@ export class CalendarEventCleanerConnectedAccountListener {
private readonly calendarQueueService: MessageQueueService,
) {}
@OnDatabaseBatchEvent('connectedAccount', DatabaseEventAction.DESTROYED)
async handleDestroyedEvent(
payload: WorkspaceEventBatch<
ObjectRecordDeleteEvent<ConnectedAccountEntity>
>,
@OnCustomBatchEvent(CONNECTED_ACCOUNT_DELETED_EVENT)
async handleDeletedEvent(
batchEvent: CustomWorkspaceEventBatch<ConnectedAccountDeletedEvent>,
) {
const { workspaceId } = batchEvent;
if (!isDefined(workspaceId)) {
return;
}
await Promise.all(
payload.events.map((eventPayload) =>
batchEvent.events.map((event) =>
this.calendarQueueService.add<DeleteConnectedAccountAssociatedCalendarDataJobData>(
DeleteConnectedAccountAssociatedCalendarDataJob.name,
{
workspaceId: payload.workspaceId,
connectedAccountId: eventPayload.recordId,
workspaceId,
connectedAccountId: event.connectedAccountId,
},
),
),
@@ -1,17 +1,17 @@
import { Injectable } from '@nestjs/common';
import { InjectRepository } from '@nestjs/typeorm';
import { type ObjectRecordDeleteEvent } from 'twenty-shared/database-events';
import { isDefined } from 'twenty-shared/utils';
import { type Repository } from 'typeorm';
import { OnDatabaseBatchEvent } from 'src/engine/api/graphql/graphql-query-runner/decorators/on-database-batch-event.decorator';
import { DatabaseEventAction } from 'src/engine/api/graphql/graphql-query-runner/enums/database-event-action';
import { OnCustomBatchEvent } from 'src/engine/api/graphql/graphql-query-runner/decorators/on-custom-batch-event.decorator';
import { UserWorkspaceEntity } from 'src/engine/core-modules/user-workspace/user-workspace.entity';
import { CONNECTED_ACCOUNT_DELETED_EVENT } from 'src/engine/metadata-modules/connected-account/constants/connected-account-deleted.constant';
import { type ConnectedAccountDeletedEvent } from 'src/engine/metadata-modules/connected-account/types/connected-account-deleted.type';
import { GlobalWorkspaceOrmManager } from 'src/engine/twenty-orm/global-workspace-datasource/global-workspace-orm.manager';
import { buildSystemAuthContext } from 'src/engine/twenty-orm/utils/build-system-auth-context.util';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import { CustomWorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/custom-workspace-batch-event.type';
import { AccountsToReconnectService } from 'src/modules/connected-account/services/accounts-to-reconnect.service';
import { type ConnectedAccountEntity } from 'src/engine/metadata-modules/connected-account/entities/connected-account.entity';
@Injectable()
export class ConnectedAccountListener {
@@ -22,35 +22,32 @@ export class ConnectedAccountListener {
private readonly userWorkspaceRepository: Repository<UserWorkspaceEntity>,
) {}
@OnDatabaseBatchEvent('connectedAccount', DatabaseEventAction.DESTROYED)
async handleDestroyedEvent(
payload: WorkspaceEventBatch<
ObjectRecordDeleteEvent<ConnectedAccountEntity>
>,
@OnCustomBatchEvent(CONNECTED_ACCOUNT_DELETED_EVENT)
async handleDeletedEvent(
batchEvent: CustomWorkspaceEventBatch<ConnectedAccountDeletedEvent>,
) {
const workspaceId = payload.workspaceId;
const { workspaceId } = batchEvent;
if (!isDefined(workspaceId)) {
return;
}
const authContext = buildSystemAuthContext(workspaceId);
await this.globalWorkspaceOrmManager.executeInWorkspaceContext(async () => {
for (const eventPayload of payload.events) {
const userWorkspaceId = eventPayload.properties.before.userWorkspaceId;
for (const event of batchEvent.events) {
const userWorkspace = await this.userWorkspaceRepository.findOne({
where: { id: userWorkspaceId },
where: { id: event.userWorkspaceId },
});
if (!userWorkspace) {
continue;
}
const userId = userWorkspace.userId;
const connectedAccountId = eventPayload.properties.before.id;
await this.accountsToReconnectService.removeAccountToReconnect(
userId,
userWorkspace.userId,
workspaceId,
connectedAccountId,
event.connectedAccountId,
);
}
}, authContext);
@@ -1,18 +1,18 @@
import { Injectable } from '@nestjs/common';
import { type ObjectRecordDeleteEvent } from 'twenty-shared/database-events';
import { isDefined } from 'twenty-shared/utils';
import { OnCustomBatchEvent } from 'src/engine/api/graphql/graphql-query-runner/decorators/on-custom-batch-event.decorator';
import { InjectMessageQueue } from 'src/engine/core-modules/message-queue/decorators/message-queue.decorator';
import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants';
import { MessageQueueService } from 'src/engine/core-modules/message-queue/services/message-queue.service';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import { type ConnectedAccountEntity } from 'src/engine/metadata-modules/connected-account/entities/connected-account.entity';
import { CONNECTED_ACCOUNT_DELETED_EVENT } from 'src/engine/metadata-modules/connected-account/constants/connected-account-deleted.constant';
import { type ConnectedAccountDeletedEvent } from 'src/engine/metadata-modules/connected-account/types/connected-account-deleted.type';
import { CustomWorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/custom-workspace-batch-event.type';
import {
MessagingConnectedAccountDeletionCleanupJob,
type MessagingConnectedAccountDeletionCleanupJobData,
} from 'src/modules/messaging/message-cleaner/jobs/messaging-connected-account-deletion-cleanup.job';
import { OnDatabaseBatchEvent } from 'src/engine/api/graphql/graphql-query-runner/decorators/on-database-batch-event.decorator';
import { DatabaseEventAction } from 'src/engine/api/graphql/graphql-query-runner/enums/database-event-action';
@Injectable()
export class MessagingMessageCleanerConnectedAccountListener {
@@ -21,19 +21,23 @@ export class MessagingMessageCleanerConnectedAccountListener {
private readonly messageQueueService: MessageQueueService,
) {}
@OnDatabaseBatchEvent('connectedAccount', DatabaseEventAction.DESTROYED)
async handleDestroyedEvent(
payload: WorkspaceEventBatch<
ObjectRecordDeleteEvent<ConnectedAccountEntity>
>,
@OnCustomBatchEvent(CONNECTED_ACCOUNT_DELETED_EVENT)
async handleDeletedEvent(
batchEvent: CustomWorkspaceEventBatch<ConnectedAccountDeletedEvent>,
) {
const { workspaceId } = batchEvent;
if (!isDefined(workspaceId)) {
return;
}
await Promise.all(
payload.events.map((eventPayload) =>
batchEvent.events.map((event) =>
this.messageQueueService.add<MessagingConnectedAccountDeletionCleanupJobData>(
MessagingConnectedAccountDeletionCleanupJob.name,
{
workspaceId: payload.workspaceId,
connectedAccountId: eventPayload.recordId,
workspaceId,
connectedAccountId: event.connectedAccountId,
},
),
),
@@ -1,18 +1,18 @@
import { Injectable } from '@nestjs/common';
import { type ObjectRecordDeleteEvent } from 'twenty-shared/database-events';
import { isDefined } from 'twenty-shared/utils';
import { OnCustomBatchEvent } from 'src/engine/api/graphql/graphql-query-runner/decorators/on-custom-batch-event.decorator';
import { InjectMessageQueue } from 'src/engine/core-modules/message-queue/decorators/message-queue.decorator';
import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants';
import { MessageQueueService } from 'src/engine/core-modules/message-queue/services/message-queue.service';
import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type';
import { type MessageChannelEntity } from 'src/engine/metadata-modules/message-channel/entities/message-channel.entity';
import { MESSAGE_CHANNEL_DELETED_EVENT } from 'src/engine/metadata-modules/message-channel/constants/message-channel-deleted.constant';
import { type MessageChannelDeletedEvent } from 'src/engine/metadata-modules/message-channel/types/message-channel-deleted.type';
import { CustomWorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/custom-workspace-batch-event.type';
import {
MessagingMessageChannelDeletionCleanupJob,
type MessagingMessageChannelDeletionCleanupJobData,
} from 'src/modules/messaging/message-cleaner/jobs/messaging-message-channel-deletion-cleanup.job';
import { OnDatabaseBatchEvent } from 'src/engine/api/graphql/graphql-query-runner/decorators/on-database-batch-event.decorator';
import { DatabaseEventAction } from 'src/engine/api/graphql/graphql-query-runner/enums/database-event-action';
@Injectable()
export class MessagingMessageCleanerMessageChannelListener {
@@ -21,17 +21,23 @@ export class MessagingMessageCleanerMessageChannelListener {
private readonly messageQueueService: MessageQueueService,
) {}
@OnDatabaseBatchEvent('messageChannel', DatabaseEventAction.DESTROYED)
async handleDestroyedEvent(
payload: WorkspaceEventBatch<ObjectRecordDeleteEvent<MessageChannelEntity>>,
@OnCustomBatchEvent(MESSAGE_CHANNEL_DELETED_EVENT)
async handleDeletedEvent(
batchEvent: CustomWorkspaceEventBatch<MessageChannelDeletedEvent>,
) {
const { workspaceId } = batchEvent;
if (!isDefined(workspaceId)) {
return;
}
await Promise.all(
payload.events.map((eventPayload) =>
batchEvent.events.map((event) =>
this.messageQueueService.add<MessagingMessageChannelDeletionCleanupJobData>(
MessagingMessageChannelDeletionCleanupJob.name,
{
workspaceId: payload.workspaceId,
messageChannelId: eventPayload.recordId,
workspaceId,
messageChannelId: event.messageChannelId,
},
),
),