spb/datacenterindex
Public
HTML 53.9%
TypeScript 44.5%
JavaScript 0.6%
SQL 0.5%
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