spb/datacenterindex
Public
HTML 53.9%
TypeScript 44.5%
JavaScript 0.6%
SQL 0.5%
1/**2 * Connector health (ConnectorHealthDTO) from the connectors row, its latest run, the last-24 h run stats, document3 * counters, the source licence and (best effort) ClickHouse crawl_log latency / premium request counts.4 */5import type { ConnectorHealthDTO } from "@dci/core";6import { chQuery } from "@dci/db/clickhouse";7import { pg } from "../../lib/sql.js";8import { bool, int, iso, json, num, reqStr, str, type Row } from "../../lib/rows.js";9import { asSourceKind } from "../../lib/dto.js";1011interface ChCost { connector_id: string; fetcher: string; n: string | number; avg_ms: string | number | null; credits: string | number | null; n24: string | number | null }12interface ChStats { avgMs: number | null; cost: ConnectorHealthDTO["cost"]; requests24h: { scrapfly: number; firecrawl: number } }1314/** Per-connector cost / latency from ClickHouse crawl_log over the last 7 days + premium request counts over 24 h (null when CH unavailable). */15async function crawlLogStats(): Promise<Map<string, ChStats> | null> {16 try {17 const rows = await Promise.race([18 chQuery<ChCost>("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"),19 new Promise<never>((_, rej) => setTimeout(() => rej(new Error("clickhouse timeout")), 3_000)),20 ]);21 const out = new Map<string, ChStats & { _n: number; _sum: number }>();22 for (const r of rows) {23 let e = out.get(r.connector_id);24 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); }25 const n = int(r.n);26 const f = r.fetcher === "firecrawl" ? "firecrawl" : r.fetcher === "scrapfly" ? "scrapfly" : "direct";27 e.cost[f] += n;28 e.cost.credits += num(r.credits) ?? 0;29 if (f !== "direct") e.requests24h[f] += int(r.n24);30 e._n += n;31 e._sum += (num(r.avg_ms) ?? 0) * n;32 e.avgMs = e._n ? Math.round(e._sum / e._n) : null;33 }34 return out;35 } catch {36 return null;37 }38}3940function runStat(stats: Record<string, unknown>, ...keys: string[]): number | null {41 for (const k of keys) { const v = num(stats[k]); if (v != null) return v; }42 return null;43}4445export function connectorHealth(r: Row, ch: Map<string, ChStats> | null): ConnectorHealthDTO {46 const stats = json<Record<string, unknown>>(r.stats, {});47 const runStats = json<Record<string, unknown>>(r.run_stats, {});48 const day = json<Record<string, unknown>>(r.day_stats, {});49 const merged = { ...stats, ...runStats };50 const paused = bool(r.paused);51 const quarantine = bool(r.quarantine);52 const consecutiveFailures = int(r.consecutive_failures);53 const blockedSince = iso(r.blocked_since);54 const healthCol = str(r.health) ?? "never_run";55 const lastError = str(r.last_error) ?? "";56 const docsFetched = int(r.doc_fetched);57 const extractTried = int(r.extract_tried);58 const extractionSuccess = extractTried ? Math.round((int(r.extract_ok) / extractTried) * 1000) / 1000 : runStat(runStats, "extractionSuccess");59 const neverRun = !r.last_run_at && !r.run_started;60 const weekRuns = int(r.week_runs);61 const weekNew = int(r.week_created) + int(r.week_changed);62 let health: ConnectorHealthDTO["health"];63 if (paused) health = "paused";64 else if (quarantine) health = "quarantine";65 else if (neverRun) health = "never_run";66 else if (blockedSince) health = "blocked";67 else if ((extractionSuccess != null && extractionSuccess < 0.2 && extractTried >= 10) || /schema|selector|parser/i.test(lastError)) health = "schema_change";68 else if (consecutiveFailures >= 3) health = "failing";69 else if (weekRuns >= 3 && weekNew === 0 && (str(r.run_status) === "ok" || str(r.last_status) === "ok")) health = "no_new_content";70 else if (healthCol === "ok" || healthCol === "degraded" || healthCol === "failing") health = healthCol;71 else health = str(r.run_status) === "failed" ? "failing" : "ok"; // stored health stale (e.g. never_run) but a run exists72 const chStats = ch?.get(reqStr(r.id));73 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 };74 const d = (k: string, ...alts: string[]) => runStat(day, k, ...alts) ?? 0;75 return {76 id: reqStr(r.id),77 sourceName: reqStr(r.source_name),78 domain: reqStr(r.domain),79 kind: asSourceKind(r.kind),80 mode: reqStr(r.mode, "hybrid"),81 enabled: bool(r.enabled) && !paused,82 health,83 parserVersion: reqStr(r.parser_version, "v1"),84 lastRunAt: iso(r.last_run_at) ?? iso(r.run_started),85 lastSuccessAt: iso(r.last_success_at),86 lastFailureAt: iso(r.last_failure_at),87 nextRunAt: iso(r.next_run_at),88 lastStatus: str(r.last_status) ?? str(r.run_status),89 discovered: runStat(runStats, "discovered", "urls") ?? int(r.doc_total),90 fetched: runStat(runStats, "fetched") ?? docsFetched,91 changed: runStat(runStats, "changed") ?? int(r.doc_changed),92 failed: runStat(runStats, "failed", "errors") ?? int(r.doc_failed),93 extracted: runStat(runStats, "extracted", "entities", "received") ?? int(r.doc_extracted),94 extractionSuccess,95 avgResponseMs: chStats?.avgMs ?? runStat(merged, "avgMs", "avg_ms", "avgResponseMs"),96 cost,97 schedule: json<Record<string, string>>(r.schedule, {}),98 quarantine,99 consecutiveFailures,100 blockedSince,101 urlsDiscovered: d("discovered"),102 urlsFetched: d("fetched"),103 newDocs: d("created"),104 changedDocs: d("changed"),105 recordsCreated: d("created"),106 recordsModified: d("updated"),107 rejectedClaims: d("unscopedClaims"),108 httpErrors: d("failed"),109 antiBotEscalations: d("blockedFetches", "robotsBlocked"),110 scrapflyRequests: chStats ? chStats.requests24h.scrapfly : d("scrapfly", "fetch_scrapfly"),111 firecrawlRequests: chStats ? chStats.requests24h.firecrawl : d("firecrawl", "fetch_firecrawl"),112 priorityScore: num(r.priority_score),113 license: str(r.src_license),114 redistribution: str(r.src_redistribution),115 };116}117118const DAY_KEYS = ["discovered", "fetched", "created", "changed", "updated", "unscopedClaims", "failed", "blockedFetches", "robotsBlocked", "scrapfly", "firecrawl", "fetch_scrapfly", "fetch_firecrawl"];119120const CONNECTOR_SELECT = (sql: ReturnType<typeof pg>) => sql`121 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,122 ls.finished_at as last_success_at, lf.finished_at as last_failure_at,123 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,124 coalesce(d.extract_ok, 0) as extract_ok, coalesce(d.extract_tried, 0) as extract_tried, coalesce(d.quarantined, 0) as doc_quarantined,125 coalesce(s24.day_stats, '{}'::jsonb) as day_stats,126 coalesce(w.runs, 0) as week_runs, coalesce(w.created, 0) as week_created, coalesce(w.changed, 0) as week_changed,127 src.license as src_license, src.redistribution as src_redistribution128 from connectors c129 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 true130 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 true131 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 true132 left join lateral (select license, redistribution from sources where connector_id = c.id order by priority, id limit 1) src on true133 left join (134 select connector_id, count(*)::int as runs, coalesce(sum((stats->>'created')::numeric), 0) as created, coalesce(sum((stats->>'changed')::numeric), 0) as changed135 from connector_runs where started_at >= now() - interval '7 days' and status <> 'running' group by connector_id136 ) w on w.connector_id = c.id137 left join (138 select connector_id, jsonb_object_agg(key, total) as day_stats139 from (140 select cr.connector_id, kv.key, sum(kv.value) as total141 from connector_runs cr142 cross join lateral (select key, (value #>> '{}')::numeric as value from jsonb_each(cr.stats) where key = any(${DAY_KEYS}) and jsonb_typeof(value) = 'number') kv143 where cr.started_at >= now() - interval '24 hours'144 group by cr.connector_id, kv.key145 ) x group by connector_id146 ) s24 on s24.connector_id = c.id147 left join (148 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,149 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,150 count(*) filter (where extract_ok is not null)::int as extract_tried, count(*) filter (where quarantined)::int as quarantined151 from documents group by connector_id152 ) d on d.connector_id = c.id`;153154export async function listConnectorHealth(): Promise<ConnectorHealthDTO[]> {155 const sql = pg();156 const [rows, ch] = await Promise.all([sql<Row[]>`${CONNECTOR_SELECT(sql)} order by c.enabled desc, c.id`, crawlLogStats()]);157 return rows.map((r) => connectorHealth(r, ch));158}159160export interface ConnectorAdminDetail { health: ConnectorHealthDTO; config: Record<string, unknown>; paused: boolean; lastError: string | null; runs: Array<Record<string, unknown>>; 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 }161162export async function getConnectorAdmin(id: string): Promise<ConnectorAdminDetail | null> {163 const sql = pg();164 const [rows, ch] = await Promise.all([sql<Row[]>`${CONNECTOR_SELECT(sql)} where c.id = ${id}`, crawlLogStats()]);165 const r = rows[0];166 if (!r) return null;167 const [runs, byType, errors] = await Promise.all([168 sql<Row[]>`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`,169 sql<Row[]>`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`,170 sql<Row[]>`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`,171 ]);172 return {173 health: connectorHealth(r, ch),174 config: json<Record<string, unknown>>(r.config, {}),175 paused: bool(r.paused),176 lastError: str(r.last_error),177 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) })),178 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) })),179 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) })),180 createdAt: iso(r.created_at),181 updatedAt: iso(r.updated_at),182 };183}184185export async function listRuns(connectorId: string | undefined, limit = 100): Promise<Array<Record<string, unknown>>> {186 const sql = pg();187 const rows = await sql<Row[]>`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}`;188 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) }));189}190191export async function getRun(id: string): Promise<Record<string, unknown> | null> {192 const sql = pg();193 const rows = await sql<Row[]>`select * from connector_runs where id = ${id}`;194 const x = rows[0];195 if (!x) return null;196 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) };197}198