/** * Connector health (ConnectorHealthDTO) from the connectors row, its latest run, the last-24 h run stats, document * counters, the source licence and (best effort) ClickHouse crawl_log latency / premium request counts. */ import type { ConnectorHealthDTO } from "@dci/core"; import { chQuery } from "@dci/db/clickhouse"; import { pg } from "../../lib/sql.js"; import { bool, int, iso, json, num, reqStr, str, type Row } from "../../lib/rows.js"; import { asSourceKind } from "../../lib/dto.js"; interface ChCost { connector_id: string; fetcher: string; n: string | number; avg_ms: string | number | null; credits: string | number | null; n24: string | number | null } interface ChStats { avgMs: number | null; cost: ConnectorHealthDTO["cost"]; requests24h: { scrapfly: number; firecrawl: number } } /** Per-connector cost / latency from ClickHouse crawl_log over the last 7 days + premium request counts over 24 h (null when CH unavailable). */ async function crawlLogStats(): Promise | null> { try { const rows = await Promise.race([ chQuery("select connector_id, fetcher, count() as n, avg(duration_ms) as avg_ms, sum(credits) as credits, countIf(ts >= now() - interval 24 hour) as n24 from crawl_log where ts >= now() - interval 7 day group by connector_id, fetcher"), new Promise((_, rej) => setTimeout(() => rej(new Error("clickhouse timeout")), 3_000)), ]); const out = new Map(); for (const r of rows) { let e = out.get(r.connector_id); if (!e) { e = { avgMs: null, cost: { direct: 0, firecrawl: 0, scrapfly: 0, credits: 0 }, requests24h: { scrapfly: 0, firecrawl: 0 }, _n: 0, _sum: 0 }; out.set(r.connector_id, e); } const n = int(r.n); const f = r.fetcher === "firecrawl" ? "firecrawl" : r.fetcher === "scrapfly" ? "scrapfly" : "direct"; e.cost[f] += n; e.cost.credits += num(r.credits) ?? 0; if (f !== "direct") e.requests24h[f] += int(r.n24); e._n += n; e._sum += (num(r.avg_ms) ?? 0) * n; e.avgMs = e._n ? Math.round(e._sum / e._n) : null; } return out; } catch { return null; } } function runStat(stats: Record, ...keys: string[]): number | null { for (const k of keys) { const v = num(stats[k]); if (v != null) return v; } return null; } export function connectorHealth(r: Row, ch: Map | null): ConnectorHealthDTO { const stats = json>(r.stats, {}); const runStats = json>(r.run_stats, {}); const day = json>(r.day_stats, {}); const merged = { ...stats, ...runStats }; const paused = bool(r.paused); const quarantine = bool(r.quarantine); const consecutiveFailures = int(r.consecutive_failures); const blockedSince = iso(r.blocked_since); const healthCol = str(r.health) ?? "never_run"; const lastError = str(r.last_error) ?? ""; const docsFetched = int(r.doc_fetched); const extractTried = int(r.extract_tried); const extractionSuccess = extractTried ? Math.round((int(r.extract_ok) / extractTried) * 1000) / 1000 : runStat(runStats, "extractionSuccess"); const neverRun = !r.last_run_at && !r.run_started; const weekRuns = int(r.week_runs); const weekNew = int(r.week_created) + int(r.week_changed); let health: ConnectorHealthDTO["health"]; if (paused) health = "paused"; else if (quarantine) health = "quarantine"; else if (neverRun) health = "never_run"; else if (blockedSince) health = "blocked"; else if ((extractionSuccess != null && extractionSuccess < 0.2 && extractTried >= 10) || /schema|selector|parser/i.test(lastError)) health = "schema_change"; else if (consecutiveFailures >= 3) health = "failing"; else if (weekRuns >= 3 && weekNew === 0 && (str(r.run_status) === "ok" || str(r.last_status) === "ok")) health = "no_new_content"; else if (healthCol === "ok" || healthCol === "degraded" || healthCol === "failing") health = healthCol; else health = str(r.run_status) === "failed" ? "failing" : "ok"; // stored health stale (e.g. never_run) but a run exists const chStats = ch?.get(reqStr(r.id)); const cost: ConnectorHealthDTO["cost"] = chStats?.cost ?? { direct: runStat(merged, "direct", "fetch_direct") ?? 0, firecrawl: runStat(merged, "firecrawl", "fetch_firecrawl") ?? 0, scrapfly: runStat(merged, "scrapfly", "fetch_scrapfly") ?? 0, credits: runStat(merged, "credits", "creditsSpent") ?? 0 }; const d = (k: string, ...alts: string[]) => runStat(day, k, ...alts) ?? 0; return { id: reqStr(r.id), sourceName: reqStr(r.source_name), domain: reqStr(r.domain), kind: asSourceKind(r.kind), mode: reqStr(r.mode, "hybrid"), enabled: bool(r.enabled) && !paused, health, parserVersion: reqStr(r.parser_version, "v1"), lastRunAt: iso(r.last_run_at) ?? iso(r.run_started), lastSuccessAt: iso(r.last_success_at), lastFailureAt: iso(r.last_failure_at), nextRunAt: iso(r.next_run_at), lastStatus: str(r.last_status) ?? str(r.run_status), discovered: runStat(runStats, "discovered", "urls") ?? int(r.doc_total), fetched: runStat(runStats, "fetched") ?? docsFetched, changed: runStat(runStats, "changed") ?? int(r.doc_changed), failed: runStat(runStats, "failed", "errors") ?? int(r.doc_failed), extracted: runStat(runStats, "extracted", "entities", "received") ?? int(r.doc_extracted), extractionSuccess, avgResponseMs: chStats?.avgMs ?? runStat(merged, "avgMs", "avg_ms", "avgResponseMs"), cost, schedule: json>(r.schedule, {}), quarantine, consecutiveFailures, blockedSince, urlsDiscovered: d("discovered"), urlsFetched: d("fetched"), newDocs: d("created"), changedDocs: d("changed"), recordsCreated: d("created"), recordsModified: d("updated"), rejectedClaims: d("unscopedClaims"), httpErrors: d("failed"), antiBotEscalations: d("blockedFetches", "robotsBlocked"), scrapflyRequests: chStats ? chStats.requests24h.scrapfly : d("scrapfly", "fetch_scrapfly"), firecrawlRequests: chStats ? chStats.requests24h.firecrawl : d("firecrawl", "fetch_firecrawl"), priorityScore: num(r.priority_score), license: str(r.src_license), redistribution: str(r.src_redistribution), }; } const DAY_KEYS = ["discovered", "fetched", "created", "changed", "updated", "unscopedClaims", "failed", "blockedFetches", "robotsBlocked", "scrapfly", "firecrawl", "fetch_scrapfly", "fetch_firecrawl"]; const CONNECTOR_SELECT = (sql: ReturnType) => sql` select c.*, r.id as run_id, r.status as run_status, r.stats as run_stats, r.started_at as run_started, r.finished_at as run_finished, ls.finished_at as last_success_at, lf.finished_at as last_failure_at, coalesce(d.total, 0) as doc_total, coalesce(d.fetched, 0) as doc_fetched, coalesce(d.changed, 0) as doc_changed, coalesce(d.failed, 0) as doc_failed, coalesce(d.extracted, 0) as doc_extracted, coalesce(d.extract_ok, 0) as extract_ok, coalesce(d.extract_tried, 0) as extract_tried, coalesce(d.quarantined, 0) as doc_quarantined, coalesce(s24.day_stats, '{}'::jsonb) as day_stats, coalesce(w.runs, 0) as week_runs, coalesce(w.created, 0) as week_created, coalesce(w.changed, 0) as week_changed, src.license as src_license, src.redistribution as src_redistribution from connectors c left join lateral (select id, status, stats, started_at, finished_at from connector_runs where connector_id = c.id order by started_at desc limit 1) r on true left join lateral (select finished_at from connector_runs where connector_id = c.id and status = 'ok' order by started_at desc limit 1) ls on true left join lateral (select finished_at from connector_runs where connector_id = c.id and status = 'failed' order by started_at desc limit 1) lf on true left join lateral (select license, redistribution from sources where connector_id = c.id order by priority, id limit 1) src on true left join ( select connector_id, count(*)::int as runs, coalesce(sum((stats->>'created')::numeric), 0) as created, coalesce(sum((stats->>'changed')::numeric), 0) as changed from connector_runs where started_at >= now() - interval '7 days' and status <> 'running' group by connector_id ) w on w.connector_id = c.id left join ( select connector_id, jsonb_object_agg(key, total) as day_stats from ( select cr.connector_id, kv.key, sum(kv.value) as total from connector_runs cr cross join lateral (select key, (value #>> '{}')::numeric as value from jsonb_each(cr.stats) where key = any(${DAY_KEYS}) and jsonb_typeof(value) = 'number') kv where cr.started_at >= now() - interval '24 hours' group by cr.connector_id, kv.key ) x group by connector_id ) s24 on s24.connector_id = c.id left join ( select connector_id, count(*)::int as total, count(*) filter (where last_fetched is not null)::int as fetched, count(*) filter (where change_count > 0)::int as changed, count(*) filter (where error_count > 0)::int as failed, coalesce(sum(extract_count), 0)::int as extracted, count(*) filter (where extract_ok)::int as extract_ok, count(*) filter (where extract_ok is not null)::int as extract_tried, count(*) filter (where quarantined)::int as quarantined from documents group by connector_id ) d on d.connector_id = c.id`; export async function listConnectorHealth(): Promise { const sql = pg(); const [rows, ch] = await Promise.all([sql`${CONNECTOR_SELECT(sql)} order by c.enabled desc, c.id`, crawlLogStats()]); return rows.map((r) => connectorHealth(r, ch)); } export interface ConnectorAdminDetail { health: ConnectorHealthDTO; config: Record; paused: boolean; lastError: string | null; runs: Array>; documentsByPageType: Array<{ pageType: string; count: number; changed: number; failed: number; quarantined: number }>; errorSamples: Array<{ id: string; url: string; error: string | null; statusCode: number | null; errorCount: number; lastChecked: string | null }>; createdAt: string | null; updatedAt: string | null } export async function getConnectorAdmin(id: string): Promise { const sql = pg(); const [rows, ch] = await Promise.all([sql`${CONNECTOR_SELECT(sql)} where c.id = ${id}`, crawlLogStats()]); const r = rows[0]; if (!r) return null; const [runs, byType, errors] = await Promise.all([ sql`select id, task, started_at, finished_at, status, stats, error, quarantined from connector_runs where connector_id = ${id} order by started_at desc limit 30`, sql`select page_type, count(*)::int as n, count(*) filter (where change_count > 0)::int as changed, count(*) filter (where error_count > 0)::int as failed, count(*) filter (where quarantined)::int as quarantined from documents where connector_id = ${id} group by page_type order by n desc`, sql`select id, url, error, status_code, error_count, last_checked from documents where connector_id = ${id} and error is not null order by last_checked desc nulls last limit 10`, ]); return { health: connectorHealth(r, ch), config: json>(r.config, {}), paused: bool(r.paused), lastError: str(r.last_error), runs: runs.map((x) => ({ id: reqStr(x.id), task: reqStr(x.task), startedAt: iso(x.started_at), finishedAt: iso(x.finished_at), status: reqStr(x.status), stats: json(x.stats, {}), error: str(x.error), quarantined: bool(x.quarantined) })), documentsByPageType: byType.map((x) => ({ pageType: reqStr(x.page_type), count: int(x.n), changed: int(x.changed), failed: int(x.failed), quarantined: int(x.quarantined) })), errorSamples: errors.map((x) => ({ id: reqStr(x.id), url: reqStr(x.url), error: str(x.error), statusCode: num(x.status_code), errorCount: int(x.error_count), lastChecked: iso(x.last_checked) })), createdAt: iso(r.created_at), updatedAt: iso(r.updated_at), }; } export async function listRuns(connectorId: string | undefined, limit = 100): Promise>> { const sql = pg(); const rows = await sql`select id, connector_id, task, started_at, finished_at, status, stats, error, quarantined from connector_runs where ${connectorId ? sql`connector_id = ${connectorId}` : sql`true`} order by started_at desc limit ${limit}`; return rows.map((x) => ({ id: reqStr(x.id), connectorId: reqStr(x.connector_id), task: reqStr(x.task), startedAt: iso(x.started_at), finishedAt: iso(x.finished_at), status: reqStr(x.status), stats: json(x.stats, {}), error: str(x.error), quarantined: bool(x.quarantined) })); } export async function getRun(id: string): Promise | null> { const sql = pg(); const rows = await sql`select * from connector_runs where id = ${id}`; const x = rows[0]; if (!x) return null; return { id: reqStr(x.id), connectorId: reqStr(x.connector_id), task: reqStr(x.task), startedAt: iso(x.started_at), finishedAt: iso(x.finished_at), status: reqStr(x.status), stats: json(x.stats, {}), error: str(x.error), log: json(x.log, []), quarantined: bool(x.quarantined) }; }