Skip to content

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.

import { EventBus } from 'veloce-ts';
const bus = new EventBus();

Or use the pre-built singleton:

import { globalEvents } from 'veloce-ts';
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.

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' });

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' });

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);
}
});
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 event
globalEvents.removeAllListeners(); // clear all events
const count = globalEvents.listenerCount('user.created');
user.service.ts
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 startup
import { 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 event
globalEvents.on('user.created', async ({ userId }) => {
await db.insert(auditLog).values({
event: 'user.created',
resourceId: userId,
timestamp: new Date(),
});
});
// 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);
}
}
// Module A
import { 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 independently
import { 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);
});
import { EventBus } from 'veloce-ts';
// Scope the bus to a plugin or module
export class NotificationPlugin {
readonly events = new EventBus();
register(app: VeloceTS) {
// internal-only events — not polluting globalEvents
this.events.on('notification.queued', this.process.bind(this));
}
}
import { globalEvents } from 'veloce-ts';
import { afterEach } from 'bun:test';
afterEach(() => {
globalEvents.removeAllListeners();
});