TypeScript 55.4%
Python 43.2%
SQL 1.2%
1import Fastify from "fastify";2import client from "prom-client";3import { db, sql } from "@websensor/db";4import { config, log } from "./config";56export const registry = new client.Registry();7client.collectDefaultMetrics({ register: registry, prefix: "websensor_engine_" });89export const m = {10 checks: new client.Counter({ name: "websensor_checks_total", help: "Sensor checks", labelNames: ["connector", "outcome"], registers: [registry] }),11 bytes: new client.Counter({ name: "websensor_fetched_bytes_total", help: "Bytes fetched", labelNames: ["connector"], registers: [registry] }),12 changes: new client.Counter({ name: "websensor_changes_total", help: "Raw changes", labelNames: ["connector", "meaningful"], registers: [registry] }),13 events: new client.Counter({ name: "websensor_events_total", help: "Events published", labelNames: ["event_type", "silent"], registers: [registry] }),14 llmCalls: new client.Counter({ name: "websensor_llm_calls_total", help: "LLM calls", labelNames: ["model", "ok"], registers: [registry] }),15 llmTokens: new client.Counter({ name: "websensor_llm_tokens_total", help: "LLM tokens", labelNames: ["model", "direction"], registers: [registry] }),16 fetchDuration: new client.Histogram({ name: "websensor_fetch_duration_seconds", help: "Fetch latency", labelNames: ["connector"], buckets: [0.1, 0.25, 0.5, 1, 2, 5, 10, 25], registers: [registry] }),17 processingLatency: new client.Histogram({ name: "websensor_processing_latency_seconds", help: "Detected → published", buckets: [0.05, 0.1, 0.25, 0.5, 1, 2, 5, 15], registers: [registry] }),18 queueDue: new client.Gauge({ name: "websensor_sensors_due", help: "Sensors due for a check", registers: [registry] }),19 inflight: new client.Gauge({ name: "websensor_inflight_checks", help: "Checks in flight", registers: [registry] }),20 httpStatus: new client.Counter({ name: "websensor_http_status_total", help: "HTTP status codes", labelNames: ["status"], registers: [registry] }),21 alertsFired: new client.Counter({ name: "websensor_alerts_fired_total", help: "Alert rules matched", labelNames: ["channel"], registers: [registry] }),22 webhookDeliveries: new client.Counter({ name: "websensor_webhook_deliveries_total", help: "Webhook deliveries", labelNames: ["ok"], registers: [registry] }),23 changeClass: new client.Counter({ name: "websensor_change_class_total", help: "Semantic class of raw changes", labelNames: ["class"], registers: [registry] }),24 suppressed: new client.Counter({ name: "websensor_suppressed_total", help: "Candidates suppressed before becoming events", labelNames: ["reason"], registers: [registry] }),25 stageLatency: new client.Histogram({ name: "websensor_stage_seconds", help: "Pipeline stage latency", labelNames: ["stage"], buckets: [0.005, 0.02, 0.05, 0.1, 0.25, 0.5, 1, 2, 5], registers: [registry] }),26 prunedBlobs: new client.Counter({ name: "websensor_pruned_blobs_total", help: "Raw snapshot bodies pruned by retention", registers: [registry] }),27};2829/** Per-source daily counters (source quality, heatmaps). */30export async function bumpSourceDaily(sourceId: string, fields: Partial<Record<"checks" | "not_modified" | "errors" | "raw_changes" | "events", number>>): Promise<void> {31 const keys = Object.keys(fields) as (keyof typeof fields)[];32 if (!keys.length) return;33 const sets = keys.map((k) => sql.raw(`${k} = source_daily.${k} + ${Number(fields[k] ?? 0)}`));34 const cols = keys.map((k) => sql.raw(k));35 const vals = keys.map((k) => sql`${Number(fields[k] ?? 0)}`);36 try {37 await db.execute(sql`insert into source_daily (source_id, day, ${sql.join(cols, sql`, `)}) values (${sourceId}, (now() at time zone 'UTC')::date, ${sql.join(vals, sql`, `)}) on conflict (source_id, day) do update set ${sql.join(sets, sql`, `)}`);38 } catch (e) {39 log.warn({ err: (e as Error).message }, "source_daily update failed");40 }41}4243/** Daily counters in Postgres (feed the public homepage stats). */44export async function bumpDaily(fields: Partial<Record<"checks" | "not_modified" | "bytes" | "raw_changes" | "events" | "silent_events" | "errors" | "llm_calls" | "llm_input_tokens" | "llm_output_tokens" | "scrapfly_calls", number>>): Promise<void> {45 const keys = Object.keys(fields) as (keyof typeof fields)[];46 if (!keys.length) return;47 const sets = keys.map((k) => sql.raw(`${k} = metrics_daily.${k} + ${Number(fields[k] ?? 0)}`));48 const cols = keys.map((k) => sql.raw(k));49 const vals = keys.map((k) => sql`${Number(fields[k] ?? 0)}`);50 try {51 await db.execute(sql`insert into metrics_daily (day, ${sql.join(cols, sql`, `)}) values (current_date, ${sql.join(vals, sql`, `)}) on conflict (day) do update set ${sql.join(sets, sql`, `)}`);52 } catch (e) {53 log.warn({ err: (e as Error).message }, "metrics_daily update failed");54 }55}5657export async function startMetricsServer(): Promise<void> {58 const app = Fastify({ logger: false });59 app.get("/metrics", async (_req, reply) => {60 reply.header("content-type", registry.contentType);61 return registry.metrics();62 });63 app.get("/health", async () => ({ status: "ok", service: "engine", version: config.version, time: new Date().toISOString() }));64 await app.listen({ port: config.metricsPort, host: config.metricsHost });65 log.info({ port: config.metricsPort }, "engine metrics listening");66}67