import { existsSync, readdirSync, readFileSync, statSync } from "node:fs"; import { join } from "node:path"; import YAML from "yaml"; import { coverageKey, coverageSectorSchema, slugify, type CoverageMember, type CoverageSector } from "@websensor/core"; import { db, factorySeeds, sql, textArray } from "@websensor/db"; import { factoryConfig, log } from "../config"; /** * Factory seeds come from two places: the coverage universes (`config/coverage/*.yaml`, every member that is not * monitored yet — or that carries hints worth expanding) and ad-hoc seed files (`config/factory/seeds/*.yaml`, * `{ seeds: [{ name, domain, categories, country, … }] }`) or the admin API. */ export interface SeedInput { id?: string; name: string; domain: string; homepage?: string; categories?: string[]; country?: string; language?: string; tier?: string; weight?: number; importance?: number; aliases?: string[]; first_party?: boolean; sector?: string; universe?: string; hints?: Record; } export function loadCoverageSectors(dir = factoryConfig.coverageDir): { sectors: CoverageSector[]; issues: string[] } { const sectors: CoverageSector[] = []; const issues: string[] = []; if (!existsSync(dir) || !statSync(dir).isDirectory()) return { sectors, issues: [`coverage dir ${dir} not found`] }; for (const f of readdirSync(dir).filter((x) => /\.ya?ml$/.test(x) && !x.startsWith("_")).sort()) { try { const parsed = coverageSectorSchema.safeParse(YAML.parse(readFileSync(join(dir, f), "utf8"))); if (!parsed.success) { issues.push(`${f}: ${parsed.error.issues.slice(0, 3).map((i) => `${i.path.join(".")}: ${i.message}`).join("; ")}`); continue; } sectors.push(parsed.data); } catch (e) { issues.push(`${f}: ${(e as Error).message}`); } } return { sectors, issues }; } /** Merge files that share a sector key. */ export function mergeSectors(sectors: CoverageSector[]): CoverageSector[] { const byKey = new Map(); for (const s of sectors) { const prev = byKey.get(s.sector); if (!prev) byKey.set(s.sector, { ...s, universes: [...s.universes] }); else prev.universes.push(...s.universes); } return [...byKey.values()]; } export function memberToSeed(sector: CoverageSector, universe: string, m: CoverageMember): SeedInput { return { name: m.name, domain: m.domain, categories: m.categories ?? sector.categories, country: m.country, language: m.language, tier: m.tier ?? sector.tier, importance: m.importance, aliases: m.aliases, first_party: m.first_party, sector: sector.sector, universe, hints: m.hints as Record, }; } function hasHints(h: Record | undefined): boolean { if (!h) return false; return Boolean(h.cik || h.github_org || (Array.isArray(h.github_repos) && h.github_repos.length) || h.hf_author || h.status_url || (Array.isArray(h.urls) && h.urls.length) || (Array.isArray(h.hosts) && h.hosts.length)); } /** Registrable domains currently monitored (source domains + hosts of enabled non-shadow sensors). */ export async function monitoredDomainKeys(): Promise> { const rows = await db.execute<{ d: string }>(sql` select domain as d from sources where enabled and kind = 'registry' union select lower(split_part(split_part(url, '/', 3), ':', 1)) as d from sensors where enabled and status <> 'SHADOW' and url like 'http%'`); const out = new Set(); for (const r of rows.rows) if (r.d) out.add(coverageKey(r.d)); return out; } export async function upsertSeeds(inputs: SeedInput[], opts: { requeue?: boolean } = {}): Promise<{ inserted: number; updated: number; skipped: number }> { let inserted = 0; let updated = 0; let skipped = 0; const existingIds = new Set((await db.execute<{ id: string }>(sql`select id from sources`)).rows.map((r) => r.id)); const existingSeedDomains = new Map((await db.execute<{ id: string; domain: string }>(sql`select id, domain from factory_seeds`)).rows.map((r) => [coverageKey(r.domain), r.id])); for (const s of inputs) { const domain = s.domain.toLowerCase().replace(/^https?:\/\//, "").replace(/^www\./, "").replace(/\/.*$/, ""); const key = coverageKey(domain); let id = s.id ?? existingSeedDomains.get(key) ?? slugify(s.name).slice(0, 60); if (!id) id = slugify(domain); // A registry source with the same id but another organization → disambiguate by domain. if (!s.id && !existingSeedDomains.has(key) && existingIds.has(id)) { const src = (await db.execute<{ domain: string }>(sql`select domain from sources where id = ${id}`)).rows[0]; if (src && coverageKey(src.domain) !== key) id = `${id}-${slugify(domain.split(".")[0] ?? domain)}`.slice(0, 70); } const r = await db.execute<{ inserted: boolean }>(sql` insert into factory_seeds (id, name, domain, homepage, categories, country, language, tier, weight, importance, aliases, first_party, sector, universe, hints, status) values (${id}, ${s.name}, ${domain}, ${s.homepage ?? null}, ${textArray(s.categories ?? [])}, ${s.country ?? null}, ${s.language ?? null}, ${s.tier ?? "B"}, ${s.weight ?? 1}, ${s.importance ?? 2}, ${textArray(s.aliases ?? [])}, ${s.first_party ?? true}, ${s.sector ?? null}, ${s.universe ?? null}, ${JSON.stringify(s.hints ?? {})}::jsonb, 'queued') on conflict (id) do update set name = excluded.name, domain = excluded.domain, categories = excluded.categories, country = coalesce(excluded.country, factory_seeds.country), language = coalesce(excluded.language, factory_seeds.language), tier = excluded.tier, importance = excluded.importance, aliases = excluded.aliases, sector = coalesce(excluded.sector, factory_seeds.sector), universe = coalesce(excluded.universe, factory_seeds.universe), hints = factory_seeds.hints || excluded.hints, updated_at = now()${opts.requeue ? sql`, status = 'queued'` : sql``} returning (xmax = 0) as inserted`); if (r.rows[0]?.inserted) inserted++; else if (r.rowCount) updated++; else skipped++; existingSeedDomains.set(key, id); } log.info({ inserted, updated, skipped }, "factory seeds upserted"); return { inserted, updated, skipped }; } /** * Turn coverage universes into seeds. `mode`: `uncovered` = members without any monitored domain; `hinted` = also * covered members that carry expansion hints; `all` = every member (re-discover everything). */ export async function seedFromCoverage(opts: { mode?: "uncovered" | "hinted" | "all"; sectors?: string[]; requeue?: boolean } = {}): Promise<{ members: number; seeds: number; inserted: number; updated: number; issues: string[] }> { const mode = opts.mode ?? "hinted"; const { sectors, issues } = loadCoverageSectors(); const monitored = await monitoredDomainKeys(); const inputs: SeedInput[] = []; const seen = new Set(); let members = 0; for (const sector of mergeSectors(sectors)) { if (opts.sectors && !opts.sectors.includes(sector.sector)) continue; for (const u of sector.universes) { for (const m of u.members) { members++; const key = coverageKey(m.domain); if (seen.has(key)) continue; const covered = monitored.has(key); if (mode === "uncovered" && covered) continue; if (mode === "hinted" && covered && !hasHints(m.hints as Record)) continue; seen.add(key); inputs.push(memberToSeed(sector, u.key, m)); } } } const r = await upsertSeeds(inputs, { requeue: opts.requeue }); return { members, seeds: inputs.length, inserted: r.inserted, updated: r.updated, issues }; } /** `config/factory/seeds/*.yaml` → `{ seeds: [...] }` */ export async function seedFromFiles(dir = factoryConfig.seedsDir): Promise<{ files: number; seeds: number }> { if (!existsSync(dir)) return { files: 0, seeds: 0 }; const inputs: SeedInput[] = []; let files = 0; for (const f of readdirSync(dir).filter((x) => /\.ya?ml$/.test(x)).sort()) { files++; const raw = YAML.parse(readFileSync(join(dir, f), "utf8")) as { seeds?: SeedInput[] }; for (const s of raw.seeds ?? []) if (s?.name && s?.domain) inputs.push(s); } await upsertSeeds(inputs); return { files, seeds: inputs.length }; } export { factorySeeds };