SPB Git forge
38commits 1branches 0releases
338.7 MBsize
maindefault branch
5 h agolast push
HTML 53.9% TypeScript 44.5% JavaScript 0.6% SQL 0.5%
6.4 KB · 110 lines typescript
Raw Blame History
1import type { FastifyInstance } from "fastify";2import { z } from "zod";3import { chQuery } from "@dci/db/clickhouse";4import { getEnv } from "../../env.js";5import { envelope, parseBody } from "../../lib/http.js";6import { pg } from "../../lib/sql.js";7import { int, iso, json, num, reqStr, str, type Row } from "../../lib/rows.js";8import { getQueueRedis } from "../../redis.js";9import { enqueueMaintenance, queueCounts } from "../../queues.js";10import { storageOk } from "../../storage.js";11import { listRuns } from "../../repositories/admin/connectors.js";12import { cacheStats } from "../../cache.js";1314async function withTimeout<T>(p: Promise<T>, ms: number, fallback: T): Promise<T> {15  return Promise.race([p, new Promise<T>((res) => setTimeout(() => res(fallback), ms))]).catch(() => fallback);16}1718async function workerStatus(): Promise<unknown> {19  try {20    const r = getQueueRedis();21    if (r.status === "wait") await r.connect().catch(() => {});22    const raw = await withTimeout(r.get("dci:worker:status"), 1500, null);23    if (raw) { try { return JSON.parse(raw); } catch { return raw; } }24    const hash = await withTimeout(r.hgetall("dci:worker:status"), 1500, {} as Record<string, string>);25    if (hash && Object.keys(hash).length) return hash;26    const keys = await withTimeout(r.keys("dci:worker:*"), 1500, [] as string[]);27    if (keys.length) {28      const out: Record<string, unknown> = {};29      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; }30      return out;31    }32    return null;33  } catch {34    return null;35  }36}3738async function budgets(): Promise<Record<string, unknown>> {39  const env = getEnv();40  const limits = { scrapfly: Number(process.env.DCI_SCRAPFLY_DAILY_BUDGET ?? 400), firecrawl: Number(process.env.DCI_FIRECRAWL_DAILY_BUDGET ?? 200) };41  // 1) keys published by the worker42  try {43    const r = getQueueRedis();44    const day = new Date().toISOString().slice(0, 10);45    const keys = await withTimeout(r.keys("dci:budget:*"), 1500, [] as string[]);46    if (keys.length) {47      const out: Record<string, unknown> = { source: "redis", day, limits };48      for (const k of keys.slice(0, 20)) out[k] = await withTimeout(r.get(k), 1000, null);49      return out;50    }51  } catch { /* fall through */ }52  // 2) ClickHouse crawl_log credits today53  try {54    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 }>);55    const used: Record<string, { credits: number; requests: number }> = {};56    for (const r of rows) used[r.fetcher] = { credits: num(r.credits) ?? 0, requests: int(r.n) };57    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(/:\/\/.*@/, "://***@") };58  } catch {59    return { source: "none", limits };60  }61}6263export async function opsAdminRoutes(app: FastifyInstance): Promise<void> {64  app.get("/ops", { schema: { summary: "Operations overview: worker, queues, budgets, DB sizes, ClickHouse/MinIO health, alerts, recent runs" } }, async () => {65    const sql = pg();66    const [worker, queues, budget, dbSize, tables, chOk, minioOk, alerts, runs, redisOk] = await Promise.all([67      workerStatus(),68      withTimeout(queueCounts(), 4000, {} as Awaited<ReturnType<typeof queueCounts>>),69      budgets(),70      sql<Row[]>`select pg_database_size(current_database()) as bytes, pg_size_pretty(pg_database_size(current_database())) as pretty`,71      sql<Row[]>`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`,72      withTimeout(chQuery("select 1 as ok").then(() => true), 3000, false),73      storageOk(),74      sql<Row[]>`select * from system_alerts where resolved_at is null order by created_at desc limit 50`,75      listRuns(undefined, 20),76      withTimeout(getQueueRedis().ping().then((p) => p === "PONG"), 1500, false),77    ]);78    return envelope({79      worker,80      queues,81      budgets: budget,82      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) })) },83      clickhouse: { ok: chOk, url: getEnv().clickhouseUrl },84      minio: { ok: minioOk, endpoint: getEnv().s3Endpoint, bucket: getEnv().s3Bucket },85      redis: { ok: redisOk },86      apiCache: cacheStats(),87      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) })),88      lastRuns: runs,89    });90  });9192  // Note for deploy: the compose healthcheck of the api service should probe GET /api/ready (Postgres ping), not /api/health.93  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();94  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) => {95    const { task } = req.params as { task: string };96    const t = z.enum(["rankings", "metrics", "refresh-stats", "quality", "snapshot", "cleanup"]).safeParse(task);97    if (!t.success) return envelope({ enqueued: false, error: "unknown task (rankings | metrics | refresh-stats | quality | snapshot | cleanup)" });98    const body = parseBody(maintenanceBody, req.body ?? {});99    const job = await enqueueMaintenance(t.data, body);100    return envelope({ enqueued: true, job });101  });102103  app.post("/alerts/:id/resolve", { schema: { summary: "Mark a system alert resolved" } }, async (req) => {104    const { id } = req.params as { id: string };105    const sql = pg();106    const rows = await sql`update system_alerts set resolved_at = now() where id = ${id} and resolved_at is null returning id`;107    return envelope({ id, resolved: rows.length > 0 });108  });109}110