Pular para o conteúdo

N-0c — Observabilidade

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


Camada transversal de monitoramento do núcleo. Instrumenta o barramento e as colônias sem interferir no fluxo de negócio. Três frentes: logs estruturados em JSON com mensagem, correlacao_id, event_id, colonia e os dados do registro; trace distribuído por correlacao_id, que encadeia os eventos de uma demanda ao longo das colônias; e métricas no formato Prometheus, com volume de eventos por tipo, latência de processamento por colônia, taxa de falha e tamanho da DLQ.

No MVP, os logs estruturados com event_id e correlacao_id já bastam para rastrear o ciclo completo de uma demanda. O trace completo é materializado pelo core.event_log da N-0a e exposto em GET /api/events/trace/:correlacaoId. A N-0c adiciona a instrumentação de tempo de processamento e as métricas de saúde.

Critério de saída: dado um correlacao_id ou um event_id, é possível recuperar o histórico completo de eventos no core.event_log e acompanhar as métricas de processamento. O endpoint /health retorna o status operacional do monolito, e /health/readiness retorna o status de banco, DLQ e storage.


A N-0c é um módulo @Global() do núcleo. Toda colônia injeta seus serviços sem importação explícita, porque observabilidade é cross-cutting. O EventBusModule (N-0a) é o único módulo que o importa, para instrumentar o wrapper de consumo. O AppModule registra o TraceInterceptor como APP_INTERCEPTOR global.

src/nucleo/n-0c-observabilidade/
├── n-0c-observabilidade.module.ts # @Global() — providers: ObservabilityService, MetricsService, HealthService, TraceInterceptor
├── n-0c-observabilidade.service.ts # Contexto de trace (ALS), spans e logs estruturados
├── n-0c-observabilidade.controller.ts # GET /api/observability/health, /health/readiness e /metrics
├── n-0c-observabilidade.interceptor.ts # TraceInterceptor — propaga correlacao_id nas requisições HTTP
├── dto/
│ └── health-response.dto.ts # Forma da resposta de health e readiness
├── health/
│ └── health.service.ts # Liveness e readiness com banco, DLQ e storage
├── metrics/
│ ├── metrics.service.ts # Métricas em memória, formato Prometheus manual
│ └── metrics.interface.ts
├── logging/
│ ├── trace-context.ts # AsyncLocalStorage
│ ├── logger.interface.ts # Interface do logger estruturado
│ └── redacao-log.ts # Redação de campos sensíveis nos logs
└── interfaces/
└── observability.interface.ts # IObservability e SpanResultado
@Global()
@Module({
providers: [ObservabilityService, MetricsService, HealthService, TraceInterceptor],
exports: [ObservabilityService, MetricsService],
controllers: [ObservabilityController],
})
export class ObservabilityModule {}
  • @Global(): toda colônia injeta ObservabilityService e MetricsService sem imports: [ObservabilityModule]. O acoplamento é zero da perspectiva da colônia.
  • O módulo não importa EventBusModule. A N-0a é quem importa a N-0c, e a dependência é unidirecional. A N-0c não conhece o Event Bus.
  • O módulo não importa RegistryModule. A N-0c não valida tipos de evento.
  • O TraceInterceptor é provido no módulo e registrado como APP_INTERCEPTOR no AppModule. Como os serviços exportados são globais, a instância resolve as dependências sem importação adicional.
  • A separação entre ObservabilityService (trace context, spans e logs) e MetricsService (métricas) é deliberada. Colônias que só precisam de métricas injetam apenas MetricsService. Colônias que precisam de contexto e log injetam ObservabilityService.
definirTraceContext(ctx: Partial<TraceContext>): void; // mescla o contexto no store atual do ALS
obterTraceContext(): TraceContext; // {} quando não há contexto
executarComTraceContext<T>(ctx: TraceContext, fn: () => Promise<T>): Promise<T>;
info(mensagem: string, dados?: Record<string, unknown>): void;
warn(mensagem: string, dados?: Record<string, unknown>): void;
error(mensagem: string, erro?: Error, dados?: Record<string, unknown>): void;
debug(mensagem: string, dados?: Record<string, unknown>): void;
iniciarSpan(name: string, metadata?: Record<string, unknown>): string; // retorna spanId
finalizarSpan(spanId: string, status?: 'success' | 'failure', erro?: Error): SpanResultado;
interface TraceContext {
correlacao_id?: string;
event_id?: string;
colonia?: string;
}
interface SpanResultado {
span_id: string;
name: string;
duration_ms: number;
status: 'success' | 'failure';
metadata: Record<string, unknown>;
}
registrarEventoProcessado(tipo: string, colonia: string, status: 'success' | 'failure', durationMs: number): void;
definirPendentesDLQ(colonia: string, quantidade: number): void;
registrarRequisicaoHttp(method: string, path: string, statusCode: number, durationMs: number): void;
obterMetricas(): Promise<string>; // formato texto do Prometheus

O método resetar() limpa todos os Map e é usado pelos testes para isolar cenários.

A N-0c não tem BFF acoplado, não processa demanda, não publica nem consome eventos de negócio. O controller REST expõe health check e métricas. O trace é exposto pela N-0a (GET /api/events/trace/:correlacaoId). A N-0c garante que os logs e as métricas de processamento carreguem correlacao_id, event_id e colonia.


