Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -161,3 +161,6 @@ LEVER_INCLUDE_ALL_JOBS=true
SCRAPER_CATALOG_LIFETIME=216h
# Cache de páginas de busca, limitado também pela próxima expiração de vaga.
JOB_SEARCH_CACHE_TTL_SECONDS=120

# Versão da aplicação para status operacional (não é label Prometheus).
APPLICATION_VERSION=local
2 changes: 2 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,8 @@ npm-debug.log
yarn-error.log
yarn-debug.log
PAV-124_REPORT.md
PAV-125_REPORT.md
PAV-126_REPORT.md

# Test coverage reports
coverage/
Expand Down
43 changes: 43 additions & 0 deletions BACKEND.md
Original file line number Diff line number Diff line change
Expand Up @@ -634,3 +634,46 @@ podem afetar uma resposta, embora uma geração modificada impeça publicar cach
obsoleto. O índice usa expiração em segundos; existe granularidade inferior a um
segundo em relação aos timestamps SQL. A meta operacional de p95 < 500 ms exige
medição com volume e concorrência representativos; testes locais não a comprovam.


## Snapshot administrativo e métricas de busca — PAV-126

`GET /admin/observability` (também sob o prefixo API existente) exige sessão,
role administrativa e permissão `observability.metrics`. Usa o padrão de resposta
administrativa direto, `Cache-Control: no-store`, com contrato:

```json
{
"status": "partial",
"timestamp": "2026-10-07T00:00:00.000Z",
"processor": null,
"availability": { "scraper": "down" }
}
```

Quando o Processor responde, `processor` contém execution, lock, concurrency,
progress, queues, resources, errors, rejectedTitles/rejectedTitlesSince,
dependencies e index. Sem histórico, execution.status é idle, os timestamps/source/
stage são null e durationSeconds é zero. Falhas de PostgreSQL/Valkey preservam o
snapshot e retornam status partial. O backend faz uma única chamada interna com
prazo de 2.5s, valida com Zod e remove campos extras antes de responder. Ausência
ou resposta inválida do Processor resulta no formato parcial acima, HTTP 200.
Prometheus não é dependência dessa rota. As rotas administrativas anteriores,
auth, rate limiting e contratos de famílias permanecem preservados.

Métricas HTTP existentes `http_request_duration_seconds`/`http_requests_total`
continuam com os mesmos nomes/labels. Route usa templates Express; caminhos sem
match caem em `__unmatched__`, salvo rotas prioritárias estáticas conhecidas.
Query strings e IDs não viram labels. Methods desconhecidos caem em OTHER.
Os buckets existentes permitem estimar p50/p95/p99, sem promessa de desempenho.

O cache PAV-125 expõe `candidate_jobs_search_cache_requests_total` com hit/miss/
stale/error, histogram de get/set/invalidate e invalidações por motivo fixo.
Cada request tem um resultado: reutilização local via singleflight conta como hit;
cache indisponível conta error; publicação recusada por geração diferente conta
stale. As keys, fingerprints, TTL, CAS e ranking não mudaram. Go observa invalidações
de catálogo/rebuild; Node observa a invalidação manual já existente. Nada registra
keys, texto livre, localização privada ou dados de usuário em labels.

Ver `observability/PAV-126_REPORT.md` para inventário, regras, testes, limitações e
rollback. A meta de busca p95 < 500ms requer validação em staging representativo.
54 changes: 54 additions & 0 deletions SCRAPER.md
Original file line number Diff line number Diff line change
Expand Up @@ -582,3 +582,57 @@ forward-only: não existe down migration destrutiva automática. Não remover a
para reverter código; preservar os dados e planejar eventual arquivamento separado.
Nenhuma operação usa FLUSH ou o comando KEYS; cleanup só alcança o manifesto do
namespace selecionado e recusa a versão ativa e chaves externas.


## Observabilidade operacional — PAV-126

