SPB Git forge
38commits 1branches 0releases
338.7 MBsize
maindefault branch
3 h agolast push
HTML 53.9% TypeScript 44.5% JavaScript 0.6% SQL 0.5%
8.5 KB · 142 lines typescript
Raw Blame History
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