A N-0c não possui tabelas próprias no MVP. A decisão se apoia em três fatores:

  1. O trace completo por correlacao_id já existe no core.event_log da N-0a. Cada evento publicado registra correlacao_id, origem, tipo e timestamp. Consultar SELECT * FROM core.event_log WHERE correlacao_id = ? ORDER BY sequence_number reconstrói o ciclo completo de uma demanda. A N-0c não duplica essa responsabilidade.

  2. As métricas de processamento residem em memória no MetricsService, em Map por combinação de labels, e são renderizadas sob demanda no endpoint /metrics. Um Prometheus externo coleta periodicamente e armazena no próprio time-series database. Não há persistência de métricas no PostgreSQL do monolito.

  3. Os logs estruturados são emitidos em JSON de uma linha para stdout pelo Logger do NestJS. Um coletor externo (Fluentd, Vector, Loki) ingere e persiste. O monolito não retém logs em banco. A retenção e a busca são responsabilidade da infraestrutura de observabilidade externa.

Configuração de alertas usa variáveis de ambiente, sem tabela de configuração:

Variável Descrição Default MVP
OBS_DLQ_WARN_THRESHOLD Entradas pendentes na DLQ que disparam mensagem de aviso no readiness 50
OBS_DLQ_CRITICAL_THRESHOLD Entradas pendentes que disparam status degraded 200

São as duas únicas variáveis de ambiente lidas pela N-0c.

Na Fase 2, quando o sistema operar em múltiplos processos e o volume de eventos exigir retenção de métricas históricas, a N-0c pode ganhar uma tabela core.alert_history para registro de alertas disparados e uma tabela core.metrics_snapshot para snapshots periódicos de métricas agregadas, úteis para dashboards de longo prazo sem dependência de Prometheus externo. Essas tabelas entram no schema core junto com as demais tabelas do núcleo.


A N-0c não publica nem consome eventos de negócio no barramento. Ela instrumenta o barramento, com hooks acoplados ao ciclo de vida dos eventos sem interferir no payload.

Quando a N-0a publica um evento, a N-0c é acionada em dois momentos.

Contexto de trace antes da persistência. A N-0a envolve a persistência e a entrega em observabilityService.executarComTraceContext({ correlacao_id, event_id, colonia: origem }, fn). Todo log emitido dentro da função carrega os três campos a partir do AsyncLocalStorage.

Métrica depois da entrega. A N-0a chama metricsService.registrarEventoProcessado(tipo, origem, 'success', duration_ms) após finalizarSpan. A métrica de publish registra apenas sucesso, porque uma falha na publicação lança exceção antes da emissão. A falha individual de consumidor é registrada pela instrumentação de subscribe.

Fluxo no EventBusService.publicar():

1. Validar a forma do payload, o event_id, o correlacao_id e o timestamp
2. Validar tipo, versão e payload contra o Registry
3. executarComTraceContext({ correlacao_id, event_id, colonia: origem }, ...)
4. Persistir no event_log com o payload redigido
5. iniciarSpan('event_bus.publish', { tipo, origem, event_id })
6. Emitir para os consumidores com o payload completo
7. finalizarSpan(spanId, 'success')
8. registrarEventoProcessado(tipo, origem, 'success', duration_ms)
9. Retornar o registro entregue

O log estruturado emitido nesses passos carrega automaticamente correlacao_id, event_id e colonia porque formatarMensagem() lê o AsyncLocalStorage. Nenhum console.log manual é necessário.

O wrapper executarProtegido() da N-0a é o ponto de instrumentação de processamento de eventos:

função executarProtegido(tipo, colonia, evento, manipulador):
spanId = iniciarSpan('event.process', { tipo, colonia, event_id: evento.event_id })
tentar:
executarComTraceContext({
correlacao_id: evento.correlacao_id,
event_id: evento.event_id,
colonia: colonia,
}, () => manipulador(evento))
span = finalizarSpan(spanId, 'success')
registrarEventoProcessado(tipo, colonia, 'success', span.duration_ms)
capturar erro:
span = finalizarSpan(spanId, 'failure', erro)
registrarEventoProcessado(tipo, colonia, 'failure', span.duration_ms)
logar erro "Falha em manipulador de evento"
inserir na DLQ { event_id, sequence_number, colonia, handler_nome, mensagem_erro, pilha_erro }
atualizar gauges da DLQ
propagar o erro

O handler executa dentro do contexto de trace. O erro é capturado, registrado na DLQ e propagado, e o cursor da colônia não avança. A instrumentação acrescenta a medição de duração, a métrica Prometheus e o log estruturado com os campos de trace. A reentrega programática (reentregarParaColonia) usa o mesmo wrapper, com span, DLQ e métricas.

A N-0c registra o TraceInterceptor como APP_INTERCEPTOR no AppModule. Toda requisição HTTP passa por ele:

intercept(ctx, next):
request = ctx.switchToHttp().getRequest()
response = ctx.switchToHttp().getResponse()
corrId = x-correlation-id, se UUID válido
?? x-request-id, se UUID válido
?? uuidv4()
definirTraceContext({ correlacao_id: corrId })
response.setHeader('x-correlation-id', corrId)
spanId = iniciarSpan('http.request', { method, path: request.route?.path ?? request.url })
aguardar next.handle() e, ao final:
sucesso: finalizarSpan(spanId, 'success')
registrarRequisicaoHttp(method, path, response.statusCode, duração)
erro: finalizarSpan(spanId, 'failure', erro)
registrarRequisicaoHttp(method, path, erro.status ?? 500, duração)

