rhythmjs

Search documentation

Search guides, the tutorial and every package.

On this page

emit() and dispatch order#

emit(event, payload) dispatches synchronously. Exact-name subscribers run first, then wildcard patterns in the order their pattern first appeared in the registry. Every handler receives (payload, event), where event is always the concrete emitted name; a "order.*" subscriber sees "order.created", not the pattern.

TypeScript
eventBus.on("order.created", audit); // runs first (exact)
eventBus.on("order.*", project); // then wildcards, registration order
eventBus.on("order.**", replicate);

eventBus.emit("order.created", { orderId: "o1", total: 99 });

Dispatch snapshots the registry per pattern, so a handler may unsubscribe itself or its siblings mid-emit without skipping anyone already snapshotted. Once-subscriptions are removed before their handler runs.

Errors never cross listeners#

A listener that throws never stops its siblings. With emit, a synchronous throw is caught on the spot and a returned promise gets a rejection handler; both route to the bus's onError(error, event, payload) hook. The default hook rethrows on a microtask, so failures surface as unhandled errors instead of vanishing.

TypeScript
import { createEventBus } from "@rhythmjs/events";

const bus = createEventBus<AppEvents>({
  onError: (error, event, payload) =>
    log.error(`listener failed for ${event}`, { error, payload }),
});

emitAsync(): await every listener#

emitAsync(event, payload) is for emitters that need the outcome. It collects every matching handler (removing once-subscriptions first), awaits them all with Promise.allSettled, so every listener runs to completion regardless of failures, and then, if any rejected, throws an AggregateError whose errors array holds every rejection reason, with the message N listener(s) failed for event "…".

TypeScript
try {
  await bus.emitAsync("order.created", { orderId: "o1", total: 99 });
} catch (error) {
  // AggregateError: every listener rejection, none skipped
  for (const cause of (error as AggregateError).errors) report(cause);
}

Note the difference: emit is fire-and-forget with errors funneled to onError; emitAsync bypasses onError and hands failures back to the caller.

The kernel module#

eventsModule.forRoot<TEvents>(options?) wraps createEventBus in an encapsulated Rhythm module that provides eventBus; options.onError is the same hook.

TypeScript
import { Rhythm } from "@rhythmjs/rhythm";
import { eventsModule } from "@rhythmjs/events";

const app = new Rhythm().register(
  eventsModule.forRoot<AppEvents>(),
  ({ eventBus }) => ({ eventBus }),
);
await app.setup();
const { eventBus } = await app.run({});

Announce, then obligate#

Events pair naturally with @rhythmjs/queue: an event announces (ephemeral, every current listener), a queue job obligates (at-least-once, exactly one worker). Bridging one into the other is one line, and because both are typed by name-to-payload maps, the bridge type-checks end to end.

TypeScript
eventBus.on("order.**", (payload, event) => {
  void queueService.add(event, payload); // fan announcements into durable work
});

See the API reference for every export and type.