Pular para o conteúdo

N-0a — Event Bus (Feromônios)

Parte do Núcleo — O Chão do Formigueiro


O barramento de eventos é o canal único de comunicação do sistema. Toda colônia publica e consome exclusivamente por aqui. Nenhuma colônia se comunica diretamente com outra.

Não executa lógica de negócio e não decide quem pode publicar o quê. A validação de tipo e de payload é responsabilidade do Registry (N-0b). Suas responsabilidades são: receber eventos publicados, validar tipo, versão e payload contra o Registry, atribuir sequence number global monotônico, redigir o payload para a persistência conforme as regras de LGPD, persistir o histórico de forma append-only, entregar aos consumidores registrados e expor o histórico completo para auditoria e replay.

O log de eventos é a memória do sistema. Um evento publicado é permanente. Pode ser reprocessado, auditado e comparado com qualquer versão anterior.


O Event Bus é um módulo NestJS do núcleo, separado do Registry (N-0b) e da Observabilidade (N-0c). É um dos módulos declarados como @Global(): toda colônia o injeta sem importá-lo explicitamente no array imports do seu módulo. O ObservabilityModule (N-0c) também é @Global(). O EventBusModule exporta EventBusService e PrismaService, a instância única de banco injetada por todas as colônias.

src/nucleo/n-0a-event-bus/
├── event-bus.module.ts # @Global(); imports EventEmitterModule + ObservabilityModule + RegistryModule
├── event-bus.service.ts # publicar(), inscrever(), consultas e replay
├── event-bus.controller.ts # rotas internas de operação em /api/events
├── prisma.service.ts # PrismaClient com adapter pg, instância única do monolito
├── dto/
│ ├── publicar-evento.dto.ts # Contrato de publicação recebido das colônias
│ ├── consultar-eventos.dto.ts # Filtros da consulta ao event_log
│ └── reprocessar-dlq.dto.ts # Corpo de POST /api/events/dlq/reprocessar
├── redacao/
│ └── redacao-eventos.ts # Redação do payload na persistência (allowlist por tipo)
└── repositories/
├── event-log.repository.ts # insert, consulta paginada, replay e maior sequence
└── dead-letter.repository.ts # insert, transições de status e estatísticas da DLQ
@Global()
@Module({
imports: [
EventEmitterModule.forRoot({
wildcard: false,
delimiter: '.',
maxListeners: 30,
verboseMemoryLeak: true,
ignoreErrors: false,
}),
ObservabilityModule, // N-0c — instrumentação cross-cutting
RegistryModule, // N-0b — validação de tipo, versão e payload
],
controllers: [EventBusController],
providers: [
EventBusService,
PrismaService,
EventLogRepository,
DeadLetterRepository,
PapelOperadorGuard, // proteção das rotas internas
],
exports: [EventBusService, PrismaService], // PrismaService único do monolito (conexão global)
})
export class EventBusModule {}
  • wildcard: false — cada manipulador escuta um tipo exato (demanda.recebida), não padrões com wildcard. A entrega é precisa.
  • maxListeners: 30 — limite definido no módulo, sem variável de ambiente. Ajustes de escala passam pelo código.
  • ignoreErrors: false — erros em manipuladores não são silenciados. O wrapper do inscrever() captura a exceção, registra na DLQ e não deixa a rejeição subir.
  • O módulo importa RegistryModule para validar tipo, versão e payload na publicação. É a única dependência síncrona de validação no caminho de publicação. O ObservabilityModule também é importado, para trace e métricas.

Os métodos implementados, em português:

publicar(dto: PublicarEventoDto): Promise<EventoConsultado>; // valida tipo, versão e payload contra o Registry, persiste e emite
inscrever(tipo: string, colonia: string, manipulador: ManipuladorEvento): void;
consultarEventos(filtros: ConsultarEventosDto): Promise<ResultadoPaginado<EventoConsultado>>; // limite máximo 500
obterEventoPorId(eventId: string): Promise<EventoConsultado | null>;
replayDeSequence(deSequence: number | bigint, tipos?: string[]): Promise<EventoConsultado[]>; // paginado internamente
obterMaiorSequence(): Promise<bigint>; // seed de cursor das colônias no boot
obterUltimoEventId(): Promise<string | null>; // seed do sorteio (D-6a)
obterEstatisticasDLQ(): Promise<DLQEstatisticas>;
obterPendentesDLQ(colonia: string): Promise<DLQEntrada[]>;
reentregarParaColonia(tipo: string, colonia: string, evento: EventoReentrega): Promise<void>; // mesma proteção do consumo ao vivo
reprocessarPendentesDLQ(colonia: string): Promise<{ reprocessados: number }>; // usado pela rota POST /api/events/dlq/reprocessar

O registro de consumo é programático: eventBus.inscrever(tipo, colonia, manipulador) dentro do iniciar() do service da colônia, após o replay de eventos perdidos (protocolo da seção 3.6). Não há decorator de consumo. A reentrega e o reprocessamento da DLQ usam a mesma lista de ouvintes por tipo e colônia.

A interface é o contrato que isola o transporte. Trocar EventEmitter2 por Kafka no futuro substitui a implementação interna sem alterar uma linha nas colônias.

A N-0a não tem BFF acoplado, não processa demanda e não altera payload de negócio. O controller REST expõe consulta ao log e o reprocessamento manual da DLQ, ambos restritos a operadores. Nenhuma rota publica eventos de negócio.


Todas as tabelas do núcleo residem no schema core do PostgreSQL. Este schema é de uso exclusivo dos módulos N-0a, N-0b e N-0c. Nenhuma colônia de negócio escreve ou lê diretamente tabelas deste schema. Todo acesso passa pelo EventBusService ou pelo RegistryService.

Registro append-only de todo evento que já transitou pelo sistema. Nenhuma linha é atualizada ou removida. É a memória imutável do sistema.

