chore: Treat internal event differently, primarely do not display info logs for those events (#8767)

* chore: Treat internal event differently, primarely do not display info log for those events

* revert doc

* add few tests

* only set internal option if present

* revert to previous condition

* start including feedback after discussion

* include feedback

* fix modules integration tests

* fix modules integration tests

* fix event bus local
This commit is contained in:
Adrien de Peretti
2024-08-28 16:46:40 +02:00
committed by GitHub
parent 6cfe9bd874
commit 5bec38538a
27 changed files with 915 additions and 465 deletions
@@ -307,6 +307,8 @@ describe("RedisEventBusService", () => {
JSON.stringify(testGroup2Event2),
])
}
return
})
queue = (eventBus as any).queue_
@@ -337,7 +339,7 @@ describe("RedisEventBusService", () => {
})
describe("worker_", () => {
let result
let result!: any
describe("Successfully processes the jobs", () => {
beforeEach(async () => {
@@ -421,12 +423,16 @@ describe("RedisEventBusService", () => {
})
it("should retry processing when subcribers fail, if configured - final attempt", async () => {
eventBus.subscribe("eventName", async () => Promise.resolve(), {
eventBus.subscribe("eventName", async () => await Promise.resolve(), {
subscriberId: "1",
})
eventBus.subscribe("eventName", async () => Promise.reject("fail1"), {
subscriberId: "2",
})
eventBus.subscribe(
"eventName",
async () => await Promise.reject("fail1"),
{
subscriberId: "2",
}
)
result = await eventBus
.worker_({
@@ -456,12 +462,16 @@ describe("RedisEventBusService", () => {
})
it("should retry processing when subcribers fail, if configured", async () => {
eventBus.subscribe("eventName", async () => Promise.resolve(), {
eventBus.subscribe("eventName", async () => await Promise.resolve(), {
subscriberId: "1",
})
eventBus.subscribe("eventName", async () => Promise.reject("fail1"), {
subscriberId: "2",
})
eventBus.subscribe(
"eventName",
async () => await Promise.reject("fail1"),
{
subscriberId: "2",
}
)
result = await eventBus
.worker_({
@@ -11,7 +11,7 @@ import {
} from "@medusajs/utils"
import { BulkJobOptions, Queue, Worker } from "bullmq"
import { Redis } from "ioredis"
import { BullJob, EventBusRedisModuleOptions } from "../types"
import { BullJob, EventBusRedisModuleOptions, Options } from "../types"
type InjectedDependencies = {
logger: Logger
@@ -87,7 +87,7 @@ export default class RedisEventBusService extends AbstractEventBusModuleService
private buildEvents<T>(
eventsData: Message<T>[],
options: BulkJobOptions = {}
options: Options = {}
): IORedisEventType<T>[] {
const opts = {
// default options
@@ -127,7 +127,7 @@ export default class RedisEventBusService extends AbstractEventBusModuleService
*/
async emit<T = unknown>(
eventsData: Message<T> | Message<T>[],
options: BulkJobOptions & { groupedEventsTTL?: number } = {}
options: Options = {}
): Promise<void> {
let eventsDataArray = Array.isArray(eventsData) ? eventsData : [eventsData]
@@ -169,7 +169,7 @@ export default class RedisEventBusService extends AbstractEventBusModuleService
// This will be helpful in preventing stale data from staying in redis for too long
// in the event the module fails to cleanup events. For long running workflows, setting a much higher
// TTL or even skipping the TTL would be required
this.setExpire(groupId, groupedEventsTTL)
void this.setExpire(groupId, groupedEventsTTL)
const eventsData = this.buildEvents(events, options)
@@ -229,7 +229,7 @@ export default class RedisEventBusService extends AbstractEventBusModuleService
* @return resolves to the results of the subscriber calls.
*/
worker_ = async <T>(job: BullJob<T>): Promise<unknown> => {
const { data, name } = job
const { data, name, opts } = job
const eventSubscribers = this.eventToSubscribersMap.get(name) || []
const wildcardSubscribers = this.eventToSubscribersMap.get("*") || []
@@ -250,18 +250,20 @@ export default class RedisEventBusService extends AbstractEventBusModuleService
const isFinalAttempt = currentAttempt === configuredAttempts
if (isRetry) {
if (isFinalAttempt) {
this.logger_.info(`Final retry attempt for ${name}`)
}
if (!opts.internal) {
if (isRetry) {
if (isFinalAttempt) {
this.logger_.info(`Final retry attempt for ${name}`)
}
this.logger_.info(
`Retrying ${name} which has ${eventSubscribers.length} subscribers (${subscribersInCurrentAttempt.length} of them failed)`
)
} else {
this.logger_.info(
`Processing ${name} which has ${eventSubscribers.length} subscribers`
)
this.logger_.info(
`Retrying ${name} which has ${eventSubscribers.length} subscribers (${subscribersInCurrentAttempt.length} of them failed)`
)
} else {
this.logger_.info(
`Processing ${name} which has ${eventSubscribers.length} subscribers`
)
}
}
const completedSubscribersInCurrentAttempt: string[] = []
@@ -315,7 +317,7 @@ export default class RedisEventBusService extends AbstractEventBusModuleService
this.logger_.warn(errorMessage)
return Promise.reject(Error(errorMessage))
throw Error(errorMessage)
}
if (didSubscribersFail && !isFinalAttempt) {
@@ -325,6 +327,6 @@ export default class RedisEventBusService extends AbstractEventBusModuleService
)
}
return Promise.resolve(subscribersResult)
return subscribersResult
}
}