|
1 | | -import { IEventBusModuleService } from "@medusajs/framework/types" |
2 | 1 | import { EventEmitter } from "events" |
3 | 2 |
|
4 | | -// Allows you to wait for all subscribers to execute for a given event. Only works with the local event bus. |
5 | | -export const waitSubscribersExecution = ( |
6 | | - eventName: string, |
7 | | - eventBus: IEventBusModuleService, |
8 | | - { |
9 | | - timeout = 5000, |
10 | | - }: { |
11 | | - timeout?: number |
12 | | - } = {} |
13 | | -) => { |
14 | | - const eventEmitter: EventEmitter = (eventBus as any).eventEmitter_ |
15 | | - const subscriberPromises: Promise<any>[] = [] |
16 | | - const originalListeners = eventEmitter.listeners(eventName) |
17 | | - let timeoutId: NodeJS.Timeout | null = null |
| 3 | +type EventBus = { |
| 4 | + eventEmitter_: EventEmitter |
| 5 | +} |
| 6 | + |
| 7 | +type WaitSubscribersExecutionOptions = { |
| 8 | + timeout?: number |
| 9 | +} |
| 10 | + |
| 11 | +// Map to hold pending promises for each event. |
| 12 | +const waits = new Map<string | symbol, Promise<any>>() |
18 | 13 |
|
19 | | - // Create a promise that rejects after the timeout |
20 | | - const timeoutPromise = new Promise((_, reject) => { |
| 14 | +/** |
| 15 | + * Creates a promise that rejects after a specified timeout. |
| 16 | + * @param timeout - The timeout in milliseconds. |
| 17 | + * @param eventName - The name of the event being waited on. |
| 18 | + * @returns A tuple containing the timeout promise and a function to clear the timeout. |
| 19 | + */ |
| 20 | +const createTimeoutPromise = ( |
| 21 | + timeout: number, |
| 22 | + eventName: string | symbol |
| 23 | +): [Promise<never>, () => void] => { |
| 24 | + let timeoutId: NodeJS.Timeout | null = null |
| 25 | + const promise = new Promise<never>((_, reject) => { |
21 | 26 | timeoutId = setTimeout(() => { |
22 | 27 | reject( |
23 | 28 | new Error( |
24 | | - `Timeout of ${timeout}ms exceeded while waiting for event "${eventName}"` |
| 29 | + `Timeout of ${timeout}ms exceeded while waiting for event "${String( |
| 30 | + eventName |
| 31 | + )}"` |
25 | 32 | ) |
26 | 33 | ) |
27 | 34 | }, timeout) |
28 | 35 | timeoutId.unref() |
29 | 36 | }) |
| 37 | + return [promise, () => timeoutId && clearTimeout(timeoutId)] |
| 38 | +} |
| 39 | + |
| 40 | +// Core logic to wait for subscribers. |
| 41 | +const doWaitSubscribersExecution = ( |
| 42 | + eventName: string | symbol, |
| 43 | + eventBus: EventBus, |
| 44 | + { timeout = 15000 }: WaitSubscribersExecutionOptions = {} |
| 45 | +): Promise<any> => { |
| 46 | + const eventEmitter = eventBus.eventEmitter_ |
| 47 | + const subscriberPromises: Promise<any>[] = [] |
| 48 | + const [timeoutPromise, clearTimeout] = createTimeoutPromise( |
| 49 | + timeout, |
| 50 | + eventName |
| 51 | + ) |
30 | 52 |
|
31 | | - // If there are no existing listeners, resolve once the event happens. Otherwise, wrap the existing subscribers in a promise and resolve once they are done. |
32 | 53 | if (!eventEmitter.listeners(eventName).length) { |
33 | | - let ok |
| 54 | + let ok: (value?: any) => void |
34 | 55 | const promise = new Promise((resolve) => { |
35 | 56 | ok = resolve |
36 | 57 | }) |
37 | | - |
38 | 58 | subscriberPromises.push(promise) |
39 | | - eventEmitter.on(eventName, ok) |
| 59 | + |
| 60 | + const newListener = async () => { |
| 61 | + eventEmitter.removeListener(eventName, newListener) |
| 62 | + ok() |
| 63 | + } |
| 64 | + |
| 65 | + Object.defineProperty(newListener, "__isSubscribersExecutionWrapper", { |
| 66 | + value: true, |
| 67 | + configurable: true, |
| 68 | + enumerable: false, |
| 69 | + }) |
| 70 | + |
| 71 | + eventEmitter.on(eventName, newListener) |
40 | 72 | } else { |
41 | 73 | eventEmitter.listeners(eventName).forEach((listener: any) => { |
| 74 | + if (listener.__isSubscribersExecutionWrapper) { |
| 75 | + return |
| 76 | + } |
| 77 | + |
42 | 78 | eventEmitter.removeListener(eventName, listener) |
43 | 79 |
|
44 | | - let ok, nok |
| 80 | + let ok: (value?: any) => void, nok: (reason?: any) => void |
45 | 81 | const promise = new Promise((resolve, reject) => { |
46 | 82 | ok = resolve |
47 | 83 | nok = reject |
48 | 84 | }) |
49 | 85 | subscriberPromises.push(promise) |
50 | 86 |
|
51 | | - const newListener = async (...args2) => { |
| 87 | + const newListener = async (...args2: any[]) => { |
| 88 | + // As soon as the subscriber is executed, we restore the original listener |
| 89 | + eventEmitter.removeListener(eventName, newListener) |
| 90 | + let listenerToAdd = listener |
| 91 | + while (listenerToAdd.originalListener) { |
| 92 | + listenerToAdd = listenerToAdd.originalListener |
| 93 | + } |
| 94 | + eventEmitter.on(eventName, listenerToAdd) |
| 95 | + |
52 | 96 | try { |
53 | 97 | const res = await listener.apply(eventBus, args2) |
54 | | - |
55 | 98 | ok(res) |
56 | | - |
57 | | - return res |
58 | 99 | } catch (error) { |
59 | 100 | nok(error) |
60 | 101 | } |
61 | 102 | } |
62 | 103 |
|
| 104 | + Object.defineProperty(newListener, "__isSubscribersExecutionWrapper", { |
| 105 | + value: true, |
| 106 | + configurable: true, |
| 107 | + enumerable: false, |
| 108 | + }) |
| 109 | + Object.defineProperty(newListener, "originalListener", { |
| 110 | + value: listener, |
| 111 | + configurable: true, |
| 112 | + enumerable: false, |
| 113 | + }) |
63 | 114 | eventEmitter.on(eventName, newListener) |
64 | 115 | }) |
65 | 116 | } |
66 | 117 |
|
67 | 118 | const subscribersPromise = Promise.all(subscriberPromises).finally(() => { |
68 | 119 | // Clear the timeout since events have been fired and handled |
69 | | - if (timeoutId !== null) { |
70 | | - clearTimeout(timeoutId) |
71 | | - } |
72 | | - |
73 | | - // Restore original event listeners |
74 | | - eventEmitter.removeAllListeners(eventName) |
75 | | - originalListeners.forEach((listener) => { |
76 | | - eventEmitter.on(eventName, listener as (...args: any) => void) |
77 | | - }) |
| 120 | + clearTimeout() |
78 | 121 | }) |
79 | 122 |
|
80 | 123 | // Race between the subscribers and the timeout |
81 | 124 | return Promise.race([subscribersPromise, timeoutPromise]) |
82 | 125 | } |
| 126 | + |
| 127 | +/** |
| 128 | + * Allows you to wait for all subscribers to execute for a given event. |
| 129 | + * It ensures that concurrent waits for the same event are queued and executed sequentially. |
| 130 | + * |
| 131 | + * @param eventName - The name of the event to wait for. |
| 132 | + * @param eventBus - The event bus instance. |
| 133 | + * @param options - Options including timeout. |
| 134 | + */ |
| 135 | +export const waitSubscribersExecution = ( |
| 136 | + eventName: string | symbol, |
| 137 | + eventBus: EventBus, |
| 138 | + options?: WaitSubscribersExecutionOptions |
| 139 | +): Promise<any> => { |
| 140 | + const chain = waits.get(eventName) || Promise.resolve() |
| 141 | + |
| 142 | + const runner = () => { |
| 143 | + return doWaitSubscribersExecution(eventName, eventBus, options) |
| 144 | + } |
| 145 | + |
| 146 | + const newPromise = chain |
| 147 | + .then(runner) |
| 148 | + .catch(runner) // Still execute the runner on error to prevent cascading tests failing because the previous wait failed |
| 149 | + .finally(() => { |
| 150 | + waits.delete(eventName) |
| 151 | + }) |
| 152 | + |
| 153 | + waits.set(eventName, newPromise) |
| 154 | + |
| 155 | + return newPromise |
| 156 | +} |
0 commit comments