import type { FastifyInstance } from "fastify"; import { z } from "zod"; import { chQuery } from "@dci/db/clickhouse"; import { getEnv } from "../../env.js"; import { envelope, parseBody } from "../../lib/http.js"; import { pg } from "../../lib/sql.js"; import { int, iso, json, num, reqStr, str, type Row } from "../../lib/rows.js"; import { getQueueRedis } from "../../redis.js"; import { enqueueMaintenance, queueCounts } from "../../queues.js"; import { storageOk } from "../../storage.js"; import { listRuns } from "../../repositories/admin/connectors.js"; import { cacheStats } from "../../cache.js"; async function withTimeout(p: Promise, ms: number, fallback: T): Promise { return Promise.race([p, new Promise((res) => setTimeout(() => res(fallback), ms))]).catch(() => fallback); } async function workerStatus(): Promise { try { const r = getQueueRedis(); if (r.status === "wait") await r.connect().catch(() => {}); const raw = await withTimeout(r.get("dci:worker:status"), 1500, null); if (raw) { try { return JSON.parse(raw); } catch { return raw; } } const hash = await withTimeout(r.hgetall("dci:worker:status"), 1500, {} as Record); if (hash && Object.keys(hash).length) return hash; const keys = await withTimeout(r.keys("dci:worker:*"), 1500, [] as string[]); if (keys.length) { const out: Record = {}; for (const k of keys.slice(0, 20)) { const v = await withTimeout(r.get(k), 1000, null); out[k] = v ? (() => { try { return JSON.parse(v); } catch { return v; } })() : null; } return out; } return null; } catch { return null; } } async function budgets(): Promise> { const env = getEnv(); const limits = { scrapfly: Number(process.env.DCI_SCRAPFLY_DAILY_BUDGET ?? 400), firecrawl: Number(process.env.DCI_FIRECRAWL_DAILY_BUDGET ?? 200) }; // 1) keys published by the worker try { const r = getQueueRedis(); const day = new Date().toISOString().slice(0, 10); const keys = await withTimeout(r.keys("dci:budget:*"), 1500, [] as string[]); if (keys.length) { const out: Record = { source: "redis", day, limits }; for (const k of keys.slice(0, 20)) out[k] = await withTimeout(r.get(k), 1000, null); return out; } } catch { /* fall through */ } // 2) ClickHouse crawl_log credits today try { const rows = await withTimeout(chQuery<{ fetcher: string; credits: number | string; n: number | string }>("select fetcher, sum(credits) as credits, count() as n from crawl_log where ts >= toStartOfDay(now()) group by fetcher"), 3000, [] as Array<{ fetcher: string; credits: number | string; n: number | string }>); const used: Record = {}; for (const r of rows) used[r.fetcher] = { credits: num(r.credits) ?? 0, requests: int(r.n) }; return { source: "clickhouse", day: new Date().toISOString().slice(0, 10), limits, used, remaining: { scrapfly: Math.max(0, limits.scrapfly - (used.scrapfly?.credits ?? 0)), firecrawl: Math.max(0, limits.firecrawl - (used.firecrawl?.credits ?? 0)) }, redis: env.redisUrl.replace(/:\/\/.*@/, "://***@") }; } catch { return { source: "none", limits }; } } export async function opsAdminRoutes(app: FastifyInstance): Promise { app.get("/ops", { schema: { summary: "Operations overview: worker, queues, budgets, DB sizes, ClickHouse/MinIO health, alerts, recent runs" } }, async () => { const sql = pg(); const [worker, queues, budget, dbSize, tables, chOk, minioOk, alerts, runs, redisOk] = await Promise.all([ workerStatus(), withTimeout(queueCounts(), 4000, {} as Awaited>), budgets(), sql`select pg_database_size(current_database()) as bytes, pg_size_pretty(pg_database_size(current_database())) as pretty`, sql`select relname as table, n_live_tup as rows, pg_total_relation_size(relid) as bytes from pg_stat_user_tables order by pg_total_relation_size(relid) desc limit 40`, withTimeout(chQuery("select 1 as ok").then(() => true), 3000, false), storageOk(), sql`select * from system_alerts where resolved_at is null order by created_at desc limit 50`, listRuns(undefined, 20), withTimeout(getQueueRedis().ping().then((p) => p === "PONG"), 1500, false), ]); return envelope({ worker, queues, budgets: budget, db: { sizeBytes: num(dbSize[0]?.bytes), sizePretty: str(dbSize[0]?.pretty), tables: tables.map((t) => ({ table: reqStr(t.table), rows: int(t.rows), bytes: num(t.bytes) })) }, clickhouse: { ok: chOk, url: getEnv().clickhouseUrl }, minio: { ok: minioOk, endpoint: getEnv().s3Endpoint, bucket: getEnv().s3Bucket }, redis: { ok: redisOk }, apiCache: cacheStats(), alerts: alerts.map((a) => ({ id: reqStr(a.id), level: reqStr(a.level), component: reqStr(a.component), message: reqStr(a.message), details: json(a.details, null), createdAt: iso(a.created_at) })), lastRuns: runs, }); }); // Note for deploy: the compose healthcheck of the api service should probe GET /api/ready (Postgres ping), not /api/health. const maintenanceBody = z.object({ limit: z.number().int().min(1).max(100000).optional(), dryRun: z.boolean().optional(), connector: z.string().max(64).optional(), reason: z.string().max(500).optional() }).strict(); app.post("/maintenance/:task", { schema: { summary: "Enqueue a maintenance job (rankings | metrics | refresh-stats | quality | snapshot | cleanup) on dci:maintenance; body {limit?, dryRun?, connector?, reason?} (strict)" } }, async (req) => { const { task } = req.params as { task: string }; const t = z.enum(["rankings", "metrics", "refresh-stats", "quality", "snapshot", "cleanup"]).safeParse(task); if (!t.success) return envelope({ enqueued: false, error: "unknown task (rankings | metrics | refresh-stats | quality | snapshot | cleanup)" }); const body = parseBody(maintenanceBody, req.body ?? {}); const job = await enqueueMaintenance(t.data, body); return envelope({ enqueued: true, job }); }); app.post("/alerts/:id/resolve", { schema: { summary: "Mark a system alert resolved" } }, async (req) => { const { id } = req.params as { id: string }; const sql = pg(); const rows = await sql`update system_alerts set resolved_at = now() where id = ${id} and resolved_at is null returning id`; return envelope({ id, resolved: rows.length > 0 }); }); }