/** * Pure scheduling / change-detection math (no I/O) — unit-tested in scheduling.test.ts. * * Adaptive re-check interval: base interval of the document's schedule group, scaled by an EMA of * "did this page change when fetched" (change_frequency_score). Pages that keep changing are checked * twice as often; static pages back off up to 4× the base interval. 404/410 back off ×4 and after three * consecutive ones the document is quarantined (no more scheduling). */ import { intervalMs, type ConnectorSchedule } from "@dci/connectors"; export const DEFAULT_INTERVAL = "weekly"; /** EMA smoothing factor for change_frequency_score (0.5 = neutral prior). */ export const EMA_ALPHA = 0.3; export const MAX_CONSECUTIVE_GONE = 3; export const MIN_INTERVAL_MS = 15 * 60_000; /** "group|https://from" ⇄ { group, from } — the `documents.discovered_from` column carries both. */ export function encodeDiscoveredFrom(group: string, from: string | null | undefined): string { return `${group.replace(/\|/g, "_")}|${from ?? ""}`; } export function parseDiscoveredFrom(s: string | null | undefined): { group: string; from: string | null } { if (!s) return { group: "default", from: null }; const i = s.indexOf("|"); if (i < 0) return { group: "default", from: s || null }; return { group: s.slice(0, i) || "default", from: s.slice(i + 1) || null }; } /** Interval for a group: schedule[group] → schedule.default → weekly. `never` → Infinity. */ export function baseIntervalMs(schedule: ConnectorSchedule, group: string): number { const spec = schedule[group] ?? schedule.default ?? schedule["*"] ?? DEFAULT_INTERVAL; try { return intervalMs(spec); } catch { return intervalMs(DEFAULT_INTERVAL); } } /** Shortest finite interval across groups (used for connectors.next_run_at fallback and discovery cadence). */ export function minIntervalMs(schedule: ConnectorSchedule): number { let min = Number.POSITIVE_INFINITY; for (const spec of Object.values(schedule)) { try { min = Math.min(min, intervalMs(spec)); } catch { /* ignore bad spec */ } } return Number.isFinite(min) ? min : intervalMs(DEFAULT_INTERVAL); } /** EMA update of the change score: 1 = changed on every fetch, 0 = never changes. */ export function updateChangeScore(prev: number, changed: boolean, alpha = EMA_ALPHA): number { const p = Number.isFinite(prev) ? Math.min(1, Math.max(0, prev)) : 0.5; const v = p * (1 - alpha) + (changed ? 1 : 0) * alpha; return Math.round(v * 1000) / 1000; } /** * Interval multiplier from the change score: * ≥ 0.5 → 0.5 (frequent changer: check twice as often) * 0.25–0.5 → 1 * 0.1–0.25 → 2 (static: back off) * < 0.1 → 4 (very static: cap at 4× base) */ export function intervalScale(score: number): number { if (score >= 0.5) return 0.5; if (score >= 0.25) return 1; if (score >= 0.1) return 2; return 4; } export interface NextCheckInput { now: Date; baseMs: number; changeScore: number; statusCode: number | null; /** consecutive error count AFTER this fetch (0 on success) */ consecutiveErrors: number; /** transport error (dns, timeout, blocked…) on this fetch */ hadError: boolean; /** server-requested wait from `Retry-After` on a 429 (ms), when the fetcher exposed it (`doc.meta.retryAfterMs`) */ retryAfterMs?: number | null; } export interface NextCheckResult { /** null = stop scheduling (quarantined or `never`) */ nextCheck: Date | null; quarantined: boolean; reason: string; } export function isGone(status: number | null): boolean { return status === 404 || status === 410; } /** Decide the next check time for a document after a fetch. */ export function nextCheckAfter(i: NextCheckInput): NextCheckResult { if (!Number.isFinite(i.baseMs)) return { nextCheck: null, quarantined: false, reason: "never" }; if (isGone(i.statusCode)) { if (i.consecutiveErrors >= MAX_CONSECUTIVE_GONE) return { nextCheck: null, quarantined: true, reason: `gone ×${i.consecutiveErrors}` }; return { nextCheck: new Date(i.now.getTime() + Math.max(MIN_INTERVAL_MS, i.baseMs * 4)), quarantined: false, reason: "gone backoff ×4" }; } if (i.statusCode === 429 && i.retryAfterMs != null && Number.isFinite(i.retryAfterMs) && i.retryAfterMs > 0) { // the server told us when to come back: honour it (never sooner than the floor, never later than the error backoff cap) const wait = Math.min(Math.max(MIN_INTERVAL_MS, i.retryAfterMs), Math.max(MIN_INTERVAL_MS, i.baseMs * 8)); return { nextCheck: new Date(i.now.getTime() + wait), quarantined: false, reason: `retry-after ${Math.round(i.retryAfterMs / 1000)}s` }; } if (i.hadError || (i.statusCode != null && i.statusCode >= 400)) { // exponential backoff on errors: base × 2^(n-1), capped at 8× base const mult = Math.min(8, 2 ** Math.max(0, i.consecutiveErrors - 1)); return { nextCheck: new Date(i.now.getTime() + Math.max(MIN_INTERVAL_MS, i.baseMs * mult)), quarantined: false, reason: `error backoff ×${mult}` }; } const scale = intervalScale(i.changeScore); return { nextCheck: new Date(i.now.getTime() + Math.max(MIN_INTERVAL_MS, i.baseMs * scale)), quarantined: false, reason: `adaptive ×${scale}` }; } /** Text-only significance when no field-level changes were detected (10–30 by diff ratio). */ export function ratioSignificance(ratio: number): number { if (ratio <= 0) return 0; if (ratio < 0.02) return 10; if (ratio < 0.1) return 20; return 30; } /** Overall significance for a document version: max of field-level changes, else ratio-based. */ export function versionSignificance(detected: Array<{ significance: number }>, ratio: number): number { const fromFields = detected.reduce((m, c) => Math.max(m, Math.round(c.significance)), 0); return Math.max(0, Math.min(100, fromFields > 0 ? fromFields : ratioSignificance(ratio))); } /** * Skip decision: extraction is skipped when the content fingerprint equals the stored one and the extractor * version has not changed and the caller did not force it. */ export function shouldSkipExtraction(i: { newHash: string; storedHash: string | null; force: boolean; extractorVersion: string; storedExtractorVersion: string | null; notModified: boolean }): boolean { if (i.force) return false; if (i.notModified) return true; if (!i.storedHash || i.storedHash !== i.newHash) return false; return (i.storedExtractorVersion ?? i.extractorVersion) === i.extractorVersion; } /** * Health label from a run: failure rate < 20 % ok, < 50 % degraded, else failing. A run that attempted nothing is * `ok` only when it had nothing to do (no due documents); a run whose discovery yielded nothing for a connector that * used to discover is reported by the caller as `no_new_content` / `blocked` (see pipeline connector health). */ export function healthFrom(fetched: number, failed: number): "ok" | "degraded" | "failing" { if (fetched <= 0) return failed > 0 ? "failing" : "ok"; const rate = failed / fetched; if (rate < 0.2) return "ok"; if (rate < 0.5) return "degraded"; return "failing"; } /** Premium levels (3/4) are allowed only while the run's credit budget is not exhausted. */ export function maxLevelForBudget(cfgMaxLevel: number, creditsSpent: number, maxCreditsPerRun: number): 1 | 2 | 3 | 4 { const lvl = Math.max(1, Math.min(4, cfgMaxLevel)) as 1 | 2 | 3 | 4; if (lvl <= 2) return lvl; return creditsSpent >= maxCreditsPerRun ? 2 : lvl; } /** Discovery fetches (sitemaps, RSS feeds) never use premium fetchers: direct transport only (L1, L2 = browser identity). */ export const DISCOVERY_GROUPS: ReadonlySet = new Set(["sitemap", "rss"]); export function isDiscoveryGroup(group: string | undefined): boolean { return group !== undefined && DISCOVERY_GROUPS.has(group); } /** * Premium escalation for a document that keeps failing: after two consecutive failed attempts (which may already * have cost credits), only every 4th attempt may escalate again — the error backoff (×2^(n−1), cap ×8) spaces those * out to at most one premium retry per ~week for a daily group. A successful fetch resets `errorCount` to 0. */ export function premiumAllowedAfterErrors(errorCount: number): boolean { if (!Number.isFinite(errorCount) || errorCount < 2) return true; return errorCount % 4 === 0; } /** Connector health label (admin health center): HEALTHY | DEGRADED | BLOCKED | SCHEMA_CHANGE | NO_NEW_CONTENT | FAILED. */ export type ConnectorHealth = "ok" | "degraded" | "failing" | "blocked" | "schema_change" | "no_new_content" | "never_run"; export 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 { if (i.status === "failed") return "failing"; if (i.fetched > 0 && i.blocked >= Math.max(1, Math.ceil(i.fetched * 0.5))) return "blocked"; // pages fetched fine but the parsers found nothing on pages that used to yield entities → the site's HTML changed if (i.fetched >= 5 && i.extracted >= 5 && i.entities === 0 && i.docsTotal > 0) return "schema_change"; if (i.previouslyDiscovered != null && i.previouslyDiscovered > 0 && i.discovered === 0 && i.fetched === 0) return "no_new_content"; return i.runHealth; }