/** * One connector run: discover → register URLs → select due documents → fetch (robots, rate limit, conditional * GET, escalation, budgets) → archive raw → fingerprint → skip unchanged → extract → normalize → validate → * ingest → document bookkeeping (versions, detected changes, adaptive next_check) → run + connector health. * * Every per-document step is wrapped: one bad page never kills a run. SIGTERM finishes the current documents * and marks the run `aborted`. */ import type { DiscoveredUrl, RawDocument } from "@dci/connectors"; import { entityIsValid, pageTitle } from "@dci/connectors"; import type { NormalizedEntity, ValidationIssue } from "@dci/core"; import { contentFingerprint, sha256, newId } from "@dci/core"; import { getDb, connectorRuns, connectors as connectorsTable, eq, inArray, sql } from "@dci/db"; import { loadAllConnectors, requireConnector, type LoadedConnector } from "./configs.js"; import { createRunContext, type LogLevel, type LogLine, type RunContext } from "./context.js"; import { connectorDocStats, documentIdFor, dueDocuments, nextCheckByGroup, recordFetch, recordVersion, registerDiscovered, textProjection, updateExtraction, type DocumentRow } from "./documents.js"; import { ingestEntities } from "./ingest/index.js"; import type { IngestDocRef, IngestRun, IngestStats } from "./ingest/contract.js"; import { metrics } from "./prom.js"; import { connectorHealthFrom, healthFrom, minIntervalMs, parseDiscoveredFrom, shouldSkipExtraction } from "./scheduling.js"; import { getRaw, putRaw } from "./storage.js"; export type RunTask = "discover" | "crawl" | "full" | "reprocess"; export interface RunOptions { task: RunTask; group?: string; limit?: number; dryRun?: boolean; /** crawl exactly these URLs (registered on the fly, forced) */ urls?: string[]; /** ignore next_check and content-hash skip */ force?: boolean; onLog?: (line: LogLine, extra?: Record) => void; logLevel?: LogLevel; /** collect up to N normalized entities for display (dry runs) */ sampleEntities?: number; /** quarantine: fetch, archive and extract normally but roll the ingest back (set from connectors.quarantine) */ quarantine?: boolean; /** reprocess only documents whose stored extractor version differs from the current one */ staleOnly?: boolean; } export interface RunStats extends Record { discovered: number; registered: number; selected: number; fetched: number; notModified: number; changed: number; unchanged: number; markupOnly: number; failed: number; robotsBlocked: number; quarantined: number; extracted: number; entities: number; valid: number; rejected: number; created: number; updated: number; merged: number; pendingMatches: number; events: number; provenanceRows: number; credits: number; avgMs: number; totalMs: number; } export type RunStatus = "ok" | "partial" | "failed" | "aborted"; export interface RunResult { runId: string; connectorId: string; task: RunTask; status: RunStatus; stats: RunStats; error: string | null; startedAt: string; finishedAt: string; durationMs: number; samples: NormalizedEntity[]; issues: ValidationIssue[]; health: "ok" | "degraded" | "failing"; } function emptyStats(): RunStats { return { discovered: 0, registered: 0, selected: 0, fetched: 0, notModified: 0, changed: 0, unchanged: 0, markupOnly: 0, failed: 0, robotsBlocked: 0, quarantined: 0, extracted: 0, entities: 0, valid: 0, rejected: 0, created: 0, updated: 0, merged: 0, pendingMatches: 0, events: 0, provenanceRows: 0, credits: 0, avgMs: 0, totalMs: 0 }; } /* ---------- graceful abort ---------- */ let abortReason: string | null = null; const activeRuns = new Set(); export function requestAbort(reason = "SIGTERM"): void { abortReason = reason; } export function isAbortRequested(): boolean { return abortReason !== null; } export function abortReasonText(): string | null { return abortReason; } export function activeRunIds(): string[] { return [...activeRuns]; } /** For tests / long-lived processes that want to resume after a handled abort. */ export function clearAbort(): void { abortReason = null; } /** * Forced shutdown (graceful deadline reached): mark the runs still in flight as aborted in Postgres so they do not * linger as `running` until the orphan reconciliation. Best effort; returns the number of runs updated. */ export async function abortActiveRunsInDb(reason: string): Promise { const ids = [...activeRuns]; if (!ids.length) return 0; await getDb().update(connectorRuns).set({ status: "aborted", finishedAt: new Date().toISOString(), error: `aborted: ${reason}`.slice(0, 2000) }).where(inArray(connectorRuns.id, ids)); return ids.length; } /* ---------- inline p-limit ---------- */ export function pLimit(concurrency: number): (fn: () => Promise) => Promise { const n = Math.max(1, Math.floor(concurrency)); let active = 0; const queue: Array<() => void> = []; const next = () => { active--; queue.shift()?.(); }; return (fn: () => Promise) => new Promise((resolve, reject) => { const start = () => { active++; fn().then(resolve, reject).finally(next); }; if (active < n) start(); else queue.push(start); }); } /* ---------- helpers ---------- */ /** * Runtime-level noise on top of core's contentFingerprint: epoch-like cache busters in URLs (`?t=1789110041117`, * `&_=1789110041`, `?v=1700000000000`) that some CMSs stamp into og:url / share links on every render. */ export function stripRuntimeNoise(text: string): string { return text.replace(/([?&](?:t|ts|_t|_|v|cb|cache|nocache|timestamp|time|rnd|r)=)\d{9,13}(?=[&"'\s<)]|$)/gi, "$1").replace(/\b1[5-9]\d{11}\b/g, ""); } export function fingerprintOf(raw: RawDocument): string { return raw.text ? contentFingerprint(stripRuntimeNoise(raw.text), raw.contentType) : sha256(raw.body); } function docToDiscovered(doc: DocumentRow): DiscoveredUrl { const { group, from } = parseDiscoveredFrom(doc.discoveredFrom); return { url: doc.url, group, pageType: doc.pageType !== "unknown" ? (doc.pageType as DiscoveredUrl["pageType"]) : undefined, priority: doc.priority, discoveredFrom: from, minLevel: doc.fetchLevel > 1 ? (doc.fetchLevel as DiscoveredUrl["minLevel"]) : undefined }; } function sumIngest(stats: RunStats, s: IngestStats): void { stats.created += s.created; stats.updated += s.updated; stats.merged += s.merged; stats.pendingMatches += s.pendingMatches; stats.events += s.events; stats.provenanceRows += s.provenanceRows; stats.claims = (stats.claims ?? 0) + (s.claims ?? 0); stats.qualityFlags = (stats.qualityFlags ?? 0) + (s.qualityFlags ?? 0); stats.unscopedClaims = (stats.unscopedClaims ?? 0) + (s.unscopedClaims ?? 0); stats.projectsVetoed = (stats.projectsVetoed ?? 0) + (s.projectsVetoed ?? 0); stats.campusLinks = (stats.campusLinks ?? 0) + (s.campusLinks ?? 0); } interface DocOutcome { kind: "fetched" | "notModified" | "failed" | "unchanged" | "changed" | "aborted"; ms: number } /** Extract → normalize → validate → ingest for one raw document. Returns ingest stats (or null when nothing valid). */ async function extractAndIngest(ctx: RunContext, loaded: LoadedConnector, doc: DocumentRow, raw: RawDocument, stats: RunStats, result: RunResult, opts: RunOptions): Promise<{ ingest: IngestStats | null; valid: number; rejected: number; classifier: string | null; pageType: string | null }> { const { connector } = loaded; const records = await connector.extract(ctx, raw); const entities = await connector.normalize(ctx, records); const report = await connector.validate(ctx, entities); const valid = entities.filter((e) => entityIsValid(report, e.key)); stats.extracted++; stats.entities += entities.length; stats.valid += valid.length; stats.rejected += report.rejected; for (const i of report.issues) { if (i.level === "error") ctx.log("warn", `validation ${i.key}${i.field ? "." + i.field : ""}: ${i.message}`); if (result.issues.length < 200) result.issues.push(i); } const sampleCap = opts.sampleEntities ?? (opts.dryRun ? 20 : 0); for (const e of valid) if (result.samples.length < sampleCap) result.samples.push(e); const pageType = (raw.meta?.pageType as string | undefined) ?? doc.pageType; const classifier = (raw.meta?.classifier as string | undefined) ?? null; metrics.ingest.inc({ connector: loaded.cfg.id, result: "rejected" }, report.rejected); if (!valid.length) return { ingest: null, valid: 0, rejected: report.rejected, classifier, pageType }; const run: IngestRun = { connectorId: loaded.cfg.id, sourceId: loaded.sourceId, runId: ctx.runId, sourceKind: loaded.cfg.kind, sourcePriority: loaded.cfg.priority, dryRun: Boolean(opts.dryRun) || Boolean(opts.quarantine) }; const ref: IngestDocRef = { documentId: doc.id, url: raw.finalUrl || doc.url, pageType, fetchedAt: raw.fetchedAt }; const ingest = await ingestEntities(run, valid, ref); sumIngest(stats, ingest); if (!opts.dryRun && !opts.quarantine) { metrics.ingest.inc({ connector: loaded.cfg.id, result: "created" }, ingest.created); metrics.ingest.inc({ connector: loaded.cfg.id, result: "updated" }, ingest.updated); metrics.ingest.inc({ connector: loaded.cfg.id, result: "merged" }, ingest.merged); metrics.ingest.inc({ connector: loaded.cfg.id, result: "unchanged" }, Math.max(0, valid.length - ingest.created - ingest.updated - ingest.merged)); metrics.events.inc(undefined, ingest.events); } return { ingest, valid: valid.length, rejected: report.rejected, classifier, pageType }; } /** Fetch + archive + fingerprint + (maybe) extract one due document. Never throws. */ async function processDocument(ctx: RunContext, loaded: LoadedConnector, doc: DocumentRow, stats: RunStats, result: RunResult, opts: RunOptions): Promise { const t0 = Date.now(); const { connector, cfg } = loaded; const dryRun = Boolean(opts.dryRun); const force = Boolean(opts.force); try { if (!force && (doc.etag || doc.lastModified)) ctx.validators.set(doc.url, { etag: doc.etag, lastModified: doc.lastModified }); ctx.hints.set(doc.url, { errorCount: doc.errorCount ?? 0 }); const raw = await connector.fetch(ctx, docToDiscovered(doc)); ctx.validators.delete(doc.url); ctx.hints.delete(doc.url); stats.fetched++; if (raw.error || raw.status >= 400) { const plan = await recordFetch(doc, raw, { contentHash: null, changed: false, storageKey: null, title: null }, cfg.schedule, { dryRun }); stats.failed++; if (plan.quarantined) stats.quarantined++; ctx.log(raw.error?.code === "robots_disallow" ? "debug" : "warn", `${doc.url} → ${raw.error ? raw.error.code : "HTTP " + raw.status}${plan.quarantined ? " (quarantined)" : ""} next=${plan.nextCheck ?? "never"}`); return { kind: "failed", ms: Date.now() - t0 }; } if (raw.notModified) { const plan = await recordFetch(doc, raw, { contentHash: doc.contentHash, changed: false, storageKey: doc.storageKey, title: null }, cfg.schedule, { dryRun }); stats.notModified++; ctx.log("debug", `${doc.url} → 304 next=${plan.nextCheck} (${plan.reason})`); return { kind: "notModified", ms: Date.now() - t0 }; } const hash = fingerprintOf(raw); let changed = hash !== doc.contentHash; const isHtml = !raw.contentType || /html|xml/i.test(raw.contentType); const title = isHtml && raw.text ? pageTitle(raw.text) : null; // Markup-only change guard: the fingerprint moved but the visible text did not → no re-extraction, no version. let markupOnly = false; if (changed && doc.contentHash && doc.storageKey && raw.text) { try { const old = await getRaw(doc.storageKey); if (old) { const oldText = textProjection(old.body, old.contentType ?? doc.contentType); if (oldText.length > 0 && oldText === textProjection(raw.text, raw.contentType)) { markupOnly = true; changed = false; } } } catch { /* storage hiccup: fall through and treat as changed */ } } // archive raw body (dedupe by hash; unchanged pages reuse the existing key). Pages carrying an X-Robots-Tag / // `noarchive` (or `noindex`) directive are never archived — facts are still extracted below. let storageKey: string | null = doc.storageKey; const noarchive = raw.meta?.noarchive === true; if (!dryRun && (changed || !doc.storageKey) && !noarchive) { try { storageKey = (await putRaw(cfg.id, doc.id, hash, raw.body, raw.contentType, raw.finalUrl)).key; } catch (e) { ctx.log("warn", `archive failed for ${doc.url}: ${(e as Error).message}`); } } else if (noarchive && (changed || !doc.storageKey)) { ctx.log("debug", `${doc.url} not archived (${String(raw.meta?.robotsDirective ?? "noarchive")})`); } const skip = markupOnly || shouldSkipExtraction({ newHash: hash, storedHash: doc.contentHash, force, extractorVersion: loaded.extractorVersion, storedExtractorVersion: doc.extractorVersion, notModified: false }); // markup-only: adopt the new fingerprint (so a stabilised page re-syncs after one fetch) but keep the archived // body — its text projection is identical, which is what versions/diffs are built from. const plan = await recordFetch(doc, raw, { contentHash: hash, changed, storageKey, title }, cfg.schedule, { dryRun }); if (skip) { stats.unchanged++; if (markupOnly) stats.markupOnly++; ctx.log("debug", `${doc.url} unchanged${markupOnly ? " (markup-only)" : ""} (${hash.slice(0, 12)}) next=${plan.nextCheck} (${plan.reason})`); return { kind: "unchanged", ms: Date.now() - t0 }; } if (changed) stats.changed++; let ingestStats: IngestStats | null = null; let extractError: string | null = null; let classifier: string | null = null; let pageType: string | null = null; let validCount = 0; try { const r = await extractAndIngest(ctx, loaded, doc, raw, stats, result, opts); ingestStats = r.ingest; classifier = r.classifier; pageType = r.pageType; validCount = r.valid; } catch (e) { extractError = `extract: ${(e as Error).message}`.slice(0, 500); ctx.log("error", `${doc.url} ${extractError}`); } await updateExtraction(doc.id, { entityRefs: ingestStats?.refs ?? [], extractOk: extractError === null, extractCount: validCount, extractorVersion: loaded.extractorVersion, error: extractError, pageType, classifier }, dryRun); if (changed || !doc.contentHash || force) { await recordVersion(doc, raw, hash, storageKey, ingestStats?.changes ?? [], { dryRun, runId: ctx.runId, extractorVersion: loaded.extractorVersion }); } ctx.log("info", `${doc.url} → ${raw.status} ${raw.fetcher} L${raw.level} ${raw.durationMs}ms ${changed ? "CHANGED" : "same"} · ${validCount} entities${ingestStats ? ` (+${ingestStats.created} new, ${ingestStats.updated} upd, ${ingestStats.events} events)` : ""} next=${plan.nextCheck ?? "never"}`); return { kind: changed ? "changed" : "fetched", ms: Date.now() - t0 }; } catch (e) { ctx.validators.delete(doc.url); ctx.hints.delete(doc.url); stats.failed++; ctx.log("error", `${doc.url} unexpected: ${(e as Error).stack ?? (e as Error).message}`); try { if (!dryRun) await updateExtraction(doc.id, { entityRefs: [], extractOk: false, extractCount: 0, extractorVersion: loaded.extractorVersion, error: `pipeline: ${(e as Error).message}`.slice(0, 500) }); } catch { /* ignore */ } return { kind: "failed", ms: Date.now() - t0 }; } } /** Re-run extraction from the archived body (no network). */ async function reprocessDocument(ctx: RunContext, loaded: LoadedConnector, doc: DocumentRow, stats: RunStats, result: RunResult, opts: RunOptions): Promise { const t0 = Date.now(); const dryRun = Boolean(opts.dryRun); try { if (!doc.storageKey) { ctx.log("debug", `${doc.url}: no archived body`); return { kind: "unchanged", ms: Date.now() - t0 }; } const stored = await getRaw(doc.storageKey); if (!stored) { stats.failed++; ctx.log("warn", `${doc.url}: archived body ${doc.storageKey} missing`); return { kind: "failed", ms: Date.now() - t0 }; } const contentType = stored.contentType ?? doc.contentType; const isText = !contentType || /text|json|xml|javascript|html|csv|markdown/i.test(contentType); const { group } = parseDiscoveredFrom(doc.discoveredFrom); const raw: RawDocument = { url: doc.url, finalUrl: doc.url, fetchedAt: doc.lastFetched ?? new Date().toISOString(), status: doc.statusCode ?? 200, contentType, body: stored.body, text: isText ? stored.body.toString("utf8") : "", headers: {}, etag: doc.etag, lastModified: doc.lastModified, notModified: false, fetcher: "cache", level: doc.fetchLevel as RawDocument["level"], durationMs: 0, credits: 0, group, pageType: doc.pageType !== "unknown" ? (doc.pageType as RawDocument["pageType"]) : undefined }; const r = await extractAndIngest(ctx, loaded, doc, raw, stats, result, opts); await updateExtraction(doc.id, { entityRefs: r.ingest?.refs ?? [], extractOk: true, extractCount: r.valid, extractorVersion: loaded.extractorVersion, error: null, pageType: r.pageType, classifier: r.classifier }, dryRun); ctx.log("info", `${doc.url} reprocessed → ${r.valid} entities${r.ingest ? ` (+${r.ingest.created} new, ${r.ingest.updated} upd, ${r.ingest.events} events)` : ""}`); return { kind: "fetched", ms: Date.now() - t0 }; } catch (e) { stats.failed++; ctx.log("error", `${doc.url} reprocess failed: ${(e as Error).message}`); return { kind: "failed", ms: Date.now() - t0 }; } } /* ---------- the run ---------- */ export async function runConnector(connectorId: string, opts: RunOptions): Promise { const loaded = requireConnector(connectorId); const { cfg, connector } = loaded; const dryRun = Boolean(opts.dryRun); const runId = newId("run"); const startedAt = new Date().toISOString(); const stats = emptyStats(); const result: RunResult = { runId, connectorId, task: opts.task, status: "ok", stats, error: null, startedAt, finishedAt: startedAt, durationMs: 0, samples: [], issues: [], health: "ok" }; const ctx = createRunContext(loaded, { runId, dryRun, onLog: opts.onLog, logLevel: opts.logLevel }); const db = getDb(); activeRuns.add(runId); let fatal: string | null = null; const durations: number[] = []; // quarantine mode (connectors.quarantine): everything runs, nothing is published let quarantine = Boolean(opts.quarantine); let previouslyDiscovered: number | null = null; if (!dryRun) { const row = (await db.execute<{ quarantine: boolean; last_discovered: number | null }>(sql`select quarantine, last_discovered from connectors where id = ${connectorId}`))[0]; if (row?.quarantine) quarantine = true; previouslyDiscovered = row?.last_discovered == null ? null : Number(row.last_discovered); } opts = { ...opts, quarantine }; if (!dryRun) await db.insert(connectorRuns).values({ id: runId, connectorId, task: opts.task, startedAt, status: "running", quarantined: quarantine }); ctx.log("info", `run ${runId} task=${opts.task}${opts.group ? ` group=${opts.group}` : ""}${opts.limit ? ` limit=${opts.limit}` : ""}${dryRun ? " DRY-RUN" : ""}${quarantine ? " QUARANTINE (ingest rolled back)" : ""}${opts.force ? " force" : ""}`); try { // ---- discover let targetIds: string[] | undefined; if (opts.urls?.length) { const manual: DiscoveredUrl[] = opts.urls.map((u) => ({ url: u, group: opts.group ?? "manual", priority: 100, discoveredFrom: "manual" })); const reg = await registerDiscovered(loaded, manual, { dryRun, preserveGroup: true }); stats.discovered += manual.length; stats.registered += reg.inserted + reg.updated; targetIds = reg.ids; } else if (opts.task === "discover" || opts.task === "full") { const urls = await connector.discover(ctx); stats.discovered = urls.length; const filtered = opts.group ? urls.filter((u) => u.group === opts.group) : urls; const reg = await registerDiscovered(loaded, filtered, { dryRun }); stats.registered = reg.inserted + reg.updated; ctx.log("info", `discovered ${urls.length} urls (${reg.inserted} new, ${reg.updated} known, ${reg.skipped} skipped)${opts.group ? ` · group ${opts.group}: ${filtered.length}` : ""}`); if (!dryRun) await ctx.setState("lastDiscoverAt", new Date().toISOString()); if (dryRun && opts.task === "full") { // nothing was persisted: crawl the freshly discovered urls directly (priority order) targetIds = undefined; const sorted = [...filtered].sort((a, b) => (b.priority ?? 0) - (a.priority ?? 0)).slice(0, opts.limit ?? 10); await crawlVirtual(ctx, loaded, sorted, stats, result, opts, durations); } } // ---- crawl / reprocess if (opts.task === "crawl" || opts.task === "reprocess" || (opts.task === "full" && !dryRun) || (opts.urls?.length && opts.task !== "discover")) { let docs = await dueDocuments(connectorId, { group: opts.group, limit: opts.limit, force: Boolean(opts.force) || Boolean(targetIds) || opts.task === "reprocess", ids: targetIds }); // `reprocess --stale`: only documents extracted with an older parser version (parser fix → controlled re-extraction) if (opts.task === "reprocess" && opts.staleOnly) docs = docs.filter((d) => d.extractorVersion !== loaded.extractorVersion); stats.selected = docs.length; ctx.log("info", `${docs.length} document(s) selected${opts.task === "reprocess" ? " for reprocessing" : ""}`); const limit = pLimit(cfg.fetch.concurrency); let processed = 0; await Promise.all( docs.map((doc) => limit(async () => { if (isAbortRequested()) return; const out = opts.task === "reprocess" ? await reprocessDocument(ctx, loaded, doc, stats, result, opts) : await processDocument(ctx, loaded, doc, stats, result, opts); durations.push(out.ms); if (++processed % 25 === 0) await ctx.flushLog(); }), ), ); } } catch (e) { fatal = (e as Error).stack ?? (e as Error).message; ctx.log("error", `run failed: ${(e as Error).message}`); } // ---- finish stats.credits = Math.round(ctx.credits * 100) / 100; stats.robotsBlocked = ctx.robotsBlocked; stats.totalMs = durations.reduce((a, b) => a + b, 0); stats.avgMs = durations.length ? Math.round(stats.totalMs / durations.length) : 0; const finishedAt = new Date().toISOString(); result.finishedAt = finishedAt; result.durationMs = Date.parse(finishedAt) - Date.parse(startedAt); const attempted = stats.fetched; result.health = healthFrom(attempted, stats.failed); if (isAbortRequested()) result.status = "aborted"; else if (fatal) result.status = "failed"; else if (attempted > 0 && stats.failed >= attempted) result.status = "failed"; else if (stats.failed > 0 || stats.rejected > 0) result.status = "partial"; else result.status = "ok"; result.error = fatal ? fatal.slice(0, 2000) : isAbortRequested() ? `aborted: ${abortReasonText()}` : null; if (!dryRun) metrics.runs.inc({ status: result.status }); ctx.log("info", `run ${result.status} in ${result.durationMs}ms · fetched=${stats.fetched} 304=${stats.notModified} changed=${stats.changed} unchanged=${stats.unchanged} failed=${stats.failed} entities=${stats.valid}/${stats.entities} created=${stats.created} updated=${stats.updated} events=${stats.events} credits=${stats.credits}${ctx.budgetCapped ? ` budgetCapped=${ctx.budgetCapped}` : ""}`); if (!dryRun) { try { await ctx.flushLog(); await db.update(connectorRuns).set({ finishedAt, status: result.status, stats, error: result.error, log: ctx.logs.slice(-500) }).where(eq(connectorRuns.id, runId)); const groups = await nextCheckByGroup(connectorId); const nexts = groups.map((g) => g.nextCheck).filter((x): x is string => Boolean(x)).map((x) => Date.parse(x)); const nextRunAt = new Date(nexts.length ? Math.min(...nexts) : Date.now() + minIntervalMs(cfg.schedule)).toISOString(); const docStats = await connectorDocStats(connectorId); const blocked = Object.values(ctx.costByFetcher).reduce((n, c) => n + (c.errors ?? 0), 0); const health = connectorHealthFrom({ runHealth: result.health, status: result.status, fetched: stats.fetched, blocked: Math.min(blocked, stats.failed), discovered: stats.discovered, previouslyDiscovered: opts.task === "discover" || opts.task === "full" ? previouslyDiscovered : null, extracted: stats.extracted, entities: stats.entities, docsTotal: docStats.total }); const failedRun = result.status === "failed" || health === "blocked"; const prev = (await db.execute<{ consecutive_failures: number; blocked_since: string | null }>(sql`select consecutive_failures, blocked_since from connectors where id = ${connectorId}`))[0]; const consecutive = failedRun ? Number(prev?.consecutive_failures ?? 0) + 1 : 0; const blockedSince = health === "blocked" ? (prev?.blocked_since ?? finishedAt) : null; // guard rails: a connector failing 5 runs in a row, or a non-full run that suddenly creates hundreds of records, is quarantined const suspiciousSpike = opts.task !== "full" && !quarantine && stats.created >= 500; const autoQuarantine = consecutive >= 5 || suspiciousSpike; await db .update(connectorsTable) .set({ health, lastRunAt: finishedAt, lastStatus: result.status, lastError: result.error, nextRunAt, consecutiveFailures: consecutive, blockedSince, lastDiscovered: opts.task === "discover" || opts.task === "full" ? stats.discovered : undefined, ...(autoQuarantine ? { quarantine: true } : {}), stats: { ...stats, docsTotal: docStats.total, docsDue: docStats.due, docsQuarantined: docStats.quarantined, docsErrors: docStats.errors, docsExtracted: docStats.extracted, lastRunMs: result.durationMs, blockedFetches: blocked, unscopedClaims: stats.unscopedClaims ?? 0, projectsVetoed: stats.projectsVetoed ?? 0 }, updatedAt: finishedAt }) .where(eq(connectorsTable.id, connectorId)); if (autoQuarantine) { const why = suspiciousSpike ? `run ${runId} created ${stats.created} records in a ${opts.task} run — quarantined pending review` : `${consecutive} consecutive failed/blocked runs — quarantined`; ctx.log("error", why); await db.execute(sql`insert into system_alerts (id, level, component, message, details) values (${`alr_${Date.now().toString(36)}_q_${connectorId.replace(/[^a-z0-9]/gi, "")}`}, 'error', ${`connector:${connectorId}`}, ${why}, ${JSON.stringify({ runId, stats })}::jsonb)`).catch(() => undefined); } } catch (e) { console.error(`[${connectorId}] finishing run ${runId} failed: ${(e as Error).message}`); } } activeRuns.delete(runId); return result; } /** Dry-run helper: crawl discovered URLs that were never registered (no DocumentRow exists). */ async function crawlVirtual(ctx: RunContext, loaded: LoadedConnector, urls: DiscoveredUrl[], stats: RunStats, result: RunResult, opts: RunOptions, durations: number[]): Promise { const now = new Date().toISOString(); const limit = pLimit(loaded.cfg.fetch.concurrency); stats.selected = urls.length; await Promise.all( urls.map((u) => limit(async () => { if (isAbortRequested()) return; let id: string; try { id = documentIdFor(u.url); } catch { stats.failed++; return; } const virtual: DocumentRow = { id, connectorId: loaded.cfg.id, sourceId: loaded.sourceId, url: u.url, canonicalUrl: u.url, urlFingerprint: "", pageType: u.pageType ?? "unknown", classifier: null, fetchLevel: u.minLevel ?? 1, priority: u.priority ?? 50, contentHash: null, etag: null, lastModified: null, statusCode: null, contentType: null, sizeBytes: null, title: null, storageKey: null, extractorVersion: null, extractOk: null, extractCount: 0, error: null, errorCount: 0, firstSeen: now, lastFetched: null, lastChanged: null, lastChecked: null, nextCheck: now, changeFrequencyScore: 0.5, fetchCount: 0, changeCount: 0, discoveredFrom: `${u.group}|${u.discoveredFrom ?? ""}`, quarantined: false, entityRefs: [], updatedAt: now }; const out = await processDocument(ctx, loaded, virtual, stats, result, { ...opts, dryRun: true }); durations.push(out.ms); }), ), ); } /** Run every enabled connector sequentially (CLI `run-all`). */ export async function runAll(opts: Omit & { task?: RunTask; ids?: string[] }): Promise { const out: RunResult[] = []; for (const l of loadAllConnectors()) { if (!l.cfg.enabled) continue; if (opts.ids?.length && !opts.ids.includes(l.cfg.id)) continue; if (isAbortRequested()) break; out.push(await runConnector(l.cfg.id, { ...opts, task: opts.task ?? "full" })); } return out; }