feat(core-shared/jobs): InMemoryJobQueue with register()
This commit is contained in:
54
packages/core-shared/src/jobs/in-memory-job-queue.test.ts
Normal file
54
packages/core-shared/src/jobs/in-memory-job-queue.test.ts
Normal file
@@ -0,0 +1,54 @@
|
||||
import { describe, it, expect, vi } from "vitest";
|
||||
import { InMemoryJobQueue } from "@/jobs/in-memory-job-queue";
|
||||
import type { IJobQueue } from "@/jobs/job-queue.interface";
|
||||
|
||||
describe("InMemoryJobQueue", () => {
|
||||
it("returns a synthetic jobId on enqueue", async () => {
|
||||
const handler = vi.fn();
|
||||
const queue: IJobQueue = new InMemoryJobQueue({
|
||||
"test.task": handler,
|
||||
});
|
||||
const result = await queue.enqueue("test.task", { x: 1 });
|
||||
expect(result.jobId).toMatch(/^in-memory-/);
|
||||
});
|
||||
|
||||
it("invokes the registered handler asynchronously with the input", async () => {
|
||||
const handler = vi.fn();
|
||||
const queue = new InMemoryJobQueue({ "test.task": handler });
|
||||
await queue.enqueue("test.task", { x: 42 });
|
||||
await new Promise((r) => setImmediate(r));
|
||||
expect(handler).toHaveBeenCalledWith({ x: 42 });
|
||||
});
|
||||
|
||||
it("throws if the task slug has no registered handler", async () => {
|
||||
const queue = new InMemoryJobQueue({});
|
||||
await expect(queue.enqueue("missing.task", {})).rejects.toThrow(
|
||||
/no handler registered for task slug: missing\.task/,
|
||||
);
|
||||
});
|
||||
|
||||
it("delays execution when runAt is in the future", async () => {
|
||||
vi.useFakeTimers();
|
||||
try {
|
||||
const handler = vi.fn();
|
||||
const queue = new InMemoryJobQueue({ "test.task": handler });
|
||||
const future = new Date(Date.now() + 1000);
|
||||
await queue.enqueue("test.task", {}, { runAt: future });
|
||||
expect(handler).not.toHaveBeenCalled();
|
||||
vi.advanceTimersByTime(1000);
|
||||
await Promise.resolve();
|
||||
expect(handler).toHaveBeenCalledTimes(1);
|
||||
} finally {
|
||||
vi.useRealTimers();
|
||||
}
|
||||
});
|
||||
|
||||
it("register adds a handler that can be enqueued against", async () => {
|
||||
const queue = new InMemoryJobQueue();
|
||||
const handler = vi.fn();
|
||||
queue.register("late.task", handler);
|
||||
await queue.enqueue("late.task", { z: 1 });
|
||||
await new Promise((r) => setImmediate(r));
|
||||
expect(handler).toHaveBeenCalledWith({ z: 1 });
|
||||
});
|
||||
});
|
||||
36
packages/core-shared/src/jobs/in-memory-job-queue.ts
Normal file
36
packages/core-shared/src/jobs/in-memory-job-queue.ts
Normal file
@@ -0,0 +1,36 @@
|
||||
import type { IJobQueue } from "./job-queue.interface";
|
||||
|
||||
export type InMemoryHandler = (input: unknown) => Promise<void> | void;
|
||||
|
||||
export class InMemoryJobQueue implements IJobQueue {
|
||||
private counter = 0;
|
||||
private readonly handlers: Record<string, InMemoryHandler>;
|
||||
|
||||
constructor(handlers: Record<string, InMemoryHandler> = {}) {
|
||||
this.handlers = { ...handlers };
|
||||
}
|
||||
|
||||
register(slug: string, handler: InMemoryHandler): void {
|
||||
this.handlers[slug] = handler;
|
||||
}
|
||||
|
||||
async enqueue<T>(
|
||||
taskSlug: string,
|
||||
input: T,
|
||||
options?: { runAt?: Date },
|
||||
): Promise<{ jobId: string }> {
|
||||
const handler = this.handlers[taskSlug];
|
||||
if (!handler) {
|
||||
throw new Error(`no handler registered for task slug: ${taskSlug}`);
|
||||
}
|
||||
this.counter += 1;
|
||||
const jobId = `in-memory-${this.counter}`;
|
||||
const delay = options?.runAt ? options.runAt.getTime() - Date.now() : 0;
|
||||
if (delay > 0) {
|
||||
setTimeout(() => void handler(input), delay);
|
||||
} else {
|
||||
setImmediate(() => void handler(input));
|
||||
}
|
||||
return { jobId };
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user