CREATE SCHEMA IF NOT EXISTS core;
CREATE TABLE core.event_log (
sequence_number BIGSERIAL NOT NULL,
event_id UUID NOT NULL,
tipo VARCHAR(255) NOT NULL,
versao_schema VARCHAR(20) NOT NULL DEFAULT '1.0.0',
timestamp TIMESTAMPTZ(2) NOT NULL DEFAULT CURRENT_TIMESTAMP,
origem VARCHAR(100) NOT NULL,
correlacao_id UUID,
payload JSONB NOT NULL,
criado_em TIMESTAMPTZ(2) NOT NULL DEFAULT CURRENT_TIMESTAMP,
CONSTRAINT event_log_pkey PRIMARY KEY (sequence_number)
);
CREATE UNIQUE INDEX event_log_event_id_key ON core.event_log (event_id);
CREATE INDEX event_log_tipo_idx ON core.event_log (tipo);
CREATE INDEX event_log_correlacao_id_idx ON core.event_log (correlacao_id);
CREATE INDEX event_log_timestamp_idx ON core.event_log (timestamp);
CREATE INDEX event_log_origem_idx ON core.event_log (origem);
CREATE INDEX event_log_tipo_timestamp_idx ON core.event_log (tipo, timestamp);
Coluna Tipo Descrição
sequence_number BIGSERIAL Sequência global monotônica. Atribuída pelo PostgreSQL no INSERT. É a ordem canônica do sistema. No Prisma o campo é BigInt.
event_id UUID UNIQUE Identificador único do evento. Gerado pela colônia publicadora ou pelo Event Bus se ausente.
tipo VARCHAR(255) Nome completo do tipo conforme Registry (ex: demanda.recebida). Validado contra o Registry antes do INSERT.
versao_schema VARCHAR(20) Versão do schema do payload no Registry (semver). Padrão '1.0.0' quando a colônia não informa.
timestamp TIMESTAMPTZ(2) Momento da publicação. Fornecido pela colônia ou preenchido pelo Event Bus com NOW().
origem VARCHAR(100) Identificador da colônia que publicou (ex: D-1a, D-3). Campo auditável.
correlacao_id UUID nullable ID de rastreamento de ponta a ponta. Vincula eventos do mesmo fluxo (ex: mesma sessão de captura). Pode diferir do event_id.
payload JSONB Corpo do evento conforme schema definido no Registry. A persistência passa pela redação de LGPD: textos viram hash, campos pessoais saem e coordenadas exatas são arredondadas. O payload completo segue em memória para os consumidores.
criado_em TIMESTAMPTZ(2) Timestamp de persistência no banco (pode diferir do timestamp do evento se houver atraso na rede ou fila).

No Prisma, sequence_number é BigInt (o obterMaiorSequence() retorna bigint), payload é Json e os timestamps usam Timestamptz(2). correlacao_id é nullable e pode receber valor diferente do event_id.

  • sequence_number BIGSERIAL — monotônico para MVP monolítico. Suficiente para ordenação causal dentro do processo único. Na migração para microsserviços, substituir por ULID ou Snowflake com partição temporal.
  • payload JSONB — permite consulta futura indexada via GIN para busca em campos específicos do payload sem desserializar. Útil para queries de auditoria (ex: “todos os eventos com demanda_id = X”).
  • event_id UNIQUE — garante idempotência: se a mesma colônia republicar o mesmo evento (por retry), o INSERT falha com violação de unicidade e o Event Bus retorna o registro existente.
  • Sem foreign key para tabelas do Registry. A validação é lógica via RegistryService, não via constraint de banco. O Registry pode ter tipos removidos ou descontinuados que já tiveram eventos publicados, e o log histórico permanece íntegro.
  • Schema core como namespace — isola as tabelas do núcleo das tabelas de colônias de negócio. Cada colônia tem seu próprio schema.

Registro de eventos cujo processamento por um consumidor falhou. Cada linha representa uma falha em uma colônia específica para um evento específico.

CREATE TABLE core.dead_letter_queue (
id BIGSERIAL NOT NULL,
event_id UUID NOT NULL,
sequence_number BIGINT NOT NULL,
colonia VARCHAR(100) NOT NULL,
handler_nome VARCHAR(255),
mensagem_erro TEXT NOT NULL,
pilha_erro TEXT,
tentativas INTEGER NOT NULL DEFAULT 1,
ultima_tentativa TIMESTAMPTZ(2) NOT NULL DEFAULT CURRENT_TIMESTAMP,
status VARCHAR(20) NOT NULL DEFAULT 'pendente',
criado_em TIMESTAMPTZ(2) NOT NULL DEFAULT CURRENT_TIMESTAMP,
resolvido_em TIMESTAMPTZ(2),
CONSTRAINT dead_letter_queue_pkey PRIMARY KEY (id),
CONSTRAINT dead_letter_queue_event_id_fkey
FOREIGN KEY (event_id) REFERENCES core.event_log(event_id)
);
CREATE UNIQUE INDEX dead_letter_queue_event_id_colonia_key ON core.dead_letter_queue (event_id, colonia);
CREATE INDEX dead_letter_queue_status_idx ON core.dead_letter_queue (status);
CREATE INDEX dead_letter_queue_colonia_idx ON core.dead_letter_queue (colonia);
CREATE INDEX dead_letter_queue_status_colonia_idx ON core.dead_letter_queue (status, colonia);

O índice único dead_letter_queue_event_id_colonia_key sustenta a deduplicação por (event_id, colonia). As transições reais são pendente → retrying → resolved ou failed_permanent, com reabertura para pendente em nova falha. resolvido_em é preenchido também em failed_permanent.

