import { z } from "zod"; import { parse as parseYaml } from "yaml"; import { readFileSync, readdirSync } from "node:fs"; import { join } from "node:path"; import { SOURCE_KINDS, PAGE_TYPES } from "@dci/core"; import { getImplementation, listImplementations, listParsers } from "./registry.js"; /** * Connector configuration (YAML in config/connectors/*.yaml). Declarative first; `parser:` names a * code-backed parser from the registry when selectors are not enough. */ const selectorRule = z.union([ z.string(), z.object({ selector: z.string().optional(), attr: z.string().optional(), regex: z.string().optional(), regexFlags: z.string().optional(), group: z.number().int().optional(), all: z.boolean().optional(), jsonPath: z.string().optional(), // for embedded JSON / JSON-LD ("$.address.addressLocality") jsonld: z.string().optional(), // JSON-LD @type to look in (e.g. "Place", "LocalBusiness") meta: z.string().optional(), // transform: z.enum(["mw", "sqm", "ha", "money", "date", "status", "type", "country", "int", "float", "trim", "lower", "url", "bool"]).optional(), default: z.unknown().optional(), }), ]); export const extractorSchema = z.object({ /** page types this extractor handles */ pageTypes: z.array(z.enum(PAGE_TYPES)).optional(), /** which URL patterns select this extractor */ match: z.array(z.string()).optional(), kind: z.enum(["facility", "operator", "campus", "cloud_region", "ixp", "project", "news_event", "country"]).default("facility"), parser: z.string().optional(), params: z.record(z.string(), z.unknown()).optional(), /** declarative field rules */ fields: z.record(z.string(), selectorRule).optional(), /** repeat over a list container: each match yields one record */ each: z.string().optional(), /** stable key template, e.g. "equinix:{code}" or "{url}" */ key: z.string().optional(), minCertainty: z.number().min(0).max(1).optional(), }); export const connectorConfigSchema = z.object({ id: z.string().regex(/^[a-z0-9][a-z0-9-]*$/), name: z.string(), domain: z.string(), kind: z.enum(SOURCE_KINDS), mode: z.enum(["html", "sitemap", "rss", "pdf", "json", "hybrid", "api", "dataset"]).default("hybrid"), enabled: z.boolean().default(true), priority: z.number().int().min(1).max(5).default(3), license: z.string().optional(), attribution: z.string().optional(), homepage: z.string().url().optional(), notes: z.string().optional(), /** which countries / operators this source mostly covers (for admin + scheduling) */ coverage: z.object({ countries: z.array(z.string()).optional(), operators: z.array(z.string()).optional(), global: z.boolean().optional() }).optional(), fetch: z .object({ /** starting level; the runtime escalates on 403/429/JS-shell detection if `maxLevel` allows */ level: z.number().int().min(1).max(4).default(1), maxLevel: z.number().int().min(1).max(4).default(2), renderJs: z.boolean().optional(), country: z.string().optional(), waitForSelector: z.string().optional(), userAgent: z.enum(["bot", "browser"]).default("bot"), accept: z.string().optional(), headers: z.record(z.string(), z.string()).optional(), timeoutMs: z.number().int().optional(), maxBytes: z.number().int().optional(), /** requests per minute against this domain */ rpm: z.number().min(1).max(600).default(30), concurrency: z.number().int().min(1).max(8).default(2), respectRobots: z.boolean().default(true), /** premium credit budget per run */ maxCreditsPerRun: z.number().int().default(200), /** premium credit budget per UTC day for this connector (all runs, all workers — Redis counter); unset = only the global daily budgets apply */ maxCreditsPerDay: z.number().int().min(0).optional(), }) .prefault({}), discovery: z .object({ sitemap: z.union([z.boolean(), z.array(z.string())]).default(false), rss: z.array(z.string()).default([]), seeds: z.array(z.union([z.string(), z.object({ url: z.string(), group: z.string().default("seed"), pageType: z.enum(PAGE_TYPES).optional(), minLevel: z.number().int().min(1).max(4).optional() })])).default([]), /** follow internal links matching these patterns (regex), from index pages only */ follow: z.array(z.string()).default([]), include: z.array(z.string()).default([]), exclude: z.array(z.string()).default([]), /** map URL regex → group / pageType */ classify: z .array(z.object({ pattern: z.string(), group: z.string(), pageType: z.enum(PAGE_TYPES).optional(), priority: z.number().int().optional(), minLevel: z.number().int().min(1).max(4).optional() })) .default([]), maxUrlsPerRun: z.number().int().default(2000), /** pagination: "?page={n}" or "/page/{n}" appended to index seeds, until no new links */ pagination: z.object({ template: z.string(), start: z.number().int().default(1), max: z.number().int().default(20) }).optional(), }) .prefault({}), schedule: z.record(z.string(), z.string()).default({ seed: "weekly" }), extractors: z.record(z.string(), extractorSchema).default({}), /** defaults merged into every normalized entity from this connector */ defaults: z .object({ operatorName: z.string().optional(), operatorKind: z.string().optional(), countryIso2: z.string().optional(), facilityType: z.string().optional(), status: z.string().optional(), isHyperscale: z.boolean().optional(), providerName: z.string().optional() }) .optional(), /** code-backed connector class overriding the generic pipeline (e.g. "peeringdb", "overpass") */ implementation: z.string().optional(), params: z.record(z.string(), z.unknown()).optional(), parserVersion: z.string().default("v1"), }).superRefine((cfg, ctx) => { // cross-field rules — a config that passes the schema must also be internally consistent if (cfg.fetch.level > cfg.fetch.maxLevel) ctx.addIssue({ code: "custom", path: ["fetch", "level"], message: `fetch.level ${cfg.fetch.level} exceeds fetch.maxLevel ${cfg.fetch.maxLevel}` }); if (cfg.fetch.maxCreditsPerDay !== undefined && cfg.fetch.maxCreditsPerDay < cfg.fetch.maxCreditsPerRun && cfg.fetch.maxLevel >= 3) ctx.addIssue({ code: "custom", path: ["fetch", "maxCreditsPerRun"], message: `fetch.maxCreditsPerRun ${cfg.fetch.maxCreditsPerRun} exceeds fetch.maxCreditsPerDay ${cfg.fetch.maxCreditsPerDay}` }); for (const s of cfg.discovery.seeds) { const lvl = typeof s === "object" ? s.minLevel : undefined; if (lvl !== undefined && lvl > cfg.fetch.maxLevel) ctx.addIssue({ code: "custom", path: ["discovery", "seeds"], message: `seed ${typeof s === "object" ? s.url : s} minLevel ${lvl} exceeds fetch.maxLevel ${cfg.fetch.maxLevel}` }); } cfg.discovery.classify.forEach((r, i) => { if (r.minLevel !== undefined && r.minLevel > cfg.fetch.maxLevel) ctx.addIssue({ code: "custom", path: ["discovery", "classify", i, "minLevel"], message: `minLevel ${r.minLevel} exceeds fetch.maxLevel ${cfg.fetch.maxLevel}` }); }); // every regex is compiled once here so a bad pattern fails `dci validate` instead of throwing mid-discovery const tryRe = (pattern: string, path: Array) => { try { new RegExp(pattern, "i"); } catch (e) { ctx.addIssue({ code: "custom", path, message: `invalid regex /${pattern}/: ${(e as Error).message}` }); } }; cfg.discovery.include.forEach((p, i) => tryRe(p, ["discovery", "include", i])); cfg.discovery.exclude.forEach((p, i) => tryRe(p, ["discovery", "exclude", i])); cfg.discovery.follow.forEach((p, i) => tryRe(p, ["discovery", "follow", i])); cfg.discovery.classify.forEach((r, i) => tryRe(r.pattern, ["discovery", "classify", i, "pattern"])); for (const [name, ex] of Object.entries(cfg.extractors)) { ex.match?.forEach((p, i) => tryRe(p, ["extractors", name, "match", i])); if (!ex.parser && !ex.fields) ctx.addIssue({ code: "custom", path: ["extractors", name], message: "extractor needs either `parser` or `fields`" }); for (const [field, rule] of Object.entries(ex.fields ?? {})) if (typeof rule === "object" && rule.regex) { try { new RegExp(rule.regex, rule.regexFlags ?? "i"); } catch (e) { ctx.addIssue({ code: "custom", path: ["extractors", name, "fields", field, "regex"], message: `invalid regex: ${(e as Error).message}` }); } } } for (const [group, spec] of Object.entries(cfg.schedule)) { try { intervalMs(spec); } catch { ctx.addIssue({ code: "custom", path: ["schedule", group], message: `bad interval "${spec}" (daily | weekly | monthly | 6h | 30m …)` }); } } }); export type ConnectorConfig = z.infer; export type ExtractorConfig = z.infer; export type SelectorRule = z.infer; export function parseConnectorConfig(text: string, file = ""): ConnectorConfig { const raw = parseYaml(text); const res = connectorConfigSchema.safeParse(raw); if (!res.success) throw new Error(`Invalid connector config ${file}: ${res.error.issues.map((i) => `${i.path.join(".")}: ${i.message}`).join("; ")}`); return res.data; } /** * Registry-dependent checks, run lazily once parsers / implementations are registered (the YAML schema cannot know * them): every `extractors.*.parser` and the `implementation` must exist. Returns the list of problems (empty = ok); * `assertAgainstRegistry` throws one clear error listing them. When the registry is still empty the check is skipped * (returns []), so configs can be parsed before the worker groups are loaded. */ export function validateAgainstRegistry(cfg: ConnectorConfig): string[] { if (!listParsers().length && !listImplementations().length) return []; const problems: string[] = []; const known = new Set(listParsers().map((p) => p.name)); for (const [name, ex] of Object.entries(cfg.extractors)) if (ex.parser && !known.has(ex.parser)) problems.push(`extractors.${name}.parser "${ex.parser}" is not registered`); if (cfg.implementation && !getImplementation(cfg.implementation)) problems.push(`implementation "${cfg.implementation}" is not registered`); return problems; } export function assertAgainstRegistry(cfgs: ConnectorConfig[]): void { const all = cfgs.flatMap((c) => validateAgainstRegistry(c).map((p) => `${c.id}: ${p}`)); if (all.length) throw new Error(`connector config(s) reference unknown parsers/implementations — ${all.join("; ")}. Registered parsers: ${listParsers().map((p) => p.name).sort().join(", ") || "none"}`); } export function loadConnectorConfigs(dir: string): ConnectorConfig[] { const out: ConnectorConfig[] = []; const seen = new Set(); for (const f of readdirSync(dir).filter((x) => /\.ya?ml$/.test(x)).sort()) { const cfg = parseConnectorConfig(readFileSync(join(dir, f), "utf8"), f); if (seen.has(cfg.id)) throw new Error(`duplicate connector id ${cfg.id} in ${f}`); seen.add(cfg.id); out.push(cfg); } return out; } /** "daily" | "weekly" | "monthly" | "hourly" | "3h" | "12h" | "30m" → milliseconds */ export function intervalMs(spec: string): number { const s = spec.trim().toLowerCase(); const named: Record = { hourly: 3_600_000, daily: 86_400_000, weekly: 7 * 86_400_000, biweekly: 14 * 86_400_000, monthly: 30 * 86_400_000, quarterly: 91 * 86_400_000, never: Number.POSITIVE_INFINITY }; if (named[s]) return named[s]!; const m = s.match(/^(\d+)\s*(m|min|h|d|w)$/); if (!m) throw new Error(`bad interval: ${spec}`); const n = Number(m[1]); return { m: 60_000, min: 60_000, h: 3_600_000, d: 86_400_000, w: 7 * 86_400_000 }[m[2]!]! * n; }