spb/datacenterindex
Public
HTML 53.9%
TypeScript 44.5%
JavaScript 0.6%
SQL 0.5%
1/**2 * One connector run: discover → register URLs → select due documents → fetch (robots, rate limit, conditional3 * GET, escalation, budgets) → archive raw → fingerprint → skip unchanged → extract → normalize → validate →4 * ingest → document bookkeeping (versions, detected changes, adaptive next_check) → run + connector health.5 *6 * Every per-document step is wrapped: one bad page never kills a run. SIGTERM finishes the current documents7 * and marks the run `aborted`.8 */9import type { DiscoveredUrl, RawDocument } from "@dci/connectors";10import { entityIsValid, pageTitle } from "@dci/connectors";11import type { NormalizedEntity, ValidationIssue } from "@dci/core";12import { contentFingerprint, sha256, newId } from "@dci/core";13import { getDb, connectorRuns, connectors as connectorsTable, eq, inArray, sql } from "@dci/db";14import { loadAllConnectors, requireConnector, type LoadedConnector } from "./configs.js";15import { createRunContext, type LogLevel, type LogLine, type RunContext } from "./context.js";16import { connectorDocStats, documentIdFor, dueDocuments, nextCheckByGroup, recordFetch, recordVersion, registerDiscovered, textProjection, updateExtraction, type DocumentRow } from "./documents.js";17import { ingestEntities } from "./ingest/index.js";18import type { IngestDocRef, IngestRun, IngestStats } from "./ingest/contract.js";19import { metrics } from "./prom.js";20import { connectorHealthFrom, healthFrom, minIntervalMs, parseDiscoveredFrom, shouldSkipExtraction } from "./scheduling.js";21import { getRaw, putRaw } from "./storage.js";2223export type RunTask = "discover" | "crawl" | "full" | "reprocess";2425export interface RunOptions {26 task: RunTask;27 group?: string;28 limit?: number;29 dryRun?: boolean;30 /** crawl exactly these URLs (registered on the fly, forced) */31 urls?: string[];32 /** ignore next_check and content-hash skip */33 force?: boolean;34 onLog?: (line: LogLine, extra?: Record<string, unknown>) => void;35 logLevel?: LogLevel;36 /** collect up to N normalized entities for display (dry runs) */37 sampleEntities?: number;38 /** quarantine: fetch, archive and extract normally but roll the ingest back (set from connectors.quarantine) */39 quarantine?: boolean;40 /** reprocess only documents whose stored extractor version differs from the current one */41 staleOnly?: boolean;42}4344export interface RunStats extends Record<string, number> {45 discovered: number;46 registered: number;47 selected: number;48 fetched: number;49 notModified: number;50 changed: number;51 unchanged: number;52 markupOnly: number;53 failed: number;54 robotsBlocked: number;55 quarantined: number;56 extracted: number;57 entities: number;58 valid: number;59 rejected: number;60 created: number;61 updated: number;62 merged: number;63 pendingMatches: number;64 events: number;65 provenanceRows: number;66 credits: number;67 avgMs: number;68 totalMs: number;69}7071export type RunStatus = "ok" | "partial" | "failed" | "aborted";7273export interface RunResult {74 runId: string;75 connectorId: string;76 task: RunTask;77 status: RunStatus;78 stats: RunStats;79 error: string | null;80 startedAt: string;81 finishedAt: string;82 durationMs: number;83 samples: NormalizedEntity[];84 issues: ValidationIssue[];85 health: "ok" | "degraded" | "failing";86}8788function emptyStats(): RunStats {89 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 };90}9192/* ---------- graceful abort ---------- */93let abortReason: string | null = null;94const activeRuns = new Set<string>();95export function requestAbort(reason = "SIGTERM"): void { abortReason = reason; }96export function isAbortRequested(): boolean { return abortReason !== null; }97export function abortReasonText(): string | null { return abortReason; }98export function activeRunIds(): string[] { return [...activeRuns]; }99/** For tests / long-lived processes that want to resume after a handled abort. */100export function clearAbort(): void { abortReason = null; }101/**102 * Forced shutdown (graceful deadline reached): mark the runs still in flight as aborted in Postgres so they do not103 * linger as `running` until the orphan reconciliation. Best effort; returns the number of runs updated.104 */105export async function abortActiveRunsInDb(reason: string): Promise<number> {106 const ids = [...activeRuns];107 if (!ids.length) return 0;108 await getDb().update(connectorRuns).set({ status: "aborted", finishedAt: new Date().toISOString(), error: `aborted: ${reason}`.slice(0, 2000) }).where(inArray(connectorRuns.id, ids));109 return ids.length;110}111112/* ---------- inline p-limit ---------- */113export function pLimit(concurrency: number): <T>(fn: () => Promise<T>) => Promise<T> {114 const n = Math.max(1, Math.floor(concurrency));115 let active = 0;116 const queue: Array<() => void> = [];117 const next = () => { active--; queue.shift()?.(); };118 return <T>(fn: () => Promise<T>) =>119 new Promise<T>((resolve, reject) => {120 const start = () => { active++; fn().then(resolve, reject).finally(next); };121 if (active < n) start(); else queue.push(start);122 });123}124125/* ---------- helpers ---------- */126/**127 * Runtime-level noise on top of core's contentFingerprint: epoch-like cache busters in URLs (`?t=1789110041117`,128 * `&_=1789110041`, `?v=1700000000000`) that some CMSs stamp into og:url / share links on every render.129 */130export function stripRuntimeNoise(text: string): string {131 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, "");132}133export function fingerprintOf(raw: RawDocument): string {134 return raw.text ? contentFingerprint(stripRuntimeNoise(raw.text), raw.contentType) : sha256(raw.body);135}136137function docToDiscovered(doc: DocumentRow): DiscoveredUrl {138 const { group, from } = parseDiscoveredFrom(doc.discoveredFrom);139 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 };140}141142function sumIngest(stats: RunStats, s: IngestStats): void {143 stats.created += s.created; stats.updated += s.updated; stats.merged += s.merged; stats.pendingMatches += s.pendingMatches; stats.events += s.events; stats.provenanceRows += s.provenanceRows;144 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);145}146147interface DocOutcome { kind: "fetched" | "notModified" | "failed" | "unchanged" | "changed" | "aborted"; ms: number }148149/** Extract → normalize → validate → ingest for one raw document. Returns ingest stats (or null when nothing valid). */150async 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 }> {151 const { connector } = loaded;152 const records = await connector.extract(ctx, raw);153 const entities = await connector.normalize(ctx, records);154 const report = await connector.validate(ctx, entities);155 const valid = entities.filter((e) => entityIsValid(report, e.key));156 stats.extracted++;157 stats.entities += entities.length;158 stats.valid += valid.length;159 stats.rejected += report.rejected;160 for (const i of report.issues) {161 if (i.level === "error") ctx.log("warn", `validation ${i.key}${i.field ? "." + i.field : ""}: ${i.message}`);162 if (result.issues.length < 200) result.issues.push(i);163 }164 const sampleCap = opts.sampleEntities ?? (opts.dryRun ? 20 : 0);165 for (const e of valid) if (result.samples.length < sampleCap) result.samples.push(e);166 const pageType = (raw.meta?.pageType as string | undefined) ?? doc.pageType;167 const classifier = (raw.meta?.classifier as string | undefined) ?? null;168 metrics.ingest.inc({ connector: loaded.cfg.id, result: "rejected" }, report.rejected);169 if (!valid.length) return { ingest: null, valid: 0, rejected: report.rejected, classifier, pageType };170 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) };171 const ref: IngestDocRef = { documentId: doc.id, url: raw.finalUrl || doc.url, pageType, fetchedAt: raw.fetchedAt };172 const ingest = await ingestEntities(run, valid, ref);173 sumIngest(stats, ingest);174 if (!opts.dryRun && !opts.quarantine) {175 metrics.ingest.inc({ connector: loaded.cfg.id, result: "created" }, ingest.created);176 metrics.ingest.inc({ connector: loaded.cfg.id, result: "updated" }, ingest.updated);177 metrics.ingest.inc({ connector: loaded.cfg.id, result: "merged" }, ingest.merged);178 metrics.ingest.inc({ connector: loaded.cfg.id, result: "unchanged" }, Math.max(0, valid.length - ingest.created - ingest.updated - ingest.merged));179 metrics.events.inc(undefined, ingest.events);180 }181 return { ingest, valid: valid.length, rejected: report.rejected, classifier, pageType };182}183184/** Fetch + archive + fingerprint + (maybe) extract one due document. Never throws. */185async function processDocument(ctx: RunContext, loaded: LoadedConnector, doc: DocumentRow, stats: RunStats, result: RunResult, opts: RunOptions): Promise<DocOutcome> {186 const t0 = Date.now();187 const { connector, cfg } = loaded;188 const dryRun = Boolean(opts.dryRun);189 const force = Boolean(opts.force);190 try {191 if (!force && (doc.etag || doc.lastModified)) ctx.validators.set(doc.url, { etag: doc.etag, lastModified: doc.lastModified });192 ctx.hints.set(doc.url, { errorCount: doc.errorCount ?? 0 });193 const raw = await connector.fetch(ctx, docToDiscovered(doc));194 ctx.validators.delete(doc.url);195 ctx.hints.delete(doc.url);196 stats.fetched++;197198 if (raw.error || raw.status >= 400) {199 const plan = await recordFetch(doc, raw, { contentHash: null, changed: false, storageKey: null, title: null }, cfg.schedule, { dryRun });200 stats.failed++;201 if (plan.quarantined) stats.quarantined++;202 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"}`);203 return { kind: "failed", ms: Date.now() - t0 };204 }205 if (raw.notModified) {206 const plan = await recordFetch(doc, raw, { contentHash: doc.contentHash, changed: false, storageKey: doc.storageKey, title: null }, cfg.schedule, { dryRun });207 stats.notModified++;208 ctx.log("debug", `${doc.url} → 304 next=${plan.nextCheck} (${plan.reason})`);209 return { kind: "notModified", ms: Date.now() - t0 };210 }211212 const hash = fingerprintOf(raw);213 let changed = hash !== doc.contentHash;214 const isHtml = !raw.contentType || /html|xml/i.test(raw.contentType);215 const title = isHtml && raw.text ? pageTitle(raw.text) : null;216 // Markup-only change guard: the fingerprint moved but the visible text did not → no re-extraction, no version.217 let markupOnly = false;218 if (changed && doc.contentHash && doc.storageKey && raw.text) {219 try {220 const old = await getRaw(doc.storageKey);221 if (old) {222 const oldText = textProjection(old.body, old.contentType ?? doc.contentType);223 if (oldText.length > 0 && oldText === textProjection(raw.text, raw.contentType)) { markupOnly = true; changed = false; }224 }225 } catch { /* storage hiccup: fall through and treat as changed */ }226 }227228 // archive raw body (dedupe by hash; unchanged pages reuse the existing key). Pages carrying an X-Robots-Tag /229 // <meta name="robots"> `noarchive` (or `noindex`) directive are never archived — facts are still extracted below.230 let storageKey: string | null = doc.storageKey;231 const noarchive = raw.meta?.noarchive === true;232 if (!dryRun && (changed || !doc.storageKey) && !noarchive) {233 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}`); }234 } else if (noarchive && (changed || !doc.storageKey)) {235 ctx.log("debug", `${doc.url} not archived (${String(raw.meta?.robotsDirective ?? "noarchive")})`);236 }237238 const skip = markupOnly || shouldSkipExtraction({ newHash: hash, storedHash: doc.contentHash, force, extractorVersion: loaded.extractorVersion, storedExtractorVersion: doc.extractorVersion, notModified: false });239 // markup-only: adopt the new fingerprint (so a stabilised page re-syncs after one fetch) but keep the archived240 // body — its text projection is identical, which is what versions/diffs are built from.241 const plan = await recordFetch(doc, raw, { contentHash: hash, changed, storageKey, title }, cfg.schedule, { dryRun });242 if (skip) {243 stats.unchanged++;244 if (markupOnly) stats.markupOnly++;245 ctx.log("debug", `${doc.url} unchanged${markupOnly ? " (markup-only)" : ""} (${hash.slice(0, 12)}) next=${plan.nextCheck} (${plan.reason})`);246 return { kind: "unchanged", ms: Date.now() - t0 };247 }248 if (changed) stats.changed++;249250 let ingestStats: IngestStats | null = null;251 let extractError: string | null = null;252 let classifier: string | null = null;253 let pageType: string | null = null;254 let validCount = 0;255 try {256 const r = await extractAndIngest(ctx, loaded, doc, raw, stats, result, opts);257 ingestStats = r.ingest; classifier = r.classifier; pageType = r.pageType; validCount = r.valid;258 } catch (e) {259 extractError = `extract: ${(e as Error).message}`.slice(0, 500);260 ctx.log("error", `${doc.url} ${extractError}`);261 }262 await updateExtraction(doc.id, { entityRefs: ingestStats?.refs ?? [], extractOk: extractError === null, extractCount: validCount, extractorVersion: loaded.extractorVersion, error: extractError, pageType, classifier }, dryRun);263 if (changed || !doc.contentHash || force) {264 await recordVersion(doc, raw, hash, storageKey, ingestStats?.changes ?? [], { dryRun, runId: ctx.runId, extractorVersion: loaded.extractorVersion });265 }266 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"}`);267 return { kind: changed ? "changed" : "fetched", ms: Date.now() - t0 };268 } catch (e) {269 ctx.validators.delete(doc.url);270 ctx.hints.delete(doc.url);271 stats.failed++;272 ctx.log("error", `${doc.url} unexpected: ${(e as Error).stack ?? (e as Error).message}`);273 try {274 if (!dryRun) await updateExtraction(doc.id, { entityRefs: [], extractOk: false, extractCount: 0, extractorVersion: loaded.extractorVersion, error: `pipeline: ${(e as Error).message}`.slice(0, 500) });275 } catch { /* ignore */ }276 return { kind: "failed", ms: Date.now() - t0 };277 }278}279280/** Re-run extraction from the archived body (no network). */281async function reprocessDocument(ctx: RunContext, loaded: LoadedConnector, doc: DocumentRow, stats: RunStats, result: RunResult, opts: RunOptions): Promise<DocOutcome> {282 const t0 = Date.now();283 const dryRun = Boolean(opts.dryRun);284 try {285 if (!doc.storageKey) { ctx.log("debug", `${doc.url}: no archived body`); return { kind: "unchanged", ms: Date.now() - t0 }; }286 const stored = await getRaw(doc.storageKey);287 if (!stored) { stats.failed++; ctx.log("warn", `${doc.url}: archived body ${doc.storageKey} missing`); return { kind: "failed", ms: Date.now() - t0 }; }288 const contentType = stored.contentType ?? doc.contentType;289 const isText = !contentType || /text|json|xml|javascript|html|csv|markdown/i.test(contentType);290 const { group } = parseDiscoveredFrom(doc.discoveredFrom);291 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 };292 const r = await extractAndIngest(ctx, loaded, doc, raw, stats, result, opts);293 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);294 ctx.log("info", `${doc.url} reprocessed → ${r.valid} entities${r.ingest ? ` (+${r.ingest.created} new, ${r.ingest.updated} upd, ${r.ingest.events} events)` : ""}`);295 return { kind: "fetched", ms: Date.now() - t0 };296 } catch (e) {297 stats.failed++;298 ctx.log("error", `${doc.url} reprocess failed: ${(e as Error).message}`);299 return { kind: "failed", ms: Date.now() - t0 };300 }301}302303/* ---------- the run ---------- */304export async function runConnector(connectorId: string, opts: RunOptions): Promise<RunResult> {305 const loaded = requireConnector(connectorId);306 const { cfg, connector } = loaded;307 const dryRun = Boolean(opts.dryRun);308 const runId = newId("run");309 const startedAt = new Date().toISOString();310 const stats = emptyStats();311 const result: RunResult = { runId, connectorId, task: opts.task, status: "ok", stats, error: null, startedAt, finishedAt: startedAt, durationMs: 0, samples: [], issues: [], health: "ok" };312 const ctx = createRunContext(loaded, { runId, dryRun, onLog: opts.onLog, logLevel: opts.logLevel });313 const db = getDb();314 activeRuns.add(runId);315 let fatal: string | null = null;316 const durations: number[] = [];317318 // quarantine mode (connectors.quarantine): everything runs, nothing is published319 let quarantine = Boolean(opts.quarantine);320 let previouslyDiscovered: number | null = null;321 if (!dryRun) {322 const row = (await db.execute<{ quarantine: boolean; last_discovered: number | null }>(sql`select quarantine, last_discovered from connectors where id = ${connectorId}`))[0];323 if (row?.quarantine) quarantine = true;324 previouslyDiscovered = row?.last_discovered == null ? null : Number(row.last_discovered);325 }326 opts = { ...opts, quarantine };327 if (!dryRun) await db.insert(connectorRuns).values({ id: runId, connectorId, task: opts.task, startedAt, status: "running", quarantined: quarantine });328 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" : ""}`);329330 try {331 // ---- discover332 let targetIds: string[] | undefined;333 if (opts.urls?.length) {334 const manual: DiscoveredUrl[] = opts.urls.map((u) => ({ url: u, group: opts.group ?? "manual", priority: 100, discoveredFrom: "manual" }));335 const reg = await registerDiscovered(loaded, manual, { dryRun, preserveGroup: true });336 stats.discovered += manual.length; stats.registered += reg.inserted + reg.updated;337 targetIds = reg.ids;338 } else if (opts.task === "discover" || opts.task === "full") {339 const urls = await connector.discover(ctx);340 stats.discovered = urls.length;341 const filtered = opts.group ? urls.filter((u) => u.group === opts.group) : urls;342 const reg = await registerDiscovered(loaded, filtered, { dryRun });343 stats.registered = reg.inserted + reg.updated;344 ctx.log("info", `discovered ${urls.length} urls (${reg.inserted} new, ${reg.updated} known, ${reg.skipped} skipped)${opts.group ? ` · group ${opts.group}: ${filtered.length}` : ""}`);345 if (!dryRun) await ctx.setState("lastDiscoverAt", new Date().toISOString());346 if (dryRun && opts.task === "full") {347 // nothing was persisted: crawl the freshly discovered urls directly (priority order)348 targetIds = undefined;349 const sorted = [...filtered].sort((a, b) => (b.priority ?? 0) - (a.priority ?? 0)).slice(0, opts.limit ?? 10);350 await crawlVirtual(ctx, loaded, sorted, stats, result, opts, durations);351 }352 }353354 // ---- crawl / reprocess355 if (opts.task === "crawl" || opts.task === "reprocess" || (opts.task === "full" && !dryRun) || (opts.urls?.length && opts.task !== "discover")) {356 let docs = await dueDocuments(connectorId, { group: opts.group, limit: opts.limit, force: Boolean(opts.force) || Boolean(targetIds) || opts.task === "reprocess", ids: targetIds });357 // `reprocess --stale`: only documents extracted with an older parser version (parser fix → controlled re-extraction)358 if (opts.task === "reprocess" && opts.staleOnly) docs = docs.filter((d) => d.extractorVersion !== loaded.extractorVersion);359 stats.selected = docs.length;360 ctx.log("info", `${docs.length} document(s) selected${opts.task === "reprocess" ? " for reprocessing" : ""}`);361 const limit = pLimit(cfg.fetch.concurrency);362 let processed = 0;363 await Promise.all(364 docs.map((doc) =>365 limit(async () => {366 if (isAbortRequested()) return;367 const out = opts.task === "reprocess" ? await reprocessDocument(ctx, loaded, doc, stats, result, opts) : await processDocument(ctx, loaded, doc, stats, result, opts);368 durations.push(out.ms);369 if (++processed % 25 === 0) await ctx.flushLog();370 }),371 ),372 );373 }374 } catch (e) {375 fatal = (e as Error).stack ?? (e as Error).message;376 ctx.log("error", `run failed: ${(e as Error).message}`);377 }378379 // ---- finish380 stats.credits = Math.round(ctx.credits * 100) / 100;381 stats.robotsBlocked = ctx.robotsBlocked;382 stats.totalMs = durations.reduce((a, b) => a + b, 0);383 stats.avgMs = durations.length ? Math.round(stats.totalMs / durations.length) : 0;384 const finishedAt = new Date().toISOString();385 result.finishedAt = finishedAt;386 result.durationMs = Date.parse(finishedAt) - Date.parse(startedAt);387 const attempted = stats.fetched;388 result.health = healthFrom(attempted, stats.failed);389 if (isAbortRequested()) result.status = "aborted";390 else if (fatal) result.status = "failed";391 else if (attempted > 0 && stats.failed >= attempted) result.status = "failed";392 else if (stats.failed > 0 || stats.rejected > 0) result.status = "partial";393 else result.status = "ok";394 result.error = fatal ? fatal.slice(0, 2000) : isAbortRequested() ? `aborted: ${abortReasonText()}` : null;395 if (!dryRun) metrics.runs.inc({ status: result.status });396 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}` : ""}`);397398 if (!dryRun) {399 try {400 await ctx.flushLog();401 await db.update(connectorRuns).set({ finishedAt, status: result.status, stats, error: result.error, log: ctx.logs.slice(-500) }).where(eq(connectorRuns.id, runId));402 const groups = await nextCheckByGroup(connectorId);403 const nexts = groups.map((g) => g.nextCheck).filter((x): x is string => Boolean(x)).map((x) => Date.parse(x));404 const nextRunAt = new Date(nexts.length ? Math.min(...nexts) : Date.now() + minIntervalMs(cfg.schedule)).toISOString();405 const docStats = await connectorDocStats(connectorId);406 const blocked = Object.values(ctx.costByFetcher).reduce((n, c) => n + (c.errors ?? 0), 0);407 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 });408 const failedRun = result.status === "failed" || health === "blocked";409 const prev = (await db.execute<{ consecutive_failures: number; blocked_since: string | null }>(sql`select consecutive_failures, blocked_since from connectors where id = ${connectorId}`))[0];410 const consecutive = failedRun ? Number(prev?.consecutive_failures ?? 0) + 1 : 0;411 const blockedSince = health === "blocked" ? (prev?.blocked_since ?? finishedAt) : null;412 // guard rails: a connector failing 5 runs in a row, or a non-full run that suddenly creates hundreds of records, is quarantined413 const suspiciousSpike = opts.task !== "full" && !quarantine && stats.created >= 500;414 const autoQuarantine = consecutive >= 5 || suspiciousSpike;415 await db416 .update(connectorsTable)417 .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 })418 .where(eq(connectorsTable.id, connectorId));419 if (autoQuarantine) {420 const why = suspiciousSpike ? `run ${runId} created ${stats.created} records in a ${opts.task} run — quarantined pending review` : `${consecutive} consecutive failed/blocked runs — quarantined`;421 ctx.log("error", why);422 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);423 }424 } catch (e) {425 console.error(`[${connectorId}] finishing run ${runId} failed: ${(e as Error).message}`);426 }427 }428 activeRuns.delete(runId);429 return result;430}431432/** Dry-run helper: crawl discovered URLs that were never registered (no DocumentRow exists). */433async function crawlVirtual(ctx: RunContext, loaded: LoadedConnector, urls: DiscoveredUrl[], stats: RunStats, result: RunResult, opts: RunOptions, durations: number[]): Promise<void> {434 const now = new Date().toISOString();435 const limit = pLimit(loaded.cfg.fetch.concurrency);436 stats.selected = urls.length;437 await Promise.all(438 urls.map((u) =>439 limit(async () => {440 if (isAbortRequested()) return;441 let id: string;442 try { id = documentIdFor(u.url); } catch { stats.failed++; return; }443 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 };444 const out = await processDocument(ctx, loaded, virtual, stats, result, { ...opts, dryRun: true });445 durations.push(out.ms);446 }),447 ),448 );449}450451/** Run every enabled connector sequentially (CLI `run-all`). */452export async function runAll(opts: Omit<RunOptions, "task"> & { task?: RunTask; ids?: string[] }): Promise<RunResult[]> {453 const out: RunResult[] = [];454 for (const l of loadAllConnectors()) {455 if (!l.cfg.enabled) continue;456 if (opts.ids?.length && !opts.ids.includes(l.cfg.id)) continue;457 if (isAbortRequested()) break;458 out.push(await runConnector(l.cfg.id, { ...opts, task: opts.task ?? "full" }));459 }460 return out;461}462