spb/datacenterindex
Public
HTML 53.9%
TypeScript 44.5%
JavaScript 0.6%
SQL 0.5%
1/**2 * Tiny in-process Prometheus registry (text exposition format 0.0.4) — no prom-client dependency.3 * Counters, gauges and a fixed-bucket histogram with string labels. Shared by the worker (`main.ts`) and the4 * standalone scheduler (`scheduler.ts`); `renderMetrics()` produces the `/metrics` body.5 *6 * Metric contract (docs/DEPLOY.md, Grafana "DataCenterIndex — overview", deploy/monitoring/alerts.yml):7 * dci_crawl_fetches_total{connector,level,outcome=ok|not_modified|error|blocked}8 * dci_crawl_credits_total{provider=scrapfly|firecrawl}9 * dci_crawl_daily_budget{provider} dci_crawl_daily_credits_used{provider}10 * dci_crawl_fetch_duration_seconds (histogram)11 * dci_ingest_entities_total{connector,result=created|updated|unchanged|rejected}12 * dci_events_total13 * dci_queue_jobs{queue,state} dci_worker_running_jobs dci_worker_up14 */1516type Labels = Record<string, string | number>;1718function esc(v: string | number): string { return String(v).replace(/\\/g, "\\\\").replace(/"/g, '\\"').replace(/\n/g, "\\n"); }19function fmt(v: number): string { return Number.isFinite(v) ? (Number.isInteger(v) ? String(v) : v.toString()) : "0"; }2021function key(names: readonly string[], labels: Labels | undefined): string {22 if (!names.length) return "";23 return names.map((n) => `${n}="${esc(labels?.[n] ?? "")}"`).join(",");24}2526interface Metric { render(): string[]; reset(): void }27const registry = new Map<string, Metric>();2829function register<T extends Metric>(name: string, m: T): T {30 if (registry.has(name)) throw new Error(`metric ${name} already registered`);31 registry.set(name, m);32 return m;33}3435export class Counter implements Metric {36 private readonly values = new Map<string, number>();37 constructor(readonly name: string, readonly help: string, readonly labelNames: readonly string[] = []) {}38 inc(labels?: Labels, n = 1): void {39 if (!Number.isFinite(n) || n === 0) return;40 const k = key(this.labelNames, labels);41 this.values.set(k, (this.values.get(k) ?? 0) + n);42 }43 get(labels?: Labels): number { return this.values.get(key(this.labelNames, labels)) ?? 0; }44 reset(): void { this.values.clear(); }45 render(): string[] {46 const out = [`# HELP ${this.name} ${this.help}`, `# TYPE ${this.name} counter`];47 // an unlabelled, never-incremented counter renders as 0; a labelled one renders no sample (Prometheus convention)48 if (!this.values.size && !this.labelNames.length) out.push(`${this.name} 0`);49 for (const [k, v] of this.values) out.push(`${this.name}${k ? `{${k}}` : ""} ${fmt(v)}`);50 return out;51 }52}5354export class Gauge implements Metric {55 private readonly values = new Map<string, number>();56 constructor(readonly name: string, readonly help: string, readonly labelNames: readonly string[] = []) {}57 set(value: number, labels?: Labels): void { this.values.set(key(this.labelNames, labels), Number.isFinite(value) ? value : 0); }58 inc(labels?: Labels, n = 1): void { const k = key(this.labelNames, labels); this.values.set(k, (this.values.get(k) ?? 0) + n); }59 dec(labels?: Labels, n = 1): void { this.inc(labels, -n); }60 get(labels?: Labels): number { return this.values.get(key(this.labelNames, labels)) ?? 0; }61 /** drop every series (used before re-publishing a full snapshot such as queue depths) */62 reset(): void { this.values.clear(); }63 render(): string[] {64 const out = [`# HELP ${this.name} ${this.help}`, `# TYPE ${this.name} gauge`];65 if (!this.values.size && !this.labelNames.length) out.push(`${this.name} 0`);66 for (const [k, v] of this.values) out.push(`${this.name}${k ? `{${k}}` : ""} ${fmt(v)}`);67 return out;68 }69}7071export class Histogram implements Metric {72 private readonly series = new Map<string, { counts: number[]; sum: number; count: number }>();73 constructor(readonly name: string, readonly help: string, readonly buckets: readonly number[], readonly labelNames: readonly string[] = []) {74 if (!buckets.length || buckets.some((b, i) => i > 0 && b <= buckets[i - 1]!)) throw new Error(`histogram ${name}: buckets must be strictly increasing`);75 }76 observe(value: number, labels?: Labels): void {77 if (!Number.isFinite(value)) return;78 const k = key(this.labelNames, labels);79 let s = this.series.get(k);80 if (!s) { s = { counts: new Array<number>(this.buckets.length).fill(0), sum: 0, count: 0 }; this.series.set(k, s); }81 for (let i = 0; i < this.buckets.length; i++) if (value <= this.buckets[i]!) s.counts[i]!++;82 s.sum += value; s.count++;83 }84 reset(): void { this.series.clear(); }85 render(): string[] {86 const out = [`# HELP ${this.name} ${this.help}`, `# TYPE ${this.name} histogram`];87 const emit = (k: string, s: { counts: number[]; sum: number; count: number }) => {88 const sep = k ? "," : "";89 for (let i = 0; i < this.buckets.length; i++) out.push(`${this.name}_bucket{${k}${sep}le="${fmt(this.buckets[i]!)}"} ${s.counts[i]}`);90 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}`);91 };92 if (!this.series.size && !this.labelNames.length) emit("", { counts: new Array<number>(this.buckets.length).fill(0), sum: 0, count: 0 });93 for (const [k, s] of this.series) emit(k, s);94 return out;95 }96}9798export function counter(name: string, help: string, labelNames: readonly string[] = []): Counter { return register(name, new Counter(name, help, labelNames)); }99export function gauge(name: string, help: string, labelNames: readonly string[] = []): Gauge { return register(name, new Gauge(name, help, labelNames)); }100export function histogram(name: string, help: string, buckets: readonly number[], labelNames: readonly string[] = []): Histogram { return register(name, new Histogram(name, help, buckets, labelNames)); }101102/** Full exposition body (trailing newline included). */103export function renderMetrics(): string {104 const lines: string[] = [];105 for (const m of registry.values()) lines.push(...m.render());106 return lines.join("\n") + "\n";107}108109/** Tests only. */110export function resetMetrics(): void { for (const m of registry.values()) m.reset(); }111112export type FetchOutcome = "ok" | "not_modified" | "error" | "blocked";113114/* ---------- the worker's metrics ---------- */115export const metrics = {116 fetches: counter("dci_crawl_fetches_total", "fetch attempts that returned a document, by connector, final level and outcome", ["connector", "level", "outcome"]),117 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]),118 credits: counter("dci_crawl_credits_total", "premium credits spent since process start", ["provider"]),119 dailyBudget: gauge("dci_crawl_daily_budget", "daily premium credit limit (DCI_<PROVIDER>_DAILY_BUDGET)", ["provider"]),120 dailyUsed: gauge("dci_crawl_daily_credits_used", "premium credits used today, shared across workers (Redis dci:budget:<provider>:<day>)", ["provider"]),121 ingest: counter("dci_ingest_entities_total", "normalized entities handed to ingest, by result", ["connector", "result"]),122 events: counter("dci_events_total", "change events emitted by ingest"),123 runs: counter("dci_crawl_runs_total", "connector runs finished, by status", ["status"]),124 jobs: counter("dci_worker_jobs_total", "BullMQ jobs processed by this worker", ["queue", "result"]),125 queueJobs: gauge("dci_queue_jobs", "jobs per queue and state (snapshot at scrape time)", ["queue", "state"]),126 runningJobs: gauge("dci_worker_running_jobs", "jobs currently processed by this worker"),127 up: gauge("dci_worker_up", "1 while the process accepts work (0 while shutting down)"),128 uptime: gauge("dci_worker_uptime_seconds", "seconds since process start"),129 schedulerTicks: counter("dci_scheduler_ticks_total", "scheduler passes, by result", ["result"]),130 schedulerEnqueued: counter("dci_scheduler_enqueued_total", "crawl jobs enqueued by the scheduler"),131 schedulerLastTick: gauge("dci_scheduler_last_tick_timestamp_seconds", "unix time of the last successful scheduler pass"),132};133134/** Outcome label for a fetched RawDocument-like object. */135export function fetchOutcome(doc: { notModified: boolean; status: number; error?: { code: string } | null }): FetchOutcome {136 if (doc.notModified) return "not_modified";137 if (doc.error) return doc.error.code === "robots_disallow" || doc.error.code === "ssrf_blocked" ? "blocked" : "error";138 if (doc.status === 401 || doc.status === 403 || doc.status === 429 || doc.status === 503) return "blocked";139 if (doc.status >= 400 || doc.status === 0) return "error";140 return "ok";141}142