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%
15.7 KB · 234 lines typescript
Raw Blame History
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