Coluna Tipo Descrição
id BIGSERIAL PK interna da DLQ.
event_id UUID FK Referência ao evento no event_log.
sequence_number BIGINT Redundância intencional para consultas de intervalo sem JOIN com event_log.
colonia VARCHAR(100) Colônia cujo manipulador falhou (ex: D-4).
handler_nome VARCHAR(255) Nome do manipulador registrado. Para diagnóstico.
mensagem_erro TEXT Mensagem da exceção capturada.
pilha_erro TEXT Stack trace completo.
tentativas INT Número de tentativas de reprocessamento. Incrementado a cada retry.
ultima_tentativa TIMESTAMPTZ Timestamp da última tentativa.
status VARCHAR(20) pendente (aguardando replay), retrying (em processamento), resolved (reprocessado com sucesso), failed_permanent (desistência após N tentativas).
resolvido_em TIMESTAMPTZ Preenchido quando status muda para resolved ou failed_permanent.

Duas migrations criam o esquema do núcleo:

  1. 20260808112047_init_core — cria o schema core, as tabelas event_log e dead_letter_queue, os índices e a chave estrangeira entre elas.
  2. 20260814154200_add_dlq_unique_event_colonia — adiciona o índice único (event_id, colonia) que sustenta a deduplicação da DLQ.

Migrations futuras (Fase 2+): partição por timestamp no event_log, índices GIN no payload, índices parciais para queries de auditoria frequentes.

A única relação dentro do schema core é a FK de dead_letter_queue.event_id → event_log.event_id. Não há outras FKs dentro do schema. O Registry (N-0b) mantém suas próprias tabelas no mesmo schema, mas sem FKs entre as tabelas do Event Bus e do Registry. O acoplamento é lógico, não estrutural.


O Event Bus não publica nem consome eventos de negócio. É o canal. Os contratos descritos aqui são o que o Event Bus impõe a quem publica e o que entrega a quem consome.

Toda colônia que publica fornece ao EventBusService.publicar() um PublicarEventoDto:

interface PublicarEventoDto {
tipo: string; // obrigatório — validado contra o Registry
payload: Record<string, unknown>; // obrigatório — validado contra o schema do Registry
origem: string; // obrigatório — código da colônia (ex: 'D-1a')
event_id?: string; // opcional — UUID v4. Se ausente, gerado pelo Event Bus
versao_schema?: string; // opcional — se ausente, literal '1.0.0'
timestamp?: string; // opcional — ISO-8601. Se ausente, NOW()
correlacao_id?: string; // opcional — se ausente, = event_id
}

Toda colônia que consome recebe do manipulador registrado via inscrever() um EventoConsultado:

interface EventoConsultado {
sequence_number: bigint; // BIGSERIAL — atribuído pelo banco no INSERT
event_id: string; // UUID v4
tipo: string; // ex: 'demanda.recebida'
versao_schema: string; // ex: '1.0.0'
timestamp: string; // ISO-8601
origem: string; // colônia publicadora
correlacao_id: string | null; // UUID de rastreamento; nullable no banco
payload: Record<string, unknown>; // corpo do evento conforme entregue ao consumidor
criado_em: string; // ISO-8601 de persistência
}

O consumo ao vivo entrega o payload completo em memória. O replay e a reentrega da DLQ leem o core.event_log e entregam o payload já redigido. Manipuladores que dependem de conteúdo devem tolerar campos ausentes, conforme o README.md do repo api.

O fluxo dentro do EventBusService.publicar() segue esta ordem exata:

1. Validar o formato do payload
→ obrigatório ser objeto não nulo
2. Normalizar os campos
→ event_id: dto.event_id ?? uuidv4()
→ versao_schema: dto.versao_schema ?? '1.0.0'
→ timestamp: dto.timestamp ?? new Date().toISOString()
→ correlacao_id: dto.correlacao_id ?? event_id
3. Validar formatos
→ event_id e correlacao_id como UUID
→ timestamp como ISO-8601
4. Validar tipo e versão contra o Registry
→ RegistryService.validar({ tipo, versao_schema }) cobre tipo inexistente, tipo deprecated e versão sem schema
5. Validar o payload contra o schema do Registry
→ RegistryService.validarPayload(tipo, versao, payload) com Ajv draft 2020-12; retorna a lista de erros
6. Redigir o payload para a persistência
→ redigirPayload(tipo, payload): textos viram hash, campos pessoais são removidos, coordenadas exatas são arredondadas
7. PERSISTIR no event_log
→ INSERT via Prisma
→ Se violação de UNIQUE(event_id) detectada por código P2002: busca o registro existente e retorna sem reemitir (idempotência)
8. ENTREGAR aos consumidores
→ emit do EventEmitter2 com o payload completo em memória
→ cada manipulador registrado é invocado sob o wrapper protegido (span, DLQ, métricas)
→ se o manipulador lança exceção: registra na dead_letter_queue e o cursor da colônia NÃO avança
9. Retornar o EventoConsultado ao publicador

Persistir antes de entregar garante que o evento esteja no log antes que qualquer consumidor o processe. Se o banco falhar na persistência, o evento não é publicado e o publicador recebe o erro. Se um manipulador falhar na entrega, o evento já está no log e a DLQ garante a reentrega. A validação e a redação acontecem antes da persistência; a entrega usa o payload completo.

Se uma colônia publicar o mesmo event_id duas vezes (retry de rede, timeout):

  • Primeira chamada: INSERT bem-sucedido, o evento é emitido e o EventoConsultado é retornado.
  • Segunda chamada: o INSERT falha por violação de unicidade em event_id, a exceção é capturada, o registro existente é buscado e retornado sem nova emissão.

O mesmo event_id nunca é entregue duas vezes aos consumidores por erro de retry do publicador.

Quando um manipulador registrado via inscrever() lança exceção:

