/** * Maintenance jobs owned by the runtime: cleanup (old runs, resolved alerts, expired admin sessions) and the * `doctor` diagnostics (infrastructure probes, budgets, stale connectors, error rates → system_alerts). */ import { getDb, closeDb, sql } from "@dci/db"; import { chQuery } from "@dci/db/clickhouse"; import { budgetSnapshot, creditsToday } from "./context.js"; import { getEnv } from "./env.js"; import { getRedis } from "./scheduler.js"; import { storageProbe } from "./storage.js"; import { loadAllConnectors } from "./configs.js"; import { registerAllConnectors } from "./connectors/index.js"; export async function cleanupMaintenance(): Promise<{ runsDeleted: number; alertsDeleted: number; sessionsDeleted: number; orphanRunsAborted: number }> { const db = getDb(); 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`); 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`); 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`); const orphans = await reconcileOrphanRuns(); return { runsDeleted: Number(runs[0]?.n ?? 0), alertsDeleted: Number(alerts[0]?.n ?? 0), sessionsDeleted: Number(sessions[0]?.n ?? 0), orphanRunsAborted: orphans.length }; } /** * Runs left in `running` by a crashed / hard-killed process never finish on their own. Mark those older than * `maxAgeHours` as aborted so health, next_run_at and the admin UI stay truthful. Called at worker start, * by `doctor` and by the cleanup job. */ export async function reconcileOrphanRuns(maxAgeHours = 6): Promise { const rows = await getDb().execute<{ id: string; connector_id: string }>(sql` update connector_runs set status = 'aborted', finished_at = now(), error = coalesce(error, 'orphaned: worker stopped without finishing this run') where status = 'running' and started_at < now() - (${maxAgeHours} || ' hours')::interval returning id, connector_id`); 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')`); return rows.map((r) => `${r.id} (${r.connector_id})`); } export interface Check { name: string; ok: boolean; level: "ok" | "warn" | "error"; message: string; ms?: number } async function timed(name: string, fn: () => Promise, level: Check["level"] = "error"): Promise { const t0 = Date.now(); 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 }; } } export interface DoctorReport { checks: Check[]; ok: boolean; alertsWritten: number } /** Probe infrastructure and data health. Writes `system_alerts` rows for warn/error findings (unless dryRun). */ export async function doctor(opts: { dryRun?: boolean } = {}): Promise { const env = getEnv(); const checks: Check[] = []; 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]}`; })); 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}`; })); 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")); checks.push(await timed("minio", async () => { const p = await storageProbe(); if (!p.ok) throw new Error(p.message); return p.message; })); // budgets const b = budgetSnapshot(); const spentToday = await creditsToday().catch(() => -1); const sf = process.env.SCRAPFLY_API_KEY ? `scrapfly ${b.scrapfly.used}/${b.scrapfly.limit}` : "scrapfly: no key"; const fc = process.env.FIRECRAWL_API_KEY ? `firecrawl ${b.firecrawl.used}/${b.firecrawl.limit}` : "firecrawl: no key"; const budgetWarn = spentToday >= 0 && spentToday > 0.8 * (b.scrapfly.limit + b.firecrawl.limit); checks.push({ name: "budgets", ok: !budgetWarn, level: budgetWarn ? "warn" : "ok", message: `${sf} · ${fc} · credits recorded today: ${spentToday < 0 ? "?" : spentToday}` }); // configs (code-backed connectors need their groups registered before they can be built) const reg = await registerAllConnectors(); if (reg.failed.length) checks.push({ name: "connector-groups", ok: false, level: "error", message: reg.failed.map((f) => `${f.group}: ${f.error}`).join("; ") }); 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 }); } // stale / failing connectors if (checks[0]?.ok) { 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 | null }>(sql`select id, health, last_run_at, next_run_at, enabled, paused, last_error, stats from connectors order by id`); const stale: string[] = []; const failing: string[] = []; for (const r of rows) { if (!r.enabled || r.paused) continue; if (r.health === "failing") failing.push(`${r.id}${r.last_error ? ` (${r.last_error.slice(0, 80)})` : ""}`); 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)})`); if (!r.last_run_at && r.next_run_at == null) stale.push(`${r.id} (never run)`); } 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` }); 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" }); 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'`); const fetched = Number(err[0]?.fetched ?? 0), failed = Number(err[0]?.failed ?? 0), runs = Number(err[0]?.runs ?? 0); const rate = fetched ? failed / fetched : 0; 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)} %)` }); const orphans = await reconcileOrphanRuns(); const stillRunning = await getDb().execute<{ n: number }>(sql`select count(*)::int as n from connector_runs where status = 'running'`); 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(", ")}` : ""}` }); const q = await getDb().execute<{ n: number }>(sql`select count(*)::int as n from documents where quarantined`); const pend = await getDb().execute<{ n: number }>(sql`select count(*)::int as n from entity_matches where status = 'pending'`); checks.push({ name: "documents", ok: true, level: "ok", message: `${Number(q[0]?.n ?? 0)} quarantined · ${Number(pend[0]?.n ?? 0)} pending entity matches` }); } let alertsWritten = 0; if (!opts.dryRun && checks[0]?.ok) { for (const c of checks) { if (c.level === "ok") continue; try { 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)`); alertsWritten++; } catch { /* alerts are best effort */ } } // auto-resolve older doctor alerts for components that are OK now const okNames = checks.filter((c) => c.level === "ok").map((c) => `doctor:${c.name}`); 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); } return { checks, ok: checks.every((c) => c.level !== "error"), alertsWritten }; } export async function closeAll(): Promise { await closeDb().catch(() => undefined); }