import Fastify from "fastify"; import client from "prom-client"; import { db, sql } from "@websensor/db"; import { config, log } from "./config"; export const registry = new client.Registry(); client.collectDefaultMetrics({ register: registry, prefix: "websensor_engine_" }); export const m = { checks: new client.Counter({ name: "websensor_checks_total", help: "Sensor checks", labelNames: ["connector", "outcome"], registers: [registry] }), bytes: new client.Counter({ name: "websensor_fetched_bytes_total", help: "Bytes fetched", labelNames: ["connector"], registers: [registry] }), changes: new client.Counter({ name: "websensor_changes_total", help: "Raw changes", labelNames: ["connector", "meaningful"], registers: [registry] }), events: new client.Counter({ name: "websensor_events_total", help: "Events published", labelNames: ["event_type", "silent"], registers: [registry] }), llmCalls: new client.Counter({ name: "websensor_llm_calls_total", help: "LLM calls", labelNames: ["model", "ok"], registers: [registry] }), llmTokens: new client.Counter({ name: "websensor_llm_tokens_total", help: "LLM tokens", labelNames: ["model", "direction"], registers: [registry] }), 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] }), 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] }), queueDue: new client.Gauge({ name: "websensor_sensors_due", help: "Sensors due for a check", registers: [registry] }), inflight: new client.Gauge({ name: "websensor_inflight_checks", help: "Checks in flight", registers: [registry] }), httpStatus: new client.Counter({ name: "websensor_http_status_total", help: "HTTP status codes", labelNames: ["status"], registers: [registry] }), alertsFired: new client.Counter({ name: "websensor_alerts_fired_total", help: "Alert rules matched", labelNames: ["channel"], registers: [registry] }), webhookDeliveries: new client.Counter({ name: "websensor_webhook_deliveries_total", help: "Webhook deliveries", labelNames: ["ok"], registers: [registry] }), changeClass: new client.Counter({ name: "websensor_change_class_total", help: "Semantic class of raw changes", labelNames: ["class"], registers: [registry] }), suppressed: new client.Counter({ name: "websensor_suppressed_total", help: "Candidates suppressed before becoming events", labelNames: ["reason"], registers: [registry] }), 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] }), prunedBlobs: new client.Counter({ name: "websensor_pruned_blobs_total", help: "Raw snapshot bodies pruned by retention", registers: [registry] }), }; /** Per-source daily counters (source quality, heatmaps). */ export async function bumpSourceDaily(sourceId: string, fields: Partial>): Promise { const keys = Object.keys(fields) as (keyof typeof fields)[]; if (!keys.length) return; const sets = keys.map((k) => sql.raw(`${k} = source_daily.${k} + ${Number(fields[k] ?? 0)}`)); const cols = keys.map((k) => sql.raw(k)); const vals = keys.map((k) => sql`${Number(fields[k] ?? 0)}`); try { 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`, `)}`); } catch (e) { log.warn({ err: (e as Error).message }, "source_daily update failed"); } } /** Daily counters in Postgres (feed the public homepage stats). */ export async function bumpDaily(fields: Partial>): Promise { const keys = Object.keys(fields) as (keyof typeof fields)[]; if (!keys.length) return; const sets = keys.map((k) => sql.raw(`${k} = metrics_daily.${k} + ${Number(fields[k] ?? 0)}`)); const cols = keys.map((k) => sql.raw(k)); const vals = keys.map((k) => sql`${Number(fields[k] ?? 0)}`); try { 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`, `)}`); } catch (e) { log.warn({ err: (e as Error).message }, "metrics_daily update failed"); } } export async function startMetricsServer(): Promise { const app = Fastify({ logger: false }); app.get("/metrics", async (_req, reply) => { reply.header("content-type", registry.contentType); return registry.metrics(); }); app.get("/health", async () => ({ status: "ok", service: "engine", version: config.version, time: new Date().toISOString() })); await app.listen({ port: config.metricsPort, host: config.metricsHost }); log.info({ port: config.metricsPort }, "engine metrics listening"); }