spb/datacenterindex
Public
HTML 53.9%
TypeScript 44.5%
JavaScript 0.6%
SQL 0.5%
1/**2 * Pure scheduling / change-detection math (no I/O) — unit-tested in scheduling.test.ts.3 *4 * Adaptive re-check interval: base interval of the document's schedule group, scaled by an EMA of5 * "did this page change when fetched" (change_frequency_score). Pages that keep changing are checked6 * twice as often; static pages back off up to 4× the base interval. 404/410 back off ×4 and after three7 * consecutive ones the document is quarantined (no more scheduling).8 */9import { intervalMs, type ConnectorSchedule } from "@dci/connectors";1011export const DEFAULT_INTERVAL = "weekly";12/** EMA smoothing factor for change_frequency_score (0.5 = neutral prior). */13export const EMA_ALPHA = 0.3;14export const MAX_CONSECUTIVE_GONE = 3;15export const MIN_INTERVAL_MS = 15 * 60_000;1617/** "group|https://from" ⇄ { group, from } — the `documents.discovered_from` column carries both. */18export function encodeDiscoveredFrom(group: string, from: string | null | undefined): string {19 return `${group.replace(/\|/g, "_")}|${from ?? ""}`;20}21export function parseDiscoveredFrom(s: string | null | undefined): { group: string; from: string | null } {22 if (!s) return { group: "default", from: null };23 const i = s.indexOf("|");24 if (i < 0) return { group: "default", from: s || null };25 return { group: s.slice(0, i) || "default", from: s.slice(i + 1) || null };26}2728/** Interval for a group: schedule[group] → schedule.default → weekly. `never` → Infinity. */29export function baseIntervalMs(schedule: ConnectorSchedule, group: string): number {30 const spec = schedule[group] ?? schedule.default ?? schedule["*"] ?? DEFAULT_INTERVAL;31 try {32 return intervalMs(spec);33 } catch {34 return intervalMs(DEFAULT_INTERVAL);35 }36}3738/** Shortest finite interval across groups (used for connectors.next_run_at fallback and discovery cadence). */39export function minIntervalMs(schedule: ConnectorSchedule): number {40 let min = Number.POSITIVE_INFINITY;41 for (const spec of Object.values(schedule)) {42 try { min = Math.min(min, intervalMs(spec)); } catch { /* ignore bad spec */ }43 }44 return Number.isFinite(min) ? min : intervalMs(DEFAULT_INTERVAL);45}4647/** EMA update of the change score: 1 = changed on every fetch, 0 = never changes. */48export function updateChangeScore(prev: number, changed: boolean, alpha = EMA_ALPHA): number {49 const p = Number.isFinite(prev) ? Math.min(1, Math.max(0, prev)) : 0.5;50 const v = p * (1 - alpha) + (changed ? 1 : 0) * alpha;51 return Math.round(v * 1000) / 1000;52}5354/**55 * Interval multiplier from the change score:56 * ≥ 0.5 → 0.5 (frequent changer: check twice as often)57 * 0.25–0.5 → 158 * 0.1–0.25 → 2 (static: back off)59 * < 0.1 → 4 (very static: cap at 4× base)60 */61export function intervalScale(score: number): number {62 if (score >= 0.5) return 0.5;63 if (score >= 0.25) return 1;64 if (score >= 0.1) return 2;65 return 4;66}6768export interface NextCheckInput {69 now: Date;70 baseMs: number;71 changeScore: number;72 statusCode: number | null;73 /** consecutive error count AFTER this fetch (0 on success) */74 consecutiveErrors: number;75 /** transport error (dns, timeout, blocked…) on this fetch */76 hadError: boolean;77 /** server-requested wait from `Retry-After` on a 429 (ms), when the fetcher exposed it (`doc.meta.retryAfterMs`) */78 retryAfterMs?: number | null;79}8081export interface NextCheckResult {82 /** null = stop scheduling (quarantined or `never`) */83 nextCheck: Date | null;84 quarantined: boolean;85 reason: string;86}8788export function isGone(status: number | null): boolean { return status === 404 || status === 410; }8990/** Decide the next check time for a document after a fetch. */91export function nextCheckAfter(i: NextCheckInput): NextCheckResult {92 if (!Number.isFinite(i.baseMs)) return { nextCheck: null, quarantined: false, reason: "never" };93 if (isGone(i.statusCode)) {94 if (i.consecutiveErrors >= MAX_CONSECUTIVE_GONE) return { nextCheck: null, quarantined: true, reason: `gone ×${i.consecutiveErrors}` };95 return { nextCheck: new Date(i.now.getTime() + Math.max(MIN_INTERVAL_MS, i.baseMs * 4)), quarantined: false, reason: "gone backoff ×4" };96 }97 if (i.statusCode === 429 && i.retryAfterMs != null && Number.isFinite(i.retryAfterMs) && i.retryAfterMs > 0) {98 // the server told us when to come back: honour it (never sooner than the floor, never later than the error backoff cap)99 const wait = Math.min(Math.max(MIN_INTERVAL_MS, i.retryAfterMs), Math.max(MIN_INTERVAL_MS, i.baseMs * 8));100 return { nextCheck: new Date(i.now.getTime() + wait), quarantined: false, reason: `retry-after ${Math.round(i.retryAfterMs / 1000)}s` };101 }102 if (i.hadError || (i.statusCode != null && i.statusCode >= 400)) {103 // exponential backoff on errors: base × 2^(n-1), capped at 8× base104 const mult = Math.min(8, 2 ** Math.max(0, i.consecutiveErrors - 1));105 return { nextCheck: new Date(i.now.getTime() + Math.max(MIN_INTERVAL_MS, i.baseMs * mult)), quarantined: false, reason: `error backoff ×${mult}` };106 }107 const scale = intervalScale(i.changeScore);108 return { nextCheck: new Date(i.now.getTime() + Math.max(MIN_INTERVAL_MS, i.baseMs * scale)), quarantined: false, reason: `adaptive ×${scale}` };109}110111/** Text-only significance when no field-level changes were detected (10–30 by diff ratio). */112export function ratioSignificance(ratio: number): number {113 if (ratio <= 0) return 0;114 if (ratio < 0.02) return 10;115 if (ratio < 0.1) return 20;116 return 30;117}118119/** Overall significance for a document version: max of field-level changes, else ratio-based. */120export function versionSignificance(detected: Array<{ significance: number }>, ratio: number): number {121 const fromFields = detected.reduce((m, c) => Math.max(m, Math.round(c.significance)), 0);122 return Math.max(0, Math.min(100, fromFields > 0 ? fromFields : ratioSignificance(ratio)));123}124125/**126 * Skip decision: extraction is skipped when the content fingerprint equals the stored one and the extractor127 * version has not changed and the caller did not force it.128 */129export function shouldSkipExtraction(i: { newHash: string; storedHash: string | null; force: boolean; extractorVersion: string; storedExtractorVersion: string | null; notModified: boolean }): boolean {130 if (i.force) return false;131 if (i.notModified) return true;132 if (!i.storedHash || i.storedHash !== i.newHash) return false;133 return (i.storedExtractorVersion ?? i.extractorVersion) === i.extractorVersion;134}135136/**137 * Health label from a run: failure rate < 20 % ok, < 50 % degraded, else failing. A run that attempted nothing is138 * `ok` only when it had nothing to do (no due documents); a run whose discovery yielded nothing for a connector that139 * used to discover is reported by the caller as `no_new_content` / `blocked` (see pipeline connector health).140 */141export function healthFrom(fetched: number, failed: number): "ok" | "degraded" | "failing" {142 if (fetched <= 0) return failed > 0 ? "failing" : "ok";143 const rate = failed / fetched;144 if (rate < 0.2) return "ok";145 if (rate < 0.5) return "degraded";146 return "failing";147}148149/** Premium levels (3/4) are allowed only while the run's credit budget is not exhausted. */150export function maxLevelForBudget(cfgMaxLevel: number, creditsSpent: number, maxCreditsPerRun: number): 1 | 2 | 3 | 4 {151 const lvl = Math.max(1, Math.min(4, cfgMaxLevel)) as 1 | 2 | 3 | 4;152 if (lvl <= 2) return lvl;153 return creditsSpent >= maxCreditsPerRun ? 2 : lvl;154}155156/** Discovery fetches (sitemaps, RSS feeds) never use premium fetchers: direct transport only (L1, L2 = browser identity). */157export const DISCOVERY_GROUPS: ReadonlySet<string> = new Set(["sitemap", "rss"]);158export function isDiscoveryGroup(group: string | undefined): boolean { return group !== undefined && DISCOVERY_GROUPS.has(group); }159160/**161 * Premium escalation for a document that keeps failing: after two consecutive failed attempts (which may already162 * have cost credits), only every 4th attempt may escalate again — the error backoff (×2^(n−1), cap ×8) spaces those163 * out to at most one premium retry per ~week for a daily group. A successful fetch resets `errorCount` to 0.164 */165export function premiumAllowedAfterErrors(errorCount: number): boolean {166 if (!Number.isFinite(errorCount) || errorCount < 2) return true;167 return errorCount % 4 === 0;168}169170/** Connector health label (admin health center): HEALTHY | DEGRADED | BLOCKED | SCHEMA_CHANGE | NO_NEW_CONTENT | FAILED. */171export type ConnectorHealth = "ok" | "degraded" | "failing" | "blocked" | "schema_change" | "no_new_content" | "never_run";172export function connectorHealthFrom(i: { runHealth: "ok" | "degraded" | "failing"; status: string; fetched: number; blocked: number; discovered: number; previouslyDiscovered: number | null; extracted: number; entities: number; docsTotal: number }): ConnectorHealth {173 if (i.status === "failed") return "failing";174 if (i.fetched > 0 && i.blocked >= Math.max(1, Math.ceil(i.fetched * 0.5))) return "blocked";175 // pages fetched fine but the parsers found nothing on pages that used to yield entities → the site's HTML changed176 if (i.fetched >= 5 && i.extracted >= 5 && i.entities === 0 && i.docsTotal > 0) return "schema_change";177 if (i.previouslyDiscovered != null && i.previouslyDiscovered > 0 && i.discovered === 0 && i.fetched === 0) return "no_new_content";178 return i.runHealth;179}180