1. O wrapper interno de executarProtegido() captura o erro
2. Finaliza o span com status 'failure' e registra a métrica de falha
3. Loga o erro com nível ERROR (para consumo pela Observabilidade N-0c)
4. Registra na dead_letter_queue:
→ event_id, sequence_number, colonia, handler_nome, mensagem_erro, pilha_erro
→ status assume o default 'pendente' e tentativas assume 1
5. Atualiza os gauges da DLQ da N-0c
→ Falha ao gravar na DLQ é logada e não interrompe o fluxo
6. NÃO propaga o erro para outros manipuladores
→ Cada manipulador é independente. Um falhar não interrompe os demais

O replay é iniciado pela colônia consumidora, não pelo Event Bus. A colônia decide quando reprocessar:

1. Colônia consulta seu último sequence_number processado (tabela própria <id>.consumer_offset)
2. No boot, faz seed dos cursors ausentes com eventBus.obterMaiorSequence() — nunca sobrescreve cursor existente
3. Chama eventBus.replayDeSequence(cursor, [tipos de interesse])
4. O replay pagina internamente em lotes de 1000, com teto de 100 iterações (100.000 eventos); a colônia recebe o conjunto acumulado
5. Colônia processa cada evento chamando seu próprio manipulador
6. Colônia atualiza seu cursor para o último sequence_number processado com sucesso

O seed com obterMaiorSequence() e o protocolo completo (seed, replay e registro de consumidores) são obrigatórios em toda colônia consumidora, conforme a seção “Protocolo de consumo com replay” do AGENTS.md do repo api.

A DLQ é um caso especial de reentrega: POST /api/events/dlq/reprocessar (corpo { "colonia": "D-5" }) reprocessa os eventos pendentes de uma colônia, e EventBusService.reentregarParaColonia(tipo, colonia, evento) reentrega programaticamente um evento com a mesma proteção do consumo ao vivo (span, DLQ, métricas).


função publicar(dto: PublicarEventoDto) → EventoConsultado:
// 1. Forma do payload
se dto.payload não é objeto ou é nulo:
lançar erro "Payload deve ser um objeto não nulo"
// 2. Normalização de campos
eventId = dto.event_id ?? uuidv4()
versao = dto.versao_schema ?? '1.0.0'
ts = dto.timestamp ?? new Date().toISOString()
corrId = dto.correlacao_id ?? eventId
// 3. Validação de formatos
se eventId ou corrId não são UUID:
lançar erro
se ts não é ISO-8601:
lançar erro
// 4. Validação contra o Registry
se não registry.validar({ tipo: dto.tipo, versao_schema: versao }):
lançar erro "Tipo de evento não registrado no Registry"
erros = registry.validarPayload(dto.tipo, versao, dto.payload)
se erros não está vazio:
lançar erro "Payload inválido para <tipo>@<versao>"
// 5. Persistência e entrega (dentro do contexto de trace)
payloadPersistido = redigirPayload(dto.tipo, dto.payload)
tentar:
registro = eventLogRepo.inserir({
event_id: eventId,
tipo: dto.tipo,
versao_schema: versao,
timestamp: new Date(ts),
origem: dto.origem,
correlacao_id: corrId,
payload: payloadPersistido,
})
capturar violação de unicidade (código P2002):
existente = eventLogRepo.buscarPorEventId(eventId)
se existente não é nulo:
logar warning "Evento duplicado detectado, idempotência aplicada"
retornar existente // sem reemitir
senão:
propagar o erro
// 6. Entrega aos consumidores com o payload completo
registroEntrega = { ...registro, payload: dto.payload }
emitir(registroEntrega.tipo, registroEntrega)
// 7. Span e métrica de sucesso
finalizarSpan(spanId, 'success')
registrar métrica de evento processado
retornar registroEntrega

O registro é inscrever(tipo, colonia, manipulador), sem decorator, dentro do iniciar() do service da colônia, após o replay de eventos perdidos.

função inscrever(tipo, colonia, manipulador):
acrescentar { colonia, manipulador } à lista de ouvintes do tipo
eventEmitter.on(tipo, (evento) => {
executarProtegido(tipo, colonia, evento, manipulador)
})
função executarProtegido(tipo, colonia, evento, manipulador):
spanId = iniciarSpan('event.process', { tipo, colonia, event_id: evento.event_id })
tentar:
executarComTraceContext({ correlacao_id, event_id, colonia }, () => manipulador(evento))
finalizarSpan(spanId, 'success')
registrar métrica 'success'
capturar erro:
mensagem = erro.message
pilha = erro.stack
finalizarSpan(spanId, 'failure', erro)
registrar métrica 'failure'
logar erro "Falha em manipulador de evento"
tentar:
deadLetterRepo.inserir({
event_id: evento.event_id,
sequence_number: evento.sequence_number,
colonia: colonia,
handler_nome: manipulador.name || 'anônimo',
mensagem_erro: mensagem,
pilha_erro: pilha,
})
atualizarGaugesDLQ()
capturar erro de DLQ:
logar "Falha ao registrar na DLQ"

O dedup da DLQ fica por conta do índice UNIQUE (event_id, colonia). A reentrega programática e o reprocessamento usam a mesma lista de ouvintes por tipo e colônia.

função consultarEventos(filtros: ConsultarEventosDto) → ResultadoPaginado<EventoConsultado>:
limite = min(filtros.limite ?? 100, 500) // cap de segurança
deslocamento = filtros.deslocamento ?? 0
where = filtros aplicáveis:
tipo, origem, correlacao_id
timestamp entre timestamp_inicio e timestamp_fim
sequence_number entre sequence_inicio e sequence_fim
[total, linhas] = em paralelo:
COUNT(where)
findMany(where, orderBy sequence_number asc, take limite, skip deslocamento)
retornar { linhas, total, limite, deslocamento }

A paginação é interna: o Event Bus busca em lotes de 1000 até esvaziar, com teto de 100 iterações e log de progresso, e retorna o conjunto acumulado.