Para requisições que iniciam um fluxo de negócio, como o POST de demanda na D-1a, o front-end propaga o x-correlation-id devolvido na resposta nas chamadas seguintes. Se nenhum header válido estiver presente, o interceptor gera um UUID novo.

As métricas vivem em Map na memória do MetricsService e são renderizadas manualmente no formato texto do Prometheus em obterMetricas(), com linhas HELP e TYPE por métrica, _count, _sum, _bucket{le=...} e +Inf nos histogramas. Não há prom-client no MVP. Catálogo:

Nome Tipo Labels Descrição
event_processing_total Counter tipo, colonia, status Total de eventos processados, particionado por resultado
event_processing_duration_seconds Histogram (manual) tipo, colonia Duração do processamento, buckets [0.01, 0.05, 0.1, 0.5, 1, 5, 10, 30]
dead_letter_queue_pending Gauge colonia Entradas pendentes na DLQ, atualizado pela N-0a após registrar falha, resolver ou reabrir; a contagem inclui pendente e retrying
http_requests_total Counter method, path, status_code Requisições HTTP recebidas pelo monolito
http_request_duration_seconds Histogram (manual) method, path Duração de requisições HTTP

O histograma de duração de evento só observa status success. O histograma HTTP observa todas as requisições, de sucesso ou erro. Não existem event_publish_total nem métricas nodejs_* no MVP.

Os buckets correspondem a 10ms, 50ms, 100ms, 500ms, 1s, 5s, 10s e 30s. O bucket de 10ms captura processamentos rápidos, como categorização leve. O bucket de 30s captura anomalias, como handler travado aguardando I/O externo sem timeout. Os buckets cobrem desde normalização e categorização, na casa de dezenas de milissegundos, até georreferenciamento com geocodificação externa, que pode levar segundos. O histograma http_request_duration_seconds usa os mesmos buckets.

O log do span é emitido somente no finalizarSpan(), com duration_ms. Não há log de início do span.

Log de subscribe: DEPOIS. O log só é emitido quando o handler completa, com sucesso ou falha. Em falha, o log do span vem antes do registro na DLQ e carrega o campo duration_ms e o erro.

Log HTTP: DEPOIS. O log do span é emitido quando a resposta é enviada, com duration_ms. O status code não entra no log; ele entra na métrica http_requests_total.


classe ObservabilityService:
logger = new Logger(ObservabilityService.name)
spansAtivos = new Map<string, { span_id, name, inicio: bigint, metadata }>()
TAMANHO_MAX_SPANS = 10000
função definirTraceContext(ctx):
existente = traceContext.getStore()
traceContext.enterWith({ ...existente, ...ctx })
função obterTraceContext():
retornar traceContext.getStore() ?? {}
função executarComTraceContext(ctx, fn):
retornar traceContext.run({ ...ctx }, fn)
função iniciarSpan(name, metadata):
se spansAtivos.size >= TAMANHO_MAX_SPANS:
remover a chave mais antiga do Map
logar warning "Limite de spans ativos excedido"
spanId = uuidv4()
spansAtivos.set(spanId, {
span_id: spanId,
name: name,
inicio: process.hrtime.bigint(),
metadata: metadata ?? {},
})
retornar spanId
função finalizarSpan(spanId, status, erro):
span = spansAtivos.get(spanId)
se !span: lançar ErroSpanNaoEncontrado(spanId)
spansAtivos.delete(spanId)
duracaoMs = Number(process.hrtime.bigint() - span.inicio) / 1_000_000
logData = {
span_id: spanId,
span_name: span.name,
duration_ms: Math.round(duracaoMs * 100) / 100,
status: status,
error: erro ? { message: erro.message, name: erro.name } : undefined,
correlacao_id, event_id, colonia, // do contexto, quando presentes
}
se status == 'failure': logger.error(JSON.stringify(logData))
senão: logger.log(JSON.stringify(logData))
retornar { span_id: spanId, name: span.name, duration_ms: duracaoMs, status, metadata: span.metadata }

O metadata do span fica no resultado da finalização e não entra no log. O log do span carrega apenas os identificadores de instrumentação, a duração, o status e o erro.

classe HealthService:
construtor(prisma: PrismaService)
função verificarLiveness():
retornar { status: 'ok', timestamp: new Date().toISOString(), checks: {} }
função verificarReadiness():
checks.database = verificarBancoDados() // SELECT 1 via $queryRawUnsafe
checks.dead_letter_queue = verificarDLQ() // count de status 'pendente' com thresholds
checks.minio = verificarStorage() // /minio/health/live ou bucketExists S3
retornar {
status: resolverStatusGeral(checks), // 'down' > 'degraded' > 'ok'
timestamp: new Date().toISOString(),
checks: checks,
}
função verificarDLQ():
limiarAviso = OBS_DLQ_WARN_THRESHOLD ?? 50
limiarCritico = OBS_DLQ_CRITICAL_THRESHOLD ?? 200
contagem = prisma.dead_letter_queue.count({ where: { status: 'pendente' } })
se contagem > limiarCritico:
retornar { status: 'degraded', message: "DLQ com N eventos pendentes (crítico > N)" }
senão se contagem > limiarAviso:
retornar { status: 'ok', message: "DLQ com N eventos pendentes (aviso > N)" }
senão:
retornar { status: 'ok', message: "DLQ com N eventos pendentes" }

O check de storage tem dois caminhos:

