/** * Dry-run a connector against LIVE sources without persisting anything. * pnpm tsx scripts/try-connector.ts config/connectors/.yaml [--limit 5] [--url ] [--verbose] [--json out.json] * Registers every apps/worker/src/connectors//index.ts (export function register()) found on disk. */ import { readFileSync, readdirSync, existsSync } from "node:fs"; import { join, resolve } from "node:path"; import { parseConnectorConfig, GenericConnector, getImplementation, createTestContext, type Connector, type DiscoveredUrl } from "@dci/connectors"; const args = process.argv.slice(2); const file = args.find((a) => !a.startsWith("--")); if (!file) { console.error("usage: try-connector.ts [--limit N] [--url U] [--verbose] [--json out]"); process.exit(2); } const opt = (k: string) => { const i = args.indexOf(`--${k}`); return i >= 0 ? args[i + 1] : undefined; }; const limit = Number(opt("limit") ?? 5); const only = opt("url"); const verbose = args.includes("--verbose"); const jsonOut = opt("json"); async function registerAll() { const dir = resolve("apps/worker/src/connectors"); if (!existsSync(dir)) return; for (const g of readdirSync(dir, { withFileTypes: true })) { if (!g.isDirectory()) continue; const idx = join(dir, g.name, "index.ts"); if (!existsSync(idx)) continue; const mod = (await import(idx)) as { register?: () => void }; mod.register?.(); } } async function main() { await registerAll(); const cfg = parseConnectorConfig(readFileSync(file!, "utf8"), file); const impl = cfg.implementation ? getImplementation(cfg.implementation) : undefined; if (cfg.implementation && !impl) throw new Error(`implementation "${cfg.implementation}" not registered`); const connector: Connector = impl ? impl(cfg) : new GenericConnector(cfg); const ctx = createTestContext(cfg, { verbose }); const t0 = Date.now(); let urls: DiscoveredUrl[] = only ? [{ url: only, group: "manual", priority: 100 }] : await connector.discover(ctx); console.error(`discovered ${urls.length} urls in ${Date.now() - t0}ms`); const groups = new Map(); for (const u of urls) groups.set(u.group, (groups.get(u.group) ?? 0) + 1); console.error(`groups: ${[...groups.entries()].map(([g, n]) => `${g}=${n}`).join(", ")}`); urls = urls.sort((a, b) => (b.priority ?? 0) - (a.priority ?? 0)).slice(0, limit); const all: unknown[] = []; let totalValid = 0, totalRejected = 0; for (const u of urls) { const doc = await connector.fetch(ctx, u); if (doc.error) { console.error(`✗ ${u.url} → ${doc.error.code}: ${doc.error.message}`); continue; } const recs = await connector.extract(ctx, doc); const ents = await connector.normalize(ctx, recs); const report = await connector.validate(ctx, ents); totalValid += report.valid; totalRejected += report.rejected; console.error(`✓ ${u.url} [${doc.fetcher} L${doc.level} ${doc.status} ${doc.durationMs}ms] → ${recs.length} records, ${ents.length} entities, ${report.valid} valid, ${report.rejected} rejected${(doc.meta?.pageType as string | undefined) ? " · " + doc.meta!.pageType : ""}`); for (const i of report.issues) console.error(` ${i.level === "error" ? "ERR " : "warn"} ${i.key} ${i.field ?? ""}: ${i.message}`); for (const e of ents.slice(0, 3)) console.log(JSON.stringify(e, null, verbose ? 1 : 0).slice(0, verbose ? 4000 : 700)); all.push(...ents); } console.error(`\nTOTAL entities=${all.length} valid=${totalValid} rejected=${totalRejected} credits=${ctx.credits} fetched=${ctx.fetched.length}`); if (jsonOut) { const { writeFileSync } = await import("node:fs"); writeFileSync(jsonOut, JSON.stringify(all, null, 1)); console.error(`wrote ${jsonOut}`); } process.exit(totalRejected > 0 && totalValid === 0 ? 1 : 0); } main().catch((e) => { console.error(e); process.exit(1); });