This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

semble / src / shared / infrastructure / events / BullMQEventPublisher.ts
3.6 kB 110 lines
1import { Queue } from 'bullmq'; 2import Redis from 'ioredis'; 3import { IEventPublisher } from '../../application/events/IEventPublisher'; 4import { IDomainEvent } from '../../domain/events/IDomainEvent'; 5import { Result, ok, err } from '../../core/Result'; 6import { QueueNames, QueueOptions, QueueName } from './QueueConfig'; 7import { EventMapper } from './EventMapper'; 8import { EventName, EventNames } from './EventConfig'; 9 10export class BullMQEventPublisher implements IEventPublisher { 11 private queues: Map<string, Queue> = new Map(); 12 13 constructor(private redisConnection: Redis) {} 14 15 async publishEvents(events: IDomainEvent[]): Promise<Result<void>> { 16 // Fire and forget - don't await, just catch errors silently 17 this.publishEventsAsync(events).catch((error) => { 18 console.error('[BullMQEventPublisher] Failed to publish events:', error); 19 }); 20 21 return ok(undefined); 22 } 23 24 private async publishEventsAsync(events: IDomainEvent[]): Promise<void> { 25 // Group events by queue to enable batch publishing 26 const eventsByQueue = new Map<QueueName, IDomainEvent[]>(); 27 28 for (const event of events) { 29 const targetQueues = this.getTargetQueues(event.eventName); 30 31 for (const queueName of targetQueues) { 32 if (!eventsByQueue.has(queueName)) { 33 eventsByQueue.set(queueName, []); 34 } 35 eventsByQueue.get(queueName)!.push(event); 36 } 37 } 38 39 // Publish all events to each queue in batch 40 const publishPromises: Promise<void>[] = []; 41 42 for (const [queueName, queueEvents] of eventsByQueue.entries()) { 43 publishPromises.push(this.publishBatchToQueue(queueName, queueEvents)); 44 } 45 46 await Promise.all(publishPromises); 47 } 48 49 private async publishBatchToQueue( 50 queueName: QueueName, 51 events: IDomainEvent[], 52 ): Promise<void> { 53 if (!this.queues.has(queueName)) { 54 this.queues.set( 55 queueName, 56 new Queue(queueName, { 57 connection: this.redisConnection, 58 defaultJobOptions: QueueOptions[queueName], 59 }), 60 ); 61 } 62 63 const queue = this.queues.get(queueName)!; 64 65 // Use addBulk for batch publishing 66 const jobs = events.map((event) => { 67 const serializedEvent = EventMapper.toSerialized(event); 68 return { 69 name: serializedEvent.eventType, 70 data: serializedEvent, 71 }; 72 }); 73 74 await queue.addBulk(jobs); 75 } 76 77 private getTargetQueues(eventName: EventName): QueueName[] { 78 switch (eventName) { 79 case EventNames.CARD_ADDED_TO_LIBRARY: 80 return [ 81 QueueNames.FEEDS, 82 QueueNames.SEARCH, 83 QueueNames.NOTIFICATIONS, 84 QueueNames.SYNC, 85 ]; 86 case EventNames.CARD_ADDED_TO_COLLECTION: 87 return [QueueNames.FEEDS, QueueNames.NOTIFICATIONS]; 88 case EventNames.CARD_REMOVED_FROM_LIBRARY: 89 return [QueueNames.NOTIFICATIONS]; 90 case EventNames.CARD_REMOVED_FROM_COLLECTION: 91 return [QueueNames.NOTIFICATIONS]; 92 case EventNames.USER_FOLLOWED_TARGET: 93 return [QueueNames.NOTIFICATIONS]; 94 case EventNames.USER_UNFOLLOWED_TARGET: 95 return [QueueNames.NOTIFICATIONS]; 96 case EventNames.CONNECTION_CREATED: 97 return [QueueNames.FEEDS, QueueNames.NOTIFICATIONS, QueueNames.SEARCH]; 98 case EventNames.CONNECTION_REMOVED: 99 return [QueueNames.NOTIFICATIONS]; 100 default: 101 return [QueueNames.FEEDS]; 102 } 103 } 104 105 async close(): Promise<void> { 106 await Promise.all( 107 Array.from(this.queues.values()).map((queue) => queue.close()), 108 ); 109 } 110}