import { describe, expect, it } from "vitest"; import type { Observation } from "@market-atlas/market-model"; import { ConsensusEngine, weightedMedian } from "./consensus.js"; const NOW = 1_789_200_000_000; const profiles: Record = { a: { family: "fa", reliability: 0.9, isOfficial: false }, b: { family: "fb", reliability: 0.9, isOfficial: false }, c: { family: "fc", reliability: 0.9, isOfficial: false }, d: { family: "fd", reliability: 0.9, isOfficial: false }, e: { family: "fa", reliability: 0.9, isOfficial: false }, // same upstream as a eod: { family: "official", reliability: 0.95, isOfficial: true }, }; const engine = () => new ConsensusEngine((id) => ({ ...(profiles[id] ?? { family: null, reliability: 0.6, isOfficial: false }), realtimeStatus: "REALTIME" })); function obs(source: string, value: number, ageMs = 500, field: Observation["field"] = "LAST_PRICE", rt: Observation["realtimeStatus"] = "REALTIME", trust: Observation["timestampTrust"] = "EXCHANGE"): Observation { return { observationId: `${source}-${field}-${value}-${ageMs}`, instrumentId: "crypto_btc_usd", symbol: "BTC-USD", field, value, currency: "USD", sourceTimestamp: NOW - ageMs, timestampTrust: trust, rightsStatus: "PUBLIC_ATTRIBUTED", realtimeStatus: rt, sourceId: source, connectorId: `${source}-ws`, receivedAt: NOW - ageMs + 50, latencyMs: 50, rawRef: null, normalizerVersion: "1.0", }; } describe("consensus engine", () => { it("all sources agree → canonical price with high confidence", () => { const e = engine(); for (const [s, v] of [["a", 100], ["b", 100], ["c", 100]] as const) e.ingest(obs(s, v)); const q = e.compute("crypto_btc_usd", "BTC-USD", NOW)!; expect(q.price).toBe(100); expect(q.sourceCount).toBe(3); expect(q.dispersionBps).toBe(0); expect(q.confidence).toBeGreaterThan(0.85); expect(q.confidence).toBeLessThanOrEqual(0.995); }); it("one stale source is excluded and labelled", () => { const e = engine(); e.ingest(obs("a", 100)); e.ingest(obs("b", 100.1)); e.ingest(obs("c", 90, 60_000)); // 60 s old realtime → stale const q = e.compute("crypto_btc_usd", "BTC-USD", NOW)!; expect(q.price).toBeGreaterThanOrEqual(100); expect(q.contributions.find((c) => c.sourceId === "c")).toMatchObject({ included: false, reason: "stale" }); expect(q.sourceCount).toBe(2); }); it("one extreme outlier among ≥3 is rejected", () => { const e = engine(); e.ingest(obs("a", 100)); e.ingest(obs("b", 100.2)); e.ingest(obs("c", 99.9)); e.ingest(obs("d", 130)); const q = e.compute("crypto_btc_usd", "BTC-USD", NOW)!; expect(q.price).toBeLessThan(101); expect(q.contributions.find((c) => c.sourceId === "d")).toMatchObject({ included: false, reason: "outlier" }); }); it("two source clusters disagree → low confidence, large dispersion", () => { const e = engine(); e.ingest(obs("a", 100)); e.ingest(obs("b", 101.5)); const q = e.compute("crypto_btc_usd", "BTC-USD", NOW)!; expect(q.dispersionBps).toBeGreaterThan(100); expect(q.confidence).toBeLessThan(0.7); }); it("only one source available → modest confidence, count 1", () => { const e = engine(); e.ingest(obs("a", 100)); const q = e.compute("crypto_btc_usd", "BTC-USD", NOW)!; expect(q.price).toBe(100); expect(q.sourceCount).toBe(1); expect(q.confidence).toBeLessThan(0.85); }); it("sources sharing an upstream family count once", () => { const e = engine(); e.ingest(obs("a", 100)); e.ingest(obs("e", 100)); const q = e.compute("crypto_btc_usd", "BTC-USD", NOW)!; expect(q.sourceCount).toBe(1); }); it("out-of-order older observation does not overwrite the newer one", () => { const e = engine(); e.ingest(obs("a", 100, 100)); e.ingest(obs("a", 90, 5000)); expect(e.compute("crypto_btc_usd", "BTC-USD", NOW)!.price).toBe(100); }); it("end-of-day value is superseded by a live source, but used when alone", () => { const e = engine(); e.ingest(obs("eod", 95, 3_600_000, "LAST_PRICE", "END_OF_DAY", "SOURCE")); expect(e.compute("crypto_btc_usd", "BTC-USD", NOW)!.price).toBe(95); expect(e.compute("crypto_btc_usd", "BTC-USD", NOW)!.realtimeStatus).toBe("END_OF_DAY"); e.ingest(obs("a", 100)); const q = e.compute("crypto_btc_usd", "BTC-USD", NOW)!; expect(q.price).toBe(100); expect(q.realtimeStatus).toBe("REALTIME"); expect(q.contributions.find((c) => c.sourceId === "eod")).toMatchObject({ included: false, reason: "not_comparable" }); }); it("market closed / everything stale → STALE status with last known value and zero confidence", () => { const e = engine(); e.ingest(obs("a", 100, 10 * 60_000)); const q = e.compute("crypto_btc_usd", "BTC-USD", NOW)!; expect(q.realtimeStatus).toBe("STALE"); expect(q.price).toBe(100); expect(q.confidence).toBe(0); }); it("derives change from previous close and tracks session high/low", () => { const e = engine(); e.ingest(obs("a", 100)); e.ingest(obs("a", 95, 500, "PREVIOUS_CLOSE")); let q = e.compute("crypto_btc_usd", "BTC-USD", NOW)!; expect(q.change).toBeCloseTo(5); expect(q.changePercent).toBeCloseTo(5.263, 2); e.ingest(obs("a", 102, 100)); q = e.compute("crypto_btc_usd", "BTC-USD", NOW)!; expect(q.sessionHigh).toBe(102); expect(q.sessionLow).toBe(100); }); it("weighted median", () => { expect(weightedMedian([{ v: 1, w: 1 }, { v: 2, w: 1 }, { v: 100, w: 1 }])).toBe(2); expect(weightedMedian([{ v: 1, w: 10 }, { v: 2, w: 1 }, { v: 100, w: 1 }])).toBe(1); }); }); describe("comparability, roles and coverage inputs", () => { const obsT = (source: string, value: number, ageMs: number, type: Observation["observationType"], rt: Observation["realtimeStatus"], rights: Observation["rightsStatus"] = "PUBLIC_ATTRIBUTED"): Observation => ({ ...obs(source, value, ageMs, "LAST_PRICE", rt, "SOURCE"), observationType: type, rightsStatus: rights, observationId: `${source}-${type}-${value}`, }); it("an official fixing is never compared to a live market value (no false divergence)", () => { const e = engine(); e.ingest(obsT("a", 1.1725, 500, "TRADE", "REALTIME")); // Kraken EUR/USD live e.ingest(obsT("eod", 1.1592, 20 * 3_600_000, "OFFICIAL_FIX", "END_OF_DAY", "OFFICIAL_OPEN_DATA")); // ECB fixing yesterday const q = e.compute("crypto_btc_usd", "EURUSD", NOW)!; expect(q.price).toBe(1.1725); expect(q.comparability).toBe("LIVE"); expect(q.dispersionBps).toBe(0); expect(q.contributions.find((c) => c.sourceId === "eod")).toMatchObject({ included: false, reason: "not_comparable" }); }); it("official fixing is canonical when no real live market exists; stablecoin proxies alone are INDICATIVE", () => { const e = engine(); e.ingest(obsT("eod", 1.1592, 20 * 3_600_000, "OFFICIAL_FIX", "END_OF_DAY", "OFFICIAL_OPEN_DATA")); e.ingest(obsT("b", 1.1601, 500, "STABLECOIN_PROXY", "REALTIME")); let q = e.compute("crypto_btc_usd", "EURUSD", NOW)!; expect(q.comparability).toBe("FIX"); expect(q.price).toBe(1.1592); const e2 = engine(); e2.ingest(obsT("b", 1.1601, 500, "STABLECOIN_PROXY", "REALTIME")); q = e2.compute("crypto_btc_usd", "EURUSD", NOW)!; expect(q.realtimeStatus).toBe("INDICATIVE"); expect(q.sourceCount).toBe(0); // proxies are not independent families expect(q.price).toBe(1.1601); }); it("proxies confirm a real market without setting the price; validators never vote", () => { const e = engine(); e.ingest(obsT("a", 100, 300, "TRADE", "REALTIME")); e.ingest(obsT("b", 100.02, 300, "STABLECOIN_PROXY", "REALTIME")); e.ingest(obsT("c", 100.01, 300, "TRADE", "DELAYED", "PUBLIC_RESTRICTED_REDISTRIBUTION")); const q = e.compute("crypto_btc_usd", "EURUSD", NOW)!; expect(q.price).toBe(100); expect(q.sourceCount).toBe(1); expect(q.proxyCount).toBe(1); expect(q.validatorCount).toBe(1); expect(q.contributions.find((c) => c.sourceId === "c")).toMatchObject({ included: false, reason: "validation_only" }); expect(q.contributions.find((c) => c.sourceId === "c")!.deltaBps).toBeCloseTo(1, 0); expect(q.rightsStatus).toBe("PUBLIC_ATTRIBUTED"); // restricted validator does not taint the canonical value }); it("same-class values far apart in time are temporal mismatches, not divergence", () => { const e = engine(); e.ingest(obsT("h", 91.01, 8 * 86_400_000, "EOD_CLOSE", "END_OF_DAY", "LICENSED")); // last week's close e.ingest(obsT("g", 95.5, 20 * 3_600_000, "EOD_CLOSE", "END_OF_DAY", "LICENSED")); // yesterday's close const q = e.compute("crypto_btc_usd", "CL=F", NOW)!; expect(q.price).toBe(95.5); expect(q.contributions.find((c) => c.sourceId === "h")).toMatchObject({ included: false, reason: "temporal_mismatch" }); }); });