função replayDeSequence(deSequence, tipos?) → EventoConsultado[]:
acumulados = []
cursor = deSequence
para iteracao de 1 até 100:
lote = eventLogRepo.replayDeSequence(cursor, tipos, 1000)
se lote vazio: parar
acumulados.push(...lote)
cursor = último sequence_number do lote
se tamanho do lote < 1000: parar
logar progresso do replay
se acumulados atingiram o teto de 100.000 eventos:
logar warning "Teto de iterações do replay atingido"
retornar acumulados

O limite de 1000 eventos por lote evita sobrecarga de memória. O teto de 100 iterações protege contra loops infinitos.

função obterEstatisticasDLQ() → DLQEstatisticas:
registros = carregar entradas com status 'pendente' ou 'retrying', selecionando a colônia
pendentes_por_colonia = {}
total_pendentes = 0
para cada registro:
pendentes_por_colonia[colonia] += 1
total_pendentes += 1
retornar { pendentes_por_colonia, total_pendentes }

A N-0a atualiza o gauge dead_letter_queue_pending da N-0c depois de registrar uma falha e depois de resolver ou reabrir uma entrada. A contagem inclui pendente e retrying, e o índice dead_letter_queue_status_colonia_idx cobre a consulta.

Caso Comportamento
Evento com tipo não registrado publicar() lança erro na validação RegistryService.validar, antes de persistir. Nada entra no log.
Evento duplicado (mesmo event_id) INSERT falha com violação de unicidade (P2002). O registro existente é retornado e o evento não é reemitido.
Evento fora de ordem temporal O sequence_number do banco define a ordem, não o timestamp do evento. O timestamp é metadado informativo do publicador.
Manipulador lança exceção Erro capturado pelo wrapper. DLQ registrada. Demais manipuladores do mesmo evento não são afetados.
EventEmitter2 atinge maxListeners Configurado com verboseMemoryLeak: true, que loga aviso no console. O limite é ajustado no módulo, sem variável de ambiente.
Timeout de manipulador MVP: manipuladores são síncronos ou async com await. Se um manipulador nunca resolve, bloqueia a thread. Solução Fase 2: timeout wrapper com Promise.race.
DLQ cresce sem reprocessamento Reentrega manual via POST /api/events/dlq/reprocessar (por colônia) e reentrega programática via reentregarParaColonia. Em falha de manipulador, o cursor da colônia fica parado e o replay do boot seguinte retenta. Alerta de Observabilidade (N-0c) sinaliza DLQ acima de limiar configurável.
Correlação entre eventos correlacao_id é propagado pelas colônias. Se uma colônia publica um evento derivado, inclui o correlacao_id do evento de origem. O Event Bus não altera, apenas transporta.
INSERT concorrente mesmo event_id PostgreSQL garante atomicidade do UNIQUE constraint. A primeira transação vence, a segunda recebe o erro e cai no caminho de idempotência.
Payload com estruturas aninhadas profundas JSONB suporta até 255 níveis no PostgreSQL. Limite prático: payloads de negócio raramente excedem 5 níveis. Se ocorrer, erro de validação no INSERT.

Persistir antes de entregar. Se o Event Bus emitisse antes de persistir e o banco falhasse após a entrega, o evento teria sido processado pelos consumidores mas não estaria no log, o que quebra a rastreabilidade e impede o replay. Persistir primeiro garante que todo evento processado está no log.

BIGSERIAL e não ULID. No monolito, BIGSERIAL do PostgreSQL é a fonte mais simples e confiável de monotonicidade. Não requer biblioteca externa, não tem risco de colisão, é nativo do banco. A migração para ULID será necessária na transição para microsserviços, quando o sequence_number global perde o sentido e cada serviço passa a ter ordenação local. O tipo do campo é detalhe interno do barramento.

Validação completa de payload já no MVP. A validação de payload contra os JSON Schemas do Registry (Ajv, draft 2020-12) está ligada no publicar() desde a segunda auditoria. O MVP não usa validação relaxada e não existe flag de validação parcial.

Redação do payload na persistência. O payload gravado no core.event_log passa pela redação de LGPD definida em redacao/redacao-eventos.ts, com allowlist por tipo de evento: textos pessoais viram hash, campos pessoais são removidos e coordenadas exatas são arredondadas para o nível da unidade cívica. O payload completo segue em memória para os consumidores ao vivo. O log registra o processo, nunca o dado pessoal.

Wrapper de inscrever em vez de modificar o EventEmitter2. Modificar o comportamento interno do EventEmitter2 (subclass, monkey-patch) é frágil e quebra com atualizações de versão da biblioteca. O wrapper explícito no inscrever() é transparente, testável e desacoplado da implementação interna do emissor.


5. Integração com o Barramento e Outras Colônias

Seção intitulada “5. Integração com o Barramento e Outras Colônias”
// Exemplo na D-1a (Captura)
@Injectable()
export class CapturaService {
constructor(private readonly eventBus: EventBusService) {}
async receberDemanda(input: DemandaInput): Promise<{ demanda_id: string }> {
const demandaId = uuidv4();
await this.eventBus.publicar({
tipo: 'demanda.recebida',
origem: 'D-1a',
versao_schema: '1.1.0',
event_id: demandaId,
correlacao_id: demandaId,
payload: {
demanda_id: demandaId,
texto_bruto: input.texto,
localizacao_bruta: { lat: input.lat, lng: input.lng },
cidadao_id: input.cidadaoId,
canal: 'app',
},
});
return { demanda_id: demandaId };
}
}

O cursor é seedado com a maior sequência do log, nunca 0. O protocolo completo está na seção “Protocolo de consumo com replay” do AGENTS.md do repo api.

