Usar o event bus
Publicar um evento de domínio a partir de um aggregate, inscrever um listener que reage em outro módulo.
Última atualização 19 de jul. de 2026
O event bus é in-process, respaldado por um outbox durável. Existe para manter os módulos desacoplados no nível de código — iam não importa notifications; em vez disso, ele levanta user.created e um listener em notifications reage. Desde #fix-02, publish() persiste cada evento na tabela event_outbox antes de despachar, então handlers que falham são retentados em vez de se perderem em silêncio.
Implementação
apps/server/src/infrastructure/events/event-bus.ts (interface + uma implementação in-memory leve mantida para tests) e apps/server/src/infrastructure/events/outbox-event-bus.ts (o bus de produção ligado no container):
export interface IEventBus {
publish(events: ReadonlyArray<DomainEvent>): Promise<void>;
subscribe(eventName: string, handler: EventHandler): void;
}
export class InMemoryEventBus implements IEventBus { /* ... */ }OutboxEventBus.publish() insere cada evento como linha pending em event_outbox e depois despacha inline. Um handler que dá throw não faz o aggregate publicador falhar — mas também não é mais engolido em silêncio: o erro vai para o Sentry via captureError, a linha fica pending, e um poller in-process a retenta com backoff exponencial (30s dobrando, teto de 15m). Após 8 tentativas a linha vira dead-letter (status = failed, lastError preenchido) — consulte event_outbox para inspecionar. A entrega é at-least-once por evento: um retry re-executa todos os handlers daquele evento, então os handlers devem ser idempotentes (os listeners do read model de billing são upserts). Só uma falha do insert do outbox rejeita publish() — no path do webhook de billing isso vira um 5xx e o provider reenvia.
1. Levantar um evento a partir de um aggregate
Dentro do método de domínio do aggregate:
this.addDomainEvent({
name: 'user.created',
aggregateId: this.id,
occurredAt: new Date(),
payload: { email: this.email.value, locale: this.locale },
});Os eventos se acumulam no aggregate mas ainda não são publicados.
2. Publicar no use case
Depois de persistir o 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() retorna os eventos acumulados e os limpa do aggregate, então re-salvar nunca re-publica.
3. Inscrever um listener
No bootstrap do módulo receptor (tipicamente chamado de bootstrap/container.ts ou de um 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,
});
});
};O listener faz o trabalho síncrono mínimo — tipicamente: enfileirar um job BullMQ. Trabalho pesado (enviar o email, chamar um terceiro) fica no worker.
4. Tipe os nomes dos eventos
Adicione o nome do evento a uma union compartilhada se quiser que o sistema de tipos force você a se inscrever apenas em eventos existentes:
type ApplicationEvent =
| { name: 'user.created'; payload: { /* ... */ } }
| { name: 'subscription.activated'; payload: { /* ... */ } };Quando usar eventos vs chamada direta
| Use um evento quando | Use chamada direta quando |
|---|---|
| Múltiplos módulos poderiam reagir | Exatamente um módulo é dono do próximo passo |
| A reação é best-effort (sem garantia transacional necessária) | O próximo passo precisa ter sucesso para a operação ser considerada feita |
| A reação é async / pode ser adiada | O resultado é necessário para a response HTTP |
O bus não é uma message queue completa, mas é durável: os eventos são persistidos em event_outbox antes do dispatch, e um processo que crashou retoma as linhas pendentes no próximo poll (intervalo de 30s; tanto o processo API quanto o worker fazem poll, e um claim otimista torna pollers concorrentes seguros). O gap restante é um crash exatamente entre o save() do aggregate e o publish() — essa janela dura microssegundos e deliberadamente não é fechada (exigiria atravessar um unit-of-work por todos os repositórios). Trabalho pesado continua indo em jobs BullMQ; para fan-out cross-instance, sobreponha Redis pub/sub. Uma exceção deliberada: o email de aceite de convite não viaja pelo bus — o token plaintext nunca deve ser persistido (o outbox guarda os payloads), então InviteMemberUseCase o envia direto via um callback SendInvitationEmail injetado.
Testing
Injete InMemoryEventBus a partir do código de produção ou construa um double que grava:
const recorded: DomainEvent[] = [];
const bus: IEventBus = {
publish: async (events) => { recorded.push(...events); },
subscribe: () => {},
};Faça asserts contra recorded depois que o use case roda. Se o use case esquecer de publicar, o teste falha.