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;
}