import { existsSync, readFileSync, writeFileSync, copyFileSync } from "node:fs"; import { join } from "node:path"; import type { FastifyInstance } from "fastify"; import { z } from "zod"; import { parseDocument } from "yaml"; import { intervalMs } from "@dci/connectors"; import { getEnv } from "../../env.js"; import { envelope, notFound, parseBody, parseQuery, badRequest, HttpError } from "../../lib/http.js"; import { strParam, intParam } from "../../lib/params.js"; import { pg } from "../../lib/sql.js"; import { enqueueCrawl } from "../../queues.js"; import { getConnectorAdmin, getRun, listConnectorHealth, listRuns } from "../../repositories/admin/connectors.js"; import { rollbackRun, runChanges } from "../../repositories/admin/runs.js"; import { invalidate } from "../../cache.js"; const 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() }); const scheduleBody = z.record(z.string().regex(/^[a-z0-9_*-]+$/i), z.string().min(1).max(32)); export function safeConnectorId(id: string): string { if (!/^[a-z0-9][a-z0-9-]*$/.test(id)) throw badRequest("invalid connector id"); return id; } export function configPath(id: string): string { return join(getEnv().configDir, `${safeConnectorId(id)}.yaml`); } export function backupConfig(path: string): string | null { if (!existsSync(path)) return null; const bak = `${path}.bak-${new Date().toISOString().replace(/[-:.TZ]/g, "").slice(0, 14)}`; copyFileSync(path, bak); return bak; } export async function connectorAdminRoutes(app: FastifyInstance): Promise { app.get("/connectors", { schema: { summary: "Connector health (ConnectorHealthDTO[])" } }, async () => envelope(await listConnectorHealth())); app.get("/connectors/:id", { schema: { summary: "Connector detail: config, last 30 runs, documents by page type, error samples" } }, async (req) => { const { id } = req.params as { id: string }; const d = await getConnectorAdmin(id); if (!d) throw notFound("connector"); return envelope(d); }); app.post("/connectors/:id/run", { schema: { summary: "Enqueue a crawl job on dci:crawl (jobId manual____; task crawl|discover|full|reprocess)" } }, async (req) => { const { id } = req.params as { id: string }; const body = parseBody(runBody, req.body); const sql = pg(); const exists = await sql`select id from connectors where id = ${id}`; if (!exists.length) throw notFound("connector"); const job = await enqueueCrawl({ connectorId: id, task: body.task, group: body.group, limit: body.limit, force: body.force, requestedBy: "admin" }); return envelope({ enqueued: true, job }); }); for (const action of ["pause", "resume"] as const) { app.post(`/connectors/:id/${action}`, { schema: { summary: `${action} a connector (connectors.paused)` } }, async (req) => { const { id } = req.params as { id: string }; const sql = pg(); 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`; if (!rows.length) throw notFound("connector"); return envelope(rows[0]); }); } app.patch("/connectors/:id/schedule", { schema: { summary: "Update schedule {group: interval}: DB row + YAML file (backup kept)" } }, async (req) => { const { id } = req.params as { id: string }; const body = parseBody(scheduleBody, req.body); for (const [g, spec] of Object.entries(body)) { try { intervalMs(spec); } catch { throw badRequest(`bad interval for group ${g}: ${spec}`); } } const sql = pg(); const rows = await sql`update connectors set schedule = schedule || ${JSON.stringify(body)}::jsonb, updated_at = now() where id = ${id} returning id, schedule`; if (!rows.length) throw notFound("connector"); let file: { path: string; updated: boolean; backup: string | null; error?: string } = { path: configPath(id), updated: false, backup: null }; try { if (existsSync(file.path)) { const doc = parseDocument(readFileSync(file.path, "utf8")); const backup = backupConfig(file.path); for (const [g, spec] of Object.entries(body)) doc.setIn(["schedule", g], spec); writeFileSync(file.path, doc.toString(), "utf8"); file = { ...file, updated: true, backup }; } else file.error = "yaml file not found"; } catch (e) { file.error = (e as Error).message; } return envelope({ id, schedule: rows[0]!.schedule, file }); }); app.get("/runs", { schema: { summary: "Recent connector runs (?connector=&limit=)" } }, async (req) => { const q = parseQuery(z.object({ connector: strParam, limit: intParam }), req.query); const items = await listRuns(q.connector, Math.min(500, q.limit ?? 100)); return envelope(items, { total: items.length }); }); app.post("/connectors/:id/quarantine", { schema: { summary: "Quarantine {on: boolean}: the connector extracts and previews but publishes nothing; health = quarantine while on" } }, async (req) => { const { id } = req.params as { id: string }; const body = parseBody(z.object({ on: z.boolean() }), req.body); const sql = pg(); 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`; if (!rows.length) throw notFound("connector"); return envelope(rows[0]); }); app.get("/runs/:id", { schema: { summary: "Run detail with log" } }, async (req) => { const { id } = req.params as { id: string }; const r = await getRun(id); if (!r) throw notFound("run"); return envelope(r); }); 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) => { const { id } = req.params as { id: string }; const r = await runChanges(id); if (!r) throw notFound("run"); return envelope(r, { ...r.counts }); }); 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) => { const { id } = req.params as { id: string }; const r = await rollbackRun(id); if (!r) throw notFound("run"); await Promise.all([invalidate("/datacenters"), invalidate("/projects"), invalidate("/events"), invalidate("/dashboard")]); return envelope(r); }); app.post("/cache/invalidate", { schema: { summary: "Invalidate API cache entries by route prefix (body {prefix}) — default all" } }, async (req) => { const body = parseBody(z.object({ prefix: z.string().max(200).default("") }), req.body); const n = await invalidate(body.prefix); return envelope({ invalidated: n, prefix: body.prefix }); }); app.setErrorHandler((err, _req, reply) => { if (err instanceof HttpError) return reply.code(err.statusCode).send({ error: err.message, statusCode: err.statusCode, details: err.details }); const e = err as { statusCode?: number; message?: string }; const status = e.statusCode ?? 500; if (status >= 500) app.log.error({ err }, "admin request failed"); return reply.code(status).send({ error: status >= 500 ? `internal error: ${e.message ?? ""}` : e.message ?? "error", statusCode: status }); }); }