feat(medusa,event-bus-local,event-bus-redis): Event Bus modules (#2599)
This commit is contained in:
@@ -0,0 +1,70 @@
|
||||
import LocalEventBusService from "../event-bus-local"
|
||||
|
||||
jest.genMockFromModule("events")
|
||||
jest.mock("events")
|
||||
|
||||
const loggerMock = {
|
||||
info: jest.fn().mockReturnValue(console.log),
|
||||
warn: jest.fn().mockReturnValue(console.log),
|
||||
error: jest.fn().mockReturnValue(console.log),
|
||||
}
|
||||
|
||||
const moduleDeps = {
|
||||
logger: loggerMock,
|
||||
}
|
||||
|
||||
describe("LocalEventBusService", () => {
|
||||
let eventBus
|
||||
|
||||
describe("emit", () => {
|
||||
describe("Successfully emits events", () => {
|
||||
beforeEach(() => {
|
||||
jest.clearAllMocks()
|
||||
})
|
||||
|
||||
it("Emits an event", () => {
|
||||
eventBus = new LocalEventBusService(
|
||||
moduleDeps,
|
||||
{},
|
||||
{
|
||||
resources: "shared",
|
||||
}
|
||||
)
|
||||
|
||||
eventBus.eventEmitter_.emit.mockImplementationOnce((data) => data)
|
||||
|
||||
eventBus.emit("eventName", { hi: "1234" })
|
||||
|
||||
expect(eventBus.eventEmitter_.emit).toHaveBeenCalledTimes(1)
|
||||
expect(eventBus.eventEmitter_.emit).toHaveBeenCalledWith("eventName", {
|
||||
hi: "1234",
|
||||
})
|
||||
})
|
||||
|
||||
it("Emits multiple events", () => {
|
||||
eventBus = new LocalEventBusService(
|
||||
moduleDeps,
|
||||
{},
|
||||
{
|
||||
resources: "shared",
|
||||
}
|
||||
)
|
||||
|
||||
eventBus.eventEmitter_.emit.mockImplementationOnce((data) => data)
|
||||
|
||||
eventBus.emit([
|
||||
{ eventName: "event-1", data: { hi: "1234" } },
|
||||
{ eventName: "event-2", data: { hi: "5678" } },
|
||||
])
|
||||
|
||||
expect(eventBus.eventEmitter_.emit).toHaveBeenCalledTimes(2)
|
||||
expect(eventBus.eventEmitter_.emit).toHaveBeenCalledWith("event-1", {
|
||||
hi: "1234",
|
||||
})
|
||||
expect(eventBus.eventEmitter_.emit).toHaveBeenCalledWith("event-2", {
|
||||
hi: "5678",
|
||||
})
|
||||
})
|
||||
})
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,80 @@
|
||||
import { Logger, MedusaContainer } from "@medusajs/modules-sdk"
|
||||
import { EmitData, Subscriber } from "@medusajs/types"
|
||||
import { AbstractEventBusModuleService } from "@medusajs/utils"
|
||||
import { EventEmitter } from "events"
|
||||
|
||||
type InjectedDependencies = {
|
||||
logger: Logger
|
||||
}
|
||||
|
||||
const eventEmitter = new EventEmitter()
|
||||
|
||||
export default class LocalEventBusService extends AbstractEventBusModuleService {
|
||||
protected readonly logger_: Logger
|
||||
protected readonly eventEmitter_: EventEmitter
|
||||
|
||||
constructor({ logger }: MedusaContainer & InjectedDependencies) {
|
||||
// @ts-ignore
|
||||
super(...arguments)
|
||||
|
||||
this.logger_ = logger
|
||||
this.eventEmitter_ = eventEmitter
|
||||
}
|
||||
|
||||
async emit<T>(
|
||||
eventName: string,
|
||||
data: T,
|
||||
options: Record<string, unknown>
|
||||
): Promise<void>
|
||||
|
||||
/**
|
||||
* Emit a number of events
|
||||
* @param {EmitData} data - the data to send to the subscriber.
|
||||
*/
|
||||
async emit<T>(data: EmitData<T>[]): Promise<void>
|
||||
|
||||
async emit<T, TInput extends string | EmitData<T>[] = string>(
|
||||
eventOrData: TInput,
|
||||
data?: T,
|
||||
options: Record<string, unknown> = {}
|
||||
): Promise<void> {
|
||||
const isBulkEmit = Array.isArray(eventOrData)
|
||||
|
||||
const events: EmitData[] = isBulkEmit
|
||||
? eventOrData
|
||||
: [{ eventName: eventOrData, data }]
|
||||
|
||||
for (const event of events) {
|
||||
const eventListenersCount = this.eventEmitter_.listenerCount(
|
||||
event.eventName
|
||||
)
|
||||
|
||||
this.logger_.info(
|
||||
`Processing ${event.eventName} which has ${eventListenersCount} subscribers`
|
||||
)
|
||||
|
||||
if (eventListenersCount === 0) {
|
||||
continue
|
||||
}
|
||||
|
||||
|
||||
try {
|
||||
this.eventEmitter_.emit(event.eventName, event.data)
|
||||
} catch (error) {
|
||||
this.logger_.error(
|
||||
`An error occurred while processing ${event.eventName}: ${error}`
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
subscribe(event: string | symbol, subscriber: Subscriber): this {
|
||||
this.eventEmitter_.on(event, subscriber)
|
||||
return this
|
||||
}
|
||||
|
||||
unsubscribe(event: string | symbol, subscriber: Subscriber): this {
|
||||
this.eventEmitter_.off(event, subscriber)
|
||||
return this
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user