// Exemplo na D-4 (Priorização)
@Injectable()
export class PriorizacaoService implements OnModuleInit {
private lastProcessedSequence: bigint = 0n;
constructor(
private readonly eventBus: EventBusService,
private readonly offsetRepo: ConsumerOffsetRepository, // schema próprio da D-4
) {}
async onModuleInit() {
// 1. Seed do cursor com obterMaiorSequence() quando a linha não existe
// 2. Recupera o cursor do banco próprio
this.lastProcessedSequence = await this.offsetRepo.getLastSequence('demanda.categorizada');
// 3. Replay dos eventos acima do cursor
const missedEvents = await this.eventBus.replayDeSequence(
this.lastProcessedSequence,
['demanda.categorizada', 'parametros.atualizados'],
);
for (const event of missedEvents) {
await this.processarEvento(event);
}
// 4. Registra os consumidores ao vivo
this.eventBus.inscrever(
'demanda.categorizada', 'D-4',
this.onDemandaCategorizada.bind(this),
);
this.eventBus.inscrever(
'parametros.atualizados', 'D-4',
this.onParametrosAtualizados.bind(this),
);
}
private async onDemandaCategorizada(event: EventoConsultado): Promise<void> {
// ... lógica de ranqueamento ...
await this.offsetRepo.updateLastSequence('demanda.categorizada', event.sequence_number);
}
private async onParametrosAtualizados(event: EventoConsultado): Promise<void> {
// ... recalcular ranking com novos parâmetros ...
await this.offsetRepo.updateLastSequence('parametros.atualizados', event.sequence_number);
}
}

5.3 Tabela de consumer offset (schema de cada colônia)

Seção intitulada “5.3 Tabela de consumer offset (schema de cada colônia)”

Cada colônia mantém seu cursor no próprio schema:

-- Exemplo no schema d4
CREATE TABLE d4.consumer_offset (
tipo_evento VARCHAR(255) NOT NULL,
last_sequence BIGINT NOT NULL DEFAULT 0,
updated_at TIMESTAMPTZ(2) NOT NULL DEFAULT CURRENT_TIMESTAMP,
CONSTRAINT consumer_offset_pkey PRIMARY KEY (tipo_evento)
);

O seed do cursor usa obterMaiorSequence() do barramento, nunca 0.

O offset é por tipo de evento, permitindo replay seletivo: se apenas demanda.categorizada ficou atrasado, só ele é reprocessado. Os demais tipos mantêm seu cursor.

@Controller('events')
@UseGuards(PapelOperadorGuard)
export class EventBusController {
constructor(private readonly eventBus: EventBusService) {}
@Get()
@ApiExcludeEndpoint()
async consultarEventos(@Query() filtros: ConsultarEventosDto) {
return this.eventBus.consultarEventos(filtros);
}
@Get(':eventId')
@ApiExcludeEndpoint()
async obterEvento(@Param('eventId') eventId: string) {
const evento = await this.eventBus.obterEventoPorId(eventId);
if (!evento) throw new NotFoundException('Evento não encontrado');
return evento;
}
@Get('trace/:correlacaoId')
@ApiExcludeEndpoint()
async obterTrace(@Param('correlacaoId') correlacaoId: string) {
return this.eventBus.consultarEventos({
correlacao_id: correlacaoId,
limite: 500,
});
}
}
Método Rota Descrição
GET /api/events Consulta paginada com filtros: tipo, origem, correlacao_id, timestamp_inicio, timestamp_fim, sequence_inicio, sequence_fim, limite, deslocamento. O sequence_number é serializado como string para preservar a precisão do bigint.
GET /api/events/:eventId Evento único por event_id.
GET /api/events/trace/:correlacaoId Eventos de um mesmo fluxo (trace distribuído), com limite de 500. Retorna em ordem de sequence_number.
POST /api/events/dlq/reprocessar Reentrega manual dos eventos pendentes de uma colônia na DLQ. Corpo: { "colonia": "D-5" }. Throttle 1 req/min.

As quatro rotas são internas de operação, fora do Swagger (@ApiExcludeEndpoint) e protegidas pelo PapelOperadorGuard (JWT com papel admin e allowlist de IP). As rotas de consulta herdam o rate limiting global por IP; a rota de reprocessamento da DLQ tem limite próprio de 1/min. A lista completa está no README.md do repo api, seção “Rotas internas de operação”.

O EventBusModule importa RegistryModule. O EventBusService injeta RegistryService para duas operações na publicação:

  1. publicar() — valida tipo e versão com RegistryService.validar, que cobre tipo inexistente, tipo deprecated e versão sem schema.
  2. publicar() — valida o payload com RegistryService.validarPayload contra o JSON Schema da versão.

A dependência com o Registry é síncrona porque a validação ocorre no caminho quente de toda publicação. Um round-trip assíncrono (evento → Registry → evento de volta) inviabilizaria a latência de publicação. O ObservabilityModule é a outra dependência importada pelo módulo, restrita a trace e métricas.

O BFF da D-1a não chama o Event Bus via HTTP. A comunicação é sempre via injeção de EventBusService no mesmo processo. O controller REST do Event Bus é usado apenas para consulta de auditoria externa, nunca como mecanismo de publicação de eventos.

O Event Bus não consome projeções de leitura de outras colônias. A única consulta externa é ao Registry (validação de tipo), e o Registry é uma tabela de referência, não uma projeção de leitura.


O EventBusService não aplica rate limiting à publicação. Essa responsabilidade é das colônias de entrada (D-1a, L-1), que estão na borda do sistema e lidam diretamente com input de cidadão. As rotas internas do controller herdam o rate limiting global por IP e a rota de reprocessamento da DLQ tem limite próprio de 1/min.

