SPB Git forge
38commits 1branches 0releases
338.7 MBsize
maindefault branch
3 h agolast push
HTML 53.9% TypeScript 44.5% JavaScript 0.6% SQL 0.5%
11.4 KB · 183 lines typescript
Raw Blame History
1import { z } from "zod";2import { parse as parseYaml } from "yaml";3import { readFileSync, readdirSync } from "node:fs";4import { join } from "node:path";5import { SOURCE_KINDS, PAGE_TYPES } from "@dci/core";6import { getImplementation, listImplementations, listParsers } from "./registry.js";78/**9 * Connector configuration (YAML in config/connectors/*.yaml). Declarative first; `parser:` names a10 * code-backed parser from the registry when selectors are not enough.11 */12const selectorRule = z.union([13  z.string(),14  z.object({15    selector: z.string().optional(),16    attr: z.string().optional(),17    regex: z.string().optional(),18    regexFlags: z.string().optional(),19    group: z.number().int().optional(),20    all: z.boolean().optional(),21    jsonPath: z.string().optional(), // for embedded JSON / JSON-LD ("$.address.addressLocality")22    jsonld: z.string().optional(), // JSON-LD @type to look in (e.g. "Place", "LocalBusiness")23    meta: z.string().optional(), // <meta property|name=…>24    transform: z.enum(["mw", "sqm", "ha", "money", "date", "status", "type", "country", "int", "float", "trim", "lower", "url", "bool"]).optional(),25    default: z.unknown().optional(),26  }),27]);2829export const extractorSchema = z.object({30  /** page types this extractor handles */31  pageTypes: z.array(z.enum(PAGE_TYPES)).optional(),32  /** which URL patterns select this extractor */33  match: z.array(z.string()).optional(),34  kind: z.enum(["facility", "operator", "campus", "cloud_region", "ixp", "project", "news_event", "country"]).default("facility"),35  parser: z.string().optional(),36  params: z.record(z.string(), z.unknown()).optional(),37  /** declarative field rules */38  fields: z.record(z.string(), selectorRule).optional(),39  /** repeat over a list container: each match yields one record */40  each: z.string().optional(),41  /** stable key template, e.g. "equinix:{code}" or "{url}" */42  key: z.string().optional(),43  minCertainty: z.number().min(0).max(1).optional(),44});4546export const connectorConfigSchema = z.object({47  id: z.string().regex(/^[a-z0-9][a-z0-9-]*$/),48  name: z.string(),49  domain: z.string(),50  kind: z.enum(SOURCE_KINDS),51  mode: z.enum(["html", "sitemap", "rss", "pdf", "json", "hybrid", "api", "dataset"]).default("hybrid"),52  enabled: z.boolean().default(true),53  priority: z.number().int().min(1).max(5).default(3),54  license: z.string().optional(),55  attribution: z.string().optional(),56  homepage: z.string().url().optional(),57  notes: z.string().optional(),58  /** which countries / operators this source mostly covers (for admin + scheduling) */59  coverage: z.object({ countries: z.array(z.string()).optional(), operators: z.array(z.string()).optional(), global: z.boolean().optional() }).optional(),60  fetch: z61    .object({62      /** starting level; the runtime escalates on 403/429/JS-shell detection if `maxLevel` allows */63      level: z.number().int().min(1).max(4).default(1),64      maxLevel: z.number().int().min(1).max(4).default(2),65      renderJs: z.boolean().optional(),66      country: z.string().optional(),67      waitForSelector: z.string().optional(),68      userAgent: z.enum(["bot", "browser"]).default("bot"),69      accept: z.string().optional(),70      headers: z.record(z.string(), z.string()).optional(),71      timeoutMs: z.number().int().optional(),72      maxBytes: z.number().int().optional(),73      /** requests per minute against this domain */74      rpm: z.number().min(1).max(600).default(30),75      concurrency: z.number().int().min(1).max(8).default(2),76      respectRobots: z.boolean().default(true),77      /** premium credit budget per run */78      maxCreditsPerRun: z.number().int().default(200),79      /** premium credit budget per UTC day for this connector (all runs, all workers — Redis counter); unset = only the global daily budgets apply */80      maxCreditsPerDay: z.number().int().min(0).optional(),81    })82    .prefault({}),83  discovery: z84    .object({85      sitemap: z.union([z.boolean(), z.array(z.string())]).default(false),86      rss: z.array(z.string()).default([]),87      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([]),88      /** follow internal links matching these patterns (regex), from index pages only */89      follow: z.array(z.string()).default([]),90      include: z.array(z.string()).default([]),91      exclude: z.array(z.string()).default([]),92      /** map URL regex → group / pageType */93      classify: z94        .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() }))95        .default([]),96      maxUrlsPerRun: z.number().int().default(2000),97      /** pagination: "?page={n}" or "/page/{n}" appended to index seeds, until no new links */98      pagination: z.object({ template: z.string(), start: z.number().int().default(1), max: z.number().int().default(20) }).optional(),99    })100    .prefault({}),101  schedule: z.record(z.string(), z.string()).default({ seed: "weekly" }),102  extractors: z.record(z.string(), extractorSchema).default({}),103  /** defaults merged into every normalized entity from this connector */104  defaults: z105    .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() })106    .optional(),107  /** code-backed connector class overriding the generic pipeline (e.g. "peeringdb", "overpass") */108  implementation: z.string().optional(),109  params: z.record(z.string(), z.unknown()).optional(),110  parserVersion: z.string().default("v1"),111}).superRefine((cfg, ctx) => {112  // cross-field rules — a config that passes the schema must also be internally consistent113  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}` });114  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}` });115  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}` }); }116  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}` }); });117  // every regex is compiled once here so a bad pattern fails `dci validate` instead of throwing mid-discovery118  const tryRe = (pattern: string, path: Array<string | number>) => { try { new RegExp(pattern, "i"); } catch (e) { ctx.addIssue({ code: "custom", path, message: `invalid regex /${pattern}/: ${(e as Error).message}` }); } };119  cfg.discovery.include.forEach((p, i) => tryRe(p, ["discovery", "include", i]));120  cfg.discovery.exclude.forEach((p, i) => tryRe(p, ["discovery", "exclude", i]));121  cfg.discovery.follow.forEach((p, i) => tryRe(p, ["discovery", "follow", i]));122  cfg.discovery.classify.forEach((r, i) => tryRe(r.pattern, ["discovery", "classify", i, "pattern"]));123  for (const [name, ex] of Object.entries(cfg.extractors)) {124    ex.match?.forEach((p, i) => tryRe(p, ["extractors", name, "match", i]));125    if (!ex.parser && !ex.fields) ctx.addIssue({ code: "custom", path: ["extractors", name], message: "extractor needs either `parser` or `fields`" });126    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}` }); } }127  }128  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 …)` }); } }129});130131export type ConnectorConfig = z.infer<typeof connectorConfigSchema>;132export type ExtractorConfig = z.infer<typeof extractorSchema>;133export type SelectorRule = z.infer<typeof selectorRule>;134135export function parseConnectorConfig(text: string, file = "<inline>"): ConnectorConfig {136  const raw = parseYaml(text);137  const res = connectorConfigSchema.safeParse(raw);138  if (!res.success) throw new Error(`Invalid connector config ${file}: ${res.error.issues.map((i) => `${i.path.join(".")}: ${i.message}`).join("; ")}`);139  return res.data;140}141142/**143 * Registry-dependent checks, run lazily once parsers / implementations are registered (the YAML schema cannot know144 * them): every `extractors.*.parser` and the `implementation` must exist. Returns the list of problems (empty = ok);145 * `assertAgainstRegistry` throws one clear error listing them. When the registry is still empty the check is skipped146 * (returns []), so configs can be parsed before the worker groups are loaded.147 */148export function validateAgainstRegistry(cfg: ConnectorConfig): string[] {149  if (!listParsers().length && !listImplementations().length) return [];150  const problems: string[] = [];151  const known = new Set(listParsers().map((p) => p.name));152  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`);153  if (cfg.implementation && !getImplementation(cfg.implementation)) problems.push(`implementation "${cfg.implementation}" is not registered`);154  return problems;155}156export function assertAgainstRegistry(cfgs: ConnectorConfig[]): void {157  const all = cfgs.flatMap((c) => validateAgainstRegistry(c).map((p) => `${c.id}: ${p}`));158  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"}`);159}160161export function loadConnectorConfigs(dir: string): ConnectorConfig[] {162  const out: ConnectorConfig[] = [];163  const seen = new Set<string>();164  for (const f of readdirSync(dir).filter((x) => /\.ya?ml$/.test(x)).sort()) {165    const cfg = parseConnectorConfig(readFileSync(join(dir, f), "utf8"), f);166    if (seen.has(cfg.id)) throw new Error(`duplicate connector id ${cfg.id} in ${f}`);167    seen.add(cfg.id);168    out.push(cfg);169  }170  return out;171}172173/** "daily" | "weekly" | "monthly" | "hourly" | "3h" | "12h" | "30m" → milliseconds */174export function intervalMs(spec: string): number {175  const s = spec.trim().toLowerCase();176  const named: Record<string, number> = { 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 };177  if (named[s]) return named[s]!;178  const m = s.match(/^(\d+)\s*(m|min|h|d|w)$/);179  if (!m) throw new Error(`bad interval: ${spec}`);180  const n = Number(m[1]);181  return { m: 60_000, min: 60_000, h: 3_600_000, d: 86_400_000, w: 7 * 86_400_000 }[m[2]!]! * n;182}183