SPB Git forge
38commits 1branches 0releases
338.7 MBsize
maindefault branch
3 h agolast push
HTML 53.9% TypeScript 44.5% JavaScript 0.6% SQL 0.5%
28.4 KB · 462 lines typescript
Raw Blame History
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