se MINIO_USE_SSL == 'true':
se credenciais MINIO_ACCESS_KEY e MINIO_SECRET_KEY ausentes:
retornar { status: 'degraded', message: 'Credenciais de armazenamento ausentes' }
cliente = Minio.Client({ endPoint: MINIO_ENDPOINT, port: MINIO_PORT, useSSL: true, ... })
bucketExiste = bucketExists(MINIO_BUCKET ?? 'anexos'), com timeout de 3s
retornar 'ok' se o bucket existe; 'degraded' caso contrário
senão:
resposta = fetch("http://{MINIO_ENDPOINT ?? 'minio'}:{MINIO_PORT ?? '9000'}/minio/health/live", timeout 3s)
retornar 'ok' se a resposta é 2xx; 'degraded' caso contrário

O status geral fica down quando há check down, degraded quando há check degraded e ok quando todos passam. Banco inacessível derruba o status geral. Storage indisponível marca degraded: os anexos degradam com graça e a API segue de pé. A N-0c não acessa repositórios da N-0a; a contagem da DLQ é feita direto no banco via PrismaService.

classe MetricsService:
contadorEventosProcessados = new Map<string, number>()
histogramaDuracaoEventos = new Map<string, { valores, contagem, soma }>()
gaugesDLQ = new Map<string, number>()
contadorRequisicoesHttp = new Map<string, number>()
histogramaDuracaoHttp = new Map<string, { valores, contagem, soma }>()
BUCKETS_HISTOGRAMA = [0.01, 0.05, 0.1, 0.5, 1, 5, 10, 30]
função registrarEventoProcessado(tipo, colonia, status, durationMs):
incrementar contador "event_processing_total{tipo,colonia,status}"
se status == 'success':
observar "event_processing_duration_seconds{tipo,colonia}" em segundos
função definirPendentesDLQ(colonia, quantidade):
gaugesDLQ.set(colonia, quantidade)
função registrarRequisicaoHttp(method, path, statusCode, durationMs):
incrementar contador "http_requests_total{method,path,status_code}"
observar "http_request_duration_seconds{method,path}" em segundos
função obterMetricas():
para cada métrica, emitir linha HELP e linha TYPE
contadores: uma amostra por chave observada
histogramas: "_count", "_sum" e "_bucket{le=...}" acumulado por bucket, com "+Inf"
gauges: valor atual por colônia
retornar as linhas unidas por quebra de linha

A chave de cada Map é a série já formatada com os labels, e os valores dos histogramas ficam em segundos. O cálculo dos buckets é feito na renderização, a partir do array de observações ordenado.

O log estruturado usa o Logger do NestJS com serialização JSON manual. Não há Pino nem logger.config.ts. A função formatarMensagem() monta o registro:

função formatarMensagem(mensagem, ctx, dados):
registro = { mensagem: mensagem }
se ctx.correlacao_id: registro.correlacao_id = ctx.correlacao_id
se ctx.event_id: registro.event_id = ctx.event_id
se ctx.colonia: registro.colonia = ctx.colonia
se dados: Object.assign(registro, redigirDadosLog(dados))
retornar JSON.stringify(registro)

A função redigirDadosLog() percorre o objeto de dados e substitui por [REDIGIDO] os campos sensíveis. A lista inclui nome, email, endereco, texto, descricao, telefone, cpf, senha, token, payload, midia_urls e cargos, entre outros. A verificação é por substring do nome da chave, em minúsculas, e vale para objetos aninhados e arrays. O log do span não passa por redação porque carrega apenas identificadores de instrumentação, duração, status e erro.

4.5 Health check — HealthService.verificarReadiness()

Seção intitulada “4.5 Health check — HealthService.verificarReadiness()”
GET /api/observability/health:
retornar verificarLiveness()
// Resposta: { status: 'ok', timestamp: '...', checks: {} }
// HTTP 200 sempre que o processo responder
GET /api/observability/health/readiness:
retornar verificarReadiness()
// Resposta: {
// status: 'ok' | 'degraded' | 'down',
// timestamp: '...',
// checks: {
// database: { status: 'ok' },
// dead_letter_queue: { status: 'ok', message: 'DLQ com 3 eventos pendentes' },
// minio: { status: 'ok' },
// }
// }
// HTTP 200 se ok ou degraded, 503 se down
Caso Comportamento
AsyncLocalStorage sem contexto formatarMensagem() omite correlacao_id, event_id e colonia. O log é emitido normalmente, apenas sem metadados de trace.
Span encerrado duas vezes spansAtivos.delete(spanId) na primeira chamada. A segunda lança ErroSpanNaoEncontrado, porque o span já não existe no Map.
Span iniciado e nunca encerrado Vazamento de memória se o Map crescer indefinidamente. O limite de 10.000 spans remove o mais antigo ao exceder, com log de warning. TTL por span fica para a Fase 2.
Handler não emite log O trace continua: o wrapper da N-0a emite o log do span no término de cada processamento. A colônia não precisa logar para que a duração e o resultado sejam medidos.
Cardinalidade das métricas As chaves são combinações de tipo, colonia, status, method, path e status_code. Os tipos são limitados pelo catálogo do Registry e as rotas são estáveis, então o conjunto é finito e pequeno.
Métricas sem Prometheus externo O endpoint /metrics expõe os dados independentemente de haver servidor Prometheus coletando. Sem coleta externa, os dados acumulam em memória e são perdidos no restart do monolito. Em desenvolvimento, as métricas são consultáveis via curl localhost:3000/api/observability/metrics.
DLQ consultada com frequência O readiness faz count com filtro status: 'pendente'. O índice dead_letter_queue_status_idx cobre a consulta. O gauge da N-0a usa dead_letter_queue_status_colonia_idx.
Conflito de correlacao_id entre fluxos O correlacao_id é propagado pelas colônias. Se uma colônia publica um evento derivado, inclui o correlacao_id do evento de origem. O TraceContext do AsyncLocalStorage é sobrescrito a cada evento processado. O contexto é do processamento corrente, não do fluxo pai.

