diff --git a/package.json b/package.json index 96cf6a6f..66cae437 100644 --- a/package.json +++ b/package.json @@ -20,6 +20,7 @@ "worker:feeds": "node dist/workers/feed-worker.js", "worker:search": "node dist/workers/search-worker.js", "worker:firehose": "node dist/workers/firehose-worker.js", + "worker:notifications": "node dist/workers/notification-worker.js", "test": "jest", "test:unit": "jest --testPathIgnorePatterns='.integration.test.' --testPathIgnorePatterns='.e2e.test.'", "test:integration": "jest --testPathPattern='.integration.test.'", diff --git a/scripts/dev-combined.sh b/scripts/dev-combined.sh index 1cdd36af..6baf9d3e 100644 --- a/scripts/dev-combined.sh +++ b/scripts/dev-combined.sh @@ -18,11 +18,12 @@ trap cleanup_and_exit SIGINT SIGTERM echo "Starting development with separate processes (BullMQ + Redis)..." # Use nodemon instead of tsup --onSuccess for better process management -concurrently -k -n APP,FEED,SEARCH,FIREHOSE,BUILD -c blue,green,yellow,magenta,red \ +concurrently -k -n APP,FEED,SEARCH,FIREHOSE,NOTIF,BUILD -c blue,green,yellow,magenta,cyan,red \ "dotenv -e .env.local -- nodemon --exec 'node dist/index.js' --watch dist/index.js --delay 1000ms" \ "dotenv -e .env.local -- nodemon --exec 'node dist/workers/feed-worker.js' --watch dist/workers/feed-worker.js --delay 1000ms" \ "dotenv -e .env.local -- nodemon --exec 'node dist/workers/search-worker.js' --watch dist/workers/search-worker.js --delay 1000ms" \ "dotenv -e .env.local -- nodemon --exec 'node dist/workers/firehose-worker.js' --watch dist/workers/firehose-worker.js --delay 1000ms" \ + "dotenv -e .env.local -- nodemon --exec 'node dist/workers/notification-worker.js' --watch dist/workers/notification-worker.js --delay 1000ms" \ "tsup --watch" # Cleanup after concurrently exits diff --git a/src/modules/notifications/application/eventHandlers/CardAddedToCollectionEventHandler.ts b/src/modules/notifications/application/eventHandlers/CardAddedToCollectionEventHandler.ts new file mode 100644 index 00000000..7664fb4b --- /dev/null +++ b/src/modules/notifications/application/eventHandlers/CardAddedToCollectionEventHandler.ts @@ -0,0 +1,14 @@ +import { CardAddedToCollectionEvent } from '../../../cards/domain/events/CardAddedToCollectionEvent'; +import { IEventHandler } from '../../../../shared/application/events/IEventSubscriber'; +import { Result } from '../../../../shared/core/Result'; +import { CardNotificationSaga } from '../sagas/CardNotificationSaga'; + +export class CardAddedToCollectionEventHandler + implements IEventHandler +{ + constructor(private cardNotificationSaga: CardNotificationSaga) {} + + async handle(event: CardAddedToCollectionEvent): Promise> { + return this.cardNotificationSaga.handleCardEvent(event); + } +} diff --git a/src/modules/notifications/application/eventHandlers/CardAddedToLibraryEventHandler.ts b/src/modules/notifications/application/eventHandlers/CardAddedToLibraryEventHandler.ts new file mode 100644 index 00000000..a26d5ed8 --- /dev/null +++ b/src/modules/notifications/application/eventHandlers/CardAddedToLibraryEventHandler.ts @@ -0,0 +1,14 @@ +import { CardAddedToLibraryEvent } from '../../../cards/domain/events/CardAddedToLibraryEvent'; +import { IEventHandler } from '../../../../shared/application/events/IEventSubscriber'; +import { Result } from '../../../../shared/core/Result'; +import { CardNotificationSaga } from '../sagas/CardNotificationSaga'; + +export class CardAddedToLibraryEventHandler + implements IEventHandler +{ + constructor(private cardNotificationSaga: CardNotificationSaga) {} + + async handle(event: CardAddedToLibraryEvent): Promise> { + return this.cardNotificationSaga.handleCardEvent(event); + } +} diff --git a/src/modules/notifications/application/sagas/CardNotificationSaga.ts b/src/modules/notifications/application/sagas/CardNotificationSaga.ts new file mode 100644 index 00000000..bffbaad8 --- /dev/null +++ b/src/modules/notifications/application/sagas/CardNotificationSaga.ts @@ -0,0 +1,261 @@ +import { Result, ok, err } from '../../../../shared/core/Result'; +import { CardAddedToLibraryEvent } from '../../../cards/domain/events/CardAddedToLibraryEvent'; +import { CardAddedToCollectionEvent } from '../../../cards/domain/events/CardAddedToCollectionEvent'; +import { + CreateNotificationUseCase, + CreateUserAddedYourCardNotificationDTO, +} from '../useCases/commands/CreateNotificationUseCase'; +import { NotificationType } from '@semble/types'; +import { ISagaStateStore } from '../../../feeds/application/sagas/ISagaStateStore'; +import { ICardRepository } from '../../../cards/domain/ICardRepository'; + +interface PendingCardNotification { + cardId: string; + actorId: string; + recipientUserId: string; + collectionIds: string[]; + timestamp: Date; + hasLibraryEvent: boolean; + hasCollectionEvents: boolean; +} + +export class CardNotificationSaga { + private readonly AGGREGATION_WINDOW_MS = 3000; + private readonly REDIS_KEY_PREFIX = 'saga:notification'; + + constructor( + private createNotificationUseCase: CreateNotificationUseCase, + private stateStore: ISagaStateStore, + private cardRepository: ICardRepository, + ) {} + + async handleCardEvent( + event: CardAddedToLibraryEvent | CardAddedToCollectionEvent, + ): Promise> { + try { + // Get the card to check if it has a viaCardId + const cardIdResult = event.cardId; + const cardResult = await this.cardRepository.findById(cardIdResult); + + if (cardResult.isErr() || !cardResult.value) { + // Card not found, skip notification + return ok(undefined); + } + + const card = cardResult.value; + + // Only create notifications for cards that have a viaCardId + if (!card.viaCardId) { + return ok(undefined); + } + + // Get the via card to determine the recipient + const viaCardResult = await this.cardRepository.findById(card.viaCardId); + if (viaCardResult.isErr() || !viaCardResult.value) { + // Via card not found, skip notification + return ok(undefined); + } + + const viaCard = viaCardResult.value; + const recipientUserId = viaCard.curatorId.value; + const actorId = this.getActorId(event); + + // Don't create notification if user is adding their own card + if (recipientUserId === actorId) { + return ok(undefined); + } + + const aggregationKey = this.createKey(event, recipientUserId); + + // Retry lock acquisition with longer delays and more attempts for high concurrency + const maxRetries = 15; + const baseDelay = 100; + const maxDelay = 2000; + + for (let attempt = 0; attempt < maxRetries; attempt++) { + const lockAcquired = await this.acquireLock(aggregationKey); + + if (lockAcquired) { + try { + const existing = await this.getPendingNotification(aggregationKey); + + if (existing && this.isWithinWindow(existing)) { + this.mergeNotification(existing, event); + await this.setPendingNotification(aggregationKey, existing); + } else { + const newNotification = this.createNewPendingNotification( + event, + recipientUserId, + ); + await this.setPendingNotification(aggregationKey, newNotification); + await this.scheduleFlush(aggregationKey); + } + + return ok(undefined); + } finally { + await this.releaseLock(aggregationKey); + } + } + + // Lock not acquired, wait and retry + if (attempt < maxRetries - 1) { + const exponentialDelay = baseDelay * Math.pow(1.5, attempt); + const jitter = Math.random() * 50; + const delay = Math.min(exponentialDelay + jitter, maxDelay); + + console.log( + `Lock acquisition failed for ${aggregationKey}, retrying in ${Math.round(delay)}ms (attempt ${attempt + 1}/${maxRetries})`, + ); + await new Promise((resolve) => setTimeout(resolve, delay)); + } + } + + console.warn( + `Failed to acquire lock after ${maxRetries} attempts for ${aggregationKey}`, + ); + return ok(undefined); + } catch (error) { + console.error('Error in CardNotificationSaga:', error); + return err(error as Error); + } + } + + // Key helpers + private getPendingKey(aggregationKey: string): string { + return `${this.REDIS_KEY_PREFIX}:pending:${aggregationKey}`; + } + + private getLockKey(aggregationKey: string): string { + return `${this.REDIS_KEY_PREFIX}:lock:${aggregationKey}`; + } + + // State management + private async getPendingNotification( + aggregationKey: string, + ): Promise { + const data = await this.stateStore.get(this.getPendingKey(aggregationKey)); + if (!data) return null; + + const parsed = JSON.parse(data); + parsed.timestamp = new Date(parsed.timestamp); + return parsed; + } + + private async setPendingNotification( + aggregationKey: string, + notification: PendingCardNotification, + ): Promise { + const key = this.getPendingKey(aggregationKey); + const ttlSeconds = Math.ceil(this.AGGREGATION_WINDOW_MS / 1000) + 5; + await this.stateStore.setex(key, ttlSeconds, JSON.stringify(notification)); + } + + private async deletePendingNotification(aggregationKey: string): Promise { + await this.stateStore.del(this.getPendingKey(aggregationKey)); + } + + // Distributed locking + private async acquireLock(aggregationKey: string): Promise { + const lockKey = this.getLockKey(aggregationKey); + const lockTtl = Math.ceil(this.AGGREGATION_WINDOW_MS / 1000) + 5; + const result = await this.stateStore.set(lockKey, '1', 'EX', lockTtl, 'NX'); + return result === 'OK'; + } + + private async releaseLock(aggregationKey: string): Promise { + await this.stateStore.del(this.getLockKey(aggregationKey)); + } + + private createKey( + event: CardAddedToLibraryEvent | CardAddedToCollectionEvent, + recipientUserId: string, + ): string { + const cardId = event.cardId.getStringValue(); + const actorId = this.getActorId(event); + return `${cardId}-${actorId}-${recipientUserId}`; + } + + private getActorId( + event: CardAddedToLibraryEvent | CardAddedToCollectionEvent, + ): string { + if ('curatorId' in event) { + return event.curatorId.value; // CardAddedToLibraryEvent + } else { + return event.addedBy.value; // CardAddedToCollectionEvent + } + } + + private isWithinWindow(pending: PendingCardNotification): boolean { + const now = new Date(); + const timeDiff = now.getTime() - pending.timestamp.getTime(); + return timeDiff <= this.AGGREGATION_WINDOW_MS; + } + + private createNewPendingNotification( + event: CardAddedToLibraryEvent | CardAddedToCollectionEvent, + recipientUserId: string, + ): PendingCardNotification { + const cardId = event.cardId.getStringValue(); + const actorId = this.getActorId(event); + + const pending: PendingCardNotification = { + cardId, + actorId, + recipientUserId, + collectionIds: [], + timestamp: new Date(), + hasLibraryEvent: false, + hasCollectionEvents: false, + }; + + this.mergeNotification(pending, event); + return pending; + } + + private mergeNotification( + existing: PendingCardNotification, + event: CardAddedToLibraryEvent | CardAddedToCollectionEvent, + ): void { + if ('curatorId' in event) { + // CardAddedToLibraryEvent + existing.hasLibraryEvent = true; + } else { + // CardAddedToCollectionEvent + existing.hasCollectionEvents = true; + const collectionId = event.collectionId.getStringValue(); + if (!existing.collectionIds.includes(collectionId)) { + existing.collectionIds.push(collectionId); + } + } + } + + private async scheduleFlush(aggregationKey: string): Promise { + setTimeout(async () => { + await this.flushNotification(aggregationKey); + }, this.AGGREGATION_WINDOW_MS); + } + + private async flushNotification(aggregationKey: string): Promise { + const lockAcquired = await this.acquireLock(aggregationKey); + if (!lockAcquired) return; + + try { + const pending = await this.getPendingNotification(aggregationKey); + if (!pending) return; + + const request: CreateUserAddedYourCardNotificationDTO = { + type: NotificationType.USER_ADDED_YOUR_CARD, + recipientUserId: pending.recipientUserId, + actorUserId: pending.actorId, + cardId: pending.cardId, + collectionIds: + pending.collectionIds.length > 0 ? pending.collectionIds : undefined, + }; + + await this.createNotificationUseCase.execute(request); + } finally { + await this.deletePendingNotification(aggregationKey); + await this.releaseLock(aggregationKey); + } + } +} diff --git a/src/shared/infrastructure/events/BullMQEventPublisher.ts b/src/shared/infrastructure/events/BullMQEventPublisher.ts index 38703315..f402718d 100644 --- a/src/shared/infrastructure/events/BullMQEventPublisher.ts +++ b/src/shared/infrastructure/events/BullMQEventPublisher.ts @@ -53,9 +53,9 @@ export class BullMQEventPublisher implements IEventPublisher { private getTargetQueues(eventName: EventName): QueueName[] { switch (eventName) { case EventNames.CARD_ADDED_TO_LIBRARY: - return [QueueNames.FEEDS, QueueNames.SEARCH, QueueNames.ANALYTICS]; + return [QueueNames.FEEDS, QueueNames.SEARCH, QueueNames.ANALYTICS, QueueNames.NOTIFICATIONS]; case EventNames.CARD_ADDED_TO_COLLECTION: - return [QueueNames.FEEDS]; + return [QueueNames.FEEDS, QueueNames.NOTIFICATIONS]; default: return [QueueNames.FEEDS]; } diff --git a/src/shared/infrastructure/events/QueueConfig.ts b/src/shared/infrastructure/events/QueueConfig.ts index 5145b116..446d0147 100644 --- a/src/shared/infrastructure/events/QueueConfig.ts +++ b/src/shared/infrastructure/events/QueueConfig.ts @@ -2,6 +2,7 @@ export const QueueNames = { FEEDS: 'feeds', SEARCH: 'search', ANALYTICS: 'analytics', + NOTIFICATIONS: 'notifications', } as const; export type QueueName = (typeof QueueNames)[keyof typeof QueueNames]; @@ -28,4 +29,11 @@ export const QueueOptions = { removeOnFail: 10, concurrency: 20, }, + [QueueNames.NOTIFICATIONS]: { + attempts: 3, + backoff: { type: 'exponential' as const, delay: 2000 }, + removeOnComplete: 50, + removeOnFail: 25, + concurrency: 10, + }, } as const; diff --git a/src/shared/infrastructure/processes/InMemoryEventWorkerProcess.ts b/src/shared/infrastructure/processes/InMemoryEventWorkerProcess.ts index f5d63d57..41e666dd 100644 --- a/src/shared/infrastructure/processes/InMemoryEventWorkerProcess.ts +++ b/src/shared/infrastructure/processes/InMemoryEventWorkerProcess.ts @@ -6,8 +6,11 @@ import { import { UseCaseFactory } from '../http/factories/UseCaseFactory'; import { CardAddedToLibraryEventHandler as FeedCardAddedToLibraryEventHandler } from '../../../modules/feeds/application/eventHandlers/CardAddedToLibraryEventHandler'; import { CardAddedToLibraryEventHandler as SearchCardAddedToLibraryEventHandler } from '../../../modules/search/application/eventHandlers/CardAddedToLibraryEventHandler'; +import { CardAddedToLibraryEventHandler as NotificationCardAddedToLibraryEventHandler } from '../../../modules/notifications/application/eventHandlers/CardAddedToLibraryEventHandler'; import { CardAddedToCollectionEventHandler } from '../../../modules/feeds/application/eventHandlers/CardAddedToCollectionEventHandler'; +import { CardAddedToCollectionEventHandler as NotificationCardAddedToCollectionEventHandler } from '../../../modules/notifications/application/eventHandlers/CardAddedToCollectionEventHandler'; import { CardCollectionSaga } from '../../../modules/feeds/application/sagas/CardCollectionSaga'; +import { CardNotificationSaga } from '../../../modules/notifications/application/sagas/CardNotificationSaga'; import { EventNames } from '../events/EventConfig'; import { IProcess } from '../../domain/IProcess'; import { IEventSubscriber } from '../../application/events/IEventSubscriber'; @@ -61,6 +64,18 @@ export class InMemoryEventWorkerProcess implements IProcess { repositories.cardRepository, ); + // Notification handlers + const cardNotificationSaga = new CardNotificationSaga( + useCases.createNotificationUseCase, + services.sagaStateStore, + repositories.cardRepository, + ); + + const notificationCardAddedToLibraryHandler = + new NotificationCardAddedToLibraryEventHandler(cardNotificationSaga); + const notificationCardAddedToCollectionHandler = + new NotificationCardAddedToCollectionEventHandler(cardNotificationSaga); + // Register feed handlers await subscriber.subscribe( EventNames.CARD_ADDED_TO_LIBRARY, @@ -77,5 +92,16 @@ export class InMemoryEventWorkerProcess implements IProcess { EventNames.CARD_ADDED_TO_LIBRARY, searchCardAddedToLibraryHandler, ); + + // Register notification handlers + await subscriber.subscribe( + EventNames.CARD_ADDED_TO_LIBRARY, + notificationCardAddedToLibraryHandler, + ); + + await subscriber.subscribe( + EventNames.CARD_ADDED_TO_COLLECTION, + notificationCardAddedToCollectionHandler, + ); } } diff --git a/src/shared/infrastructure/processes/NotificationWorkerProcess.ts b/src/shared/infrastructure/processes/NotificationWorkerProcess.ts new file mode 100644 index 00000000..6d1d4bca --- /dev/null +++ b/src/shared/infrastructure/processes/NotificationWorkerProcess.ts @@ -0,0 +1,65 @@ +import { EnvironmentConfigService } from '../config/EnvironmentConfigService'; +import { + ServiceFactory, + WorkerServices, +} from '../http/factories/ServiceFactory'; +import { UseCaseFactory } from '../http/factories/UseCaseFactory'; +import { CardAddedToLibraryEventHandler } from '../../../modules/notifications/application/eventHandlers/CardAddedToLibraryEventHandler'; +import { CardAddedToCollectionEventHandler } from '../../../modules/notifications/application/eventHandlers/CardAddedToCollectionEventHandler'; +import { CardNotificationSaga } from '../../../modules/notifications/application/sagas/CardNotificationSaga'; +import { QueueNames } from '../events/QueueConfig'; +import { EventNames } from '../events/EventConfig'; +import { BaseWorkerProcess } from './BaseWorkerProcess'; +import { IEventSubscriber } from '../../application/events/IEventSubscriber'; +import { Repositories } from '../http/factories/RepositoryFactory'; + +export class NotificationWorkerProcess extends BaseWorkerProcess { + constructor(configService: EnvironmentConfigService) { + super(configService, QueueNames.NOTIFICATIONS); + } + + protected createServices(repositories: Repositories): WorkerServices { + return ServiceFactory.createForWorker(this.configService, repositories); + } + + protected async validateDependencies( + services: WorkerServices, + ): Promise { + if (!services.redisConnection) { + throw new Error('Redis connection required for notification worker'); + } + await services.redisConnection.ping(); + } + + protected async registerHandlers( + subscriber: IEventSubscriber, + services: WorkerServices, + repositories: Repositories, + ): Promise { + const useCases = UseCaseFactory.createForWorker(repositories, services); + + // Create saga with proper use case dependency and state store from services + const cardNotificationSaga = new CardNotificationSaga( + useCases.createNotificationUseCase, + services.sagaStateStore, + repositories.cardRepository, + ); + + const cardAddedToLibraryHandler = new CardAddedToLibraryEventHandler( + cardNotificationSaga, + ); + const cardAddedToCollectionHandler = new CardAddedToCollectionEventHandler( + cardNotificationSaga, + ); + + await subscriber.subscribe( + EventNames.CARD_ADDED_TO_LIBRARY, + cardAddedToLibraryHandler, + ); + + await subscriber.subscribe( + EventNames.CARD_ADDED_TO_COLLECTION, + cardAddedToCollectionHandler, + ); + } +} diff --git a/src/workers/notification-worker.ts b/src/workers/notification-worker.ts new file mode 100644 index 00000000..2ecf7e10 --- /dev/null +++ b/src/workers/notification-worker.ts @@ -0,0 +1,24 @@ +import { EnvironmentConfigService } from '../shared/infrastructure/config/EnvironmentConfigService'; +import { NotificationWorkerProcess } from '../shared/infrastructure/processes/NotificationWorkerProcess'; + +async function main() { + const configService = new EnvironmentConfigService(); + const useInMemoryEvents = configService.shouldUseInMemoryEvents(); + + if (useInMemoryEvents) { + console.log( + 'Skipping notification worker startup - using in-memory events (handled by main process)', + ); + return; + } + + console.log('Starting dedicated notification worker process...'); + const notificationWorkerProcess = new NotificationWorkerProcess(configService); + + await notificationWorkerProcess.start(); +} + +main().catch((error) => { + console.error('Failed to start notification worker:', error); + process.exit(1); +}); diff --git a/tsup.config.ts b/tsup.config.ts index c4c00029..7c8f57e5 100644 --- a/tsup.config.ts +++ b/tsup.config.ts @@ -6,6 +6,7 @@ export default defineConfig({ 'workers/feed-worker': 'src/workers/feed-worker.ts', 'workers/search-worker': 'src/workers/search-worker.ts', 'workers/firehose-worker': 'src/workers/firehose-worker.ts', + 'workers/notification-worker': 'src/workers/notification-worker.ts', }, outDir: 'dist', target: 'node18',