O `/metrics` mantém as métricas `scraper_*`, `go_*` e `process_*` existentes.
As métricas `candidate_scraper_*` acrescentam resultados por origem, lock,
concorrência, progresso, discovery mode, classificação, persistência e projeção.
`other` é exclusivamente diagnóstico; a taxonomia pública continua com 13 famílias.
Providers são os oito IDs de `ports` ou `unknown`. IDs de execução/vaga, títulos,
keywords, empresas, tokens e mensagens livres nunca são labels.

O estado operacional é atualizado nos pontos reais do pipeline e protegido por
mutex. As quatro filas refletem canais e buffers reais. Ao terminar, active,
waiting e depths voltam a zero. `keywordsTotal/Processed` contam entradas consumidas
pelas tarefas, incluindo fan-out; não contam keywords únicas. Totais de batches não
são inventados quando desconhecidos. ProvidersCompleted/AdaptersProcessed indicam
conclusão de tarefas, inclusive com erro; a situação aparece nas métricas de resultado.
O comportamento anterior de sucesso parcial de coleta permanece preservado.

A ordem continua classificação → commit PostgreSQL → projeção Valkey. Contadores
de persistência/indexação medem tentativas, inclusive retries; não são contagens de
vagas únicas no catálogo. Falhas e rollbacks possuem categorias fixas. Invalidações
são observadas somente após publicação bem-sucedida da geração. A telemetria não
transforma uma publicação válida em falha se seu hash estiver corrompido.

Manutenções CLI explícitas gravam contadores/histogramas de operação e um resumo
com TTL de 24h. O servidor amostra o hash fixo a cada 15s, com prazo de 2s, e conserva
os últimos contadores conhecidos se Valkey cair. `maintenance_telemetry_available`
e o timestamp da amostra permitem identificar defasagem. O scrape Prometheus não
faz I/O externo. A versão ativa está no snapshot administrativo, sem label dinâmica.
Não há rebuild nem reconciliação automática nova.

O GET interno `/admin/observability` é técnico e deve permanecer na rede privada
existente do Processor. Autenticação/autorização continuam no backend. O snapshot
inclui estado, última execução, próxima execução, lock/TTL, concorrência, progresso,
filas, recursos, erros, versões e última manutenção. PostgreSQL/Valkey são sondados
em paralelo com prazo compartilhado de 2s e um pool Redis de diagnóstico separado;
não há consulta ao catálogo nem ao Prometheus. Campos disponíveis permanecem na
resposta quando uma dependência falha. Configure `APPLICATION_VERSION` no deploy;
o default é `unknown` (`local` no exemplo de ambiente).

Os títulos rejeitados são um agregado administrativo em memória, normalizado,
limitado a 100 pares título/motivo e top 10 retornados. Reinicia por execução ou
após 24h; pode omitir novos títulos ao atingir o limite. Não é histórico completo.
Títulos com email/URL são descartados; não há vaga/descrição completa nem labels
por título. CPU percentual é derivada entre leituras do snapshot, inicialmente
null; RSS depende de `/proc`. CPU/container limits/GC continuam nos collectors
Go/process/cAdvisor existentes, incluindo GOMAXPROCS e GOMEMLIMIT.

