diff --git a/src/modules/notifications/application/eventHandlers/ConnectionRemovedEventHandler.ts b/src/modules/notifications/application/eventHandlers/ConnectionRemovedEventHandler.ts new file mode 100644 index 00000000..83c17221 --- /dev/null +++ b/src/modules/notifications/application/eventHandlers/ConnectionRemovedEventHandler.ts @@ -0,0 +1,47 @@ +import { IEventHandler } from '../../../../shared/application/events/IEventSubscriber'; +import { ConnectionRemovedEvent } from '../../../cards/domain/events/ConnectionRemovedEvent'; +import { Result, ok, err } from '../../../../shared/core/Result'; +import { INotificationRepository } from '../../domain/INotificationRepository'; +import { CuratorId } from '../../../cards/domain/value-objects/CuratorId'; + +export class ConnectionRemovedEventHandler + implements IEventHandler +{ + constructor(private notificationRepository: INotificationRepository) {} + + async handle(event: ConnectionRemovedEvent): Promise> { + try { + const actorUserId = CuratorId.create(event.curatorId.value); + + if (actorUserId.isErr()) { + console.error('Invalid curator ID in ConnectionRemovedEvent'); + return ok(undefined); + } + + // Find and delete any existing notifications for this connection/actor combination + const existingNotificationsResult = + await this.notificationRepository.findByConnectionAndActor( + event.connectionId.getStringValue(), + actorUserId.value, + ); + + if (existingNotificationsResult.isOk()) { + const notifications = existingNotificationsResult.value; + for (const notification of notifications) { + const deleteResult = await this.notificationRepository.delete( + notification.notificationId, + ); + if (deleteResult.isErr()) { + console.error('Failed to delete notification:', deleteResult.error); + // Continue with other notifications even if one fails + } + } + } + + return ok(undefined); + } catch (error) { + console.error('Error handling ConnectionRemovedEvent:', error); + return err(error as Error); + } + } +} diff --git a/src/modules/notifications/domain/INotificationRepository.ts b/src/modules/notifications/domain/INotificationRepository.ts index a768e5d5..f250a56e 100644 --- a/src/modules/notifications/domain/INotificationRepository.ts +++ b/src/modules/notifications/domain/INotificationRepository.ts @@ -117,6 +117,10 @@ export interface INotificationRepository { actorUserId: CuratorId, ): Promise>; findByCard(cardId: string): Promise>; + findByConnectionAndActor( + connectionId: string, + actorUserId: CuratorId, + ): Promise>; findFollowNotificationsByActorAndTarget( actorUserId: CuratorId, targetId: string, diff --git a/src/modules/notifications/infrastructure/repositories/DrizzleNotificationRepository.ts b/src/modules/notifications/infrastructure/repositories/DrizzleNotificationRepository.ts index a7534bf8..a938ce28 100644 --- a/src/modules/notifications/infrastructure/repositories/DrizzleNotificationRepository.ts +++ b/src/modules/notifications/infrastructure/repositories/DrizzleNotificationRepository.ts @@ -297,6 +297,47 @@ export class DrizzleNotificationRepository implements INotificationRepository { } } + async findByConnectionAndActor( + connectionId: string, + actorUserId: CuratorId, + ): Promise> { + try { + const result = await this.db + .select() + .from(notifications) + .where(eq(notifications.actorUserId, actorUserId.value)); + + // Filter by connectionId in metadata + const matchingNotifications: Notification[] = []; + for (const notificationData of result) { + const metadata = notificationData.metadata as any; + if (metadata.connectionId === connectionId) { + const dto: NotificationDTO = { + id: notificationData.id, + recipientUserId: notificationData.recipientUserId, + actorUserId: notificationData.actorUserId, + type: notificationData.type, + metadata: notificationData.metadata as any, + read: notificationData.read, + createdAt: notificationData.createdAt, + updatedAt: notificationData.updatedAt, + }; + + const domainResult = NotificationMapper.toDomain(dto); + if (domainResult.isErr()) { + return err(domainResult.error); + } + + matchingNotifications.push(domainResult.value); + } + } + + return ok(matchingNotifications); + } catch (error) { + return err(error as Error); + } + } + async findFollowNotificationsByActorAndTarget( actorUserId: CuratorId, targetId: string, diff --git a/src/modules/notifications/tests/infrastructure/InMemoryNotificationRepository.ts b/src/modules/notifications/tests/infrastructure/InMemoryNotificationRepository.ts index ffa85c5c..b32beaea 100644 --- a/src/modules/notifications/tests/infrastructure/InMemoryNotificationRepository.ts +++ b/src/modules/notifications/tests/infrastructure/InMemoryNotificationRepository.ts @@ -144,6 +144,23 @@ export class InMemoryNotificationRepository implements INotificationRepository { return ok(matchingNotifications); } + async findByConnectionAndActor( + connectionId: string, + actorUserId: CuratorId, + ): Promise> { + const matchingNotifications = Array.from( + this.notifications.values(), + ).filter((notification) => { + const metadata = notification.metadata as any; + return ( + metadata.connectionId === connectionId && + notification.actorUserId.equals(actorUserId) + ); + }); + + return ok(matchingNotifications); + } + async findFollowNotificationsByActorAndTarget( actorUserId: CuratorId, targetId: string, diff --git a/src/shared/infrastructure/events/BullMQEventPublisher.ts b/src/shared/infrastructure/events/BullMQEventPublisher.ts index 48767954..767ac48a 100644 --- a/src/shared/infrastructure/events/BullMQEventPublisher.ts +++ b/src/shared/infrastructure/events/BullMQEventPublisher.ts @@ -95,6 +95,8 @@ export class BullMQEventPublisher implements IEventPublisher { return [QueueNames.NOTIFICATIONS]; case EventNames.CONNECTION_CREATED: return [QueueNames.FEEDS, QueueNames.NOTIFICATIONS]; + case EventNames.CONNECTION_REMOVED: + return [QueueNames.NOTIFICATIONS]; default: return [QueueNames.FEEDS]; } diff --git a/src/shared/infrastructure/processes/InMemoryEventWorkerProcess.ts b/src/shared/infrastructure/processes/InMemoryEventWorkerProcess.ts index 765631ac..e55cfb6f 100644 --- a/src/shared/infrastructure/processes/InMemoryEventWorkerProcess.ts +++ b/src/shared/infrastructure/processes/InMemoryEventWorkerProcess.ts @@ -23,6 +23,7 @@ import { } from '../http/factories/RepositoryFactory'; import { ConnectionCreatedEventHandler } from 'src/modules/feeds/application/eventHandlers/ConnectionCreatedEventHandler'; import { ConnectionCreatedEventHandler as NotificationConnectionCreatedEventHandler } from 'src/modules/notifications/application/eventHandlers/ConnectionCreatedEventHandler'; +import { ConnectionRemovedEventHandler } from 'src/modules/notifications/application/eventHandlers/ConnectionRemovedEventHandler'; export class InMemoryEventWorkerProcess implements IProcess { constructor(private configService: EnvironmentConfigService) {} @@ -117,6 +118,10 @@ export class InMemoryEventWorkerProcess implements IProcess { services.identityResolutionService, ); + const connectionRemovedHandler = new ConnectionRemovedEventHandler( + repositories.notificationRepository, + ); + // Register feed handlers await subscriber.subscribe( EventNames.CARD_ADDED_TO_LIBRARY, @@ -178,5 +183,10 @@ export class InMemoryEventWorkerProcess implements IProcess { EventNames.CONNECTION_CREATED, notificationConnectionCreatedHandler, ); + + await subscriber.subscribe( + EventNames.CONNECTION_REMOVED, + connectionRemovedHandler, + ); } } diff --git a/src/shared/infrastructure/processes/NotificationWorkerProcess.ts b/src/shared/infrastructure/processes/NotificationWorkerProcess.ts index 8296a545..1cf3544d 100644 --- a/src/shared/infrastructure/processes/NotificationWorkerProcess.ts +++ b/src/shared/infrastructure/processes/NotificationWorkerProcess.ts @@ -18,6 +18,7 @@ import { BaseWorkerProcess } from './BaseWorkerProcess'; import { IEventSubscriber } from '../../application/events/IEventSubscriber'; import { Repositories } from '../http/factories/RepositoryFactory'; import { ConnectionCreatedEventHandler } from 'src/modules/notifications/application/eventHandlers/ConnectionCreatedEventHandler'; +import { ConnectionRemovedEventHandler } from 'src/modules/notifications/application/eventHandlers/ConnectionRemovedEventHandler'; export class NotificationWorkerProcess extends BaseWorkerProcess { constructor(configService: EnvironmentConfigService) { @@ -96,6 +97,10 @@ export class NotificationWorkerProcess extends BaseWorkerProcess { services.identityResolutionService, ); + const connectionRemovedHandler = new ConnectionRemovedEventHandler( + repositories.notificationRepository, + ); + await subscriber.subscribe( EventNames.CARD_ADDED_TO_LIBRARY, cardAddedToLibraryHandler, @@ -137,5 +142,10 @@ export class NotificationWorkerProcess extends BaseWorkerProcess { EventNames.CONNECTION_CREATED, connectionCreatedHandler, ); + + await subscriber.subscribe( + EventNames.CONNECTION_REMOVED, + connectionRemovedHandler, + ); } }