Skip to content

Commit b4843f7

Browse files
committed
allow waitSubscribersExecution to wait for X trigger
1 parent 747dadc commit b4843f7

2 files changed

Lines changed: 102 additions & 24 deletions

File tree

packages/medusa-test-utils/src/__tests__/events.spec.ts

Lines changed: 60 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,8 @@
11
import { EventEmitter } from "events"
22
import { waitSubscribersExecution } from "../events"
3+
import { setTimeout } from "timers/promises"
4+
5+
jest.setTimeout(30000)
36

47
// Mock the IEventBusModuleService
58
class MockEventBus {
@@ -31,11 +34,12 @@ describe("waitSubscribersExecution", () => {
3134
describe("with no existing listeners", () => {
3235
it("should resolve when event is fired before timeout", async () => {
3336
const waitPromise = waitSubscribersExecution(TEST_EVENT, eventBus as any)
34-
setTimeout(() => eventBus.emit(TEST_EVENT, "test-data"), 100).unref()
37+
await setTimeout(100)
38+
eventBus.emit(TEST_EVENT, "test-data")
3539

3640
jest.advanceTimersByTime(100)
3741

38-
await expect(waitPromise).resolves.toEqual(["test-data"])
42+
await expect(waitPromise).resolves.toEqual([["test-data"]])
3943
})
4044

4145
it("should reject when timeout is reached before event is fired", async () => {
@@ -70,12 +74,29 @@ describe("waitSubscribersExecution", () => {
7074
`Timeout of ${customTimeout}ms exceeded while waiting for event "${TEST_EVENT}"`
7175
)
7276
})
77+
78+
it("should resolve when event is fired multiple times", async () => {
79+
const waitPromise = waitSubscribersExecution(
80+
TEST_EVENT,
81+
eventBus as any,
82+
{ triggerCount: 2 }
83+
)
84+
eventBus.emit(TEST_EVENT, "test-data")
85+
eventBus.emit(TEST_EVENT, "test-data")
86+
87+
const promisesRes = await waitPromise
88+
const res = promisesRes.pop()
89+
expect(res).toHaveLength(2)
90+
expect(res[0]).toEqual(["test-data"])
91+
expect(res[1]).toEqual(["test-data"])
92+
})
7393
})
7494

7595
describe("with existing listeners", () => {
7696
it("should resolve when all listeners complete successfully", async () => {
77-
const listener = jest.fn().mockImplementation(() => {
78-
return new Promise((resolve) => setTimeout(resolve, 200).unref())
97+
const listener = jest.fn().mockImplementation(async () => {
98+
await setTimeout(200)
99+
return "res"
79100
})
80101

81102
eventBus.eventEmitter_.on(TEST_EVENT, listener)
@@ -132,20 +153,49 @@ describe("waitSubscribersExecution", () => {
132153

133154
expect(listener).not.toHaveBeenCalled()
134155
})
156+
157+
it("should resolve when event is fired multiple times", async () => {
158+
const listener = jest.fn().mockImplementation(async () => {
159+
await setTimeout(200)
160+
return "res"
161+
})
162+
163+
eventBus.eventEmitter_.on(TEST_EVENT, listener)
164+
165+
const waitPromise = waitSubscribersExecution(
166+
TEST_EVENT,
167+
eventBus as any,
168+
{
169+
triggerCount: 2,
170+
}
171+
)
172+
173+
eventBus.emit(TEST_EVENT, "test-data")
174+
eventBus.emit(TEST_EVENT, "test-data")
175+
176+
const promisesRes = await waitPromise
177+
const res = promisesRes.pop()
178+
expect(res).toHaveLength(2)
179+
expect(res[0]).toEqual("res")
180+
expect(res[1]).toEqual("res")
181+
})
135182
})
136183

137184
describe("with multiple listeners", () => {
138185
it("should resolve when all listeners complete", async () => {
139-
const listener1 = jest.fn().mockImplementation(() => {
140-
return new Promise((resolve) => setTimeout(resolve, 100).unref())
186+
const listener1 = jest.fn().mockImplementation(async () => {
187+
await setTimeout(100)
188+
return "res"
141189
})
142190

143-
const listener2 = jest.fn().mockImplementation(() => {
144-
return new Promise((resolve) => setTimeout(resolve, 200).unref())
191+
const listener2 = jest.fn().mockImplementation(async () => {
192+
await setTimeout(200)
193+
return "res"
145194
})
146195

147-
const listener3 = jest.fn().mockImplementation(() => {
148-
return new Promise((resolve) => setTimeout(resolve, 300).unref())
196+
const listener3 = jest.fn().mockImplementation(async () => {
197+
await setTimeout(300)
198+
return "res"
149199
})
150200

151201
eventBus.eventEmitter_.on(TEST_EVENT, listener1)

packages/medusa-test-utils/src/events.ts

Lines changed: 42 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,10 @@ type EventBus = {
55
}
66

77
type WaitSubscribersExecutionOptions = {
8+
/** Timeout in milliseconds for waiting for the event. Defaults to 15000ms. */
89
timeout?: number
10+
/** Number of times the event should be triggered before resolving. Defaults to 1. */
11+
triggerCount?: number
912
}
1013

1114
// Map to hold pending promises for each event.
@@ -41,7 +44,7 @@ const createTimeoutPromise = (
4144
const doWaitSubscribersExecution = (
4245
eventName: string | symbol,
4346
eventBus: EventBus,
44-
{ timeout = 15000 }: WaitSubscribersExecutionOptions = {}
47+
{ timeout = 15000, triggerCount = 1 }: WaitSubscribersExecutionOptions = {}
4548
): Promise<any> => {
4649
const eventEmitter = eventBus.eventEmitter_
4750
const subscriberPromises: Promise<any>[] = []
@@ -50,16 +53,25 @@ const doWaitSubscribersExecution = (
5053
eventName
5154
)
5255

56+
let currentCount = 0
57+
5358
if (!eventEmitter.listeners(eventName).length) {
5459
let ok: (value?: any) => void
5560
const promise = new Promise((resolve) => {
5661
ok = resolve
5762
})
5863
subscriberPromises.push(promise)
5964

65+
let res: any[] = []
6066
const newListener = async (...args: any[]) => {
61-
eventEmitter.removeListener(eventName, newListener)
62-
ok(...args)
67+
currentCount++
68+
res.push(args)
69+
70+
if (currentCount >= triggerCount) {
71+
eventEmitter.removeListener(eventName, newListener)
72+
const res_ = triggerCount === 1 ? res[0] : res
73+
ok(res_)
74+
}
6375
}
6476

6577
Object.defineProperty(newListener, "__isSubscribersExecutionWrapper", {
@@ -83,22 +95,38 @@ const doWaitSubscribersExecution = (
8395
nok = reject
8496
})
8597
subscriberPromises.push(promise)
98+
let res: any[] = []
8699

87100
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-
96101
try {
97-
const res = await listener.apply(eventBus, args2)
98-
ok(res)
102+
const listenerRes = listener.apply(eventBus, args2)
103+
if (typeof listenerRes?.then === "function") {
104+
await listenerRes.then((res_) => {
105+
res.push(res_)
106+
currentCount++
107+
})
108+
} else {
109+
res.push(listenerRes)
110+
currentCount++
111+
}
112+
113+
if (currentCount >= triggerCount) {
114+
const res_ = triggerCount === 1 ? res[0] : res
115+
ok(res_)
116+
}
99117
} catch (error) {
100118
nok(error)
101119
}
120+
121+
if (currentCount >= triggerCount) {
122+
// As soon as the subscriber is executed the required number of times, we restore the original listener
123+
eventEmitter.removeListener(eventName, newListener)
124+
let listenerToAdd = listener
125+
while (listenerToAdd.originalListener) {
126+
listenerToAdd = listenerToAdd.originalListener
127+
}
128+
eventEmitter.on(eventName, listenerToAdd)
129+
}
102130
}
103131

104132
Object.defineProperty(newListener, "__isSubscribersExecutionWrapper", {
@@ -130,7 +158,7 @@ const doWaitSubscribersExecution = (
130158
*
131159
* @param eventName - The name of the event to wait for.
132160
* @param eventBus - The event bus instance.
133-
* @param options - Options including timeout.
161+
* @param options - Options including timeout and triggerCount.
134162
*/
135163
export const waitSubscribersExecution = (
136164
eventName: string | symbol,

0 commit comments

Comments
 (0)