feat: Event emitting part 1/N (Fulfillment) (#7391)
**What** Add support for event emitting in the fulfillment module **NOTE** It does not include the review of the events for the abstract module factory if the method is not implemented in the module itself and rely on the default implementation
This commit is contained in:
@@ -1,39 +1,4 @@
|
||||
import { Context, EventBusTypes } from "@medusajs/types"
|
||||
import { CommonEvents } from "./common-events"
|
||||
|
||||
/**
|
||||
* Build messages from message data to be consumed by the event bus and emitted to the consumer
|
||||
* @param MessageFormat
|
||||
* @param options
|
||||
*/
|
||||
export function buildEventMessages<T>(
|
||||
messageData:
|
||||
| EventBusTypes.MessageFormat<T>
|
||||
| EventBusTypes.MessageFormat<T>[],
|
||||
options?: Record<string, unknown>
|
||||
): EventBusTypes.Message<T>[] {
|
||||
const messageData_ = Array.isArray(messageData) ? messageData : [messageData]
|
||||
const messages: EventBusTypes.Message<any>[] = []
|
||||
|
||||
messageData_.map((data) => {
|
||||
const data_ = Array.isArray(data.data) ? data.data : [data.data]
|
||||
data_.forEach((bodyData) => {
|
||||
const message = composeMessage(data.eventName, {
|
||||
data: bodyData,
|
||||
service: data.metadata.service,
|
||||
entity: data.metadata.object,
|
||||
action: data.metadata.action,
|
||||
context: {
|
||||
eventGroupId: data.metadata.eventGroupId,
|
||||
} as Context,
|
||||
options,
|
||||
})
|
||||
messages.push(message)
|
||||
})
|
||||
})
|
||||
|
||||
return messages
|
||||
}
|
||||
|
||||
/**
|
||||
* Helper function to compose and normalize a Message to be emitted by EventBus Module
|
||||
@@ -48,27 +13,29 @@ export function composeMessage(
|
||||
{
|
||||
data,
|
||||
service,
|
||||
entity,
|
||||
object,
|
||||
action,
|
||||
context = {},
|
||||
options,
|
||||
}: {
|
||||
data: unknown
|
||||
data: any
|
||||
service: string
|
||||
entity: string
|
||||
object: string
|
||||
action?: string
|
||||
context?: Context
|
||||
options?: Record<string, unknown>
|
||||
options?: Record<string, any>
|
||||
}
|
||||
): EventBusTypes.Message {
|
||||
const act = action || eventName.split(".").pop()
|
||||
if (!action && !Object.values(CommonEvents).includes(act as CommonEvents)) {
|
||||
if (
|
||||
!action /* && !Object.values(CommonEvents).includes(act as CommonEvents)*/
|
||||
) {
|
||||
throw new Error("Action is required if eventName is not a CommonEvent")
|
||||
}
|
||||
|
||||
const metadata: EventBusTypes.MessageBody["metadata"] = {
|
||||
service,
|
||||
object: entity,
|
||||
object,
|
||||
action: act!,
|
||||
}
|
||||
|
||||
|
||||
@@ -1,11 +1,12 @@
|
||||
import {
|
||||
Context,
|
||||
EventBusTypes,
|
||||
IMessageAggregator,
|
||||
Message,
|
||||
MessageAggregatorFormat,
|
||||
} from "@medusajs/types"
|
||||
|
||||
import { buildEventMessages } from "./build-event-messages"
|
||||
import { composeMessage } from "./build-event-messages"
|
||||
|
||||
export class MessageAggregator implements IMessageAggregator {
|
||||
private messages: Message[]
|
||||
@@ -28,11 +29,25 @@ export class MessageAggregator implements IMessageAggregator {
|
||||
|
||||
saveRawMessageData<T>(
|
||||
messageData:
|
||||
| EventBusTypes.MessageFormat<T>
|
||||
| EventBusTypes.MessageFormat<T>[],
|
||||
options?: Record<string, unknown>
|
||||
| EventBusTypes.RawMessageFormat<T>
|
||||
| EventBusTypes.RawMessageFormat<T>[],
|
||||
{
|
||||
options,
|
||||
sharedContext,
|
||||
}: { options?: Record<string, unknown>; sharedContext?: Context } = {}
|
||||
): void {
|
||||
this.save(buildEventMessages(messageData, options))
|
||||
const messages = Array.isArray(messageData) ? messageData : [messageData]
|
||||
const composedMessages = messages.map((message) => {
|
||||
return composeMessage(message.eventName, {
|
||||
data: message.data,
|
||||
service: message.service,
|
||||
object: message.object,
|
||||
action: message.action,
|
||||
options,
|
||||
context: sharedContext,
|
||||
})
|
||||
})
|
||||
this.save(composedMessages)
|
||||
}
|
||||
|
||||
getMessages(format?: MessageAggregatorFormat): {
|
||||
|
||||
@@ -58,7 +58,8 @@ type ReturnType<TNames extends string[]> = TNames extends [
|
||||
* @param names
|
||||
*/
|
||||
export function buildEventNamesFromEntityName<TNames extends string[]>(
|
||||
names: TNames
|
||||
names: TNames,
|
||||
prefix?: string
|
||||
): ReturnType<TNames> {
|
||||
const events = {}
|
||||
|
||||
@@ -69,16 +70,15 @@ export function buildEventNamesFromEntityName<TNames extends string[]>(
|
||||
|
||||
if (i === 0) {
|
||||
for (const event of Object.values(CommonEvents) as string[]) {
|
||||
events[event] = `${kebabCaseName}.${event}`
|
||||
events[event] = `${prefix ? prefix + "." : ""}${kebabCaseName}.${event}`
|
||||
}
|
||||
continue
|
||||
}
|
||||
|
||||
for (const event of Object.values(CommonEvents) as string[]) {
|
||||
events[`${snakedCaseName}_${event}`] =
|
||||
`${kebabCaseName}.${event}` as `${KebabCase<
|
||||
typeof name
|
||||
>}.${typeof event}`
|
||||
events[`${snakedCaseName}_${event}`] = `${
|
||||
prefix ? prefix + "." : ""
|
||||
}${kebabCaseName}.${event}` as `${KebabCase<typeof name>}.${typeof event}`
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user