/** * Connector configs: load config/connectors/*.yaml, build Connector instances (registered implementation or * GenericConnector) and mirror them into the `connectors` + `sources` tables for the admin UI / scheduler. */ import { existsSync } from "node:fs"; import { GenericConnector, assertAgainstRegistry, getImplementation, getParser, loadConnectorConfigs, type Connector, type ConnectorConfig } from "@dci/connectors"; import { sha256, stableId } from "@dci/core"; import { getDb, connectors as connectorsTable, sources as sourcesTable, sql, eq, type Db } from "@dci/db"; import { getEnv } from "./env.js"; export interface LoadedConnector { cfg: ConnectorConfig; connector: Connector; /** stableId("source", connectorId) — one source row per connector */ sourceId: string; /** * Effective extractor version = the YAML `parserVersion` + the names and versions of every parser the extractors * reference (+ the implementation's own version). Bumping a parser's `version` therefore invalidates every cached * extraction of every connector that uses it — `shouldSkipExtraction` compares this, not the YAML string alone. */ extractorVersion: string; } export function effectiveExtractorVersion(cfg: ConnectorConfig, connector: Connector): string { const parts: string[] = []; for (const ex of Object.values(cfg.extractors ?? {})) { if (!ex.parser) continue; try { const p = getParser(ex.parser); parts.push(`${p.name}@${p.version}`); } catch { parts.push(`${ex.parser}@?`); } } if (cfg.implementation) parts.push(`impl:${cfg.implementation}@${connector.parserVersion}`); parts.sort(); return parts.length ? `${cfg.parserVersion}+${sha256(parts.join("|")).slice(0, 10)}` : cfg.parserVersion; } export function sourceIdFor(connectorId: string): string { return stableId("source", connectorId); } export function buildConnector(cfg: ConnectorConfig): Connector { if (cfg.implementation) { const factory = getImplementation(cfg.implementation); if (!factory) throw new Error(`connector ${cfg.id}: implementation "${cfg.implementation}" is not registered (missing apps/worker/src/connectors//index.ts register()?)`); return factory(cfg); } return new GenericConnector(cfg); } let cache: Map | null = null; /** Load every YAML in the config directory. Invalid files throw (fail fast: a broken config must not be silently skipped). */ export function loadAllConnectors(opts: { reload?: boolean; dir?: string } = {}): LoadedConnector[] { if (cache && !opts.reload) return [...cache.values()]; const dir = opts.dir ?? getEnv().configDir; const out = new Map(); if (existsSync(dir)) { const cfgs = loadConnectorConfigs(dir); // every `extractors.*.parser` / `implementation` must be registered (no-op while the registry is still empty) assertAgainstRegistry(cfgs); for (const cfg of cfgs) { const connector = buildConnector(cfg); out.set(cfg.id, { cfg, connector, sourceId: sourceIdFor(cfg.id), extractorVersion: effectiveExtractorVersion(cfg, connector) }); } } cache = out; return [...out.values()]; } export function getLoadedConnector(id: string): LoadedConnector | undefined { if (!cache) loadAllConnectors(); return cache!.get(id); } export function requireConnector(id: string): LoadedConnector { const c = getLoadedConnector(id); if (!c) throw new Error(`unknown connector "${id}" (no config/connectors/${id}.yaml)`); return c; } export interface SyncResult { connectors: number; sources: number; disabledInDb: string[] } /** * Upsert connectors + sources from YAML. `paused`, health and run statistics are runtime columns and are left * untouched; `enabled` follows the YAML (the admin can pause without editing files). */ export async function syncConnectorsToDb(loaded: LoadedConnector[] = loadAllConnectors(), db: Db = getDb()): Promise { let n = 0, s = 0; const now = new Date().toISOString(); for (const { cfg, sourceId, extractorVersion } of loaded) { const config = JSON.parse(JSON.stringify(cfg)) as Record; await db .insert(connectorsTable) .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 }) .onConflictDoUpdate({ target: connectorsTable.id, set: { sourceName: cfg.name, domain: cfg.domain, kind: cfg.kind, mode: cfg.mode, enabled: cfg.enabled, config, parserVersion: extractorVersion, schedule: cfg.schedule, updatedAt: now }, }); n++; await db .insert(sourcesTable) .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) }) .onConflictDoUpdate({ target: sourcesTable.id, 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 }, }); s++; } // connectors that exist in the DB but no longer have a YAML: disable them (never delete — runs/documents reference them) const ids = loaded.map((l) => l.cfg.id); 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`); const disabledInDb: string[] = []; // Safety: a partial view of the YAML directory (wrong cwd, half-synced checkout) must never mass-disable connectors. if (orphans.length > 0 && orphans.length > Math.max(3, Math.floor(ids.length * 0.25))) { console.warn(`[configs] refusing to disable ${orphans.length} connectors without YAML (only ${ids.length} configs loaded) — check the config directory`); return { connectors: n, sources: s, disabledInDb }; } for (const row of orphans) { await db.update(connectorsTable).set({ enabled: false, lastError: "no YAML config found at sync", updatedAt: now }).where(eq(connectorsTable.id, row.id)); disabledInDb.push(row.id); } return { connectors: n, sources: s, disabledInDb }; } /** Redistribution policy derived from the declared licence (drives /download and the API `sources` block). */ export function redistributionOf(license: string | null | undefined): "allowed" | "attribution" | "restricted" | "unknown" { const l = (license ?? "").toLowerCase(); if (!l) return "unknown"; if (/cc0|public domain|pddl|open government|ogl|us government|unrestricted/.test(l)) return "allowed"; if (/odbl|cc[- ]by|odc-by|attribution|mit|apache|peeringdb|open data/.test(l)) return "attribution"; if (/all rights reserved|proprietary|copyright|terms of (use|service)|facts only|no redistribution|non-?commercial|nc\b/.test(l)) return "restricted"; return "unknown"; } export function attributionRequiredOf(license: string | null | undefined, attribution: string | null | undefined): boolean { return redistributionOf(license) === "attribution" || !!attribution; }