From c21407afbe69d7547df33938ae46f69868cf1432 Mon Sep 17 00:00:00 2001 From: Oli Juhl <59018053+olivermrbl@users.noreply.github.com> Date: Thu, 8 Aug 2024 14:53:42 +0200 Subject: [PATCH] fix(event-bus-redis): Consume job correct in worker (#8444) * fix(event-bus-redis): Consume job correct in worker * fix: Unit tests * fix event passed to subscriber * typo --- .../src/services/__tests__/event-bus.ts | 18 ++++---- .../src/services/event-bus-redis.ts | 44 +++++++++++++------ .../src/services/product-module-service.ts | 2 +- 3 files changed, 41 insertions(+), 23 deletions(-) diff --git a/packages/modules/event-bus-redis/src/services/__tests__/event-bus.ts b/packages/modules/event-bus-redis/src/services/__tests__/event-bus.ts index 8a788d3c7f..04b504ba37 100644 --- a/packages/modules/event-bus-redis/src/services/__tests__/event-bus.ts +++ b/packages/modules/event-bus-redis/src/services/__tests__/event-bus.ts @@ -112,7 +112,7 @@ describe("RedisEventBusService", () => { expect(queue.addBulk).toHaveBeenCalledWith([ { name: "eventName", - data: { hi: "1234" }, + data: { data: { hi: "1234" } }, opts: { attempts: 1, removeOnComplete: true, @@ -132,7 +132,7 @@ describe("RedisEventBusService", () => { expect(queue.addBulk).toHaveBeenCalledWith([ { name: "eventName", - data: { hi: "1234" }, + data: { data: { hi: "1234" } }, opts: { attempts: 3, backoff: 5000, @@ -176,7 +176,7 @@ describe("RedisEventBusService", () => { expect(queue.addBulk).toHaveBeenCalledWith([ { name: "eventName", - data: { hi: "1234" }, + data: { data: { hi: "1234" } }, opts: { attempts: 3, backoff: 5000, @@ -219,7 +219,7 @@ describe("RedisEventBusService", () => { expect(queue.addBulk).toHaveBeenCalledWith([ { name: "eventName", - data: { hi: "1234" }, + data: { data: { hi: "1234" } }, opts: { attempts: 1, removeOnComplete: 5, @@ -360,7 +360,8 @@ describe("RedisEventBusService", () => { // TODO: The typing for this is all over the place await eventBus.worker_({ - data: { eventName: "eventName", data: { test: 1 } }, + name: "eventName", + data: { data: { test: 1 } }, opts: { attempts: 1 }, } as any) @@ -393,7 +394,8 @@ describe("RedisEventBusService", () => { }) result = await eventBus.worker_({ - data: { eventName: "eventName", data: { test: 1 } }, + name: "eventName", + data: { data: { test: 1 } }, opts: { attempts: 1 }, update: (data) => data, } as any) @@ -428,8 +430,8 @@ describe("RedisEventBusService", () => { result = await eventBus .worker_({ + name: "eventName", data: { - eventName: "eventName", data: {}, completedSubscriberIds: ["1"], }, @@ -463,8 +465,8 @@ describe("RedisEventBusService", () => { result = await eventBus .worker_({ + name: "eventName", data: { - eventName: "eventName", data: {}, completedSubscriberIds: ["1"], }, diff --git a/packages/modules/event-bus-redis/src/services/event-bus-redis.ts b/packages/modules/event-bus-redis/src/services/event-bus-redis.ts index 0d2cc9c335..7edfd230fb 100644 --- a/packages/modules/event-bus-redis/src/services/event-bus-redis.ts +++ b/packages/modules/event-bus-redis/src/services/event-bus-redis.ts @@ -1,5 +1,5 @@ import { InternalModuleDeclaration } from "@medusajs/modules-sdk" -import { Logger, Message, Event } from "@medusajs/types" +import { Event, Logger, Message } from "@medusajs/types" import { AbstractEventBusModuleService, isPresent, @@ -14,7 +14,9 @@ type InjectedDependencies = { eventBusRedisConnection: Redis } -type IORedisEventType = Event & { +type IORedisEventType = { + name: string + data: Omit, "name"> // See comment in `buildEvents` method opts: BulkJobOptions } @@ -93,15 +95,22 @@ export default class RedisEventBusService extends AbstractEventBusModuleService } return eventsData.map((eventData) => { - const { options, ...eventBody } = eventData + // We want to preserve event data + metadata. However, bullmq only allows for a single data field. + // Therefore, upon adding jobs to the queue we will serialize the event data and metadata into a single field + // and upon processing the job, we will deserialize it back into the original format expected by the subscribers. + const event = { + data: eventData.data, + metadata: eventData.metadata, + } return { - ...eventBody, + data: event, + name: eventData.name, opts: { // options for event group ...opts, // options for a particular event - ...options, + ...eventData.options, }, } }) @@ -216,8 +225,8 @@ export default class RedisEventBusService extends AbstractEventBusModuleService * @return resolves to the results of the subscriber calls. */ worker_ = async (job: BullJob): Promise => { - const { eventName, data } = job.data - const eventSubscribers = this.eventToSubscribersMap.get(eventName) || [] + const { data, name } = job + const eventSubscribers = this.eventToSubscribersMap.get(name) || [] const wildcardSubscribers = this.eventToSubscribersMap.get("*") || [] const allSubscribers = eventSubscribers.concat(wildcardSubscribers) @@ -239,15 +248,15 @@ export default class RedisEventBusService extends AbstractEventBusModuleService if (isRetry) { if (isFinalAttempt) { - this.logger_.info(`Final retry attempt for ${eventName}`) + this.logger_.info(`Final retry attempt for ${name}`) } this.logger_.info( - `Retrying ${eventName} which has ${eventSubscribers.length} subscribers (${subscribersInCurrentAttempt.length} of them failed)` + `Retrying ${name} which has ${eventSubscribers.length} subscribers (${subscribersInCurrentAttempt.length} of them failed)` ) } else { this.logger_.info( - `Processing ${eventName} which has ${eventSubscribers.length} subscribers` + `Processing ${name} which has ${eventSubscribers.length} subscribers` ) } @@ -255,7 +264,14 @@ export default class RedisEventBusService extends AbstractEventBusModuleService const subscribersResult = await Promise.all( subscribersInCurrentAttempt.map(async ({ id, subscriber }) => { - return await subscriber(data) + // De-serialize the event data and metadata from a single field into the original format expected by the subscribers + const event = { + name, + data: data.data, + metadata: data.metadata, + } + + return await subscriber(event) .then(async (data) => { // For every subscriber that completes successfully, add their id to the list of completed subscribers completedSubscribersInCurrentAttempt.push(id) @@ -263,7 +279,7 @@ export default class RedisEventBusService extends AbstractEventBusModuleService }) .catch((err) => { this.logger_.warn( - `An error occurred while processing ${eventName}: ${err}` + `An error occurred while processing ${name}: ${err}` ) return err }) @@ -291,7 +307,7 @@ export default class RedisEventBusService extends AbstractEventBusModuleService await job.updateData(job.data) - const errorMessage = `One or more subscribers of ${eventName} failed. Retrying...` + const errorMessage = `One or more subscribers of ${name} failed. Retrying...` this.logger_.warn(errorMessage) @@ -301,7 +317,7 @@ export default class RedisEventBusService extends AbstractEventBusModuleService if (didSubscribersFail && !isFinalAttempt) { // If retrying is not configured, we log a warning to allow server admins to recover manually this.logger_.warn( - `One or more subscribers of ${eventName} failed. Retrying is not configured. Use 'attempts' option when emitting events.` + `One or more subscribers of ${name} failed. Retrying is not configured. Use 'attempts' option when emitting events.` ) } diff --git a/packages/modules/product/src/services/product-module-service.ts b/packages/modules/product/src/services/product-module-service.ts index b32d5e0fb3..17bedb30b4 100644 --- a/packages/modules/product/src/services/product-module-service.ts +++ b/packages/modules/product/src/services/product-module-service.ts @@ -8,10 +8,10 @@ import { ProductTypes, } from "@medusajs/types" import { - Image as ProductImage, Product, ProductCategory, ProductCollection, + Image as ProductImage, ProductOption, ProductOptionValue, ProductTag,