spb/datacenterindex
Public
HTML 53.9%
TypeScript 44.5%
JavaScript 0.6%
SQL 0.5%
1/**2 * Connector configs: load config/connectors/*.yaml, build Connector instances (registered implementation or3 * GenericConnector) and mirror them into the `connectors` + `sources` tables for the admin UI / scheduler.4 */5import { existsSync } from "node:fs";6import { GenericConnector, assertAgainstRegistry, getImplementation, getParser, loadConnectorConfigs, type Connector, type ConnectorConfig } from "@dci/connectors";7import { sha256, stableId } from "@dci/core";8import { getDb, connectors as connectorsTable, sources as sourcesTable, sql, eq, type Db } from "@dci/db";9import { getEnv } from "./env.js";1011export interface LoadedConnector {12 cfg: ConnectorConfig;13 connector: Connector;14 /** stableId("source", connectorId) — one source row per connector */15 sourceId: string;16 /**17 * Effective extractor version = the YAML `parserVersion` + the names and versions of every parser the extractors18 * reference (+ the implementation's own version). Bumping a parser's `version` therefore invalidates every cached19 * extraction of every connector that uses it — `shouldSkipExtraction` compares this, not the YAML string alone.20 */21 extractorVersion: string;22}2324export function effectiveExtractorVersion(cfg: ConnectorConfig, connector: Connector): string {25 const parts: string[] = [];26 for (const ex of Object.values(cfg.extractors ?? {})) {27 if (!ex.parser) continue;28 try { const p = getParser(ex.parser); parts.push(`${p.name}@${p.version}`); } catch { parts.push(`${ex.parser}@?`); }29 }30 if (cfg.implementation) parts.push(`impl:${cfg.implementation}@${connector.parserVersion}`);31 parts.sort();32 return parts.length ? `${cfg.parserVersion}+${sha256(parts.join("|")).slice(0, 10)}` : cfg.parserVersion;33}3435export function sourceIdFor(connectorId: string): string { return stableId("source", connectorId); }3637export function buildConnector(cfg: ConnectorConfig): Connector {38 if (cfg.implementation) {39 const factory = getImplementation(cfg.implementation);40 if (!factory) throw new Error(`connector ${cfg.id}: implementation "${cfg.implementation}" is not registered (missing apps/worker/src/connectors/<group>/index.ts register()?)`);41 return factory(cfg);42 }43 return new GenericConnector(cfg);44}4546let cache: Map<string, LoadedConnector> | null = null;4748/** Load every YAML in the config directory. Invalid files throw (fail fast: a broken config must not be silently skipped). */49export function loadAllConnectors(opts: { reload?: boolean; dir?: string } = {}): LoadedConnector[] {50 if (cache && !opts.reload) return [...cache.values()];51 const dir = opts.dir ?? getEnv().configDir;52 const out = new Map<string, LoadedConnector>();53 if (existsSync(dir)) {54 const cfgs = loadConnectorConfigs(dir);55 // every `extractors.*.parser` / `implementation` must be registered (no-op while the registry is still empty)56 assertAgainstRegistry(cfgs);57 for (const cfg of cfgs) { const connector = buildConnector(cfg); out.set(cfg.id, { cfg, connector, sourceId: sourceIdFor(cfg.id), extractorVersion: effectiveExtractorVersion(cfg, connector) }); }58 }59 cache = out;60 return [...out.values()];61}6263export function getLoadedConnector(id: string): LoadedConnector | undefined {64 if (!cache) loadAllConnectors();65 return cache!.get(id);66}6768export function requireConnector(id: string): LoadedConnector {69 const c = getLoadedConnector(id);70 if (!c) throw new Error(`unknown connector "${id}" (no config/connectors/${id}.yaml)`);71 return c;72}7374export interface SyncResult { connectors: number; sources: number; disabledInDb: string[] }7576/**77 * Upsert connectors + sources from YAML. `paused`, health and run statistics are runtime columns and are left78 * untouched; `enabled` follows the YAML (the admin can pause without editing files).79 */80export async function syncConnectorsToDb(loaded: LoadedConnector[] = loadAllConnectors(), db: Db = getDb()): Promise<SyncResult> {81 let n = 0, s = 0;82 const now = new Date().toISOString();83 for (const { cfg, sourceId, extractorVersion } of loaded) {84 const config = JSON.parse(JSON.stringify(cfg)) as Record<string, unknown>;85 await db86 .insert(connectorsTable)87 .values({ id: cfg.id, sourceName: cfg.name, domain: cfg.domain, kind: cfg.kind, mode: cfg.mode, enabled: cfg.enabled, config, parserVersion: extractorVersion, schedule: cfg.schedule })88 .onConflictDoUpdate({89 target: connectorsTable.id,90 set: { sourceName: cfg.name, domain: cfg.domain, kind: cfg.kind, mode: cfg.mode, enabled: cfg.enabled, config, parserVersion: extractorVersion, schedule: cfg.schedule, updatedAt: now },91 });92 n++;93 await db94 .insert(sourcesTable)95 .values({ id: sourceId, connectorId: cfg.id, name: cfg.name, domain: cfg.domain, kind: cfg.kind, priority: cfg.priority, url: cfg.homepage ?? `https://${cfg.domain.replace(/^https?:\/\//, "")}`, license: cfg.license ?? null, attribution: cfg.attribution ?? null, notes: cfg.notes ?? null, redistribution: redistributionOf(cfg.license), attributionRequired: attributionRequiredOf(cfg.license, cfg.attribution) })96 .onConflictDoUpdate({97 target: sourcesTable.id,98 set: { connectorId: cfg.id, name: cfg.name, domain: cfg.domain, kind: cfg.kind, priority: cfg.priority, url: cfg.homepage ?? `https://${cfg.domain.replace(/^https?:\/\//, "")}`, license: cfg.license ?? null, attribution: cfg.attribution ?? null, notes: cfg.notes ?? null, redistribution: redistributionOf(cfg.license), attributionRequired: attributionRequiredOf(cfg.license, cfg.attribution), updatedAt: now },99 });100 s++;101 }102 // connectors that exist in the DB but no longer have a YAML: disable them (never delete — runs/documents reference them)103 const ids = loaded.map((l) => l.cfg.id);104 const orphans = await db.execute<{ id: string }>(ids.length ? sql`select id from connectors where enabled = true and id not in (${sql.join(ids.map((i) => sql`${i}`), sql`, `)})` : sql`select id from connectors where enabled = true`);105 const disabledInDb: string[] = [];106 // Safety: a partial view of the YAML directory (wrong cwd, half-synced checkout) must never mass-disable connectors.107 if (orphans.length > 0 && orphans.length > Math.max(3, Math.floor(ids.length * 0.25))) {108 console.warn(`[configs] refusing to disable ${orphans.length} connectors without YAML (only ${ids.length} configs loaded) — check the config directory`);109 return { connectors: n, sources: s, disabledInDb };110 }111 for (const row of orphans) {112 await db.update(connectorsTable).set({ enabled: false, lastError: "no YAML config found at sync", updatedAt: now }).where(eq(connectorsTable.id, row.id));113 disabledInDb.push(row.id);114 }115 return { connectors: n, sources: s, disabledInDb };116}117118/** Redistribution policy derived from the declared licence (drives /download and the API `sources` block). */119export function redistributionOf(license: string | null | undefined): "allowed" | "attribution" | "restricted" | "unknown" {120 const l = (license ?? "").toLowerCase();121 if (!l) return "unknown";122 if (/cc0|public domain|pddl|open government|ogl|us government|unrestricted/.test(l)) return "allowed";123 if (/odbl|cc[- ]by|odc-by|attribution|mit|apache|peeringdb|open data/.test(l)) return "attribution";124 if (/all rights reserved|proprietary|copyright|terms of (use|service)|facts only|no redistribution|non-?commercial|nc\b/.test(l)) return "restricted";125 return "unknown";126}127export function attributionRequiredOf(license: string | null | undefined, attribution: string | null | undefined): boolean {128 return redistributionOf(license) === "attribution" || !!attribution;129}130