Event Bus
Event Bus
Section titled “Event Bus”Veloce-TS includes a built-in in-process event bus for decoupling modules without Redis or Kafka. Use it for cross-module communication, domain events, and feeding SSE streams.
EventBus
Section titled “EventBus”import { EventBus } from 'veloce-ts';
const bus = new EventBus();Or use the pre-built singleton:
import { globalEvents } from 'veloce-ts';on(event, handler) — subscribe
Section titled “on(event, handler) — subscribe”globalEvents.on('user.created', async (payload) => { await emailService.sendWelcome(payload.email);});emit() calls all handlers concurrently and resolves once they have all settled. emitSync() is the fire-and-forget variant for synchronous handlers — it does not await them.
once(event, handler) — subscribe once
Section titled “once(event, handler) — subscribe once”The handler is automatically removed after the first call.
globalEvents.once('app.ready', () => { console.log('App is ready');});emit(event, payload) — fire async (parallel)
Section titled “emit(event, payload) — fire async (parallel)”Returns a Promise<void> that resolves when all handlers have settled.
await globalEvents.emit('user.created', { userId: '123', email: 'user@example.com' });emitSync(event, payload) — fire sync
Section titled “emitSync(event, payload) — fire sync”Calls handlers in registration order and returns void — it does not await them. An async handler passed to emitSync is started but not waited for, so use emit() whenever handlers do asynchronous work.
globalEvents.emitSync('order.placed', { orderId: '456' });Error handling
Section titled “Error handling”Previously emit() used Promise.all, so the first rejection aborted the whole emit and any handler that had not been reached was silently skipped — one bad listener could stop the rest of your side effects.
Both emit() and emitSync() now run every listener, collect the failures, and rethrow them together as an AggregateError once all listeners have settled. once listeners are removed even when they throw.
globalEvents.on('user.created', async () => { throw new Error('mailer down'); });globalEvents.on('user.created', async () => { await analytics.track('signup'); });
try { await globalEvents.emit('user.created', { userId: '123' });} catch (error) { // analytics.track() still ran — only the failing listener is reported here if (error instanceof AggregateError) { for (const cause of error.errors) { logger.error('Event listener failed', cause); } }}If a failing listener should never interrupt the emitter, catch inside the handler itself:
globalEvents.on('user.created', async (payload) => { try { await emailService.sendWelcome(payload.email); } catch (error) { logger.error('Welcome email failed', error); }});off(event, handler) — unsubscribe
Section titled “off(event, handler) — unsubscribe”const handler = (payload: any) => console.log(payload);globalEvents.on('user.created', handler);// later...globalEvents.off('user.created', handler);removeAllListeners(event?) — clear listeners
Section titled “removeAllListeners(event?) — clear listeners”globalEvents.removeAllListeners('user.created'); // clear specific eventglobalEvents.removeAllListeners(); // clear all eventslistenerCount(event) — count listeners
Section titled “listenerCount(event) — count listeners”const count = globalEvents.listenerCount('user.created');Patterns
Section titled “Patterns”Domain events
Section titled “Domain events”import { globalEvents } from 'veloce-ts';
class UserService { async createUser(data: CreateUserDto) { const user = await db.insert(users).values(data).returning(); await globalEvents.emit('user.created', { userId: user.id, email: user.email, name: user.name, }); return user; }}// email.service.ts — registers at startupimport { globalEvents } from 'veloce-ts';
globalEvents.on('user.created', async ({ email, name }) => { await mailer.send({ to: email, subject: `Welcome, ${name}!`, template: 'welcome', });});// audit.service.ts — multiple handlers for the same eventglobalEvents.on('user.created', async ({ userId }) => { await db.insert(auditLog).values({ event: 'user.created', resourceId: userId, timestamp: new Date(), });});Feed SSE streams via event bus
Section titled “Feed SSE streams via event bus”// Combine EventBus + SSE streaming@Get('/events')@SSE()async* streamUserEvents(): AsyncGenerator<{ data: string; event?: string }> { yield { data: 'connected', event: 'open' };
const queue: unknown[] = []; const handler = (payload: unknown) => queue.push(payload);
globalEvents.on('user.created', handler); globalEvents.on('user.updated', handler);
try { while (true) { if (queue.length > 0) { const event = queue.shift()!; yield { data: JSON.stringify(event), event: 'update' }; } else { await new Promise(r => setTimeout(r, 100)); } } } finally { globalEvents.off('user.created', handler); globalEvents.off('user.updated', handler); }}Multiple modules, one bus
Section titled “Multiple modules, one bus”// Module Aimport { globalEvents } from 'veloce-ts';
export class OrderService { async placeOrder(order: Order) { const placed = await db.createOrder(order); await globalEvents.emit('order.placed', { orderId: placed.id, amount: order.total }); return placed; }}
// Module B — handles order events independentlyimport { globalEvents } from 'veloce-ts';
globalEvents.on('order.placed', async ({ orderId, amount }) => { await analytics.track('purchase', { orderId, amount });});
globalEvents.on('order.placed', async ({ orderId }) => { await inventory.decrementStock(orderId);});Using a local bus (not global)
Section titled “Using a local bus (not global)”import { EventBus } from 'veloce-ts';
// Scope the bus to a plugin or moduleexport class NotificationPlugin { readonly events = new EventBus();
register(app: VeloceTS) { // internal-only events — not polluting globalEvents this.events.on('notification.queued', this.process.bind(this)); }}Cleanup in tests
Section titled “Cleanup in tests”import { globalEvents } from 'veloce-ts';import { afterEach } from 'bun:test';
afterEach(() => { globalEvents.removeAllListeners();});Next Steps
Section titled “Next Steps”- Streaming — combine EventBus with SSE
- Interceptors
- Exception Filters