SaaS Starter

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 quandoUse chamada direta quando
Múltiplos módulos poderiam reagirExatamente 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 adiadaO 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.

Nesta página