Zero tabelas próprias no MVP. O trace já existe no core.event_log da N-0a. Adicionar uma tabela de métricas de processamento duplicaria dados que um time-series database externo armazena com mais eficiência. Logs estruturados no stdout são a prática padrão em sistemas cloud-native (12-Factor App). O monolito emite, a infraestrutura coleta. Adicionar tabela de log no PostgreSQL do monolito criaria contenção de I/O entre carga de negócio e carga de observabilidade.

AsyncLocalStorage e não CLS (Continuation-Local Storage). AsyncLocalStorage é nativo do Node.js desde a v14, sem dependências externas. O nestjs-cls é um wrapper, mas para o MVP o uso direto no ObservabilityService é suficiente e tem menos superfície de bug. A API usada é run(), getStore() e enterWith().

Logger do NestJS com JSON manual, sem dependência de logger. A implementação monta o JSON em formatarMensagem(), sem pacote de logger externo. Os campos de trace vêm do AsyncLocalStorage e os dados passam por redigirDadosLog() antes da serialização. A migração para um logger estruturado dedicado fica para a Fase 2, quando o volume justificar.

Buckets específicos e não uniformes. A faixa de 10ms a 30s cobre desde processamentos rápidos de categorização até georreferenciamento com geocodificação externa, sem granularidade excessiva.

Métricas de processamento no MetricsService, não na N-0a. A N-0a chama metricsService.registrarEventoProcessado() no publish e no subscribe. A chamada está na N-0a, mas a métrica é definida e mantida pela N-0c. A N-0a não sabe como a métrica é armazenada ou renderizada. Ela apenas relata o tipo, a colônia, o resultado e a duração. A separação de responsabilidades é preservada.

Health check da DLQ via PrismaService, não via repositório da N-0a. O HealthService injeta o PrismaService, a instância única exportada pela N-0a, e conta as entradas pendentes por query direta. A N-0c não acessa o DeadLetterRepository. A dependência aceita é com a conexão de banco compartilhada do núcleo.

Redação de dados nos logs. redigirDadosLog() aplica a lista de chaves sensíveis a todo log emitido via ObservabilityService, com substituição por [REDIGIDO]. A redação alcança objetos aninhados e arrays.


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

Seção intitulada “5. Integração com o Barramento e Outras Colônias”

O EventBusModule importa o ObservabilityModule e o EventBusService injeta os dois serviços:

@Global()
@Module({
imports: [
EventEmitterModule.forRoot({ ... }),
ObservabilityModule, // N-0c — trace e métricas
RegistryModule,
],
controllers: [EventBusController],
providers: [
EventBusService,
PrismaService,
EventLogRepository,
DeadLetterRepository,
PapelOperadorGuard,
],
exports: [EventBusService, PrismaService],
})
export class EventBusModule {}
@Injectable()
export class EventBusService {
constructor(
private readonly eventLogRepo: EventLogRepository,
private readonly deadLetterRepo: DeadLetterRepository,
private readonly eventEmitter: EventEmitter2,
private readonly obs: ObservabilityService, // N-0c
private readonly metrics: MetricsService, // N-0c
private readonly registry: RegistryService,
) {}

A instrumentação do publicar() e do executarProtegido() segue os pseudocódigos das seções 3.1 e 3.2.

Colônias recebem ObservabilityService e MetricsService por injeção, porque o módulo é @Global(). A N-0d, por exemplo, injeta o ObservabilityService e registra o resultado de operações de titular:

@Injectable()
export class TitularService {
constructor(
private readonly repo: TitularRepository,
private readonly armazenamento: ArmazenamentoTitularService,
private readonly obs: ObservabilityService,
) {}
// ...
this.obs.info('Eliminação de dados do titular concluída', {
titular_id: titularId,
protocolo,
total_demandas: demandaIds.length,
});
}

Os logs emitidos com obs.info() dentro de uma requisição herdam o correlacao_id definido pelo interceptor.

A colônia não precisa:

