spb/datacenterindex
Public
HTML 53.9%
TypeScript 44.5%
JavaScript 0.6%
SQL 0.5%
1/**2 * Connector dev tool: fetch a page through the escalation stack, preview extractor rules, validate / save YAML,3 * run a config against several live pages. Nothing here persists documents or entities.4 */5import { existsSync, readdirSync, readFileSync, writeFileSync } from "node:fs";6import { join } from "node:path";7import { pathToFileURL } from "node:url";8import type { FastifyInstance } from "fastify";9import { z } from "zod";10import { load } from "cheerio";11import { parse as parseYaml, parseDocument } from "yaml";12import { htmlToText, type FetchLevel } from "@dci/core";13import { connectorConfigSchema, extractorSchema, parseConnectorConfig, fetchWithEscalation, GenericConnector, createTestContext, classifyPage, extractGeo, extractAddress, jsonLd, embeddedJson, pageTitle, metaDescription, publishedDate, mainText, getImplementation, type ConnectorConfig, type RawDocument, type Connector } from "@dci/connectors";14import { getEnv } from "../../env.js";15import { envelope, parseBody, badRequest, notFound } from "../../lib/http.js";16import { pg } from "../../lib/sql.js";17import { backupConfig, configPath, safeConnectorId } from "./connectors.js";1819const HTML_CAP = 1_500_000;20const TEXT_CAP = 200_000;2122let parsersLoaded: Promise<{ registered: string[]; failed: Array<{ group: string; error: string }> } | null> | null = null;23/** Best-effort: register the worker's code-backed parsers so YAML `parser:` references resolve in previews. */24function loadWorkerParsers(): Promise<{ registered: string[]; failed: Array<{ group: string; error: string }> } | null> {25 if (parsersLoaded) return parsersLoaded;26 parsersLoaded = (async () => {27 const dir = getEnv().workerConnectorsDir;28 if (!dir) return null;29 const idx = ["index.ts", "index.js"].map((f) => join(dir, f)).find((p) => existsSync(p));30 if (!idx) return null;31 try {32 const mod = (await import(pathToFileURL(idx).href)) as { registerAllConnectors?: () => Promise<{ registered: string[]; failed: Array<{ group: string; error: string }> }> };33 return mod.registerAllConnectors ? await mod.registerAllConnectors() : null;34 } catch (e) {35 return { registered: [], failed: [{ group: "worker", error: (e as Error).message }] };36 }37 })();38 return parsersLoaded;39}4041function summarizeDoc(doc: RawDocument, includeBody = true): Record<string, unknown> {42 const isHtml = /html|xml/i.test(doc.contentType ?? "") || /<html/i.test(doc.text.slice(0, 2000));43 const html = isHtml ? doc.text : "";44 const title = isHtml ? pageTitle(html) : null;45 const text = isHtml ? mainText(html) : doc.text;46 const cls = classifyPage(doc.finalUrl, title, text);47 let links: Array<{ href: string; text: string }> = [];48 if (isHtml) {49 const $ = load(html);50 const base = doc.finalUrl;51 $("a[href]").slice(0, 400).each((_, el) => {52 const href = $(el).attr("href");53 if (!href || href.startsWith("#") || href.startsWith("javascript:") || href.startsWith("mailto:")) return;54 try { links.push({ href: new URL(href, base).toString(), text: $(el).text().replace(/\s+/g, " ").trim().slice(0, 80) }); } catch { /* skip */ }55 });56 const seen = new Set<string>();57 links = links.filter((l) => (seen.has(l.href) ? false : (seen.add(l.href), true))).slice(0, 50);58 }59 return {60 url: doc.url,61 finalUrl: doc.finalUrl,62 status: doc.status,63 contentType: doc.contentType,64 fetcher: doc.fetcher,65 level: doc.level,66 credits: doc.credits,67 durationMs: doc.durationMs,68 bytes: doc.body.length,69 error: doc.error ?? null,70 attempts: doc.meta?.attempts ?? null,71 title,72 description: isHtml ? metaDescription(html) : null,73 published: isHtml ? publishedDate(html) : null,74 classification: cls,75 geo: isHtml ? extractGeo(html) : null,76 address: isHtml ? extractAddress(html) : null,77 jsonLd: isHtml ? jsonLd(html).slice(0, 20) : [],78 embeddedJson: isHtml ? embeddedJson(html).map((b) => ({ id: b.id, preview: JSON.stringify(b.data).slice(0, 2048) })) : [],79 links,80 ...(includeBody ? { html: html.length > HTML_CAP ? html.slice(0, HTML_CAP) : html, htmlTruncated: html.length > HTML_CAP, text: (isHtml ? htmlToText(html) : text).slice(0, TEXT_CAP), markdown: doc.markdown ?? null } : {}),81 };82}8384function docFromHtml(html: string, url: string): RawDocument {85 const body = Buffer.from(html, "utf8");86 return { url, finalUrl: url, fetchedAt: new Date().toISOString(), status: 200, contentType: "text/html; charset=utf-8", body, text: html, headers: {}, etag: null, lastModified: null, notModified: false, fetcher: "cache", level: 1, durationMs: 0, credits: 0 };87}8889function buildConnector(cfg: ConnectorConfig): Connector {90 if (cfg.implementation) {91 const impl = getImplementation(cfg.implementation);92 if (!impl) throw badRequest(`implementation "${cfg.implementation}" is not registered in this process`);93 return impl(cfg);94 }95 return new GenericConnector(cfg);96}9798async function runPipeline(connector: Connector, cfg: ConnectorConfig, doc: RawDocument) {99 const ctx = createTestContext(cfg);100 const logs: Array<{ level: string; msg: string }> = [];101 const origLog = ctx.log;102 ctx.log = (level, msg, extra) => { logs.push({ level, msg: extra ? `${msg} ${JSON.stringify(extra)}` : msg }); if (level === "error") origLog(level, msg, extra); };103 const records = await connector.extract(ctx, doc);104 const entities = await connector.normalize(ctx, records);105 const report = await connector.validate(ctx, entities);106 return { records, entities, validation: report, logs, pageMeta: doc.meta ?? null };107}108109const fetchBody = z.object({ url: z.string().url(), level: z.number().int().min(1).max(4).optional(), maxLevel: z.number().int().min(1).max(4).optional(), renderJs: z.boolean().optional(), country: z.string().length(2).optional(), waitForSelector: z.string().max(200).optional(), userAgent: z.enum(["bot", "browser"]).optional() });110111export async function devtoolAdminRoutes(app: FastifyInstance): Promise<void> {112 app.post("/devtool/fetch", { schema: { summary: "Fetch a URL through the escalation stack (SSRF policy applies) and inspect it" } }, async (req) => {113 const b = parseBody(fetchBody, req.body);114 const level = (b.level ?? 1) as FetchLevel;115 const maxLevel = Math.max(level, b.maxLevel ?? level) as FetchLevel;116 const doc = await fetchWithEscalation(b.url, { level, maxLevel, renderJs: b.renderJs, country: b.country, waitForSelector: b.waitForSelector, timeoutMs: 45_000, headers: b.userAgent === "browser" ? { "user-agent": "Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/17.5 Safari/605.1.15" } : undefined });117 return envelope(summarizeDoc(doc));118 });119120 const previewBody = z.object({ url: z.string().url().optional(), html: z.string().max(6_000_000).optional(), extractor: z.unknown(), defaults: z.record(z.string(), z.unknown()).optional(), connectorId: z.string().optional(), kind: z.string().optional(), level: z.number().int().min(1).max(4).optional(), maxLevel: z.number().int().min(1).max(4).optional() }).refine((v) => v.url || v.html, "url or html is required");121 app.post("/devtool/preview", { schema: { summary: "Run one extractor (ExtractorConfig) against a page: records, normalized entities, validation report" } }, async (req) => {122 const b = parseBody(previewBody, req.body);123 const ex = extractorSchema.safeParse(b.extractor);124 if (!ex.success) throw badRequest("invalid extractor", ex.error.issues.map((i) => ({ path: i.path.join("."), message: i.message })));125 await loadWorkerParsers();126 const url = b.url ?? "https://preview.local/page";127 const host = (() => { try { return new URL(url).hostname; } catch { return "preview.local"; } })();128 const baseCfg = b.connectorId ? await loadConfigFromDiskOrDb(b.connectorId) : null;129 const cfgRaw = { ...(baseCfg ?? { id: "preview", name: "Preview", domain: host, kind: "operator" as const }), extractors: { preview: { ...ex.data, match: undefined, pageTypes: undefined } }, defaults: b.defaults ?? baseCfg?.defaults, fetch: { ...(baseCfg?.fetch ?? {}), level: b.level ?? baseCfg?.fetch?.level ?? 1, maxLevel: b.maxLevel ?? baseCfg?.fetch?.maxLevel ?? 2, respectRobots: true }, implementation: undefined };130 const parsed = connectorConfigSchema.safeParse(cfgRaw);131 if (!parsed.success) throw badRequest("could not build a connector config", parsed.error.issues.map((i) => ({ path: i.path.join("."), message: i.message })));132 const cfg = parsed.data;133 const connector = new GenericConnector(cfg);134 const doc = b.html ? docFromHtml(b.html, url) : await createTestContext(cfg).fetch(url);135 if (doc.error) return envelope({ fetch: summarizeDoc(doc, false), records: [], entities: [], validation: { total: 0, valid: 0, rejected: 0, issues: [] }, logs: [] });136 const res = await runPipeline(connector, cfg, doc);137 return envelope({ fetch: summarizeDoc(doc, false), ...res });138 });139140 app.post("/devtool/validate-config", { schema: { summary: "Validate connector YAML against connectorConfigSchema" } }, async (req) => {141 const b = parseBody(z.object({ yaml: z.string().min(1).max(500_000) }), req.body);142 let raw: unknown;143 try { raw = parseYaml(b.yaml); } catch (e) { return envelope({ ok: false, errors: [{ path: "", message: `YAML parse error: ${(e as Error).message}` }] }); }144 const res = connectorConfigSchema.safeParse(raw);145 if (!res.success) return envelope({ ok: false, errors: res.error.issues.map((i) => ({ path: i.path.join("."), message: i.message })) });146 const warnings: string[] = [];147 if (!res.data.license) warnings.push("no license recorded (sources.license)");148 if (!res.data.attribution) warnings.push("no attribution recorded");149 if (!Object.keys(res.data.extractors).length && !res.data.implementation) warnings.push("no extractors and no implementation: only newsroom classification will produce records");150 if (res.data.fetch.maxLevel >= 3 && !process.env.FIRECRAWL_API_KEY && !process.env.SCRAPFLY_API_KEY) warnings.push("maxLevel ≥ 3 but no premium fetcher key is configured in this process");151 for (const [name, ex] of Object.entries(res.data.extractors)) {152 if (ex.parser) { const { getParser } = await import("@dci/connectors"); await loadWorkerParsers(); try { getParser(ex.parser); } catch { warnings.push(`extractor ${name}: parser "${ex.parser}" is not registered in the API process (worker may still have it)`); } }153 }154 return envelope({ ok: true, config: res.data, warnings });155 });156157 app.post("/devtool/save-config", { schema: { summary: "Write config/connectors/<id>.yaml (backup previous), bump parserVersion, upsert connectors row" } }, async (req) => {158 const b = parseBody(z.object({ id: z.string().min(1), yaml: z.string().min(1).max(500_000), version: z.string().max(32).optional() }), req.body);159 const id = safeConnectorId(b.id);160 let cfg: ConnectorConfig;161 try { cfg = parseConnectorConfig(b.yaml, `${id}.yaml`); } catch (e) { throw badRequest((e as Error).message); }162 if (cfg.id !== id) throw badRequest(`yaml id "${cfg.id}" does not match "${id}"`);163 let yamlText = b.yaml;164 if (b.version) { const doc = parseDocument(b.yaml); doc.set("parserVersion", b.version); yamlText = doc.toString(); cfg = parseConnectorConfig(yamlText, `${id}.yaml`); }165 const path = configPath(id);166 const backup = backupConfig(path);167 writeFileSync(path, yamlText, "utf8");168 const sql = pg();169 const config = JSON.parse(JSON.stringify(cfg)) as Record<string, unknown>;170 await sql`insert into connectors (id, source_name, domain, kind, mode, enabled, config, parser_version, schedule) values (${id}, ${cfg.name}, ${cfg.domain}, ${cfg.kind}, ${cfg.mode}, ${cfg.enabled}, ${JSON.stringify(config)}::jsonb, ${cfg.parserVersion}, ${JSON.stringify(cfg.schedule)}::jsonb)171 on conflict (id) do update set source_name = excluded.source_name, domain = excluded.domain, kind = excluded.kind, mode = excluded.mode, enabled = excluded.enabled, config = excluded.config, parser_version = excluded.parser_version, schedule = excluded.schedule, updated_at = now()`;172 return envelope({ id, path, backup, parserVersion: cfg.parserVersion, enabled: cfg.enabled });173 });174175 app.post("/devtool/run-sample", { schema: { summary: "Run a YAML config against up to 10 live URLs: fetch + extract + normalize + validate per URL (compare pages)" } }, async (req) => {176 const b = parseBody(z.object({ yaml: z.string().min(1).max(500_000), urls: z.array(z.string().url()).min(1).max(10), level: z.number().int().min(1).max(4).optional() }), req.body);177 let cfg: ConnectorConfig;178 try { cfg = parseConnectorConfig(b.yaml, "sample.yaml"); } catch (e) { throw badRequest((e as Error).message); }179 await loadWorkerParsers();180 const connector = buildConnector(cfg);181 const ctx = createTestContext(cfg);182 const results = [];183 for (const url of b.urls) {184 const started = Date.now();185 try {186 const doc = await connector.fetch(ctx, { url, group: "sample", priority: 100, minLevel: b.level as FetchLevel | undefined });187 if (doc.error || doc.status >= 400) { results.push({ url, ok: false, fetch: summarizeDoc(doc, false), error: doc.error?.message ?? `HTTP ${doc.status}`, durationMs: Date.now() - started }); continue; }188 const res = await runPipeline(connector, cfg, doc);189 results.push({ url, ok: res.validation.rejected === 0 && res.entities.length > 0, fetch: summarizeDoc(doc, false), ...res, durationMs: Date.now() - started });190 } catch (e) {191 results.push({ url, ok: false, error: (e as Error).message, durationMs: Date.now() - started });192 }193 }194 const fieldMatrix: Record<string, number> = {};195 for (const r of results) for (const e of (r as { entities?: Array<Record<string, unknown>> }).entities ?? []) for (const [k, v] of Object.entries(e)) if (v != null && v !== "" && k !== "provenance" && k !== "entityType") fieldMatrix[k] = (fieldMatrix[k] ?? 0) + 1;196 return envelope({ results, summary: { urls: results.length, ok: results.filter((r) => r.ok).length, credits: ctx.credits, fieldCoverage: fieldMatrix } });197 });198199 app.get("/devtool/configs", { schema: { summary: "List YAML connector configs on disk (id, name, enabled, parserVersion)" } }, async () => {200 const dir = getEnv().configDir;201 const out: Array<Record<string, unknown>> = [];202 if (existsSync(dir)) {203 for (const f of readdirSync(dir).filter((x) => /\.ya?ml$/.test(x)).sort()) {204 try {205 const raw = parseYaml(readFileSync(join(dir, f), "utf8")) as Record<string, unknown>;206 const res = connectorConfigSchema.safeParse(raw);207 out.push({ file: f, id: raw?.id ?? f.replace(/\.ya?ml$/, ""), name: raw?.name ?? null, enabled: raw?.enabled ?? true, parserVersion: raw?.parserVersion ?? "v1", kind: raw?.kind ?? null, domain: raw?.domain ?? null, valid: res.success, errors: res.success ? [] : res.error.issues.slice(0, 5).map((i) => `${i.path.join(".")}: ${i.message}`) });208 } catch (e) {209 out.push({ file: f, id: f.replace(/\.ya?ml$/, ""), valid: false, errors: [(e as Error).message] });210 }211 }212 }213 return envelope(out, { total: out.length, dir });214 });215216 app.get("/devtool/configs/:id", { schema: { summary: "Raw YAML text of a connector config" } }, async (req, reply) => {217 const { id } = req.params as { id: string };218 const path = configPath(id);219 if (!existsSync(path)) throw notFound("config file");220 reply.type("text/yaml; charset=utf-8");221 return reply.send(readFileSync(path, "utf8"));222 });223}224225async function loadConfigFromDiskOrDb(id: string): Promise<ConnectorConfig | null> {226 const path = configPath(id);227 if (existsSync(path)) { try { return parseConnectorConfig(readFileSync(path, "utf8"), path); } catch { /* fall through */ } }228 const sql = pg();229 const rows = await sql`select config from connectors where id = ${id}`;230 if (!rows[0]) return null;231 const res = connectorConfigSchema.safeParse(rows[0].config);232 return res.success ? res.data : null;233}234