import { describe, expect, it } from "vitest"; import type { RawDocument } from "@dci/connectors"; import { baseIntervalMs, encodeDiscoveredFrom, healthFrom, intervalScale, maxLevelForBudget, minIntervalMs, nextCheckAfter, parseDiscoveredFrom, ratioSignificance, shouldSkipExtraction, updateChangeScore, versionSignificance } from "./scheduling.js"; import { planFetchBookkeeping } from "./documents.js"; import { fingerprintOf, pLimit, stripRuntimeNoise } from "./pipeline.js"; const DAY = 86_400_000; const now = new Date("2026-09-11T12:00:00Z"); function raw(partial: Partial): RawDocument { return { url: "https://example.com/a", finalUrl: "https://example.com/a", fetchedAt: now.toISOString(), status: 200, contentType: "text/html", body: Buffer.from("x"), text: "x", headers: {}, etag: null, lastModified: null, notModified: false, fetcher: "direct", level: 1, durationMs: 10, credits: 0, ...partial }; } describe("intervals", () => { it("resolves group → default → weekly", () => { expect(baseIntervalMs({ newsroom: "daily", default: "monthly" }, "newsroom")).toBe(DAY); expect(baseIntervalMs({ newsroom: "daily", default: "monthly" }, "other")).toBe(30 * DAY); expect(baseIntervalMs({ newsroom: "daily" }, "other")).toBe(7 * DAY); expect(baseIntervalMs({ x: "not-an-interval" }, "x")).toBe(7 * DAY); }); it("never → Infinity, min interval picks the shortest finite one", () => { expect(baseIntervalMs({ archive: "never" }, "archive")).toBe(Number.POSITIVE_INFINITY); expect(minIntervalMs({ a: "weekly", b: "6h", c: "never" })).toBe(6 * 3_600_000); }); }); describe("change score EMA + scale", () => { it("moves toward 1 on change and toward 0 when static, clamped", () => { expect(updateChangeScore(0.5, true)).toBeCloseTo(0.65, 3); expect(updateChangeScore(0.5, false)).toBeCloseTo(0.35, 3); let s = 0.5; for (let i = 0; i < 20; i++) s = updateChangeScore(s, false); expect(s).toBeLessThan(0.01); expect(updateChangeScore(Number.NaN, true)).toBeCloseTo(0.65, 3); }); it("frequent changers ×0.5, static ×2 up to ×4", () => { expect(intervalScale(0.9)).toBe(0.5); expect(intervalScale(0.5)).toBe(0.5); expect(intervalScale(0.3)).toBe(1); expect(intervalScale(0.15)).toBe(2); expect(intervalScale(0.02)).toBe(4); }); }); describe("nextCheckAfter", () => { it("scales the base interval adaptively", () => { const fast = nextCheckAfter({ now, baseMs: DAY, changeScore: 0.8, statusCode: 200, consecutiveErrors: 0, hadError: false }); expect(fast.nextCheck!.getTime() - now.getTime()).toBe(DAY / 2); const slow = nextCheckAfter({ now, baseMs: DAY, changeScore: 0.01, statusCode: 200, consecutiveErrors: 0, hadError: false }); expect(slow.nextCheck!.getTime() - now.getTime()).toBe(4 * DAY); expect(slow.quarantined).toBe(false); }); it("404 backs off ×4 then quarantines after 3 consecutive", () => { const first = nextCheckAfter({ now, baseMs: DAY, changeScore: 0.5, statusCode: 404, consecutiveErrors: 1, hadError: false }); expect(first.nextCheck!.getTime() - now.getTime()).toBe(4 * DAY); expect(first.quarantined).toBe(false); const third = nextCheckAfter({ now, baseMs: DAY, changeScore: 0.5, statusCode: 410, consecutiveErrors: 3, hadError: false }); expect(third.nextCheck).toBeNull(); expect(third.quarantined).toBe(true); }); it("transport errors back off exponentially, capped at 8×, and never → null", () => { const e1 = nextCheckAfter({ now, baseMs: DAY, changeScore: 0.5, statusCode: null, consecutiveErrors: 1, hadError: true }); expect(e1.nextCheck!.getTime() - now.getTime()).toBe(DAY); const e5 = nextCheckAfter({ now, baseMs: DAY, changeScore: 0.5, statusCode: null, consecutiveErrors: 5, hadError: true }); expect(e5.nextCheck!.getTime() - now.getTime()).toBe(8 * DAY); expect(nextCheckAfter({ now, baseMs: Number.POSITIVE_INFINITY, changeScore: 0.5, statusCode: 200, consecutiveErrors: 0, hadError: false }).nextCheck).toBeNull(); }); it("429 honours Retry-After (floor 15 min, cap 8× base), otherwise falls back to the error backoff", () => { const ra = nextCheckAfter({ now, baseMs: DAY, changeScore: 0.5, statusCode: 429, consecutiveErrors: 1, hadError: false, retryAfterMs: 3_600_000 }); expect(ra.nextCheck!.getTime() - now.getTime()).toBe(3_600_000); expect(ra.reason).toBe("retry-after 3600s"); expect(nextCheckAfter({ now, baseMs: DAY, changeScore: 0.5, statusCode: 429, consecutiveErrors: 1, hadError: false, retryAfterMs: 1_000 }).nextCheck!.getTime() - now.getTime()).toBe(15 * 60_000); expect(nextCheckAfter({ now, baseMs: DAY, changeScore: 0.5, statusCode: 429, consecutiveErrors: 1, hadError: false, retryAfterMs: 30 * DAY }).nextCheck!.getTime() - now.getTime()).toBe(8 * DAY); const noHeader = nextCheckAfter({ now, baseMs: DAY, changeScore: 0.5, statusCode: 429, consecutiveErrors: 2, hadError: false, retryAfterMs: null }); expect(noHeader.reason).toBe("error backoff ×2"); }); it("planFetchBookkeeping forwards doc.meta.retryAfterMs from a 429", () => { const p = planFetchBookkeeping({ changeFrequencyScore: 0.5, errorCount: 0, discoveredFrom: encodeDiscoveredFrom("newsroom", null) }, raw({ status: 429, body: Buffer.alloc(0), text: "", meta: { retryAfterMs: 2 * 3_600_000 } }), { contentHash: null, changed: false, storageKey: null, title: null }, { newsroom: "daily" }, now); expect(p.failed).toBe(true); expect(p.error).toBe("HTTP 429"); expect(Date.parse(p.nextCheck!) - now.getTime()).toBe(2 * 3_600_000); }); it("enforces a 15 min floor", () => { const r = nextCheckAfter({ now, baseMs: 60_000, changeScore: 0.9, statusCode: 200, consecutiveErrors: 0, hadError: false }); expect(r.nextCheck!.getTime() - now.getTime()).toBe(15 * 60_000); }); }); describe("planFetchBookkeeping (documents)", () => { const doc = { changeFrequencyScore: 0.5, errorCount: 0, discoveredFrom: encodeDiscoveredFrom("newsroom", "https://example.com/sitemap.xml") }; const schedule = { newsroom: "daily", facility_pages: "weekly" }; it("uses the group's interval and updates the score on success", () => { const p = planFetchBookkeeping(doc, raw({}), { contentHash: "h1", changed: true, storageKey: "k", title: "t" }, schedule, now); expect(p.failed).toBe(false); expect(p.changed).toBe(true); expect(p.changeScore).toBeCloseTo(0.65, 3); expect(p.errorCount).toBe(0); expect(Date.parse(p.nextCheck!) - now.getTime()).toBe(DAY / 2); // 0.65 ≥ 0.5 → ×0.5 }); it("counts consecutive 404s and quarantines on the third", () => { const p1 = planFetchBookkeeping({ ...doc, errorCount: 0 }, raw({ status: 404, body: Buffer.alloc(0), text: "" }), { contentHash: null, changed: false, storageKey: null, title: null }, schedule, now); expect(p1.failed).toBe(true); expect(p1.errorCount).toBe(1); expect(p1.error).toBe("HTTP 404"); expect(p1.changeScore).toBe(0.5); // untouched on failure const p3 = planFetchBookkeeping({ ...doc, errorCount: 2 }, raw({ status: 410, body: Buffer.alloc(0), text: "" }), { contentHash: null, changed: false, storageKey: null, title: null }, schedule, now); expect(p3.quarantined).toBe(true); expect(p3.nextCheck).toBeNull(); }); it("treats transport errors as failures with their code", () => { const p = planFetchBookkeeping({ ...doc, errorCount: 1 }, raw({ status: 0, error: { code: "timeout", message: "AbortError" } }), { contentHash: null, changed: false, storageKey: null, title: null }, schedule, now); expect(p.failed).toBe(true); expect(p.error).toMatch(/^timeout/); expect(p.errorCount).toBe(2); expect(Date.parse(p.nextCheck!) - now.getTime()).toBe(2 * DAY); }); it("304 is a success that lowers the change score", () => { const p = planFetchBookkeeping(doc, raw({ status: 304, notModified: true, body: Buffer.alloc(0), text: "" }), { contentHash: "h0", changed: false, storageKey: "k", title: null }, schedule, now); expect(p.failed).toBe(false); expect(p.notModified).toBe(true); expect(p.changeScore).toBeCloseTo(0.35, 3); }); }); describe("fingerprint / skip logic", () => { it("ignores volatile noise between two fetches of the same page", () => { const a = raw({ text: '\n

Hello

' }); const b = raw({ text: '\n\n

Hello

' }); const c = raw({ text: "

Hello world

" }); expect(fingerprintOf(a)).toBe(fingerprintOf(b)); expect(fingerprintOf(a)).not.toBe(fingerprintOf(c)); }); it("ignores epoch cache-busters stamped into URLs on every render", () => { const a = raw({ text: 'x

Body

' }); const b = raw({ text: 'x

Body

' }); expect(fingerprintOf(a)).toBe(fingerprintOf(b)); expect(stripRuntimeNoise("/a?t=1789110041117&x=1")).toBe("/a?t=&x=1"); expect(stripRuntimeNoise("/a?id=12345&page=2")).toBe("/a?id=12345&page=2"); // short ids untouched expect(stripRuntimeNoise("capacity 120 MW, opened 2024, 1,000 racks")).toBe("capacity 120 MW, opened 2024, 1,000 racks"); }); it("hashes binary bodies (PDF) by bytes", () => { const pdf1 = raw({ contentType: "application/pdf", text: "", body: Buffer.from("%PDF-1.4 one") }); const pdf2 = raw({ contentType: "application/pdf", text: "", body: Buffer.from("%PDF-1.4 two") }); expect(fingerprintOf(pdf1)).not.toBe(fingerprintOf(pdf2)); expect(fingerprintOf(pdf1)).toBe(fingerprintOf(raw({ contentType: "application/pdf", text: "", body: Buffer.from("%PDF-1.4 one") }))); }); it("skips extraction only when hash and extractor version are unchanged and not forced", () => { const base = { newHash: "h", storedHash: "h", force: false, extractorVersion: "v1", storedExtractorVersion: "v1", notModified: false }; expect(shouldSkipExtraction(base)).toBe(true); expect(shouldSkipExtraction({ ...base, force: true })).toBe(false); expect(shouldSkipExtraction({ ...base, storedHash: "other" })).toBe(false); expect(shouldSkipExtraction({ ...base, storedHash: null })).toBe(false); expect(shouldSkipExtraction({ ...base, storedExtractorVersion: "v0" })).toBe(false); expect(shouldSkipExtraction({ ...base, storedExtractorVersion: null })).toBe(true); expect(shouldSkipExtraction({ ...base, storedHash: "x", notModified: true })).toBe(true); }); }); describe("significance, health, budget, discovered_from", () => { it("significance from field changes else ratio", () => { expect(versionSignificance([{ significance: 90 }, { significance: 40 }], 0.5)).toBe(90); expect(versionSignificance([], 0.01)).toBe(10); expect(versionSignificance([], 0.05)).toBe(20); expect(versionSignificance([], 0.3)).toBe(30); expect(ratioSignificance(0)).toBe(0); }); it("health thresholds 20 % / 50 %", () => { expect(healthFrom(100, 5)).toBe("ok"); expect(healthFrom(100, 19)).toBe("ok"); expect(healthFrom(100, 20)).toBe("degraded"); expect(healthFrom(100, 50)).toBe("failing"); expect(healthFrom(0, 0)).toBe("ok"); }); it("premium levels are capped at L2 once the run budget is spent", () => { expect(maxLevelForBudget(4, 0, 100)).toBe(4); expect(maxLevelForBudget(4, 100, 100)).toBe(2); expect(maxLevelForBudget(3, 250, 200)).toBe(2); expect(maxLevelForBudget(2, 999, 0)).toBe(2); expect(maxLevelForBudget(1, 0, 0)).toBe(1); }); it("round-trips group and origin in discovered_from", () => { const enc = encodeDiscoveredFrom("facility_pages", "https://example.com/sitemap.xml"); expect(parseDiscoveredFrom(enc)).toEqual({ group: "facility_pages", from: "https://example.com/sitemap.xml" }); expect(parseDiscoveredFrom(null)).toEqual({ group: "default", from: null }); expect(parseDiscoveredFrom("newsroom|")).toEqual({ group: "newsroom", from: null }); expect(parseDiscoveredFrom("https://legacy")).toEqual({ group: "default", from: "https://legacy" }); }); }); describe("pLimit", () => { it("never exceeds the concurrency and resolves everything", async () => { const limit = pLimit(2); let active = 0, peak = 0; const results = await Promise.all( Array.from({ length: 7 }, (_, i) => limit(async () => { active++; peak = Math.max(peak, active); await new Promise((r) => setTimeout(r, 5)); active--; return i; }), ), ); expect(results).toEqual([0, 1, 2, 3, 4, 5, 6]); expect(peak).toBe(2); }); it("propagates rejections without blocking the queue", async () => { const limit = pLimit(1); await expect(limit(async () => { throw new Error("boom"); })).rejects.toThrow("boom"); expect(await limit(async () => 42)).toBe(42); }); }); describe("connectorHealthFrom", () => { it("labels blocked, schema change and no-new-content connectors", async () => { const { connectorHealthFrom } = await import("./scheduling.js"); expect(connectorHealthFrom({ runHealth: "failing", status: "partial", fetched: 10, blocked: 8, discovered: 0, previouslyDiscovered: null, extracted: 0, entities: 0, docsTotal: 50 })).toBe("blocked"); expect(connectorHealthFrom({ runHealth: "ok", status: "ok", fetched: 12, blocked: 0, discovered: 0, previouslyDiscovered: null, extracted: 12, entities: 0, docsTotal: 200 })).toBe("schema_change"); expect(connectorHealthFrom({ runHealth: "ok", status: "ok", fetched: 0, blocked: 0, discovered: 0, previouslyDiscovered: 40, extracted: 0, entities: 0, docsTotal: 40 })).toBe("no_new_content"); expect(connectorHealthFrom({ runHealth: "ok", status: "ok", fetched: 8, blocked: 0, discovered: 3, previouslyDiscovered: 40, extracted: 8, entities: 8, docsTotal: 40 })).toBe("ok"); expect(connectorHealthFrom({ runHealth: "degraded", status: "failed", fetched: 8, blocked: 0, discovered: 3, previouslyDiscovered: 40, extracted: 8, entities: 8, docsTotal: 40 })).toBe("failing"); }); });