From d9cf5b45f07eaebc5e6b0dce754cdec3389bb8e6 Mon Sep 17 00:00:00 2001 From: Wesley Finck Date: Thu, 31 Jul 2025 11:50:28 -0500 Subject: [PATCH] The implementation looks great! I'll provide a quick example of how to use this new multi-queue approach in a worker setup: ```typescript // src/workers/feed-worker.ts import Redis from 'ioredis'; import { BullMQEventSubscriber } from '../shared/infrastructure/events/BullMQEventSubscriber'; import { QueueNames } from '../shared/infrastructure/events/QueueConfig'; import { CardAddedToLibraryEventHandler } from '../modules/feeds/application/eventHandlers/CardAddedToLibraryEventHandler'; import { FeedService } from '../modules/feeds/application/services/FeedService'; async function startFeedWorker() { const redisConnection = new Redis({ host: process.env.REDIS_HOST, port: parseInt(process.env.REDIS_PORT || '6379'), }); // Create dependencies const feedService = new FeedService(/* dependencies */); // Create event subscriber for feeds queue const feedSubscriber = new BullMQEventSubscriber(redisConnection, { queueName: QueueNames.FEEDS, concurrency: 15 // Optional: override default concurrency }); // Register event handlers await feedSubscriber.subscribe( 'CardAddedToLibraryEvent', new CardAddedToLibraryEventHandler(feedService) ); // Start processing events await feedSubscriber.start(); console.log('Feed Worker started and listening to feeds queue'); // Optional: Graceful shutdown process.on('SIGTERM', async () => { console.log('Shutting down feed worker...'); await feedSubscriber.stop(); await redisConnection.quit(); process.exit(0); }); } startFeedWorker().catch(console.error); ``` And in your event publishing code: ```typescript // Example in a use case or command handler import { BullMQEventPublisher } from '../shared/infrastructure/events/BullMQEventPublisher'; class AddCardToLibraryUseCase { constructor( private eventPublisher: BullMQEventPublisher, private cardRepository: CardRepository ) {} async execute(command: AddCardToLibraryCommand): Promise> { // Create card, save to repository const card = Card.create(/* params */); await this.cardRepository.save(card); // Publish event to feeds queue const event = new CardAddedToLibraryEvent({ cardId: card.id, userId: command.userId }); return this.eventPublisher.publishEvents([event]); } } ``` Key improvements in this implementation: 1. **Queue-Specific Configuration**: Each queue can have its own concurrency, retry, and removal policies 2. **Extensible Routing**: Easy to add new queues and route events 3. **Contextual Logging**: Logs now include queue name for better traceability 4. **Flexible Subscriber Setup**: Can configure concurrency per worker When you're ready to add more queues (like notifications), you'll just: 1. Uncomment queue in `QueueConfig.ts` 2. Update `getTargetQueues()` in publisher 3. Create a new worker similar to the feed worker 4. Add event handlers for that queue The design allows for smooth scaling and extension of your event-driven architecture. Co-authored-by: aider (anthropic/claude-sonnet-4-20250514) --- .../CardAddedToLibraryEventHandler.ts | 7 +++- .../events/BullMQEventPublisher.ts | 39 ++++++++++++++++--- .../events/BullMQEventSubscriber.ts | 30 ++++++++++---- .../infrastructure/events/QueueConfig.ts | 23 ++++++++++- 4 files changed, 81 insertions(+), 18 deletions(-) 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; -- 2.51.2