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