Limite Valor MVP Justificativa
Tamanho máximo do payload 1 MB (JSONB) Suficiente para eventos de negócio. Anexos e mídia trafegam como URL na D-1c, não como base64 no payload.
maxListeners por tipo de evento 30 Cada tipo de evento tem em média 2-3 consumidores. O valor é fixo no módulo, sem variável de ambiente.
Tamanho de página na query API Máximo 500 eventos Evita vazamento de memória em queries de auditoria. Paginação obrigatória acima disso.
Limite de replay por chamada Lotes de 1000 eventos O replay interno pagina em lotes de 1000 até esvaziar, com teto de 100 iterações; a colônia recebe o conjunto acumulado.
Índice Query atendida
event_log_pkey (sequence_number) Replay: WHERE sequence_number > X ORDER BY sequence_number
event_log_event_id_key (UNIQUE) Idempotência: WHERE event_id = ? + GET /:eventId
event_log_tipo_idx WHERE tipo = ? ORDER BY sequence_number
event_log_correlacao_id_idx Trace: WHERE correlacao_id = ? ORDER BY sequence_number
event_log_timestamp_idx Auditoria por período: WHERE timestamp BETWEEN ? AND ?
event_log_origem_idx Consulta paginada por colônia de origem
event_log_tipo_timestamp_idx WHERE tipo = ? AND timestamp BETWEEN ? AND ?

O sequence_number é a PK (BIGSERIAL), naturalmente clusterizada. Scans de intervalo são eficientes. Para consultas de replay (acesso sequencial ao log), nenhum índice adicional é necessário.

Sem cache no Event Bus. O event_log é append-only com acesso majoritariamente sequencial. Cache seria contraproducente:

  • Escrita: sempre no final do log. PostgreSQL lida bem com append em B-tree clusterizada.
  • Leitura de replay: acesso sequencial a partir de um ponto. Cache de páginas do PostgreSQL (shared_buffers) já cobre.
  • Leitura de auditoria (query por tipo ou correlacao_id): esporádica, não justifica cache.

A DLQ é de baixíssimo volume e não requer cache.

Cenário Eventos/dia Tamanho estimado/dia (payload médio 2 KB) Retenção
PoC (1 bairro) ~500 ~1 MB Indefinida (append-only)
MVP (1 município) ~10.000 ~20 MB Indefinida
Fase 2 (regional) ~500.000 ~1 GB Particionamento por mês

Para Fase 2: partição por timestamp no event_log e archive de partições antigas para storage frio (S3/Parquet). A localização do log é detalhe interno: as colônias chamam replayDeSequence e a implementação resolve onde buscar.


Teste unitário do EventBusService:

O módulo de teste injeta mocks para as dependências externas: EventLogRepository, DeadLetterRepository, RegistryService, ObservabilityService e MetricsService. O EventEmitter2 é real (fornecido por EventEmitterModule.forRoot() em teste), pois é parte do comportamento que se quer testar.

// Configuração do módulo de teste
beforeEach(async () => {
const module = await Test.createTestingModule({
imports: [EventEmitterModule.forRoot()],
providers: [
EventBusService,
{ provide: EventLogRepository, useValue: mockEventLogRepo },
{ provide: DeadLetterRepository, useValue: mockDLQRepo },
{ provide: RegistryService, useValue: mockRegistryService },
{ provide: ObservabilityService, useValue: mockObs },
{ provide: MetricsService, useValue: mockMetrics },
],
}).compile();
service = module.get(EventBusService);
});
// Mock do Registry — aceita tipo, versão e payload
mockRegistryService.validar.mockResolvedValue(true);
mockRegistryService.validarPayload.mockResolvedValue([]);
// Mock do EventLogRepository — simula INSERT
mockEventLogRepo.inserir.mockImplementation((data) =>
Promise.resolve({ ...data, sequence_number: 1n, criado_em: new Date().toISOString() })
);

Happy path:

# Cenário Verificação
T1 publicar() com dados válidos Evento persistido no event_log. EventEmitter.emit() chamado com o registro. EventoConsultado retornado com sequence_number preenchido.
T2 inscrever() registra manipulador Manipulador é invocado quando emit() ocorre para o tipo.
T3 consultarEventos() com filtro de tipo Retorna apenas eventos do tipo especificado. Paginação funciona.
T4 replayDeSequence() Retorna eventos com sequence_number > deSequence, ordenados ASC.
T5 obterEventoPorId() Retorna o evento correto ou null.

Falhas e bordas:

# Cenário Verificação
T6 publicar() com tipo não registrado Lança erro na validação do Registry. Nada é inserido no event_log.
T7 publicar() com payload null Lança erro “Payload deve ser um objeto não nulo”. Nada é inserido no event_log.
T8 publicar() com payload não-objeto (string, array) Mesmo comportamento de T7. Nada é inserido no event_log.
T9 publicar() com event_id duplicado Retorna o registro existente. emit() NÃO é chamado segunda vez.
T10 Manipulador lança exceção Erro capturado. Registro criado na dead_letter_queue com status: 'pendente'. Demais manipuladores executam normalmente.
T11 Manipulador lança exceção — conteúdo da DLQ event_id, sequence_number, colonia, handler_nome, mensagem_erro e pilha_erro preenchidos corretamente.
T12 publicar() sem versao_schema Usa o literal '1.0.0'.
T13 publicar() sem correlacao_id Preenche com event_id.
T14 consultarEventos() sem filtros — limite de página Retorna no máximo 500 eventos mesmo que limite solicitado seja maior.
T15 replayDeSequence() sem eventos novos Retorna array vazio.
T16 replayDeSequence() com filtro de tipo Retorna apenas eventos dos tipos especificados, ordenados por sequence_number.

Teste de integração (com PostgreSQL de teste via Testcontainers ou banco local):

