/** * Tiny in-process Prometheus registry (text exposition format 0.0.4) — no prom-client dependency. * Counters, gauges and a fixed-bucket histogram with string labels. Shared by the worker (`main.ts`) and the * standalone scheduler (`scheduler.ts`); `renderMetrics()` produces the `/metrics` body. * * Metric contract (docs/DEPLOY.md, Grafana "DataCenterIndex — overview", deploy/monitoring/alerts.yml): * dci_crawl_fetches_total{connector,level,outcome=ok|not_modified|error|blocked} * dci_crawl_credits_total{provider=scrapfly|firecrawl} * dci_crawl_daily_budget{provider} dci_crawl_daily_credits_used{provider} * dci_crawl_fetch_duration_seconds (histogram) * dci_ingest_entities_total{connector,result=created|updated|unchanged|rejected} * dci_events_total * dci_queue_jobs{queue,state} dci_worker_running_jobs dci_worker_up */ type Labels = Record; function esc(v: string | number): string { return String(v).replace(/\\/g, "\\\\").replace(/"/g, '\\"').replace(/\n/g, "\\n"); } function fmt(v: number): string { return Number.isFinite(v) ? (Number.isInteger(v) ? String(v) : v.toString()) : "0"; } function key(names: readonly string[], labels: Labels | undefined): string { if (!names.length) return ""; return names.map((n) => `${n}="${esc(labels?.[n] ?? "")}"`).join(","); } interface Metric { render(): string[]; reset(): void } const registry = new Map(); function register(name: string, m: T): T { if (registry.has(name)) throw new Error(`metric ${name} already registered`); registry.set(name, m); return m; } export class Counter implements Metric { private readonly values = new Map(); constructor(readonly name: string, readonly help: string, readonly labelNames: readonly string[] = []) {} inc(labels?: Labels, n = 1): void { if (!Number.isFinite(n) || n === 0) return; const k = key(this.labelNames, labels); this.values.set(k, (this.values.get(k) ?? 0) + n); } get(labels?: Labels): number { return this.values.get(key(this.labelNames, labels)) ?? 0; } reset(): void { this.values.clear(); } render(): string[] { const out = [`# HELP ${this.name} ${this.help}`, `# TYPE ${this.name} counter`]; // an unlabelled, never-incremented counter renders as 0; a labelled one renders no sample (Prometheus convention) if (!this.values.size && !this.labelNames.length) out.push(`${this.name} 0`); for (const [k, v] of this.values) out.push(`${this.name}${k ? `{${k}}` : ""} ${fmt(v)}`); return out; } } export class Gauge implements Metric { private readonly values = new Map(); constructor(readonly name: string, readonly help: string, readonly labelNames: readonly string[] = []) {} set(value: number, labels?: Labels): void { this.values.set(key(this.labelNames, labels), Number.isFinite(value) ? value : 0); } inc(labels?: Labels, n = 1): void { const k = key(this.labelNames, labels); this.values.set(k, (this.values.get(k) ?? 0) + n); } dec(labels?: Labels, n = 1): void { this.inc(labels, -n); } get(labels?: Labels): number { return this.values.get(key(this.labelNames, labels)) ?? 0; } /** drop every series (used before re-publishing a full snapshot such as queue depths) */ reset(): void { this.values.clear(); } render(): string[] { const out = [`# HELP ${this.name} ${this.help}`, `# TYPE ${this.name} gauge`]; if (!this.values.size && !this.labelNames.length) out.push(`${this.name} 0`); for (const [k, v] of this.values) out.push(`${this.name}${k ? `{${k}}` : ""} ${fmt(v)}`); return out; } } export class Histogram implements Metric { private readonly series = new Map(); constructor(readonly name: string, readonly help: string, readonly buckets: readonly number[], readonly labelNames: readonly string[] = []) { if (!buckets.length || buckets.some((b, i) => i > 0 && b <= buckets[i - 1]!)) throw new Error(`histogram ${name}: buckets must be strictly increasing`); } observe(value: number, labels?: Labels): void { if (!Number.isFinite(value)) return; const k = key(this.labelNames, labels); let s = this.series.get(k); if (!s) { s = { counts: new Array(this.buckets.length).fill(0), sum: 0, count: 0 }; this.series.set(k, s); } for (let i = 0; i < this.buckets.length; i++) if (value <= this.buckets[i]!) s.counts[i]!++; s.sum += value; s.count++; } reset(): void { this.series.clear(); } render(): string[] { const out = [`# HELP ${this.name} ${this.help}`, `# TYPE ${this.name} histogram`]; const emit = (k: string, s: { counts: number[]; sum: number; count: number }) => { const sep = k ? "," : ""; for (let i = 0; i < this.buckets.length; i++) out.push(`${this.name}_bucket{${k}${sep}le="${fmt(this.buckets[i]!)}"} ${s.counts[i]}`); out.push(`${this.name}_bucket{${k}${sep}le="+Inf"} ${s.count}`, `${this.name}_sum${k ? `{${k}}` : ""} ${fmt(s.sum)}`, `${this.name}_count${k ? `{${k}}` : ""} ${s.count}`); }; if (!this.series.size && !this.labelNames.length) emit("", { counts: new Array(this.buckets.length).fill(0), sum: 0, count: 0 }); for (const [k, s] of this.series) emit(k, s); return out; } } export function counter(name: string, help: string, labelNames: readonly string[] = []): Counter { return register(name, new Counter(name, help, labelNames)); } export function gauge(name: string, help: string, labelNames: readonly string[] = []): Gauge { return register(name, new Gauge(name, help, labelNames)); } export function histogram(name: string, help: string, buckets: readonly number[], labelNames: readonly string[] = []): Histogram { return register(name, new Histogram(name, help, buckets, labelNames)); } /** Full exposition body (trailing newline included). */ export function renderMetrics(): string { const lines: string[] = []; for (const m of registry.values()) lines.push(...m.render()); return lines.join("\n") + "\n"; } /** Tests only. */ export function resetMetrics(): void { for (const m of registry.values()) m.reset(); } export type FetchOutcome = "ok" | "not_modified" | "error" | "blocked"; /* ---------- the worker's metrics ---------- */ export const metrics = { fetches: counter("dci_crawl_fetches_total", "fetch attempts that returned a document, by connector, final level and outcome", ["connector", "level", "outcome"]), fetchDuration: histogram("dci_crawl_fetch_duration_seconds", "wall time of one fetch (all escalation levels included)", [0.1, 0.25, 0.5, 1, 2.5, 5, 10, 30, 60, 120]), credits: counter("dci_crawl_credits_total", "premium credits spent since process start", ["provider"]), dailyBudget: gauge("dci_crawl_daily_budget", "daily premium credit limit (DCI__DAILY_BUDGET)", ["provider"]), dailyUsed: gauge("dci_crawl_daily_credits_used", "premium credits used today, shared across workers (Redis dci:budget::)", ["provider"]), ingest: counter("dci_ingest_entities_total", "normalized entities handed to ingest, by result", ["connector", "result"]), events: counter("dci_events_total", "change events emitted by ingest"), runs: counter("dci_crawl_runs_total", "connector runs finished, by status", ["status"]), jobs: counter("dci_worker_jobs_total", "BullMQ jobs processed by this worker", ["queue", "result"]), queueJobs: gauge("dci_queue_jobs", "jobs per queue and state (snapshot at scrape time)", ["queue", "state"]), runningJobs: gauge("dci_worker_running_jobs", "jobs currently processed by this worker"), up: gauge("dci_worker_up", "1 while the process accepts work (0 while shutting down)"), uptime: gauge("dci_worker_uptime_seconds", "seconds since process start"), schedulerTicks: counter("dci_scheduler_ticks_total", "scheduler passes, by result", ["result"]), schedulerEnqueued: counter("dci_scheduler_enqueued_total", "crawl jobs enqueued by the scheduler"), schedulerLastTick: gauge("dci_scheduler_last_tick_timestamp_seconds", "unix time of the last successful scheduler pass"), }; /** Outcome label for a fetched RawDocument-like object. */ export function fetchOutcome(doc: { notModified: boolean; status: number; error?: { code: string } | null }): FetchOutcome { if (doc.notModified) return "not_modified"; if (doc.error) return doc.error.code === "robots_disallow" || doc.error.code === "ssrf_blocked" ? "blocked" : "error"; if (doc.status === 401 || doc.status === 403 || doc.status === 429 || doc.status === 503) return "blocked"; if (doc.status >= 400 || doc.status === 0) return "error"; return "ok"; }