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%
9.2 KB · 180 lines typescript
Raw Blame History
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