  • Chamar definirTraceContext(). O interceptor HTTP ou o wrapper da N-0a já definiu o contexto.
  • Chamar iniciarSpan() e finalizarSpan(). O interceptor e o wrapper já instrumentam.
  • Chamar registrarEventoProcessado(). O wrapper da N-0a já registrou.

A colônia pode, opcionalmente, emitir logs de negócio com obs.info(), obs.warn(), obs.error() e obs.debug(). Eles herdam o contexto de trace.

5.3 Interceptor HTTP — propagação de correlacao_id

Seção intitulada “5.3 Interceptor HTTP — propagação de correlacao_id”

O interceptor global é registrado no AppModule via provider APP_INTERCEPTOR:

src/app.module.ts
@Module({
imports: [EventBusModule, /* ... demais colônias */],
providers: [
{
provide: APP_INTERCEPTOR,
useClass: TraceInterceptor, // da N-0c
},
],
})
export class AppModule {}

O TraceInterceptor:

intercept(ctx, next):
request = ctx.switchToHttp().getRequest()
response = ctx.switchToHttp().getResponse()
corrId = extrairOuGerarCorrelacaoId(request)
// x-correlation-id, se UUID válido; senão x-request-id, se UUID válido; senão uuidv4()
definirTraceContext({ correlacao_id: corrId })
response.setHeader('x-correlation-id', corrId)
spanId = iniciarSpan('http.request', { method, path: request.route?.path ?? request.url })
aguardar next.handle() e, ao final:
sucesso: finalizarSpan(spanId, 'success')
registrarRequisicaoHttp(method, path, response.statusCode, duração)
erro: finalizarSpan(spanId, 'failure', erro)
registrarRequisicaoHttp(method, path, erro.status ?? 500, duração)

O formato UUID é validado por expressão regular. Header com formato não-UUID é descartado e um UUID novo é gerado.

@ApiTags('Observabilidade')
@Controller('observability')
export class ObservabilityController {
constructor(
private readonly health: HealthService,
private readonly metricas: MetricsService,
) {}
@Get('health')
@SkipThrottle()
@ApiOperation({ summary: 'Health check', description: 'Verifica se a aplicação está rodando (liveness)' })
@ApiResponse({ status: 200, description: 'Aplicação saudável' })
async healthCheck(): Promise<HealthStatus> {
return this.health.verificarLiveness();
}
@Get('health/readiness')
@Throttle({ default: { limit: 30, ttl: 60000 } })
@UseGuards(IpInternoOuOperadorGuard)
@ApiExcludeEndpoint()
async readiness(@Res() res): Promise<void> {
const status = await this.health.verificarReadiness();
res.status(status.status === 'down' ? 503 : 200).json(status);
}
@Get('metrics')
@Throttle({ default: { limit: 120, ttl: 60000 } })
@UseGuards(IpInternoOuOperadorGuard)
@ApiExcludeEndpoint()
@Header('Content-Type', 'text/plain; charset=utf-8')
async metricaEndpoint(): Promise<string> {
return this.metricas.obterMetricas();
}
}
Método Rota Descrição Proteção
GET /api/observability/health Liveness. HTTP 200 sempre que o processo responder. Consumido pelo health check do container. Pública, sem throttle
GET /api/observability/health/readiness Readiness. Verifica banco, DLQ e storage. HTTP 200 ou 503. IpInternoOuOperadorGuard, 30/min
GET /api/observability/metrics Métricas Prometheus em formato texto (text/plain). Coletado por Prometheus externo. IpInternoOuOperadorGuard, 120/min

O liveness é público para o health check do container. O IpInternoOuOperadorGuard libera loopback, libera qualquer IP fora de produção e, em produção, exige IP na allowlist OPERADORES_ALLOWED_IPS. O Swagger documenta apenas o health; readiness e metrics ficam fora com @ApiExcludeEndpoint().

O trace por correlacao_id está no endpoint da N-0a: GET /api/events/trace/:correlacaoId. A N-0c não duplica esse endpoint. A responsabilidade do trace é compartilhada: a N-0a armazena e expõe; a N-0c garante que os logs e as métricas carreguem o correlacao_id para correlação.

O BFF da D-1a não chama a N-0c diretamente. A comunicação passa pelo interceptor HTTP global. Toda requisição que passa pelo BFF ganha x-correlation-id na resposta e contexto de trace no AsyncLocalStorage para os serviços que emitem logs pelo ObservabilityService. O Logger padrão do NestJS não lê o AsyncLocalStorage.

A N-0c não consome projeções de leitura de outras colônias. As únicas dependências de dados são:

