N-0c — Observabilidade
Parte do Núcleo — O Chão do Formigueiro
Propósito
Seção intitulada “Propósito”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.
1. Estrutura do Módulo NestJS
Seção intitulada “1. Estrutura do Módulo NestJS”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.
Árvore de diretórios
Seção intitulada “Árvore de diretórios”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 SpanResultadoModule definition
Seção intitulada “Module definition”@Global()@Module({ providers: [ObservabilityService, MetricsService, HealthService, TraceInterceptor], exports: [ObservabilityService, MetricsService], controllers: [ObservabilityController],})export class ObservabilityModule {}Pontos de atenção
Seção intitulada “Pontos de atenção”@Global(): toda colônia injetaObservabilityServiceeMetricsServicesemimports: [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 comoAPP_INTERCEPTORnoAppModule. 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) eMetricsService(métricas) é deliberada. Colônias que só precisam de métricas injetam apenasMetricsService. Colônias que precisam de contexto e log injetamObservabilityService.
Interface pública — ObservabilityService
Seção intitulada “Interface pública — ObservabilityService”definirTraceContext(ctx: Partial<TraceContext>): void; // mescla o contexto no store atual do ALSobterTraceContext(): TraceContext; // {} quando não há contextoexecutarComTraceContext<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 spanIdfinalizarSpan(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>;}Interface pública — MetricsService
Seção intitulada “Interface pública — MetricsService”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 PrometheusO método resetar() limpa todos os Map e é usado pelos testes para isolar cenários.
Colônia pura de infraestrutura
Seção intitulada “Colônia pura de infraestrutura”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.
2. Banco de Dados — Schema e Entidades
Seção intitulada “2. Banco de Dados — Schema e Entidades”A N-0c não possui tabelas próprias no MVP. A decisão se apoia em três fatores:
-
O trace completo por
correlacao_idjá existe nocore.event_logda N-0a. Cada evento publicado registracorrelacao_id,origem,tipoetimestamp. ConsultarSELECT * FROM core.event_log WHERE correlacao_id = ? ORDER BY sequence_numberreconstrói o ciclo completo de uma demanda. A N-0c não duplica essa responsabilidade. -
As métricas de processamento residem em memória no
MetricsService, emMappor 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. -
Os logs estruturados são emitidos em JSON de uma linha para stdout pelo
Loggerdo 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.
Evolução futura
Seção intitulada “Evolução futura”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.
3. Eventos — Contratos Detalhados
Seção intitulada “3. Eventos — Contratos Detalhados”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.
3.1 Instrumentação do publish (N-0a → N-0c)
Seção intitulada “3.1 Instrumentação do publish (N-0a → N-0c)”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 timestamp2. Validar tipo, versão e payload contra o Registry3. executarComTraceContext({ correlacao_id, event_id, colonia: origem }, ...)4. Persistir no event_log com o payload redigido5. iniciarSpan('event_bus.publish', { tipo, origem, event_id })6. Emitir para os consumidores com o payload completo7. finalizarSpan(spanId, 'success')8. registrarEventoProcessado(tipo, origem, 'success', duration_ms)9. Retornar o registro entregueO log estruturado emitido nesses passos carrega automaticamente correlacao_id, event_id e colonia porque formatarMensagem() lê o AsyncLocalStorage. Nenhum console.log manual é necessário.
3.2 Instrumentação do subscribe (N-0a → N-0c)
Seção intitulada “3.2 Instrumentação do subscribe (N-0a → N-0c)”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 erroO 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.
3.3 Instrumentação HTTP (interceptor global)
Seção intitulada “3.3 Instrumentação HTTP (interceptor global)”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.
3.4 Métricas expostas — catálogo
Seção intitulada “3.4 Métricas expostas — catálogo”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.
3.5 Ordem de operações — log antes ou depois?
Seção intitulada “3.5 Ordem de operações — log antes ou depois?”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.
4. Lógica de Negócio — Algoritmos e Fluxos
Seção intitulada “4. Lógica de Negócio — Algoritmos e Fluxos”4.1 ObservabilityService — pseudocódigo
Seção intitulada “4.1 ObservabilityService — pseudocódigo”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.
4.2 HealthService — pseudocódigo
Seção intitulada “4.2 HealthService — pseudocódigo”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árioO 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.
4.3 MetricsService — pseudocódigo
Seção intitulada “4.3 MetricsService — pseudocódigo”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 linhaA 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.
4.4 Configuração de logging
Seção intitulada “4.4 Configuração de logging”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 down4.6 Casos de borda
Seção intitulada “4.6 Casos de borda”| 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. |
4.7 Decisões de design com justificativa
Seção intitulada “4.7 Decisões de design com justificativa”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”5.1 Dependência da N-0a (Event Bus)
Seção intitulada “5.1 Dependência da N-0a (Event Bus)”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.
5.2 Como uma colônia usa a N-0c
Seção intitulada “5.2 Como uma colônia usa a N-0c”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()efinalizarSpan(). 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:
@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.
5.4 REST Controller — endpoints expostos
Seção intitulada “5.4 REST Controller — endpoints expostos”@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.
5.5 Chamadas síncronas via BFF
Seção intitulada “5.5 Chamadas síncronas via BFF”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.
5.6 Dependências de projeções de leitura
Seção intitulada “5.6 Dependências de projeções de leitura”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; comMINIO_USE_SSL=true, verificação de existência do bucket via S3. O readiness marcadegradedquando 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.
6. Performance e Limites
Seção intitulada “6. Performance e Limites”6.1 Rate limiting
Seção intitulada “6.1 Rate limiting”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.
6.2 Limites de tamanho e memória
Seção intitulada “6.2 Limites de tamanho e memória”| 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. |
6.3 Impacto de performance da instrumentação
Seção intitulada “6.3 Impacto de performance da instrumentação”| 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.
6.4 Estratégia de cache
Seção intitulada “6.4 Estratégia de cache”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).
7. Testabilidade
Seção intitulada “7. Testabilidade”7.1 Como testar o módulo isolado
Seção intitulada “7.1 Como testar o módulo isolado”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.
7.2 Cenários de teste críticos
Seção intitulada “7.2 Cenários de teste críticos”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. |
7.3 Dados de seed para desenvolvimento local
Seção intitulada “7.3 Dados de seed para desenvolvimento local”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:
# Health check livenesscurl http://localhost:3000/api/observability/health
# Health check readiness (requer banco e storage)curl http://localhost:3000/api/observability/health/readiness
# Métricas Prometheuscurl 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/metrics8. Alinhamento com o MVP
Seção intitulada “8. Alinhamento com o MVP”8.1 O que é MVP obrigatório
Seção intitulada “8.1 O que é MVP obrigatório”| 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 |
8.2 Simplificações válidas no MVP
Seção intitulada “8.2 Simplificações válidas no MVP”| 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. |
8.3 O que vai para a Fase 2
Seção intitulada “8.3 O que vai para a Fase 2”- OpenTelemetry SDK com export para Jaeger/Zipkin (trace distribuído entre microsserviços)
- Tabela
core.alert_historypara auditoria de alertas disparados - Tabela
core.metrics_snapshotpara 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/bullmqpara 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.
Referências
Seção intitulada “Referências”- Especificação base: Apêndice B - Colônias.md, seção “N-0c — Observabilidade”
- Colônia irmã (núcleo): N-0a - Event Bus.md
- Colônia irmã (núcleo): N-0b - Registry.md
- Stack de referência e arquitetura do MVP: Apêndice B - Colônias.md, seção “Arquitetura do MVP — Monolito Modular”
- Mapa de dependências de eventos: Apêndice B - Colônias.md, seção “Mapa de Dependências de Eventos entre Colônias”
- Princípios do Formigueiro: Apêndice B - Colônias.md, seção “Princípios herdados do Formigueiro”
- Contexto de gestão: contexto_IA.md, seção 7 (Gestão)
- Contexto de infraestrutura: contexto_IA.md, seção 10 (Infraestrutura cívica digital)
- Contexto do Formigueiro: contexto_IA.md, seção 22 (Arquitetura técnica — O Formigueiro)
Documento de especificação técnica de implementação. Aprovado e integrado.