feat(medusa): Extract cron job logic to its own service (#2821)
**What** Extract cron job logic from the `EventBusService` to its own service `JobSchedulerService` **Why** Preliminary step to extracting the event bus into a module **Tests** Tested with a local installation of Medusa Resolved CORE-918
This commit is contained in:
@@ -3,10 +3,10 @@ const inventorySync = async (container, options) => {
|
||||
return
|
||||
} else {
|
||||
const brightpearlService = container.resolve("brightpearlService")
|
||||
const eventBus = container.resolve("eventBusService")
|
||||
const jobSchedulerService = container.resolve("jobSchedulerService")
|
||||
try {
|
||||
const pattern = options.inventory_sync_cron
|
||||
eventBus.createCronJob("inventory-sync", {}, pattern, () =>
|
||||
jobSchedulerService.create("inventory-sync", {}, pattern, () =>
|
||||
brightpearlService.syncInventory()
|
||||
)
|
||||
} catch (err) {
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
import Bull from "bull"
|
||||
import { MockRepository, MockManager } from "medusa-test-utils"
|
||||
import EventBusService from "../event-bus"
|
||||
import { MockManager, MockRepository } from "medusa-test-utils"
|
||||
import config from "../../loaders/config"
|
||||
import EventBusService from "../event-bus"
|
||||
|
||||
jest.genMockFromModule("bull")
|
||||
jest.mock("bull")
|
||||
@@ -36,7 +36,7 @@ describe("EventBusService", () => {
|
||||
})
|
||||
|
||||
it("creates bull queue", () => {
|
||||
expect(Bull).toHaveBeenCalledTimes(2)
|
||||
expect(Bull).toHaveBeenCalledTimes(1)
|
||||
expect(Bull).toHaveBeenCalledWith("EventBusService:queue", {
|
||||
createClient: expect.any(Function),
|
||||
})
|
||||
@@ -97,7 +97,8 @@ describe("EventBusService", () => {
|
||||
})
|
||||
|
||||
describe("emit", () => {
|
||||
let eventBus, job
|
||||
let eventBus
|
||||
let job
|
||||
describe("successfully adds job to queue", () => {
|
||||
beforeAll(() => {
|
||||
jest.resetAllMocks()
|
||||
@@ -126,7 +127,8 @@ describe("EventBusService", () => {
|
||||
})
|
||||
|
||||
describe("worker", () => {
|
||||
let eventBus, result
|
||||
let eventBus
|
||||
let result
|
||||
describe("successfully runs the worker", () => {
|
||||
beforeAll(async () => {
|
||||
jest.resetAllMocks()
|
||||
@@ -134,11 +136,14 @@ describe("EventBusService", () => {
|
||||
find: () => Promise.resolve([]),
|
||||
})
|
||||
|
||||
eventBus = new EventBusService({
|
||||
manager: MockManager,
|
||||
stagedJobRepository,
|
||||
logger: loggerMock,
|
||||
}, {})
|
||||
eventBus = new EventBusService(
|
||||
{
|
||||
manager: MockManager,
|
||||
stagedJobRepository,
|
||||
logger: loggerMock,
|
||||
},
|
||||
{}
|
||||
)
|
||||
eventBus.subscribe("eventName", () => Promise.resolve("hi"))
|
||||
result = await eventBus.worker_({
|
||||
data: { eventName: "eventName", data: {} },
|
||||
|
||||
@@ -0,0 +1,95 @@
|
||||
import Bull from "bull"
|
||||
import config from "../../loaders/config"
|
||||
import JobSchedulerService from "../job-scheduler"
|
||||
|
||||
jest.genMockFromModule("bull")
|
||||
jest.mock("bull")
|
||||
jest.mock("../../loaders/config")
|
||||
|
||||
config.redisURI = "testhost"
|
||||
|
||||
const loggerMock = {
|
||||
info: jest.fn().mockReturnValue(console.log),
|
||||
warn: jest.fn().mockReturnValue(console.log),
|
||||
error: jest.fn().mockReturnValue(console.log),
|
||||
}
|
||||
|
||||
describe("JobSchedulerService", () => {
|
||||
describe("constructor", () => {
|
||||
let jobScheduler
|
||||
beforeAll(() => {
|
||||
jest.resetAllMocks()
|
||||
|
||||
jobScheduler = new JobSchedulerService({
|
||||
logger: loggerMock,
|
||||
})
|
||||
})
|
||||
|
||||
it("creates bull queue", () => {
|
||||
expect(Bull).toHaveBeenCalledTimes(1)
|
||||
expect(Bull).toHaveBeenCalledWith("scheduled-jobs:queue", {
|
||||
createClient: expect.any(Function),
|
||||
})
|
||||
})
|
||||
})
|
||||
|
||||
describe("create", () => {
|
||||
let jobScheduler
|
||||
describe("successfully creates scheduled job and add handler", () => {
|
||||
beforeAll(() => {
|
||||
jest.resetAllMocks()
|
||||
|
||||
jobScheduler = new JobSchedulerService({
|
||||
logger: loggerMock,
|
||||
})
|
||||
|
||||
jobScheduler.create(
|
||||
"eventName",
|
||||
{ data: "test" },
|
||||
"* * * * *",
|
||||
() => "test"
|
||||
)
|
||||
})
|
||||
|
||||
it("added the handler to the job queue", () => {
|
||||
expect(jobScheduler.handlers_.get("eventName").length).toEqual(1)
|
||||
})
|
||||
})
|
||||
|
||||
describe("scheduledJobWorker", () => {
|
||||
let jobScheduler
|
||||
let result
|
||||
describe("successfully runs the worker", () => {
|
||||
beforeAll(async () => {
|
||||
jest.resetAllMocks()
|
||||
|
||||
jobScheduler = new JobSchedulerService(
|
||||
{
|
||||
logger: loggerMock,
|
||||
},
|
||||
{}
|
||||
)
|
||||
|
||||
jobScheduler.create("eventName", { data: "test" }, "* * * * *", () =>
|
||||
Promise.resolve("hi")
|
||||
)
|
||||
|
||||
result = await jobScheduler.scheduledJobsWorker({
|
||||
data: { eventName: "eventName", data: {} },
|
||||
})
|
||||
})
|
||||
|
||||
it("calls logger", () => {
|
||||
expect(loggerMock.info).toHaveBeenCalled()
|
||||
expect(loggerMock.info).toHaveBeenCalledWith(
|
||||
"Processing scheduled job: eventName"
|
||||
)
|
||||
})
|
||||
|
||||
it("returns array with hi", async () => {
|
||||
expect(result).toEqual(["hi"])
|
||||
})
|
||||
})
|
||||
})
|
||||
})
|
||||
})
|
||||
@@ -1,15 +1,17 @@
|
||||
import Bull from "bull"
|
||||
import Redis from "ioredis"
|
||||
import { EntityManager } from "typeorm"
|
||||
import { ConfigModule, Logger } from "../types/global"
|
||||
import { StagedJobRepository } from "../repositories/staged-job"
|
||||
import { StagedJob } from "../models"
|
||||
import { StagedJobRepository } from "../repositories/staged-job"
|
||||
import { ConfigModule, Logger } from "../types/global"
|
||||
import { sleep } from "../utils/sleep"
|
||||
import JobSchedulerService from "./job-scheduler"
|
||||
|
||||
type InjectedDependencies = {
|
||||
manager: EntityManager
|
||||
logger: Logger
|
||||
stagedJobRepository: typeof StagedJobRepository
|
||||
jobSchedulerService: JobSchedulerService
|
||||
redisClient: Redis.Redis
|
||||
redisSubscriber: Redis.Redis
|
||||
}
|
||||
@@ -34,11 +36,10 @@ export default class EventBusService {
|
||||
protected readonly manager_: EntityManager
|
||||
protected readonly logger_: Logger
|
||||
protected readonly stagedJobRepository_: typeof StagedJobRepository
|
||||
protected readonly jobSchedulerService_: JobSchedulerService
|
||||
protected readonly observers_: Map<string | symbol, Subscriber[]>
|
||||
protected readonly cronHandlers_: Map<string | symbol, Subscriber[]>
|
||||
protected readonly redisClient_: Redis.Redis
|
||||
protected readonly redisSubscriber_: Redis.Redis
|
||||
protected readonly cronQueue_: Bull
|
||||
protected queue_: Bull
|
||||
protected shouldEnqueuerRun: boolean
|
||||
protected transactionManager_: EntityManager | undefined
|
||||
@@ -79,14 +80,10 @@ export default class EventBusService {
|
||||
|
||||
this.observers_ = new Map()
|
||||
this.queue_ = new Bull(`${this.constructor.name}:queue`, opts)
|
||||
this.cronHandlers_ = new Map()
|
||||
this.redisClient_ = redisClient
|
||||
this.redisSubscriber_ = redisSubscriber
|
||||
this.cronQueue_ = new Bull(`cron-jobs:queue`, opts)
|
||||
// Register our worker to handle emit calls
|
||||
this.queue_.process(this.worker_)
|
||||
// Register cron worker
|
||||
this.cronQueue_.process(this.cronWorker_)
|
||||
|
||||
if (process.env.NODE_ENV !== "test") {
|
||||
this.startEnqueuer()
|
||||
@@ -103,6 +100,7 @@ export default class EventBusService {
|
||||
{
|
||||
manager: transactionManager,
|
||||
stagedJobRepository: this.stagedJobRepository_,
|
||||
jobSchedulerService: this.jobSchedulerService_,
|
||||
logger: this.logger_,
|
||||
redisClient: this.redisClient_,
|
||||
redisSubscriber: this.redisSubscriber_,
|
||||
@@ -157,27 +155,6 @@ export default class EventBusService {
|
||||
return this
|
||||
}
|
||||
|
||||
/**
|
||||
* Adds a function to a list of event subscribers.
|
||||
* @param event - the event that the subscriber will listen for.
|
||||
* @param subscriber - the function to be called when a certain event
|
||||
* happens. Subscribers must return a Promise.
|
||||
* @return this
|
||||
*/
|
||||
protected registerCronHandler_(
|
||||
event: string | symbol,
|
||||
subscriber: Subscriber
|
||||
): this {
|
||||
if (typeof subscriber !== "function") {
|
||||
throw new Error("Handler must be a function")
|
||||
}
|
||||
|
||||
const cronHandlers = this.cronHandlers_.get(event) ?? []
|
||||
this.cronHandlers_.set(event, [...cronHandlers, subscriber])
|
||||
|
||||
return this
|
||||
}
|
||||
|
||||
/**
|
||||
* Calls all subscribers when an event occurs.
|
||||
* @param {string} eventName - the name of the event to be process.
|
||||
@@ -198,7 +175,7 @@ export default class EventBusService {
|
||||
const stagedJobInstance = stagedJobRepository.create({
|
||||
event_name: eventName,
|
||||
data,
|
||||
})
|
||||
} as StagedJob)
|
||||
return await stagedJobRepository.save(stagedJobInstance)
|
||||
} else {
|
||||
const opts: { removeOnComplete: boolean } & EmitOptions = {
|
||||
@@ -288,32 +265,9 @@ export default class EventBusService {
|
||||
)
|
||||
}
|
||||
|
||||
/**
|
||||
* Handles incoming jobs.
|
||||
* @param job The job object
|
||||
* @return resolves to the results of the subscriber calls.
|
||||
*/
|
||||
cronWorker_ = async <T>(job: {
|
||||
data: { eventName: string; data: T }
|
||||
}): Promise<unknown[]> => {
|
||||
const { eventName, data } = job.data
|
||||
const observers = this.cronHandlers_.get(eventName) || []
|
||||
this.logger_.info(`Processing cron job: ${eventName}`)
|
||||
|
||||
return await Promise.all(
|
||||
observers.map(async (subscriber) => {
|
||||
return subscriber(data, eventName).catch((err) => {
|
||||
this.logger_.warn(
|
||||
`An error occured while processing ${eventName}: ${err}`
|
||||
)
|
||||
return err
|
||||
})
|
||||
})
|
||||
)
|
||||
}
|
||||
|
||||
/**
|
||||
* Registers a cron job.
|
||||
* @deprecated All cron job logic has been refactored to the `JobSchedulerService`. This method will be removed in a future release.
|
||||
* @param eventName - the name of the event
|
||||
* @param data - the data to be sent with the event
|
||||
* @param cron - the cron pattern
|
||||
@@ -326,14 +280,6 @@ export default class EventBusService {
|
||||
cron: string,
|
||||
handler: Subscriber
|
||||
): void {
|
||||
this.logger_.info(`Registering ${eventName}`)
|
||||
this.registerCronHandler_(eventName, handler)
|
||||
return this.cronQueue_.add(
|
||||
{
|
||||
eventName,
|
||||
data,
|
||||
},
|
||||
{ repeat: { cron } }
|
||||
)
|
||||
this.jobSchedulerService_.create(eventName, data, cron, handler)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,122 @@
|
||||
import Bull from "bull"
|
||||
import Redis from "ioredis"
|
||||
import { ConfigModule, Logger } from "../types/global"
|
||||
|
||||
type InjectedDependencies = {
|
||||
logger: Logger
|
||||
redisClient: Redis.Redis
|
||||
redisSubscriber: Redis.Redis
|
||||
}
|
||||
|
||||
type ScheduledJobHandler<T = unknown> = (
|
||||
data: T,
|
||||
eventName: string
|
||||
) => Promise<void>
|
||||
|
||||
export default class JobSchedulerService {
|
||||
protected readonly config_: ConfigModule
|
||||
protected readonly logger_: Logger
|
||||
protected readonly handlers_: Map<string | symbol, ScheduledJobHandler[]> =
|
||||
new Map()
|
||||
protected readonly queue_: Bull
|
||||
|
||||
constructor(
|
||||
{ logger, redisClient, redisSubscriber }: InjectedDependencies,
|
||||
config: ConfigModule,
|
||||
singleton = true
|
||||
) {
|
||||
this.config_ = config
|
||||
this.logger_ = logger
|
||||
|
||||
if (singleton) {
|
||||
const opts = {
|
||||
createClient: (type: string): Redis.Redis => {
|
||||
switch (type) {
|
||||
case "client":
|
||||
return redisClient
|
||||
case "subscriber":
|
||||
return redisSubscriber
|
||||
default:
|
||||
if (config.projectConfig.redis_url) {
|
||||
return new Redis(config.projectConfig.redis_url)
|
||||
}
|
||||
return redisClient
|
||||
}
|
||||
},
|
||||
}
|
||||
|
||||
this.queue_ = new Bull(`scheduled-jobs:queue`, opts)
|
||||
// Register scheduled job worker
|
||||
this.queue_.process(this.scheduledJobsWorker)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Adds a function to a list of event subscribers.
|
||||
* @param event - the event that the subscriber will listen for.
|
||||
* @param subscriber - the function to be called when a certain event
|
||||
* happens. Subscribers must return a Promise.
|
||||
* @return this
|
||||
*/
|
||||
protected registerHandler(
|
||||
event: string | symbol,
|
||||
handler: ScheduledJobHandler
|
||||
): void {
|
||||
if (typeof handler !== "function") {
|
||||
throw new Error("Handler must be a function")
|
||||
}
|
||||
|
||||
const handlers = this.handlers_.get(event) ?? []
|
||||
this.handlers_.set(event, [...handlers, handler])
|
||||
}
|
||||
|
||||
/**
|
||||
* Handles incoming scheduled jobs.
|
||||
* @param job The job object
|
||||
* @return resolves to the results of the subscriber calls.
|
||||
*/
|
||||
protected scheduledJobsWorker = async <T>(job: {
|
||||
data: { eventName: string; data: T }
|
||||
}): Promise<unknown[]> => {
|
||||
const { eventName, data } = job.data
|
||||
const observers = this.handlers_.get(eventName) || []
|
||||
this.logger_.info(`Processing scheduled job: ${eventName}`)
|
||||
|
||||
return await Promise.all(
|
||||
observers.map(async (subscriber) => {
|
||||
return subscriber(data, eventName).catch((err) => {
|
||||
this.logger_.warn(
|
||||
`An error occured while processing ${eventName}: ${err}`
|
||||
)
|
||||
return err
|
||||
})
|
||||
})
|
||||
)
|
||||
}
|
||||
|
||||
/**
|
||||
* Registers a scheduled job.
|
||||
* @param eventName - the name of the event
|
||||
* @param data - the data to be sent with the event
|
||||
* @param schedule - the schedule expression
|
||||
* @param handler - the handler to call on the job
|
||||
* @return void
|
||||
*/
|
||||
create<T>(
|
||||
eventName: string,
|
||||
data: T,
|
||||
schedule: string,
|
||||
handler: ScheduledJobHandler
|
||||
): void {
|
||||
this.logger_.info(`Registering ${eventName}`)
|
||||
this.registerHandler(eventName, handler)
|
||||
|
||||
this.queue_.add(
|
||||
{
|
||||
eventName,
|
||||
data,
|
||||
},
|
||||
{ repeat: { cron: schedule } }
|
||||
)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user