  • PrismaService (instância única do núcleo). Health check de banco (SELECT 1) e contagem de DLQ pendente.
  • Storage. Liveness do MinIO via /minio/health/live; com MINIO_USE_SSL=true, verificação de existência do bucket via S3. O readiness marca degraded quando o storage está indisponível.
  • O endpoint de trace é da N-0a (GET /api/events/trace/:correlacaoId). A N-0c não chama repositórios da N-0a.

A N-0c aplica throttle em readiness (30/min) e metrics (120/min). O health de liveness usa @SkipThrottle(), porque é chamado pelo health check do container. Em produção, readiness e metrics exigem IP de operação via IpInternoOuOperadorGuard. O liveness é público. Loopback permanece liberado para o health check interno. Rate limiting dedicado por IP nos endpoints de observabilidade fica para a Fase 2.

Limite Valor MVP Justificativa
Spans ativos em memória 10.000 O Map spansAtivos tem tamanho máximo para prevenir vazamento. Spans nunca encerrados são removidos por ordem de criação quando o limite é excedido.
Resposta do /metrics ~5 KB (típico) O formato texto é compacto e cresce com as combinações de labels observadas. Com tipos finitos do Registry e rotas estáveis, o limite prático fica abaixo de 50 KB.
Log por requisição 1 linha JSON Cada requisição HTTP gera 1 log de span ao final. Cada evento processado gera 1 log de span. Cada publish gera 1 log de span. Logs de negócio emitidos pela colônia somam ao volume.
Chaves de métricas ~500 no MVP As chaves são combinações de tipo, colônia, status, método, rota e status code. Com os tipos do Registry e as rotas estáveis, o total fica bem abaixo do que a memória suporta.
Operação Overhead estimado Mitigação
AsyncLocalStorage.run() por handler < 0,01ms Operação nativa do Node.js. Um enter/exit de contexto por evento processado.
process.hrtime.bigint() para span timing < 0,001ms Chamada nativa de alta precisão. Duas por handler (início e fim).
formatarMensagem() por log < 0,01ms Monta um objeto pequeno com até três campos de contexto e serializa em JSON.
Incremento de contador no Map < 0,01ms Leitura e escrita em memória.
Observação de histograma no Map < 0,05ms Acumulação em array e soma. O cálculo dos buckets ocorre na renderização.
Health check com count na DLQ < 1ms Os índices dead_letter_queue_status_idx e dead_letter_queue_status_colonia_idx cobrem a consulta. Executado a cada 10–30s.

Total por evento processado: < 0,1ms de overhead de observabilidade. Para um pipeline de 6 colônias (captura, normalização, georreferenciamento, categorização, priorização e agenda), o overhead total fica em torno de 0,6ms, pequeno em relação à latência de I/O de banco e de processamento de negócio.

Sem cache na N-0c. As métricas vivem em Map na memória por definição. Os logs são emitidos imediatamente para stdout pelo Logger do NestJS. Health checks fazem queries leves ao banco que não justificam cache. O correlacao_id no AsyncLocalStorage é uma referência em memória, com acesso O(1).


Teste unitário do ObservabilityService:

beforeEach(async () => {
const module = await Test.createTestingModule({
providers: [ObservabilityService],
}).compile();
service = module.get(ObservabilityService);
});

O teste cobre contexto de trace, spans e níveis de log. O ObservabilityService não depende de provider externo.

Teste unitário do HealthService:

beforeEach(async () => {
const module = await Test.createTestingModule({
providers: [HealthService, { provide: PrismaService, useValue: mockPrisma }],
}).compile();
service = module.get(HealthService);
});

O mock do PrismaService responde $queryRawUnsafe e dead_letter_queue.count. O teste cobre liveness, readiness com banco, DLQ e storage, e o caminho com MINIO_USE_SSL=true.

Teste unitário do MetricsService:

beforeEach(async () => {
const module = await Test.createTestingModule({
providers: [MetricsService],
}).compile();
service = module.get(MetricsService);
service.resetar();
});

O interceptor e o controller têm specs próprios. O do interceptor usa mocks de ObservabilityService e MetricsService. O do controller verifica os metadados de throttle e os códigos de resposta.

Happy path:

# Cenário Verificação
T1 definirTraceContext() com correlacao_id e event_id obterTraceContext() retorna os campos. Logs emitidos via info() incluem os campos.
T2 iniciarSpan() + finalizarSpan() com sucesso Span retornado tem duration_ms > 0 e status: 'success'. Log emitido com span_name, duration_ms e status.
T3 executarComTraceContext() executa fn com contexto isolado Contexto dentro da fn é o passado. Contexto fora da fn não é afetado.
T4 registrarEventoProcessado() incrementa métricas Counter event_processing_total incrementado com as labels corretas. Histograma event_processing_duration_seconds observado no sucesso.
T5 verificarLiveness() Retorna { status: 'ok' }. HTTP 200.
T6 verificarReadiness() com banco ok, DLQ vazia e storage ok Retorna { status: 'ok', checks: { database: 'ok', dead_letter_queue: 'ok', minio: 'ok' } }.
T7 obterMetricas() Retorna string em formato Prometheus com todas as métricas observadas.

Falhas e bordas:

# Cenário Verificação
T8 finalizarSpan() com span inexistente Lança ErroSpanNaoEncontrado.
T9 finalizarSpan() duas vezes no mesmo spanId A primeira retorna o resultado. A segunda lança ErroSpanNaoEncontrado.
T10 spansAtivos excede 10.000 O span mais antigo é removido, com log de warning.
T11 verificarReadiness() com banco inacessível O check database retorna { status: 'down', message: '...' }. Status geral: down. HTTP 503.
T12 verificarReadiness() com DLQ acima do threshold crítico O check dead_letter_queue retorna { status: 'degraded', message: '...' }. Status geral: degraded.
T13 verificarReadiness() com erro na consulta da DLQ O check retorna { status: 'degraded', message: 'Erro ao consultar DLQ: ...' }. O erro não se propaga.
T14 obterTraceContext() sem contexto prévio Retorna {}. Não lança erro.
T15 registrarEventoProcessado() com status failure O counter recebe a label status: 'failure'. O histograma não é observado.
T16 definirPendentesDLQ() com quantidade zero O gauge fica com valor 0 e permanece visível no /metrics.
T17 Storage indisponível O check minio retorna { status: 'degraded', message: '...' }. Status geral: degraded.

Teste de integração (com módulo NestJS real):

# Cenário Verificação
T18 Handler executado com contexto de trace obterTraceContext() dentro do handler retorna { correlacao_id, event_id, colonia }.
T19 Interceptor HTTP propaga x-correlation-id Requisição com header válido inclui o mesmo header na resposta. Sem header, a resposta inclui UUID gerado.
T20 Span do wrapper da N-0a captura duração real finalizarSpan() retorna duration_ms compatível com o tempo de execução do handler.
T21 Falha no handler registra métrica de failure Counter event_processing_total{status='failure'} incrementado. O histograma de duração não é observado.

A N-0c não tem tabelas próprias, e não há seed de banco. Para desenvolvimento local, o módulo sobe com a configuração padrão e expõe health e métricas sem dependência de dados prévios.

Para simular cenários de teste:

Janela do terminal
# Health check liveness
curl http://localhost:3000/api/observability/health
# Health check readiness (requer banco e storage)
curl http://localhost:3000/api/observability/health/readiness
# Métricas Prometheus
curl http://localhost:3000/api/observability/metrics
# Simular eventos processados para popular métricas:
# - Publicar demandas via POST /api/demandas
# - Observar métricas incrementarem em /api/observability/metrics

Funcionalidade Status
Logs estruturados JSON com correlacao_id, event_id, colonia e redação de dados sensíveis MVP obrigatório
AsyncLocalStorage para propagação de contexto de trace MVP obrigatório
Spans de processamento (iniciarSpan/finalizarSpan) no interceptor HTTP e no wrapper da N-0a MVP obrigatório
Métricas Prometheus em memória: event_processing_total, event_processing_duration_seconds, dead_letter_queue_pending, http_requests_total, http_request_duration_seconds MVP obrigatório
Health check: liveness (/health) e readiness (/health/readiness) com banco, DLQ e storage MVP obrigatório
Interceptor HTTP global com x-correlation-id MVP obrigatório
@Global(): módulo disponível para todas as colônias sem import explícito MVP obrigatório
Proteção por IP e throttle em readiness e metrics MVP obrigatório
Simplificação Justificativa Quando remover
Zero tabelas no banco Trace no event_log da N-0a. Métricas em memória. Logs no stdout. Adicionar core.alert_history na Fase 2, quando alertas precisarem de persistência e auditoria.
Sem alertas ativos (apenas health check passivo) Alertas são configurados no Prometheus AlertManager externo, não no monolito. O health check da N-0c expõe o status e o AlertManager decide quando notificar. Se o monolito operar sem infraestrutura externa de monitoramento, adicionar notificação interna na Fase 2.
Configuração via env vars, sem endpoint de configuração Os thresholds do health check não mudam em runtime no MVP. Aplicar requer restart. Adicionar endpoint PATCH /api/observability/config na Fase 2, para ajuste dinâmico sem restart.
Sem OpenTelemetry OpenTelemetry adiciona complexidade de SDK, exporters e collectors. Para o MVP monolítico, o trace por correlacao_id no log e no event_log é suficiente. Migrar para OpenTelemetry quando houver múltiplos serviços (Fase 2).
Sem dashboard (Grafana fica externo) A visualização de métricas é responsabilidade da infraestrutura de monitoramento. Grafana conecta no Prometheus e exibe dashboards. Se houver necessidade de dashboard embutido no monolito, adicionar na Fase 2 como projeção de leitura da D-7.
Sem tracing distribuído (Jaeger/Zipkin) No monolito, o trace por correlacao_id no event_log é suficiente. Adicionar OpenTelemetry tracing com export para Jaeger na Fase 2.
spansAtivos em memória com limite fixo Para o MVP, 10.000 spans é suficiente. Spans órfãos são removidos quando o Map enche. Substituir por estrutura com TTL na Fase 2.
  • OpenTelemetry SDK com export para Jaeger/Zipkin (trace distribuído entre microsserviços)
  • Tabela core.alert_history para auditoria de alertas disparados
  • Tabela core.metrics_snapshot para retenção histórica de métricas agregadas
  • Endpoint de configuração dinâmica (PATCH /api/observability/config)
  • TTL por span em vez de limite fixo de tamanho do Map
  • Migração para logger estruturado com nível configurável em runtime
  • Métricas de processo do Node.js (nodejs_memory_heap_bytes, nodejs_eventloop_lag_seconds) e health check de latência p99 do event loop
  • Métricas customizáveis por colônia (ex: D-4 expõe ranking_position_distribution)
  • Integração com @nestjs/bullmq para métricas de fila (quando a N-0a migrar para fila assíncrona)

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

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

Resolução: health check da DLQ via PrismaService. A N-0c consome a conexão de banco compartilhada do núcleo e conta as entradas pendentes com dead_letter_queue.count. O gauge dead_letter_queue_pending é atualizado pela N-0a após registrar, resolver ou reabrir entradas da DLQ, com contagem de pendente e retrying. O readiness conta apenas pendente. O índice dead_letter_queue_status_idx cobre a contagem do readiness, e o índice dead_letter_queue_status_colonia_idx cobre a consulta do gauge.

Conflito potencial: x-correlation-id header já usado por outras camadas. O interceptor lê x-correlation-id e, como fallback, x-request-id. Se um proxy reverso (nginx, Envoy) injetar x-request-id, o interceptor o utiliza como correlacao_id. O formato UUID v4 é validado. Header com formato não-UUID é descartado e um UUID novo é gerado. Sem conflito.

Conflito potencial: logs fora do ObservabilityService. O Logger padrão do NestJS e bibliotecas de terceiros emitem logs em texto, sem os campos de trace. A instrumentação cobre o que passa pelo TraceInterceptor, pelo wrapper de consumo da N-0a e pelos logs emitidos via ObservabilityService. A correlação do ciclo de negócio depende desses pontos. Sem conflito com o formato JSON.

Conflito potencial: @Global() do ObservabilityModule + @Global() do EventBusModule. Múltiplos módulos @Global() coexistem sem conflito. O NestJS resolve a injeção de dependência normalmente. O EventBusModule importa ObservabilityModule explicitamente, o que é necessário para que o EventBusService injete os serviços de observabilidade. As demais colônias não precisam importar nenhum dos dois. Sem conflito.

Verificação da distribuição por decaimento e outras regras de negócio. A N-0c não implementa regras de negócio. A distribuição por decaimento, o sorteio, a priorização e as demais regras são responsabilidade das colônias D-4, D-5 e D-6a. Sem conflito.



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