Usar el event bus
Publicar un evento de dominio desde un aggregate, suscribir un listener que reaccione en otro módulo.
Última actualización 19 jul 2026
El event bus es in-process, respaldado por un outbox durable. Existe para mantener los módulos desacoplados a nivel de código — iam no importa notifications; en cambio, levanta user.created y un listener en notifications reacciona. Desde #fix-02, publish() persiste cada evento en la tabla event_outbox antes de despachar, así que los handlers que fallan se reintentan en vez de perderse en silencio.
Implementación
apps/server/src/infrastructure/events/event-bus.ts (interface + una implementación in-memory liviana que se conserva para tests) y apps/server/src/infrastructure/events/outbox-event-bus.ts (el bus de producción cableado en el container):
export interface IEventBus {
publish(events: ReadonlyArray<DomainEvent>): Promise<void>;
subscribe(eventName: string, handler: EventHandler): void;
}
export class InMemoryEventBus implements IEventBus { /* ... */ }OutboxEventBus.publish() inserta cada evento como fila pending en event_outbox y después despacha inline. Un handler que lanza no hace fallar al aggregate publicador — pero tampoco se traga en silencio: el error va a Sentry vía captureError, la fila queda pending, y un poller in-process la reintenta con backoff exponencial (30s duplicando, tope 15m). Tras 8 intentos la fila queda dead-letter (status = failed, lastError cargado) — consultá event_outbox para inspeccionar. La entrega es at-least-once por evento: un retry re-ejecuta todos los handlers de ese evento, así que los handlers deben ser idempotentes (los listeners del read model de billing son upserts). Solo un fallo del insert del outbox rechaza publish() — en el path del webhook de billing eso se convierte en un 5xx y el provider reenvía.
1. Levantar un evento desde un aggregate
Dentro del método de dominio del aggregate:
this.addDomainEvent({
name: 'user.created',
aggregateId: this.id,
occurredAt: new Date(),
payload: { email: this.email.value, locale: this.locale },
});Los eventos se acumulan en el aggregate pero todavía no se publican.
2. Publicar en el use case
Después de persistir el aggregate:
async execute(input: RegisterInput): Promise<Result<User, DomainError>> {
const userResult = User.create(input);
if (userResult.isErr()) return userResult;
await this.users.save(userResult.value);
await this.bus.publish(userResult.value.pullEvents());
return userResult;
}pullEvents() devuelve los eventos acumulados y los borra del aggregate, así que re-guardar nunca re-publica.
3. Suscribir un listener
En el bootstrap del módulo receptor (típicamente llamado desde bootstrap/container.ts o un register-listeners.ts dedicado):
// apps/server/src/modules/notifications/infrastructure/listeners.ts
export const registerNotificationListeners = (deps: {
bus: IEventBus;
jobs: JobScheduler;
}) => {
deps.bus.subscribe('user.created', async (event) => {
await deps.jobs.enqueue('emails', {
kind: 'welcome',
to: { email: event.payload.email },
recipientName: event.payload.recipientName,
appName: 'UseDeploy',
locale: event.payload.locale,
});
});
};El listener hace el trabajo síncrono mínimo — típicamente: encolar un job de BullMQ. El trabajo pesado (mandar el email, llamar a un tercero) va en el worker.
4. Tipá los nombres de eventos
Agregá el nombre del evento a un union compartido si querés que el sistema de tipos imponga que sólo te suscribas a eventos existentes:
type ApplicationEvent =
| { name: 'user.created'; payload: { /* ... */ } }
| { name: 'subscription.activated'; payload: { /* ... */ } };Cuándo usar eventos vs llamada directa
| Usá un evento cuando | Usá llamada directa cuando |
|---|---|
| Múltiples módulos podrían reaccionar | Exactamente un módulo es dueño del próximo paso |
| La reacción es best-effort (no se necesita garantía transaccional) | El próximo paso debe tener éxito para que la operación se considere hecha |
| La reacción es async / puede diferirse | El resultado se necesita para la respuesta HTTP |
El bus no es una message queue completa, pero sí es durable: los eventos se persisten en event_outbox antes del dispatch, y un proceso que crasheó retoma las filas pendientes en el próximo poll (intervalo de 30s; tanto el proceso API como el worker pollean, y un claim optimista hace seguros a los pollers concurrentes). El gap restante es un crash exactamente entre el save() del aggregate y publish() — esa ventana dura microsegundos y deliberadamente no se cierra (requeriría atravesar un unit-of-work por todos los repositorios). El trabajo pesado sigue yendo en jobs de BullMQ; para fan-out cross-instance, montá Redis pub/sub encima. Una excepción deliberada: el email de aceptación de invitación no viaja por el bus — el token plaintext no debe persistirse jamás (el outbox guarda los payloads), así que InviteMemberUseCase lo envía directo vía un callback SendInvitationEmail inyectado.
Testing
Inyectá InMemoryEventBus desde el código de producción o construí un double que registre:
const recorded: DomainEvent[] = [];
const bus: IEventBus = {
publish: async (events) => { recorded.push(...events); },
subscribe: () => {},
};Hacé asserts contra recorded después de que corre el use case. Si el use case se olvida de publicar, el test falla.