diff --git a/src/modules/feeds/application/eventHandlers/CardAddedToLibraryEventHandler.ts b/src/modules/feeds/application/eventHandlers/CardAddedToLibraryEventHandler.ts index fc042f06..8cc5591a 100644 --- a/src/modules/feeds/application/eventHandlers/CardAddedToLibraryEventHandler.ts +++ b/src/modules/feeds/application/eventHandlers/CardAddedToLibraryEventHandler.ts @@ -8,20 +8,23 @@ export class CardAddedToLibraryEventHandler implements IEventHandler> { try { + console.log(`[FEEDS] Processing CardAddedToLibraryEvent for card ${event.cardId.getStringValue()}`); + const result = await this.feedService.processCardAddedToLibrary(event); if (result.isErr()) { console.error( - 'Error processing CardAddedToLibraryEvent in feeds:', + '[FEEDS] Error processing CardAddedToLibraryEvent:', result.error, ); return err(result.error); } + console.log(`[FEEDS] Successfully processed CardAddedToLibraryEvent for card ${event.cardId.getStringValue()}`); return ok(undefined); } catch (error) { console.error( - 'Unexpected error handling CardAddedToLibraryEvent in feeds:', + '[FEEDS] Unexpected error handling CardAddedToLibraryEvent:', error, ); return err(error as Error); diff --git a/src/shared/infrastructure/events/BullMQEventPublisher.ts b/src/shared/infrastructure/events/BullMQEventPublisher.ts index 3dbb4331..6b37dcc4 100644 --- a/src/shared/infrastructure/events/BullMQEventPublisher.ts +++ b/src/shared/infrastructure/events/BullMQEventPublisher.ts @@ -3,8 +3,9 @@ import Redis from 'ioredis'; import { IEventPublisher } from '../../application/events/IEventPublisher'; import { IDomainEvent } from '../../domain/events/IDomainEvent'; import { Result, ok, err } from '../../core/Result'; -import { QueueNames, QueueOptions } from './QueueConfig'; +import { QueueNames, QueueOptions, QueueName } from './QueueConfig'; import { EventMapper } from './EventMapper'; +import { EventName } from './EventConfig'; export class BullMQEventPublisher implements IEventPublisher { private queues: Map = new Map(); @@ -23,21 +24,47 @@ export class BullMQEventPublisher implements IEventPublisher { } private async publishSingleEvent(event: IDomainEvent): Promise { - if (!this.queues.has(QueueNames.EVENTS)) { + const targetQueues = this.getTargetQueues(event.eventName); + + for (const queueName of targetQueues) { + await this.publishToQueue(queueName, event); + } + } + + private async publishToQueue(queueName: QueueName, event: IDomainEvent): Promise { + if (!this.queues.has(queueName)) { this.queues.set( - QueueNames.EVENTS, - new Queue(QueueNames.EVENTS, { + queueName, + new Queue(queueName, { connection: this.redisConnection, - defaultJobOptions: QueueOptions[QueueNames.EVENTS], + defaultJobOptions: QueueOptions[queueName], }), ); } - const queue = this.queues.get(QueueNames.EVENTS)!; + const queue = this.queues.get(queueName)!; const serializedEvent = EventMapper.toSerialized(event); await queue.add(serializedEvent.eventType, serializedEvent); } + private getTargetQueues(eventName: EventName): QueueName[] { + // Route events to appropriate queues + // For now, all events go to feeds queue + // Future: route different events to different queues + switch (eventName) { + case 'CardAddedToLibraryEvent': + return [QueueNames.FEEDS]; + // Future: return [QueueNames.FEEDS, QueueNames.NOTIFICATIONS, QueueNames.ANALYTICS]; + case 'CardAddedToCollectionEvent': + return [QueueNames.FEEDS]; + // Future: return [QueueNames.FEEDS, QueueNames.ANALYTICS]; + case 'CollectionCreatedEvent': + return [QueueNames.FEEDS]; + // Future: return [QueueNames.FEEDS, QueueNames.ANALYTICS]; + default: + return [QueueNames.FEEDS]; // Default to feeds queue + } + } async close(): Promise { await Promise.all( diff --git a/src/shared/infrastructure/events/BullMQEventSubscriber.ts b/src/shared/infrastructure/events/BullMQEventSubscriber.ts index 450ed315..d84739cd 100644 --- a/src/shared/infrastructure/events/BullMQEventSubscriber.ts +++ b/src/shared/infrastructure/events/BullMQEventSubscriber.ts @@ -5,15 +5,26 @@ import { IEventHandler, } from '../../application/events/IEventSubscriber'; import { IDomainEvent } from '../../domain/events/IDomainEvent'; -import { QueueNames } from './QueueConfig'; +import { QueueNames, QueueOptions, QueueName } from './QueueConfig'; import { EventMapper } from './EventMapper'; import { EventName } from './EventConfig'; +export interface BullMQEventSubscriberConfig { + queueName: QueueName; + concurrency?: number; +} + export class BullMQEventSubscriber implements IEventSubscriber { private workers: Worker[] = []; private handlers: Map> = new Map(); + private config: BullMQEventSubscriberConfig; - constructor(private redisConnection: Redis) {} + constructor( + private redisConnection: Redis, + config: BullMQEventSubscriberConfig, + ) { + this.config = config; + } async subscribe( eventType: EventName, @@ -23,27 +34,30 @@ export class BullMQEventSubscriber implements IEventSubscriber { } async start(): Promise { + const queueConfig = QueueOptions[this.config.queueName]; + const concurrency = this.config.concurrency || queueConfig.concurrency || 10; + const worker = new Worker( - QueueNames.EVENTS, + this.config.queueName, async (job: Job) => { await this.processJob(job); }, { connection: this.redisConnection, - concurrency: 10, + concurrency, }, ); worker.on('completed', (job) => { - console.log(`Job ${job.id} completed successfully`); + console.log(`[${this.config.queueName}] Job ${job.id} completed successfully`); }); worker.on('failed', (job, err) => { - console.error(`Job ${job?.id} failed:`, err); + console.error(`[${this.config.queueName}] Job ${job?.id} failed:`, err); }); worker.on('error', (err) => { - console.error('Worker error:', err); + console.error(`[${this.config.queueName}] Worker error:`, err); }); this.workers.push(worker); @@ -60,7 +74,7 @@ export class BullMQEventSubscriber implements IEventSubscriber { const handler = this.handlers.get(eventType); if (!handler) { - console.warn(`No handler registered for event type: ${eventType}`); + console.warn(`[${this.config.queueName}] No handler registered for event type: ${eventType}`); return; } diff --git a/src/shared/infrastructure/events/QueueConfig.ts b/src/shared/infrastructure/events/QueueConfig.ts index e0853c46..21a70bf7 100644 --- a/src/shared/infrastructure/events/QueueConfig.ts +++ b/src/shared/infrastructure/events/QueueConfig.ts @@ -1,14 +1,33 @@ export const QueueNames = { - EVENTS: 'events', + FEEDS: 'feeds', + // Future queues can be added here: + // NOTIFICATIONS: 'notifications', + // ANALYTICS: 'analytics', } as const; export type QueueName = typeof QueueNames[keyof typeof QueueNames]; export const QueueOptions = { - [QueueNames.EVENTS]: { + [QueueNames.FEEDS]: { attempts: 3, backoff: { type: 'exponential' as const, delay: 2000 }, removeOnComplete: 50, removeOnFail: 25, + concurrency: 15, }, + // Future queue configurations: + // [QueueNames.NOTIFICATIONS]: { + // attempts: 5, + // backoff: { type: 'exponential' as const, delay: 1000 }, + // removeOnComplete: 100, + // removeOnFail: 50, + // concurrency: 5, + // }, + // [QueueNames.ANALYTICS]: { + // attempts: 2, + // backoff: { type: 'exponential' as const, delay: 5000 }, + // removeOnComplete: 25, + // removeOnFail: 10, + // concurrency: 20, + // }, } as const;