spb/datacenterindex
Public
HTML 53.9%
TypeScript 44.5%
JavaScript 0.6%
SQL 0.5%
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