This repository has no description
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}