# Cenário Verificação
T17 Ciclo completo: publicar → consumir → consultar Evento no log. Manipulador executado. DLQ vazia. Consulta retorna o evento.
T18 Ciclo com falha: publicar → manipulador falha → DLQ → replay → sucesso DLQ com 1 entrada status: 'pendente'. Após o replay: DLQ status: 'resolved'.
T19 Idempotência com PostgreSQL real Dois publicar() com o mesmo event_id. O primeiro insere. O segundo retorna o mesmo registro sem erro. count(*) no event_log = 1.
T20 Replay após falha seletiva Manipulador de tipo A falha, tipo B processa. O replay de A recupera o evento pendente. O offset de B não é afetado.
-- Eventos de exemplo para testar queries e replay
INSERT INTO core.event_log (event_id, tipo, versao_schema, timestamp, origem, correlacao_id, payload)
VALUES
(
'a1b2c3d4-e5f6-7890-abcd-ef1234567890',
'demanda.recebida',
'1.0.0',
'2026-06-15T10:30:00Z',
'D-1a',
'a1b2c3d4-e5f6-7890-abcd-ef1234567890',
'{"demanda_id": "a1b2c3d4-e5f6-7890-abcd-ef1234567890", "texto_bruto": "Falta d''água na Rua das Flores", "localizacao_bruta": {"lat": -23.5505, "lng": -46.6333}, "cidadao_id": "u1", "canal": "app"}'
),
(
'b2c3d4e5-f6a7-8901-bcde-f12345678901',
'demanda.normalizada',
'1.0.0',
'2026-06-15T10:30:05Z',
'D-1b',
'a1b2c3d4-e5f6-7890-abcd-ef1234567890',
'{"demanda_id": "a1b2c3d4-e5f6-7890-abcd-ef1234567890", "titulo": "Falta de água na Rua das Flores", "descricao_limpa": "Falta dágua na Rua das Flores", "confianca_normalizacao": 0.95}'
),
(
'c3d4e5f6-a7b8-9012-cdef-123456789012',
'demanda.categorizada',
'1.0.0',
'2026-06-15T10:30:10Z',
'D-3',
'a1b2c3d4-e5f6-7890-abcd-ef1234567890',
'{"demanda_id": "a1b2c3d4-e5f6-7890-abcd-ef1234567890", "categoria_id": "1.1", "nivel_precedencia": 1, "score_horizontal": 95, "confianca_categorizacao": 0.88, "metodo": "automatico"}'
);

Funcionalidade Status
publicar() com validação de tipo, versão e payload (Registry) e persistência append-only MVP obrigatório
inscrever() com wrapper de erro → DLQ MVP obrigatório
event_log com sequence_number global monotônico MVP obrigatório
dead_letter_queue com status pendente/retrying/resolved/failed_permanent MVP obrigatório
Redação do payload na persistência (LGPD) MVP obrigatório
consultarEventos() via service e rotas internas do controller MVP obrigatório
replayDeSequence() para que colônias recuperem eventos perdidos MVP obrigatório
Idempotência por event_id MVP obrigatório
Logs estruturados em cada publicação e consumo (para consumo pela N-0c) MVP obrigatório
@Global() — módulo disponível para todas as colônias sem import explícito MVP obrigatório
Simplificação Justificativa Quando remover
EventEmitter2 in-process Único processo no monolito. Sem latência de rede entre colônias. Migrar para Kafka/Redpanda quando o monolito for particionado em microsserviços (Fase 2).
BIGSERIAL como sequence_number PostgreSQL único. Monotonicidade garantida pelo banco. Migrar para ULID quando houver múltiplas instâncias de banco ou escrita distribuída.
Sem timeout em manipuladores Manipuladores são síncronos no MVP. Se um travar, a thread trava. Problema conhecido e aceito para a fase. Adicionar Promise.race com timeout configurável na Fase 2.
Sem replay automático da DLQ Replay é manual ou via scheduled job. A colônia decide quando reprocessar. Na Fase 2, a N-0a pode opcionalmente notificar a colônia sobre eventos pendentes na DLQ.
Sem partição no event_log Volume de eventos no MVP (até ~10K/dia) não justifica particionamento. Adicionar partição por mês no event_log quando volume diário > 500K eventos.
Sem indexação GIN no payload Consultas ao log usam colunas indexadas (tipo, timestamp, correlacao_id). Busca por campo interno do payload é rara no MVP. Adicionar índice GIN no payload para queries de auditoria avançadas na Fase 2.
Sem fila de publicação assíncrona publicar() é síncrono: colônia espera INSERT + emit antes de continuar. Latência aceitável para MVP (INSERT local < 5ms). Na migração para microsserviços, publicar se torna async e o Event Bus confirma via callback/evento de confirmação.
  • Transporte distribuído (Kafka/Redpanda) com consumer groups e partições
  • Partição do event_log por timestamp
  • Índice GIN no payload para queries de auditoria profunda
  • Timeout wrapper em manipuladores (Promise.race)
  • Notificação automática de DLQ para colônias afetadas
  • Snapshot de event_log para storage frio (S3/Parquet)
  • @nestjs/bullmq ou similar para fila de publicação com backpressure

8.4 Verificação de conflitos com outras colônias

Seção intitulada “8.4 Verificação de conflitos com outras colônias”

Sem conflitos detectados. As decisões de MVP são autocontidas no Event Bus e não impõem restrições às colônias consumidoras:

  • As colônias publicam via EventBusService.publicar(), independente do transporte interno (EventEmitter2 ou Kafka).
  • As colônias consomem via EventBusService.inscrever(), com o mesmo contrato de manipulador.
  • A DLQ e o replay são mecanismos internos que não afetam o contrato público das colônias.
  • A redação do payload na persistência não impede que colônias validem e consumam os payloads internamente com seus próprios DTOs.
  • O @Global() do módulo Event Bus não conflita com a regra de isolamento. As colônias injetam o EventBusService do núcleo, que é dependência permitida.


Documento de especificação técnica de implementação. Aprovado e integrado.