Recording rules e alertas ficam em `observability/prometheus/rules`. Alertas de
CPU/memória usam janelas sustentadas e os limites atuais de 1.5 CPU/2GiB; revisar
limiares se recursos mudarem. Prometheus envia ao Alertmanager existente, cujo
receiver de exemplo não entrega notificações: configurar o destino operacional.
A validação e o inventário completo estão em `observability/PAV-126_REPORT.md`.
12 changes: 11 additions & 1 deletion backend/src/lib/cache.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,7 @@
import {
searchCacheDuration,
searchCacheInvalidations,
} from "../metrics/metrics";
import { normalizeJobTaxonomy } from "../modules/jobs/types/professionalTaxonomy";
import { randomUUID } from "node:crypto";
import { createClient, type RedisClientType } from "redis";
Expand Down Expand Up @@ -488,7 +492,13 @@ export async function cacheClearJobs(): Promise<{
if (await client.get("scraper:jobs:index-version")) {
// Catalog projections are owned by the Processor. Clearing HTTP search
// cache must not delete an active validated namespace or durable jobs.
await client.incr("jobs:search:generation");
const end = searchCacheDuration.startTimer({ operation: "invalidate" });
try {
await client.incr("jobs:search:generation");
searchCacheInvalidations.inc({ reason: "manual" });
} finally {
end();
}
return { deleted: 0, patterns: ["jobs:search:generation"] };
}
const patterns = ["scraper:job:*", "scraper:jobs:*"];
Expand Down
20 changes: 20 additions & 0 deletions backend/src/metrics/metrics.ts
Original file line number Diff line number Diff line change
Expand Up @@ -37,3 +37,23 @@ export const cacheOperationsTotal = new client.Counter({
labelNames: ["operation", "result"],
registers: [register],
});

export const searchCacheRequests = new client.Counter({
name: "candidate_jobs_search_cache_requests_total",
help: "Search cache outcomes",
labelNames: ["result"],
registers: [register],
});
export const searchCacheDuration = new client.Histogram({
name: "candidate_jobs_search_cache_operation_duration_seconds",
help: "Search cache operation duration",
labelNames: ["operation"],
buckets: [0.001, 0.005, 0.01, 0.05, 0.1, 0.5, 1, 5],
registers: [register],
});
export const searchCacheInvalidations = new client.Counter({
name: "candidate_jobs_search_cache_invalidations_total",
help: "Successful search cache generation invalidations",
labelNames: ["reason"],
registers: [register],
});
30 changes: 27 additions & 3 deletions backend/src/middleware/metrics.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,17 @@
import type { NextFunction, Request, Response } from "express";
import { httpRequestDuration, httpRequestsTotal } from "../metrics/metrics";

const priorityRoutes = new Set(
[
"/jobs/search",
"/health",
"/auth/login",
"/admin/scrapers",
"/admin/observability",
"/metrics",
].flatMap((route) => [route, `/api/v1${route}`]),
);

export function metricsMiddleware(
req: Request,
res: Response,
Expand All @@ -9,13 +20,26 @@ export function metricsMiddleware(
const end = httpRequestDuration.startTimer();

res.on("finish", () => {
// route já resolvido pelo Express (com :params), com fallback pro path cru
// Express templates only. Unmatched/auth-short-circuited requests must not
// expose arbitrary paths, query strings or dynamic IDs.
const route = req.route?.path
? `${req.baseUrl}${req.route.path}`
: req.path;
: priorityRoutes.has(req.path)
? req.path
: "__unmatched__";

const labels = {
method: req.method,
method: [
"GET",
"POST",
"PUT",
"PATCH",
"DELETE",
"HEAD",
"OPTIONS",
].includes(req.method)
? req.method
: "OTHER",
route,
status_code: String(res.statusCode),
};
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,13 @@ export class ObservabilityController {
private readonly auditService: AuditService,
) {}

async getOperationalSnapshot(req: Request, res: Response) {
const result = await this.service.getOperationalSnapshot();
this.auditService.fromRequest(req, "observability.metrics");
res.set("Cache-Control", "no-store");
return res.json(result);
}

async getHealth(req: Request, res: Response) {
try {
const result = await this.service.getHealth();
Expand Down
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
import { config } from "../../../config";
import { ProcessorSnapshotSchema } from "./observability.types";
import { HealthService } from "./health.service";
import { MetricsService } from "./metrics.service";
import type {
Expand All @@ -13,6 +15,31 @@ export class ObservabilityService {
private readonly metricsService: MetricsService,
) {}

async getOperationalSnapshot() {
const timestamp = new Date().toISOString();
try {
const response = await fetch(`${config.scraperUrl}/admin/observability`, {
signal: AbortSignal.timeout(2500),
});
if (!response.ok) throw new Error("processor unavailable");
const processor = ProcessorSnapshotSchema.parse(await response.json());
return {
status: processor.status,
timestamp,
processor,
availability: { scraper: "ok" as const },
};
} catch {
// Dependency failure is distinct from an available Processor with no runs.
return {
status: "partial" as const,
timestamp,
processor: null,
availability: { scraper: "down" as const },
};
}
}

async getHealth(): Promise<HealthcheckResult> {
return this.healthService.getHealthcheck();
}
Expand Down
112 changes: 112 additions & 0 deletions backend/src/modules/admin/observability/observability.types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -84,3 +84,115 @@ export const ObservabilityOverviewSchema = z.object({
metrics: MetricSnapshotSchema,
});
export type ObservabilityOverview = z.infer<typeof ObservabilityOverviewSchema>;

// Processor-owned runtime state. Unknown fields are stripped before forwarding,
// so an internal response cannot expose credentials or incidental payloads.
const nullableTimestamp = z.string().datetime().nullable();
export const ProcessorSnapshotSchema = z.object({
status: z.enum(["ok", "partial"]),
timestamp: z.string().datetime(),
execution: z.object({
status: z.enum(["idle", "running", "canceling", "failed", "completed"]),
source: z.enum(["manual", "cron"]).nullable(),
startedAt: nullableTimestamp,
durationSeconds: z.number().nonnegative(),
finishedAt: nullableTimestamp,
lastDurationSeconds: z.number().nonnegative(),
lastStatus: z
.enum(["success", "failed", "canceled", "skipped", "timeout"])
.nullable(),
nextRunAt: nullableTimestamp,
stage: z
.enum(["collection", "classification", "persistence", "indexing"])
.nullable(),
applicationVersion: z.string().max(128),
taxonomyVersion: z.string().max(128),
}),
lock: z.object({ held: z.boolean(), ttlSeconds: z.number().nonnegative() }),
concurrency: z.object({
configured: z.number().int().nonnegative(),
effective: z.number().int().nonnegative(),
active: z.number().int().nonnegative(),
waiting: z.number().int().nonnegative(),
}),
progress: z.object(
Object.fromEntries(
[
"providersTotal",
"providersCompleted",
"adaptersTotal",
"adaptersProcessed",
"tasksTotal",
"tasksCompleted",
"tasksCanceled",
"keywordsTotal",
"keywordsProcessed",
"batchesCompleted",
].map((k) => [k, z.number().int().nonnegative()]),
),
),
queues: z.object(
Object.fromEntries(
["collection", "classification", "persistence", "indexing"].map((k) => [
k,
z.object({
depth: z.number().int().nonnegative(),
capacity: z.number().int().nonnegative(),
}),
]),
),
),
resources: z.object({
cpuSeconds: z.number().nonnegative(),
cpuPercent: z.number().nonnegative().nullable(),
memoryBytes: z.number().nonnegative().nullable(),
heapBytes: z.number().nonnegative(),
goroutines: z.number().int().nonnegative(),
gomaxprocs: z.number().int().positive(),
gomemlimitBytes: z.number().nonnegative(),
}),
errors: z.object({
total: z.number().int().nonnegative(),
timeouts: z.number().int().nonnegative(),
}),
rejectedTitles: z
.array(
z.object({
title: z.string().max(100),
count: z.number().int().nonnegative(),
reasonCode: z.enum([
"negative_title",
"no_family_recognized",
"insufficient_title_evidence",
]),
}),
)
.max(10),
rejectedTitlesSince: z.string().datetime(),
dependencies: z.object({
postgres: z.object({ status: z.enum(["ok", "down"]) }),
valkey: z.object({ status: z.enum(["ok", "degraded", "down"]) }),
}),
index: z.object({
activeVersion: z.string().max(128),
rebuildProgress: z.number().int().nonnegative().nullable(),
maintenance: z
.object({
operation: z.enum([
"rebuild",
"reconcile",
"backfill",
"reclassify",
"expire",
"rollback",
]),
status: z.enum(["success", "failed", "canceled"]),
durationSeconds: z.number().nonnegative(),
finishedAt: z.string().datetime(),
processed: z.number().int().nonnegative(),
divergences: z.number().int().nonnegative(),
})
.nullable(),
}),
});
export type ProcessorSnapshot = z.infer<typeof ProcessorSnapshotSchema>;
Loading
Loading