/** * URL registry (`documents`) + change tracking (`document_versions`, ClickHouse page_changes). * * - registerDiscovered: upsert discovered URLs (id = stableId("document", canonicalUrl)); new rows are due now. * - dueDocuments: what to fetch next (priority order, group filter, respects quarantine / next_check). * - recordFetch: bookkeeping after a fetch — conditional-GET validators, content hash, adaptive next_check. * - recordVersion: a document_versions row when the content fingerprint changed (line diff of text * projections, significance from field-level detectedChanges or diff ratio). */ import type { DiscoveredUrl, RawDocument } from "@dci/connectors"; import type { DetectedChange } from "@dci/core"; import { canonicalUrl, urlFingerprint, stableId, sha256, htmlToText, diffLines, type LineDiff } from "@dci/core"; import { getDb, documents, documentVersions, sql, eq, and, inArray } from "@dci/db"; import { chInsert } from "@dci/db/clickhouse"; import type { LoadedConnector } from "./configs.js"; import { getRaw } from "./storage.js"; import { baseIntervalMs, encodeDiscoveredFrom, isGone, nextCheckAfter, parseDiscoveredFrom, updateChangeScore, versionSignificance } from "./scheduling.js"; export type DocumentRow = typeof documents.$inferSelect; export type DocumentVersionRow = typeof documentVersions.$inferSelect; export function documentIdFor(url: string): string { return stableId("document", canonicalUrl(url)); } export interface RegisterResult { inserted: number; updated: number; skipped: number; ids: string[] } /** Upsert discovered URLs. Existing rows keep their scheduling state; priority/pageType/fetchLevel only ratchet up. */ export 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 { const now = (opts.now ?? new Date()).toISOString(); const res: RegisterResult = { inserted: 0, updated: 0, skipped: 0, ids: [] }; type Row = typeof documents.$inferInsert; const rows = new Map(); for (const u of urls) { let canon: string; try { canon = canonicalUrl(u.url); } catch { res.skipped++; continue; } if (!/^https?:$/.test(new URL(canon).protocol)) { res.skipped++; continue; } const id = stableId("document", canon); if (rows.has(id)) continue; rows.set(id, { id, connectorId: loaded.cfg.id, sourceId: loaded.sourceId, url: u.url, canonicalUrl: canon, urlFingerprint: urlFingerprint(u.url), pageType: u.pageType ?? "unknown", fetchLevel: u.minLevel ?? 1, priority: Math.max(0, Math.min(100, Math.round(u.priority ?? 50))), discoveredFrom: encodeDiscoveredFrom(u.group ?? "default", u.discoveredFrom ?? null), nextCheck: now, firstSeen: now, updatedAt: now, lastmod: u.lastmod ?? null, }); } res.ids = [...rows.keys()]; if (opts.dryRun || !rows.size) { res.inserted = rows.size; return res; } const db = getDb(); const all = [...rows.values()]; for (let i = 0; i < all.length; i += 200) { const chunk = all.slice(i, i + 200).map(({ lastmod: _l, ...r }) => r); const ret = await db .insert(documents) .values(chunk) .onConflictDoUpdate({ target: documents.id, set: { url: sql`excluded.url`, pageType: sql`case when ${documents.pageType} = 'unknown' then excluded.page_type else ${documents.pageType} end`, priority: sql`greatest(${documents.priority}, excluded.priority)`, fetchLevel: sql`greatest(${documents.fetchLevel}, excluded.fetch_level)`, discoveredFrom: opts.preserveGroup ? sql`coalesce(${documents.discoveredFrom}, excluded.discovered_from)` : sql`excluded.discovered_from`, updatedAt: now, }, }) .returning({ id: documents.id, inserted: sql`(xmax = 0)` }); for (const r of ret) { if (r.inserted) res.inserted++; else res.updated++; } } // sitemap lastmod newer than our last fetch → the page is due now (cheap signal, saves a full cycle) const withLastmod = all.filter((r) => r.lastmod && !Number.isNaN(Date.parse(r.lastmod))); if (withLastmod.length) { const byId = new Map(withLastmod.map((r) => [r.id, Date.parse(r.lastmod!)])); 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()])); 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); if (due.length) await db.update(documents).set({ nextCheck: now, updatedAt: now }).where(inArray(documents.id, due)); } return res; } export interface DueOptions { group?: string; limit?: number; /** ignore next_check (still skips quarantined) */ force?: boolean; ids?: string[]; includeQuarantined?: boolean } export async function dueDocuments(connectorId: string, opts: DueOptions = {}): Promise { const db = getDb(); const conds = [eq(documents.connectorId, connectorId)]; if (!opts.includeQuarantined) conds.push(eq(documents.quarantined, false)); if (!opts.force) conds.push(sql`(${documents.nextCheck} is null or ${documents.nextCheck} <= now())`); if (opts.group) conds.push(sql`${documents.discoveredFrom} like ${opts.group.replace(/[%_]/g, "\\$&") + "|%"}`); if (opts.ids?.length) conds.push(inArray(documents.id, opts.ids)); 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`); return opts.limit ? q.limit(opts.limit) : q; } export async function countDue(connectorId: string, group?: string): Promise { const rows = await getDb().execute<{ n: number }>( group ? 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 + "|%"}` : sql`select count(*)::int as n from documents where connector_id = ${connectorId} and quarantined = false and (next_check is null or next_check <= now())`, ); return Number(rows[0]?.n ?? 0); } export async function getDocument(idOrUrl: string): Promise { const id = /^https?:\/\//i.test(idOrUrl) ? documentIdFor(idOrUrl) : idOrUrl; const rows = await getDb().select().from(documents).where(eq(documents.id, id)).limit(1); return rows[0] ?? null; } /** The stored content hash for a URL — powers ConnectorContext.isKnownUnchanged. */ export async function storedHashFor(url: string): Promise { let id: string; try { id = documentIdFor(url); } catch { return null; } const rows = await getDb().select({ h: documents.contentHash }).from(documents).where(eq(documents.id, id)).limit(1); return rows[0]?.h ?? null; } export interface FetchOutcome { /** content fingerprint of the new body (null on error / 304) */ contentHash: string | null; changed: boolean; storageKey: string | null; title: string | null; } export interface FetchBookkeeping { nextCheck: string | null; quarantined: boolean; changeScore: number; errorCount: number; error: string | null; reason: string; changed: boolean; notModified: boolean; failed: boolean; } /** Compute what recordFetch will write (pure; exported for tests and dry runs). */ export function planFetchBookkeeping(doc: Pick, raw: RawDocument, outcome: FetchOutcome, schedule: Record, now = new Date()): FetchBookkeeping { const group = parseDiscoveredFrom(doc.discoveredFrom).group; const baseMs = baseIntervalMs(schedule, group); const hadError = Boolean(raw.error); const httpFail = !hadError && raw.status >= 400; const failed = hadError || httpFail; const errorCount = failed ? (doc.errorCount ?? 0) + 1 : 0; const error = hadError ? `${raw.error!.code}: ${raw.error!.message}`.slice(0, 500) : httpFail ? `HTTP ${raw.status}` : null; const changeScore = failed ? doc.changeFrequencyScore : updateChangeScore(doc.changeFrequencyScore, outcome.changed); const retryAfterMs = typeof raw.meta?.retryAfterMs === "number" ? raw.meta.retryAfterMs : null; const nc = nextCheckAfter({ now, baseMs, changeScore, statusCode: hadError ? null : raw.status, consecutiveErrors: errorCount, hadError, retryAfterMs }); 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 }; } /** Persist post-fetch state. Returns the bookkeeping decisions (also in dry-run, where nothing is written). */ export async function recordFetch(doc: DocumentRow, raw: RawDocument, outcome: FetchOutcome, schedule: Record, opts: { dryRun?: boolean; now?: Date } = {}): Promise { const now = opts.now ?? new Date(); const plan = planFetchBookkeeping(doc, raw, outcome, schedule, now); if (opts.dryRun) return plan; const nowIso = now.toISOString(); const set: Partial = { lastChecked: nowIso, nextCheck: plan.nextCheck, quarantined: plan.quarantined, changeFrequencyScore: plan.changeScore, errorCount: plan.errorCount, error: plan.error, fetchCount: sql`${documents.fetchCount} + 1` as unknown as number, updatedAt: nowIso, }; if (!raw.error) { set.statusCode = raw.status; if (raw.finalUrl && raw.finalUrl !== doc.url && /^https?:\/\//.test(raw.finalUrl)) set.url = raw.finalUrl; } if (!plan.failed) { set.lastFetched = nowIso; if (!raw.notModified) { set.contentType = raw.contentType; set.sizeBytes = raw.body.length; if (raw.etag) set.etag = raw.etag; if (raw.lastModified) set.lastModified = raw.lastModified; if (outcome.title) set.title = outcome.title.slice(0, 500); if (outcome.contentHash) set.contentHash = outcome.contentHash; if (outcome.storageKey) set.storageKey = outcome.storageKey; if (plan.changed) { set.lastChanged = nowIso; set.changeCount = sql`${documents.changeCount} + 1` as unknown as number; } } else { // 304: validators may be refreshed by the server if (raw.etag) set.etag = raw.etag; if (raw.lastModified) set.lastModified = raw.lastModified; } } else if (isGone(raw.status) && plan.quarantined) { set.error = `gone: HTTP ${raw.status} ×${plan.errorCount}`; } await getDb().update(documents).set(set).where(eq(documents.id, doc.id)); return plan; } export interface ExtractionUpdate { entityRefs: Array<{ type: string; id: string }>; extractOk: boolean; extractCount: number; extractorVersion: string; error?: string | null; pageType?: string | null; classifier?: string | null } export async function updateExtraction(docId: string, u: ExtractionUpdate, dryRun = false): Promise { if (dryRun) return; const set: Partial = { entityRefs: u.entityRefs, extractOk: u.extractOk, extractCount: u.extractCount, extractorVersion: u.extractorVersion, updatedAt: new Date().toISOString() }; if (u.error !== undefined) set.error = u.error; if (u.pageType && u.pageType !== "unknown") set.pageType = u.pageType; if (u.classifier) set.classifier = u.classifier; await getDb().update(documents).set(set).where(eq(documents.id, docId)); } /** Plain-text projection used for diffs (HTML → text; PDFs/binaries are not diffed line by line). */ export function textProjection(body: Buffer | string, contentType: string | null | undefined): string { const ct = (contentType ?? "").toLowerCase(); if (/pdf|image|octet-stream|zip/.test(ct)) return ""; const s = typeof body === "string" ? body : body.toString("utf8"); if (!ct || /html|xml/.test(ct) || / { const now = opts.now ?? new Date(); let diff: LineDiff | null = null; if (prev.storageKey && prev.contentHash && prev.contentHash !== newHash) { try { const old = await getRaw(prev.storageKey); if (old) diff = diffLines(textProjection(old.body, old.contentType ?? prev.contentType), textProjection(raw.text || raw.body, raw.contentType), 200); } catch { diff = null; } } const significance = versionSignificance(detectedChanges, diff?.ratio ?? 0); const versionId = `ver_${sha256(`${prev.id}:${newHash}:${now.toISOString()}`).slice(0, 16)}`; if (opts.dryRun) return { versionId, diff, significance }; await getDb() .insert(documentVersions) .values({ id: versionId, documentId: prev.id, contentHash: newHash, fetchedAt: now.toISOString(), fetchLevel: raw.level, statusCode: raw.status, sizeBytes: raw.body.length, storageKey, 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, detectedChanges: detectedChanges as unknown as Array>, significance, runId: opts.runId ?? null, extractorVersion: opts.extractorVersion ?? null, }) .onConflictDoNothing(); if (prev.contentHash) { await chInsert("page_changes", [ { ts: now.toISOString().replace("T", " ").replace("Z", ""), connector_id: prev.connectorId, document_id: prev.id, url: prev.url, old_hash: prev.contentHash, new_hash: newHash, added_lines: diff?.addedCount ?? 0, removed_lines: diff?.removedCount ?? 0, ratio: diff?.ratio ?? 0, significance, fields: [...new Set(detectedChanges.map((c) => c.field))], event_types: [...new Set(detectedChanges.map((c) => c.eventType))], }, ]).catch(() => undefined); } return { versionId, diff, significance }; } export async function latestVersion(documentId: string): Promise { const rows = await getDb().select().from(documentVersions).where(eq(documentVersions.documentId, documentId)).orderBy(sql`${documentVersions.fetchedAt} desc`).limit(1); return rows[0] ?? null; } export async function connectorDocStats(connectorId: string): Promise<{ total: number; due: number; quarantined: number; errors: number; extracted: number }> { const rows = await getDb().execute<{ total: number; due: number; quarantined: number; errors: number; extracted: number }>(sql` select count(*)::int as total, count(*) filter (where quarantined = false and (next_check is null or next_check <= now()))::int as due, count(*) filter (where quarantined)::int as quarantined, count(*) filter (where error_count > 0)::int as errors, count(*) filter (where extract_ok)::int as extracted from documents where connector_id = ${connectorId}`); const r = rows[0]; 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) }; } /** Earliest next_check per group for a connector (drives connectors.next_run_at and the scheduler). */ export async function nextCheckByGroup(connectorId: string): Promise> { const rows = await getDb().execute<{ discovered_from: string | null; next_check: string | null; due: number }>(sql` select split_part(coalesce(discovered_from, 'default|'), '|', 1) as discovered_from, min(next_check) as next_check, count(*) filter (where next_check is null or next_check <= now())::int as due from documents where connector_id = ${connectorId} and quarantined = false group by 1`); return rows.map((r) => ({ group: r.discovered_from || "default", nextCheck: r.next_check ? new Date(r.next_check).toISOString() : null, due: Number(r.due) })); }