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%
16.6 KB · 307 lines typescript
Raw Blame History
1/**2 * URL registry (`documents`) + change tracking (`document_versions`, ClickHouse page_changes).3 *4 *  - registerDiscovered: upsert discovered URLs (id = stableId("document", canonicalUrl)); new rows are due now.5 *  - dueDocuments: what to fetch next (priority order, group filter, respects quarantine / next_check).6 *  - recordFetch: bookkeeping after a fetch — conditional-GET validators, content hash, adaptive next_check.7 *  - recordVersion: a document_versions row when the content fingerprint changed (line diff of text8 *    projections, significance from field-level detectedChanges or diff ratio).9 */10import type { DiscoveredUrl, RawDocument } from "@dci/connectors";11import type { DetectedChange } from "@dci/core";12import { canonicalUrl, urlFingerprint, stableId, sha256, htmlToText, diffLines, type LineDiff } from "@dci/core";13import { getDb, documents, documentVersions, sql, eq, and, inArray } from "@dci/db";14import { chInsert } from "@dci/db/clickhouse";15import type { LoadedConnector } from "./configs.js";16import { getRaw } from "./storage.js";17import { baseIntervalMs, encodeDiscoveredFrom, isGone, nextCheckAfter, parseDiscoveredFrom, updateChangeScore, versionSignificance } from "./scheduling.js";1819export type DocumentRow = typeof documents.$inferSelect;20export type DocumentVersionRow = typeof documentVersions.$inferSelect;2122export function documentIdFor(url: string): string { return stableId("document", canonicalUrl(url)); }2324export interface RegisterResult { inserted: number; updated: number; skipped: number; ids: string[] }2526/** Upsert discovered URLs. Existing rows keep their scheduling state; priority/pageType/fetchLevel only ratchet up. */27export async function registerDiscovered(loaded: LoadedConnector, urls: DiscoveredUrl[], opts: { dryRun?: boolean; now?: Date; /** keep the existing group/origin of known documents (manual --url runs) */ preserveGroup?: boolean } = {}): Promise<RegisterResult> {28  const now = (opts.now ?? new Date()).toISOString();29  const res: RegisterResult = { inserted: 0, updated: 0, skipped: 0, ids: [] };30  type Row = typeof documents.$inferInsert;31  const rows = new Map<string, Row & { lastmod: string | null }>();32  for (const u of urls) {33    let canon: string;34    try { canon = canonicalUrl(u.url); } catch { res.skipped++; continue; }35    if (!/^https?:$/.test(new URL(canon).protocol)) { res.skipped++; continue; }36    const id = stableId("document", canon);37    if (rows.has(id)) continue;38    rows.set(id, {39      id,40      connectorId: loaded.cfg.id,41      sourceId: loaded.sourceId,42      url: u.url,43      canonicalUrl: canon,44      urlFingerprint: urlFingerprint(u.url),45      pageType: u.pageType ?? "unknown",46      fetchLevel: u.minLevel ?? 1,47      priority: Math.max(0, Math.min(100, Math.round(u.priority ?? 50))),48      discoveredFrom: encodeDiscoveredFrom(u.group ?? "default", u.discoveredFrom ?? null),49      nextCheck: now,50      firstSeen: now,51      updatedAt: now,52      lastmod: u.lastmod ?? null,53    });54  }55  res.ids = [...rows.keys()];56  if (opts.dryRun || !rows.size) { res.inserted = rows.size; return res; }57  const db = getDb();58  const all = [...rows.values()];59  for (let i = 0; i < all.length; i += 200) {60    const chunk = all.slice(i, i + 200).map(({ lastmod: _l, ...r }) => r);61    const ret = await db62      .insert(documents)63      .values(chunk)64      .onConflictDoUpdate({65        target: documents.id,66        set: {67          url: sql`excluded.url`,68          pageType: sql`case when ${documents.pageType} = 'unknown' then excluded.page_type else ${documents.pageType} end`,69          priority: sql`greatest(${documents.priority}, excluded.priority)`,70          fetchLevel: sql`greatest(${documents.fetchLevel}, excluded.fetch_level)`,71          discoveredFrom: opts.preserveGroup ? sql`coalesce(${documents.discoveredFrom}, excluded.discovered_from)` : sql`excluded.discovered_from`,72          updatedAt: now,73        },74      })75      .returning({ id: documents.id, inserted: sql<boolean>`(xmax = 0)` });76    for (const r of ret) { if (r.inserted) res.inserted++; else res.updated++; }77  }78  // sitemap lastmod newer than our last fetch → the page is due now (cheap signal, saves a full cycle)79  const withLastmod = all.filter((r) => r.lastmod && !Number.isNaN(Date.parse(r.lastmod)));80  if (withLastmod.length) {81    const byId = new Map(withLastmod.map((r) => [r.id, Date.parse(r.lastmod!)]));82    const existing = await db.select({ id: documents.id, lastFetched: documents.lastFetched, nextCheck: documents.nextCheck, quarantined: documents.quarantined }).from(documents).where(inArray(documents.id, [...byId.keys()]));83    const due = existing.filter((e) => !e.quarantined && e.lastFetched && byId.get(e.id)! > Date.parse(e.lastFetched) && (!e.nextCheck || Date.parse(e.nextCheck) > Date.now())).map((e) => e.id);84    if (due.length) await db.update(documents).set({ nextCheck: now, updatedAt: now }).where(inArray(documents.id, due));85  }86  return res;87}8889export interface DueOptions { group?: string; limit?: number; /** ignore next_check (still skips quarantined) */ force?: boolean; ids?: string[]; includeQuarantined?: boolean }9091export async function dueDocuments(connectorId: string, opts: DueOptions = {}): Promise<DocumentRow[]> {92  const db = getDb();93  const conds = [eq(documents.connectorId, connectorId)];94  if (!opts.includeQuarantined) conds.push(eq(documents.quarantined, false));95  if (!opts.force) conds.push(sql`(${documents.nextCheck} is null or ${documents.nextCheck} <= now())`);96  if (opts.group) conds.push(sql`${documents.discoveredFrom} like ${opts.group.replace(/[%_]/g, "\\$&") + "|%"}`);97  if (opts.ids?.length) conds.push(inArray(documents.id, opts.ids));98  const q = db.select().from(documents).where(and(...conds)).orderBy(sql`${documents.priority} desc`, sql`${documents.nextCheck} asc nulls first`, sql`${documents.lastChecked} asc nulls first`);99  return opts.limit ? q.limit(opts.limit) : q;100}101102export async function countDue(connectorId: string, group?: string): Promise<number> {103  const rows = await getDb().execute<{ n: number }>(104    group105      ? sql`select count(*)::int as n from documents where connector_id = ${connectorId} and quarantined = false and (next_check is null or next_check <= now()) and discovered_from like ${group + "|%"}`106      : sql`select count(*)::int as n from documents where connector_id = ${connectorId} and quarantined = false and (next_check is null or next_check <= now())`,107  );108  return Number(rows[0]?.n ?? 0);109}110111export async function getDocument(idOrUrl: string): Promise<DocumentRow | null> {112  const id = /^https?:\/\//i.test(idOrUrl) ? documentIdFor(idOrUrl) : idOrUrl;113  const rows = await getDb().select().from(documents).where(eq(documents.id, id)).limit(1);114  return rows[0] ?? null;115}116117/** The stored content hash for a URL — powers ConnectorContext.isKnownUnchanged. */118export async function storedHashFor(url: string): Promise<string | null> {119  let id: string;120  try { id = documentIdFor(url); } catch { return null; }121  const rows = await getDb().select({ h: documents.contentHash }).from(documents).where(eq(documents.id, id)).limit(1);122  return rows[0]?.h ?? null;123}124125export interface FetchOutcome {126  /** content fingerprint of the new body (null on error / 304) */127  contentHash: string | null;128  changed: boolean;129  storageKey: string | null;130  title: string | null;131}132133export interface FetchBookkeeping {134  nextCheck: string | null;135  quarantined: boolean;136  changeScore: number;137  errorCount: number;138  error: string | null;139  reason: string;140  changed: boolean;141  notModified: boolean;142  failed: boolean;143}144145/** Compute what recordFetch will write (pure; exported for tests and dry runs). */146export function planFetchBookkeeping(doc: Pick<DocumentRow, "changeFrequencyScore" | "errorCount" | "discoveredFrom">, raw: RawDocument, outcome: FetchOutcome, schedule: Record<string, string>, now = new Date()): FetchBookkeeping {147  const group = parseDiscoveredFrom(doc.discoveredFrom).group;148  const baseMs = baseIntervalMs(schedule, group);149  const hadError = Boolean(raw.error);150  const httpFail = !hadError && raw.status >= 400;151  const failed = hadError || httpFail;152  const errorCount = failed ? (doc.errorCount ?? 0) + 1 : 0;153  const error = hadError ? `${raw.error!.code}: ${raw.error!.message}`.slice(0, 500) : httpFail ? `HTTP ${raw.status}` : null;154  const changeScore = failed ? doc.changeFrequencyScore : updateChangeScore(doc.changeFrequencyScore, outcome.changed);155  const retryAfterMs = typeof raw.meta?.retryAfterMs === "number" ? raw.meta.retryAfterMs : null;156  const nc = nextCheckAfter({ now, baseMs, changeScore, statusCode: hadError ? null : raw.status, consecutiveErrors: errorCount, hadError, retryAfterMs });157  return { nextCheck: nc.nextCheck ? nc.nextCheck.toISOString() : null, quarantined: nc.quarantined, changeScore, errorCount, error, reason: nc.reason, changed: outcome.changed && !failed, notModified: raw.notModified, failed };158}159160/** Persist post-fetch state. Returns the bookkeeping decisions (also in dry-run, where nothing is written). */161export async function recordFetch(doc: DocumentRow, raw: RawDocument, outcome: FetchOutcome, schedule: Record<string, string>, opts: { dryRun?: boolean; now?: Date } = {}): Promise<FetchBookkeeping> {162  const now = opts.now ?? new Date();163  const plan = planFetchBookkeeping(doc, raw, outcome, schedule, now);164  if (opts.dryRun) return plan;165  const nowIso = now.toISOString();166  const set: Partial<typeof documents.$inferInsert> = {167    lastChecked: nowIso,168    nextCheck: plan.nextCheck,169    quarantined: plan.quarantined,170    changeFrequencyScore: plan.changeScore,171    errorCount: plan.errorCount,172    error: plan.error,173    fetchCount: sql`${documents.fetchCount} + 1` as unknown as number,174    updatedAt: nowIso,175  };176  if (!raw.error) {177    set.statusCode = raw.status;178    if (raw.finalUrl && raw.finalUrl !== doc.url && /^https?:\/\//.test(raw.finalUrl)) set.url = raw.finalUrl;179  }180  if (!plan.failed) {181    set.lastFetched = nowIso;182    if (!raw.notModified) {183      set.contentType = raw.contentType;184      set.sizeBytes = raw.body.length;185      if (raw.etag) set.etag = raw.etag;186      if (raw.lastModified) set.lastModified = raw.lastModified;187      if (outcome.title) set.title = outcome.title.slice(0, 500);188      if (outcome.contentHash) set.contentHash = outcome.contentHash;189      if (outcome.storageKey) set.storageKey = outcome.storageKey;190      if (plan.changed) { set.lastChanged = nowIso; set.changeCount = sql`${documents.changeCount} + 1` as unknown as number; }191    } else {192      // 304: validators may be refreshed by the server193      if (raw.etag) set.etag = raw.etag;194      if (raw.lastModified) set.lastModified = raw.lastModified;195    }196  } else if (isGone(raw.status) && plan.quarantined) {197    set.error = `gone: HTTP ${raw.status} ×${plan.errorCount}`;198  }199  await getDb().update(documents).set(set).where(eq(documents.id, doc.id));200  return plan;201}202203export interface ExtractionUpdate { entityRefs: Array<{ type: string; id: string }>; extractOk: boolean; extractCount: number; extractorVersion: string; error?: string | null; pageType?: string | null; classifier?: string | null }204205export async function updateExtraction(docId: string, u: ExtractionUpdate, dryRun = false): Promise<void> {206  if (dryRun) return;207  const set: Partial<typeof documents.$inferInsert> = { entityRefs: u.entityRefs, extractOk: u.extractOk, extractCount: u.extractCount, extractorVersion: u.extractorVersion, updatedAt: new Date().toISOString() };208  if (u.error !== undefined) set.error = u.error;209  if (u.pageType && u.pageType !== "unknown") set.pageType = u.pageType;210  if (u.classifier) set.classifier = u.classifier;211  await getDb().update(documents).set(set).where(eq(documents.id, docId));212}213214/** Plain-text projection used for diffs (HTML → text; PDFs/binaries are not diffed line by line). */215export function textProjection(body: Buffer | string, contentType: string | null | undefined): string {216  const ct = (contentType ?? "").toLowerCase();217  if (/pdf|image|octet-stream|zip/.test(ct)) return "";218  const s = typeof body === "string" ? body : body.toString("utf8");219  if (!ct || /html|xml/.test(ct) || /<html/i.test(s.slice(0, 2000))) return htmlToText(s);220  return s;221}222223export interface VersionResult { versionId: string; diff: LineDiff | null; significance: number }224225/**226 * Record a new content version. `prev` is the document row as it was BEFORE recordFetch (old hash / storage key).227 * The old text is loaded from object storage when available; `detectedChanges` come from the ingest layer.228 */229export async function recordVersion(prev: DocumentRow, raw: RawDocument, newHash: string, storageKey: string | null, detectedChanges: DetectedChange[], opts: { dryRun?: boolean; now?: Date; runId?: string; extractorVersion?: string } = {}): Promise<VersionResult> {230  const now = opts.now ?? new Date();231  let diff: LineDiff | null = null;232  if (prev.storageKey && prev.contentHash && prev.contentHash !== newHash) {233    try {234      const old = await getRaw(prev.storageKey);235      if (old) diff = diffLines(textProjection(old.body, old.contentType ?? prev.contentType), textProjection(raw.text || raw.body, raw.contentType), 200);236    } catch { diff = null; }237  }238  const significance = versionSignificance(detectedChanges, diff?.ratio ?? 0);239  const versionId = `ver_${sha256(`${prev.id}:${newHash}:${now.toISOString()}`).slice(0, 16)}`;240  if (opts.dryRun) return { versionId, diff, significance };241  await getDb()242    .insert(documentVersions)243    .values({244      id: versionId,245      documentId: prev.id,246      contentHash: newHash,247      fetchedAt: now.toISOString(),248      fetchLevel: raw.level,249      statusCode: raw.status,250      sizeBytes: raw.body.length,251      storageKey,252      diffSummary: diff ? { addedCount: diff.addedCount, removedCount: diff.removedCount, ratio: Math.round(diff.ratio * 10_000) / 10_000, added: diff.added.slice(0, 60), removed: diff.removed.slice(0, 60) } : null,253      detectedChanges: detectedChanges as unknown as Array<Record<string, unknown>>,254      significance,255      runId: opts.runId ?? null,256      extractorVersion: opts.extractorVersion ?? null,257    })258    .onConflictDoNothing();259  if (prev.contentHash) {260    await chInsert("page_changes", [261      {262        ts: now.toISOString().replace("T", " ").replace("Z", ""),263        connector_id: prev.connectorId,264        document_id: prev.id,265        url: prev.url,266        old_hash: prev.contentHash,267        new_hash: newHash,268        added_lines: diff?.addedCount ?? 0,269        removed_lines: diff?.removedCount ?? 0,270        ratio: diff?.ratio ?? 0,271        significance,272        fields: [...new Set(detectedChanges.map((c) => c.field))],273        event_types: [...new Set(detectedChanges.map((c) => c.eventType))],274      },275    ]).catch(() => undefined);276  }277  return { versionId, diff, significance };278}279280export async function latestVersion(documentId: string): Promise<DocumentVersionRow | null> {281  const rows = await getDb().select().from(documentVersions).where(eq(documentVersions.documentId, documentId)).orderBy(sql`${documentVersions.fetchedAt} desc`).limit(1);282  return rows[0] ?? null;283}284285export async function connectorDocStats(connectorId: string): Promise<{ total: number; due: number; quarantined: number; errors: number; extracted: number }> {286  const rows = await getDb().execute<{ total: number; due: number; quarantined: number; errors: number; extracted: number }>(sql`287    select count(*)::int as total,288           count(*) filter (where quarantined = false and (next_check is null or next_check <= now()))::int as due,289           count(*) filter (where quarantined)::int as quarantined,290           count(*) filter (where error_count > 0)::int as errors,291           count(*) filter (where extract_ok)::int as extracted292    from documents where connector_id = ${connectorId}`);293  const r = rows[0];294  return { total: Number(r?.total ?? 0), due: Number(r?.due ?? 0), quarantined: Number(r?.quarantined ?? 0), errors: Number(r?.errors ?? 0), extracted: Number(r?.extracted ?? 0) };295}296297/** Earliest next_check per group for a connector (drives connectors.next_run_at and the scheduler). */298export async function nextCheckByGroup(connectorId: string): Promise<Array<{ group: string; nextCheck: string | null; due: number }>> {299  const rows = await getDb().execute<{ discovered_from: string | null; next_check: string | null; due: number }>(sql`300    select split_part(coalesce(discovered_from, 'default|'), '|', 1) as discovered_from,301           min(next_check) as next_check,302           count(*) filter (where next_check is null or next_check <= now())::int as due303    from documents where connector_id = ${connectorId} and quarantined = false304    group by 1`);305  return rows.map((r) => ({ group: r.discovered_from || "default", nextCheck: r.next_check ? new Date(r.next_check).toISOString() : null, due: Number(r.due) }));306}307