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%
7.7 KB · 138 lines typescript
Raw Blame History
1import { existsSync, readFileSync, writeFileSync, copyFileSync } from "node:fs";2import { join } from "node:path";3import type { FastifyInstance } from "fastify";4import { z } from "zod";5import { parseDocument } from "yaml";6import { intervalMs } from "@dci/connectors";7import { getEnv } from "../../env.js";8import { envelope, notFound, parseBody, parseQuery, badRequest, HttpError } from "../../lib/http.js";9import { strParam, intParam } from "../../lib/params.js";10import { pg } from "../../lib/sql.js";11import { enqueueCrawl } from "../../queues.js";12import { getConnectorAdmin, getRun, listConnectorHealth, listRuns } from "../../repositories/admin/connectors.js";13import { rollbackRun, runChanges } from "../../repositories/admin/runs.js";14import { invalidate } from "../../cache.js";1516const runBody = z.object({ task: z.enum(["crawl", "discover", "full", "reprocess"]).default("crawl"), group: z.string().min(1).max(64).optional(), limit: z.number().int().min(1).max(100000).optional(), force: z.boolean().optional() });17const scheduleBody = z.record(z.string().regex(/^[a-z0-9_*-]+$/i), z.string().min(1).max(32));1819export function safeConnectorId(id: string): string {20  if (!/^[a-z0-9][a-z0-9-]*$/.test(id)) throw badRequest("invalid connector id");21  return id;22}2324export function configPath(id: string): string {25  return join(getEnv().configDir, `${safeConnectorId(id)}.yaml`);26}2728export function backupConfig(path: string): string | null {29  if (!existsSync(path)) return null;30  const bak = `${path}.bak-${new Date().toISOString().replace(/[-:.TZ]/g, "").slice(0, 14)}`;31  copyFileSync(path, bak);32  return bak;33}3435export async function connectorAdminRoutes(app: FastifyInstance): Promise<void> {36  app.get("/connectors", { schema: { summary: "Connector health (ConnectorHealthDTO[])" } }, async () => envelope(await listConnectorHealth()));3738  app.get("/connectors/:id", { schema: { summary: "Connector detail: config, last 30 runs, documents by page type, error samples" } }, async (req) => {39    const { id } = req.params as { id: string };40    const d = await getConnectorAdmin(id);41    if (!d) throw notFound("connector");42    return envelope(d);43  });4445  app.post("/connectors/:id/run", { schema: { summary: "Enqueue a crawl job on dci:crawl (jobId manual__<id>__<ts>; task crawl|discover|full|reprocess)" } }, async (req) => {46    const { id } = req.params as { id: string };47    const body = parseBody(runBody, req.body);48    const sql = pg();49    const exists = await sql`select id from connectors where id = ${id}`;50    if (!exists.length) throw notFound("connector");51    const job = await enqueueCrawl({ connectorId: id, task: body.task, group: body.group, limit: body.limit, force: body.force, requestedBy: "admin" });52    return envelope({ enqueued: true, job });53  });5455  for (const action of ["pause", "resume"] as const) {56    app.post(`/connectors/:id/${action}`, { schema: { summary: `${action} a connector (connectors.paused)` } }, async (req) => {57      const { id } = req.params as { id: string };58      const sql = pg();59      const rows = await sql`update connectors set paused = ${action === "pause"}, health = case when ${action === "pause"} then 'paused' else case when last_run_at is null then 'never_run' else 'ok' end end, updated_at = now() where id = ${id} returning id, paused, health`;60      if (!rows.length) throw notFound("connector");61      return envelope(rows[0]);62    });63  }6465  app.patch("/connectors/:id/schedule", { schema: { summary: "Update schedule {group: interval}: DB row + YAML file (backup kept)" } }, async (req) => {66    const { id } = req.params as { id: string };67    const body = parseBody(scheduleBody, req.body);68    for (const [g, spec] of Object.entries(body)) { try { intervalMs(spec); } catch { throw badRequest(`bad interval for group ${g}: ${spec}`); } }69    const sql = pg();70    const rows = await sql`update connectors set schedule = schedule || ${JSON.stringify(body)}::jsonb, updated_at = now() where id = ${id} returning id, schedule`;71    if (!rows.length) throw notFound("connector");72    let file: { path: string; updated: boolean; backup: string | null; error?: string } = { path: configPath(id), updated: false, backup: null };73    try {74      if (existsSync(file.path)) {75        const doc = parseDocument(readFileSync(file.path, "utf8"));76        const backup = backupConfig(file.path);77        for (const [g, spec] of Object.entries(body)) doc.setIn(["schedule", g], spec);78        writeFileSync(file.path, doc.toString(), "utf8");79        file = { ...file, updated: true, backup };80      } else file.error = "yaml file not found";81    } catch (e) {82      file.error = (e as Error).message;83    }84    return envelope({ id, schedule: rows[0]!.schedule, file });85  });8687  app.get("/runs", { schema: { summary: "Recent connector runs (?connector=&limit=)" } }, async (req) => {88    const q = parseQuery(z.object({ connector: strParam, limit: intParam }), req.query);89    const items = await listRuns(q.connector, Math.min(500, q.limit ?? 100));90    return envelope(items, { total: items.length });91  });9293  app.post("/connectors/:id/quarantine", { schema: { summary: "Quarantine {on: boolean}: the connector extracts and previews but publishes nothing; health = quarantine while on" } }, async (req) => {94    const { id } = req.params as { id: string };95    const body = parseBody(z.object({ on: z.boolean() }), req.body);96    const sql = pg();97    const rows = await sql`update connectors set quarantine = ${body.on}, health = case when ${body.on} then 'quarantine' when paused then 'paused' when last_run_at is null then 'never_run' else 'ok' end, updated_at = now() where id = ${id} returning id, quarantine, health`;98    if (!rows.length) throw notFound("connector");99    return envelope(rows[0]);100  });101102  app.get("/runs/:id", { schema: { summary: "Run detail with log" } }, async (req) => {103    const { id } = req.params as { id: string };104    const r = await getRun(id);105    if (!r) throw notFound("run");106    return envelope(r);107  });108109  app.get("/runs/:id/changes", { schema: { summary: "Everything a run wrote: provenance rows, claims, events, document versions (run_id = id; 2 000 rows per list)" } }, async (req) => {110    const { id } = req.params as { id: string };111    const r = await runChanges(id);112    if (!r) throw notFound("run");113    return envelope(r, { ...r.counts });114  });115116  app.post("/runs/:id/rollback", { schema: { summary: "Roll a run back: its claims → rejected (rollback), its provenance rows → not current (winner restored to the latest remaining row per field), its events → rejected. Entity columns are re-derived by the worker's next reconciliation" } }, async (req) => {117    const { id } = req.params as { id: string };118    const r = await rollbackRun(id);119    if (!r) throw notFound("run");120    await Promise.all([invalidate("/datacenters"), invalidate("/projects"), invalidate("/events"), invalidate("/dashboard")]);121    return envelope(r);122  });123124  app.post("/cache/invalidate", { schema: { summary: "Invalidate API cache entries by route prefix (body {prefix}) — default all" } }, async (req) => {125    const body = parseBody(z.object({ prefix: z.string().max(200).default("") }), req.body);126    const n = await invalidate(body.prefix);127    return envelope({ invalidated: n, prefix: body.prefix });128  });129130  app.setErrorHandler((err, _req, reply) => {131    if (err instanceof HttpError) return reply.code(err.statusCode).send({ error: err.message, statusCode: err.statusCode, details: err.details });132    const e = err as { statusCode?: number; message?: string };133    const status = e.statusCode ?? 500;134    if (status >= 500) app.log.error({ err }, "admin request failed");135    return reply.code(status).send({ error: status >= 500 ? `internal error: ${e.message ?? ""}` : e.message ?? "error", statusCode: status });136  });137}138