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