SPB Git forge
38commits 1branches 0releases
338.7 MBsize
maindefault branch
3 h agolast push
HTML 53.9% TypeScript 44.5% JavaScript 0.6% SQL 0.5%
9.4 KB · 114 lines typescript
Raw Blame History
1/**2 * Maintenance jobs owned by the runtime: cleanup (old runs, resolved alerts, expired admin sessions) and the3 * `doctor` diagnostics (infrastructure probes, budgets, stale connectors, error rates → system_alerts).4 */5import { getDb, closeDb, sql } from "@dci/db";6import { chQuery } from "@dci/db/clickhouse";7import { budgetSnapshot, creditsToday } from "./context.js";8import { getEnv } from "./env.js";9import { getRedis } from "./scheduler.js";10import { storageProbe } from "./storage.js";11import { loadAllConnectors } from "./configs.js";12import { registerAllConnectors } from "./connectors/index.js";1314export async function cleanupMaintenance(): Promise<{ runsDeleted: number; alertsDeleted: number; sessionsDeleted: number; orphanRunsAborted: number }> {15  const db = getDb();16  const runs = await db.execute<{ n: number }>(sql`with d as (delete from connector_runs where started_at < now() - interval '90 days' returning 1) select count(*)::int as n from d`);17  const alerts = await db.execute<{ n: number }>(sql`with d as (delete from system_alerts where resolved_at is not null and resolved_at < now() - interval '30 days' returning 1) select count(*)::int as n from d`);18  const sessions = await db.execute<{ n: number }>(sql`with d as (delete from admin_sessions where expires_at < now() returning 1) select count(*)::int as n from d`);19  const orphans = await reconcileOrphanRuns();20  return { runsDeleted: Number(runs[0]?.n ?? 0), alertsDeleted: Number(alerts[0]?.n ?? 0), sessionsDeleted: Number(sessions[0]?.n ?? 0), orphanRunsAborted: orphans.length };21}2223/**24 * Runs left in `running` by a crashed / hard-killed process never finish on their own. Mark those older than25 * `maxAgeHours` as aborted so health, next_run_at and the admin UI stay truthful. Called at worker start,26 * by `doctor` and by the cleanup job.27 */28export async function reconcileOrphanRuns(maxAgeHours = 6): Promise<string[]> {29  const rows = await getDb().execute<{ id: string; connector_id: string }>(sql`30    update connector_runs set status = 'aborted', finished_at = now(), error = coalesce(error, 'orphaned: worker stopped without finishing this run')31    where status = 'running' and started_at < now() - (${maxAgeHours} || ' hours')::interval32    returning id, connector_id`);33  for (const r of rows) await getDb().execute(sql`update connectors set last_status = 'aborted', health = case when health = 'never_run' then 'degraded' else health end, updated_at = now() where id = ${r.connector_id} and (last_run_at is null or last_run_at < now() - interval '1 hour')`);34  return rows.map((r) => `${r.id} (${r.connector_id})`);35}3637export interface Check { name: string; ok: boolean; level: "ok" | "warn" | "error"; message: string; ms?: number }3839async function timed(name: string, fn: () => Promise<string>, level: Check["level"] = "error"): Promise<Check> {40  const t0 = Date.now();41  try { const message = await fn(); return { name, ok: true, level: "ok", message, ms: Date.now() - t0 }; } catch (e) { return { name, ok: false, level, message: (e as Error).message, ms: Date.now() - t0 }; }42}4344export interface DoctorReport { checks: Check[]; ok: boolean; alertsWritten: number }4546/** Probe infrastructure and data health. Writes `system_alerts` rows for warn/error findings (unless dryRun). */47export async function doctor(opts: { dryRun?: boolean } = {}): Promise<DoctorReport> {48  const env = getEnv();49  const checks: Check[] = [];50  checks.push(await timed("postgres", async () => { const r = await getDb().execute<{ v: string }>(sql`select version() as v`); return `${env.databaseUrl.replace(/\/\/.*@/, "//…@")} · ${String(r[0]?.v ?? "").split(" on ")[0]}`; }));51  checks.push(await timed("redis", async () => { const pong = await getRedis().ping(); const info = await getRedis().info("server"); const ver = /redis_version:(\S+)/.exec(info)?.[1] ?? "?"; return `${env.redisUrl} · ${pong} · redis ${ver}`; }));52  checks.push(await timed("clickhouse", async () => { const r = await chQuery<{ v: string }>("select version() as v"); return `${env.clickhouseUrl}/${env.clickhouseDb} · ${r[0]?.v ?? "?"}`; }, "warn"));53  checks.push(await timed("minio", async () => { const p = await storageProbe(); if (!p.ok) throw new Error(p.message); return p.message; }));5455  // budgets56  const b = budgetSnapshot();57  const spentToday = await creditsToday().catch(() => -1);58  const sf = process.env.SCRAPFLY_API_KEY ? `scrapfly ${b.scrapfly.used}/${b.scrapfly.limit}` : "scrapfly: no key";59  const fc = process.env.FIRECRAWL_API_KEY ? `firecrawl ${b.firecrawl.used}/${b.firecrawl.limit}` : "firecrawl: no key";60  const budgetWarn = spentToday >= 0 && spentToday > 0.8 * (b.scrapfly.limit + b.firecrawl.limit);61  checks.push({ name: "budgets", ok: !budgetWarn, level: budgetWarn ? "warn" : "ok", message: `${sf} · ${fc} · credits recorded today: ${spentToday < 0 ? "?" : spentToday}` });6263  // configs (code-backed connectors need their groups registered before they can be built)64  const reg = await registerAllConnectors();65  if (reg.failed.length) checks.push({ name: "connector-groups", ok: false, level: "error", message: reg.failed.map((f) => `${f.group}: ${f.error}`).join("; ") });66  try { const loaded = loadAllConnectors({ reload: true }); const enabled = loaded.filter((l) => l.cfg.enabled).length; checks.push({ name: "configs", ok: true, level: enabled ? "ok" : "warn", message: `${loaded.length} YAML connectors (${enabled} enabled) · groups ${reg.registered.join(", ") || "none"} · ${env.configDir}` }); } catch (e) { checks.push({ name: "configs", ok: false, level: "error", message: (e as Error).message }); }6768  // stale / failing connectors69  if (checks[0]?.ok) {70    const rows = await getDb().execute<{ id: string; health: string; last_run_at: string | null; next_run_at: string | null; enabled: boolean; paused: boolean; last_error: string | null; stats: Record<string, number> | null }>(sql`select id, health, last_run_at, next_run_at, enabled, paused, last_error, stats from connectors order by id`);71    const stale: string[] = [];72    const failing: string[] = [];73    for (const r of rows) {74      if (!r.enabled || r.paused) continue;75      if (r.health === "failing") failing.push(`${r.id}${r.last_error ? ` (${r.last_error.slice(0, 80)})` : ""}`);76      if (r.next_run_at && Date.parse(r.next_run_at) < Date.now() - 24 * 3_600_000) stale.push(`${r.id} (due ${r.next_run_at.slice(0, 16)})`);77      if (!r.last_run_at && r.next_run_at == null) stale.push(`${r.id} (never run)`);78    }79    checks.push({ name: "connectors:stale", ok: stale.length === 0, level: stale.length ? "warn" : "ok", message: stale.length ? `${stale.length} overdue: ${stale.slice(0, 8).join(", ")}` : `${rows.length} connectors, none overdue` });80    checks.push({ name: "connectors:failing", ok: failing.length === 0, level: failing.length ? "error" : "ok", message: failing.length ? failing.slice(0, 8).join("; ") : "no failing connector" });8182    const err = await getDb().execute<{ fetched: number; failed: number; runs: number }>(sql`select coalesce(sum((stats->>'fetched')::int),0)::int as fetched, coalesce(sum((stats->>'failed')::int),0)::int as failed, count(*)::int as runs from connector_runs where started_at > now() - interval '24 hours'`);83    const fetched = Number(err[0]?.fetched ?? 0), failed = Number(err[0]?.failed ?? 0), runs = Number(err[0]?.runs ?? 0);84    const rate = fetched ? failed / fetched : 0;85    checks.push({ name: "error-rate:24h", ok: rate < 0.2, level: rate >= 0.5 ? "error" : rate >= 0.2 ? "warn" : "ok", message: `${runs} runs · ${fetched} fetches · ${failed} failed (${(rate * 100).toFixed(1)} %)` });8687    const orphans = await reconcileOrphanRuns();88    const stillRunning = await getDb().execute<{ n: number }>(sql`select count(*)::int as n from connector_runs where status = 'running'`);89    checks.push({ name: "runs", ok: true, level: orphans.length ? "warn" : "ok", message: `${Number(stillRunning[0]?.n ?? 0)} running${orphans.length ? ` · aborted ${orphans.length} orphaned run(s): ${orphans.slice(0, 5).join(", ")}` : ""}` });90    const q = await getDb().execute<{ n: number }>(sql`select count(*)::int as n from documents where quarantined`);91    const pend = await getDb().execute<{ n: number }>(sql`select count(*)::int as n from entity_matches where status = 'pending'`);92    checks.push({ name: "documents", ok: true, level: "ok", message: `${Number(q[0]?.n ?? 0)} quarantined · ${Number(pend[0]?.n ?? 0)} pending entity matches` });93  }9495  let alertsWritten = 0;96  if (!opts.dryRun && checks[0]?.ok) {97    for (const c of checks) {98      if (c.level === "ok") continue;99      try {100        await getDb().execute(sql`insert into system_alerts (id, level, component, message, details) values (${`alr_${Date.now().toString(36)}_${c.name.replace(/[^a-z0-9]/gi, "")}`}, ${c.level}, ${`doctor:${c.name}`}, ${c.message.slice(0, 1000)}, ${JSON.stringify({ ms: c.ms ?? null })}::jsonb)`);101        alertsWritten++;102      } catch { /* alerts are best effort */ }103    }104    // auto-resolve older doctor alerts for components that are OK now105    const okNames = checks.filter((c) => c.level === "ok").map((c) => `doctor:${c.name}`);106    if (okNames.length) await getDb().execute(sql`update system_alerts set resolved_at = now() where resolved_at is null and component in (${sql.join(okNames.map((n) => sql`${n}`), sql`, `)})`).catch(() => undefined);107  }108  return { checks, ok: checks.every((c) => c.level !== "error"), alertsWritten };109}110111export async function closeAll(): Promise<void> {112  await closeDb().catch(() => undefined);113}114