feat: restructure events payload (#8143)

* refactor: restructure events payload

Breaking change: This PR changes the event payload accepted by the event
listeners

* refactor: fix failing tests and implement feedback

* add integration tests

* fix timeout

---------

Co-authored-by: Adrien de Peretti <adrien.deperetti@gmail.com>
This commit is contained in:
Harminder Virk
2024-07-16 17:09:16 +01:00
committed by GitHub
co-authored by Adrien de Peretti
parent 5813216c88
commit f579f0b3be
26 changed files with 194 additions and 111 deletions
@@ -30,14 +30,14 @@ describe("LocalEventBusService", () => {
eventEmitter.emit = jest.fn((data) => data)
await eventBus.emit({
eventName: "eventName",
name: "eventName",
data: { hi: "1234" },
})
expect(eventEmitter.emit).toHaveBeenCalledTimes(1)
expect(eventEmitter.emit).toHaveBeenCalledWith("eventName", {
data: { hi: "1234" },
eventName: "eventName",
name: "eventName",
})
})
@@ -45,18 +45,18 @@ describe("LocalEventBusService", () => {
eventEmitter.emit = jest.fn((data) => data)
await eventBus.emit([
{ eventName: "event-1", data: { hi: "1234" } },
{ eventName: "event-2", data: { hi: "5678" } },
{ name: "event-1", data: { hi: "1234" } },
{ name: "event-2", data: { hi: "5678" } },
])
expect(eventEmitter.emit).toHaveBeenCalledTimes(2)
expect(eventEmitter.emit).toHaveBeenCalledWith("event-1", {
data: { hi: "1234" },
eventName: "event-1",
name: "event-1",
})
expect(eventEmitter.emit).toHaveBeenCalledWith("event-2", {
data: { hi: "5678" },
eventName: "event-2",
name: "event-2",
})
})
@@ -65,7 +65,7 @@ describe("LocalEventBusService", () => {
eventEmitter.emit = jest.fn((data) => data)
await eventBus.emit({
eventName: "test-event",
name: "test-event",
data: {
test: "1234",
},
@@ -79,7 +79,7 @@ describe("LocalEventBusService", () => {
expect(groupEventFn).toHaveBeenCalledWith("test", {
data: { test: "1234" },
metadata: { eventGroupId: "test" },
eventName: "test-event",
name: "test-event",
})
jest.clearAllMocks()
@@ -89,12 +89,12 @@ describe("LocalEventBusService", () => {
eventBus.emit([
{
eventName: "test-event",
name: "test-event",
data: { test: "1234" },
metadata: { eventGroupId: "test" },
},
{
eventName: "test-event",
name: "test-event",
data: { test: "test-1" },
},
])
@@ -102,18 +102,18 @@ describe("LocalEventBusService", () => {
expect(groupEventFn).toHaveBeenCalledTimes(1)
expect((eventBus as any).groupedEventsMap_.get("test")).toEqual([
expect.objectContaining({ eventName: "test-event" }),
expect.objectContaining({ eventName: "test-event" }),
expect.objectContaining({ name: "test-event" }),
expect.objectContaining({ name: "test-event" }),
])
await eventBus.emit({
eventName: "test-event",
name: "test-event",
data: { test: "1234" },
metadata: { eventGroupId: "test-2" },
})
expect((eventBus as any).groupedEventsMap_.get("test-2")).toEqual([
expect.objectContaining({ eventName: "test-event" }),
expect.objectContaining({ name: "test-event" }),
])
})
@@ -122,32 +122,32 @@ describe("LocalEventBusService", () => {
await eventBus.emit([
{
eventName: "event-1",
name: "event-1",
data: { test: "1" },
metadata: { eventGroupId: "group-1" },
},
{
eventName: "event-2",
name: "event-2",
data: { test: "2" },
metadata: { eventGroupId: "group-1" },
},
{
eventName: "event-1",
name: "event-1",
data: { test: "1" },
metadata: { eventGroupId: "group-2" },
},
{
eventName: "event-2",
name: "event-2",
data: { test: "2" },
metadata: { eventGroupId: "group-2" },
},
{ eventName: "event-1", data: { test: "1" } },
{ name: "event-1", data: { test: "1" } },
])
expect(eventEmitter.emit).toHaveBeenCalledTimes(1)
expect(eventEmitter.emit).toHaveBeenCalledWith("event-1", {
data: { test: "1" },
eventName: "event-1",
name: "event-1",
})
expect((eventBus as any).groupedEventsMap_.get("group-1")).toHaveLength(
@@ -171,12 +171,12 @@ describe("LocalEventBusService", () => {
expect(eventEmitter.emit).toHaveBeenCalledTimes(2)
expect(eventEmitter.emit).toHaveBeenCalledWith("event-1", {
data: { test: "1" },
eventName: "event-1",
name: "event-1",
metadata: { eventGroupId: "group-1" },
})
expect(eventEmitter.emit).toHaveBeenCalledWith("event-2", {
data: { test: "2" },
eventName: "event-2",
name: "event-2",
metadata: { eventGroupId: "group-1" },
})
})
@@ -187,12 +187,12 @@ describe("LocalEventBusService", () => {
await eventBus.emit([
{
eventName: "event-1",
name: "event-1",
data: { test: "1" },
metadata: { eventGroupId: "group-1" },
},
{
eventName: "event-1",
name: "event-1",
data: { test: "1" },
metadata: { eventGroupId: "group-2" },
},
@@ -3,7 +3,7 @@ import {
EventBusTypes,
Logger,
Message,
MessageBody,
Event,
Subscriber,
} from "@medusajs/types"
import { AbstractEventBusModuleService } from "@medusajs/utils"
@@ -45,11 +45,11 @@ export default class LocalEventBusService extends AbstractEventBusModuleService
for (const eventData of normalizedEventsData) {
const eventListenersCount = this.eventEmitter_.listenerCount(
eventData.eventName
eventData.name
)
this.logger_?.info(
`Processing ${eventData.eventName} which has ${eventListenersCount} subscribers`
`Processing ${eventData.name} which has ${eventListenersCount} subscribers`
)
if (eventListenersCount === 0) {
@@ -73,7 +73,7 @@ export default class LocalEventBusService extends AbstractEventBusModuleService
await this.groupEvent(eventGroupId, eventData)
} else {
const { options, ...eventBody } = eventData
this.eventEmitter_.emit(eventData.eventName, eventBody)
this.eventEmitter_.emit(eventData.name, eventBody)
}
}
@@ -95,7 +95,7 @@ export default class LocalEventBusService extends AbstractEventBusModuleService
for (const event of groupedEvents) {
const { options, ...eventBody } = event
this.eventEmitter_.emit(event.eventName, eventBody)
this.eventEmitter_.emit(event.name, eventBody)
}
this.clearGroupedEvents(eventGroupId)
@@ -108,7 +108,7 @@ export default class LocalEventBusService extends AbstractEventBusModuleService
subscribe(event: string | symbol, subscriber: Subscriber): this {
const randId = ulid()
this.storeSubscribers({ event, subscriberId: randId, subscriber })
this.eventEmitter_.on(event, async (data: MessageBody) => {
this.eventEmitter_.on(event, async (data: Event) => {
try {
await subscriber(data)
} catch (e) {