Tutorial · Step 11 of 14
Background jobs
Events, cron schedules, and durable queues for work beyond the request.
Announce changes with events#
The event bus decouples "a note was created" from everything that reacts to it. Define the event map once, register the module, and emit from handlers — subscribers use exact names or wildcard patterns:
import { eventsModule } from "@rhythmjs/events";
interface Events {
"note.created": { id: string; title: string };
"note.removed": { id: string };
}
export const appModule = new Rhythm<RhythmHttpContext>({ name: "app", type: "module" })
.register(eventsModule.forRoot<Events>(), ({ eventBus }) => ({ eventBus }))
/* ... */;
// in the controller:
ctx.eventBus.emit("note.created", { id: note.id, title: note.title });
// anywhere with the bus:
eventBus.on("note.*", ({ event, payload }) => {
console.log(`${event}:`, payload);
});Pass forRoot({ onError }) in real apps: the default rethrows a failing listener's error as an uncaught exception. Handler errors never break the bus or sibling listeners either way.
Run on a schedule#
@rhythmjs/schedule runs cron, interval, and timeout jobs on an in-house cron engine (no dependencies). Define jobs, register the module, and start the scheduler alongside the server. forRoot builds the service eagerly; when main.ts needs the same instance, build it with createScheduleService and assign it to the context:
import { createScheduleService, cronJob, cronPatterns, intervalJob } from "@rhythmjs/schedule";
import { startScheduler } from "@rhythmjs/schedule/scheduler";
const jobs = [
cronJob("digest", cronPatterns.daily, async () => {
await sendDailyDigest();
}),
cronJob("cleanup", "*/30 * * * *", async () => {
await purgeExpiredNotes();
}),
intervalJob("heartbeat", 30_000, () => pingUptimeMonitor()),
];
export const scheduleService = createScheduleService(...jobs);
// or, as a module: appModule.register(scheduleModule.forRoot(...jobs), ({ scheduleService }) => ({ scheduleService }))
// in main.ts, after Bun.serve:
const scheduler = startScheduler(scheduleService);
// scheduler.stop() in your shutdown hookOverlap protection defaults to "skip" (a still-running job is not started again), handler errors are recorded on scheduleService.state(name) rather than thrown, and a removed job's timer cancels itself.
Durable work with queues#
For work that must survive the request — emails, exports, webhooks — use a typed queue. In-memory by default; pass redis for durability across processes. The job map types both producers and processors. Build the service once at startup and put it on the context; close() is yours to call on shutdown:
import { createQueueService } from "@rhythmjs/queue";
interface Jobs {
"email.digest": { userId: string };
"notes.export": { noteIds: string[] };
}
export const queueService = createQueueService<Jobs>({ redis: process.env.REDIS_URL });
// declared on the app: new Rhythm<RhythmHttpContext, { queueService: typeof queueService }>(...)
appModule.context.queueService = queueService;
// (queueModule.forRoot<Jobs>(...) does the same inside a module you register)
// produce (from a handler):
await ctx.queueService.add("email.digest", { userId }, { attempts: 3, backoff: 5_000 });
// process (in a worker entrypoint):
queueService.process(
{
"email.digest": async (payload) => sendDigest(payload.userId),
"notes.export": async (payload) => exportNotes(payload.noteIds),
},
{
concurrency: 4,
onFailed: (name, id, error) => log.error(`job ${name}#${id} exhausted retries`, error),
onError: (error) => log.warn("queue engine hiccup", error), // backend errors; the loop survives
},
);
// on shutdown: stops the workers, then releases the engine
await queueService.close();schedule(name, { pattern }, payload) adds repeatable jobs on the same cron engine, and counts() exposes queue depth for your health checks — which is exactly where the next step picks up.