From e6b7219dd278de0f266eaad6b992af4cc737b9d7 Mon Sep 17 00:00:00 2001 From: Investor Flow Build Date: Mon, 29 Jun 2026 22:09:12 -0400 Subject: [PATCH] slice 4d (omlx/ornith-35): backfill integration test (first-track 4 keys -> drain -> candles+adjustments; rerun no-op) Slice 4 (yfinance-backfill-permanent-ohlcv) COMPLETE: parseAdjustments + adjustments kind + candles 10y backfill + adjustments fetchOne + subscribe queues backfill. 102/102 tests green. Implemented by local ornith-35, reviewed by orchestrator. --- .../src/adapters/__tests__/backfill.test.ts | 191 ++++++++++++++++++ 1 file changed, 191 insertions(+) create mode 100644 app/server/src/adapters/__tests__/backfill.test.ts diff --git a/app/server/src/adapters/__tests__/backfill.test.ts b/app/server/src/adapters/__tests__/backfill.test.ts new file mode 100644 index 0000000..a0818bc --- /dev/null +++ b/app/server/src/adapters/__tests__/backfill.test.ts @@ -0,0 +1,191 @@ +// Investor Flow — slice-4 backfill integration test. +// +// Verifies the full subscribe → queue → drain → cache-write path end-to-end: +// 1. CacheRepository.subscribe('NVDA','equity') enqueues 4 keys into the AdapterQueue. +// 2. AdapterQueue.drain() processes them; FakeSourceAdapter serves canned responses. +// 3. Re-subscribing increments refcount but does NOT re-queue (no duplicate fetches). +// +// Conventions: Node 26 native TS; verbatimModuleSyntax; relative imports end in .ts; +// node --test + node:assert; ADR-0007 (no imperative trade verbs). +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { DatabaseSync } from "node:sqlite"; + +import type { CacheKey, SourceKind } from "../../cache/CacheRepository.ts"; +import { createCacheRepository } from "../../cache/CacheRepository.ts"; +import { AdapterQueue } from "../../queue/AdapterQueue.ts"; +import { FakeSourceAdapter } from "../SourceAdapter.ts"; +import { createDb, initSchema } from "../../db/client.ts"; + +test("subscribe enqueues 4 pending rows for NVDA/equity", async () => { + const db = createDb({ path: ":memory:" }); + initSchema(db); + + const fake = new FakeSourceAdapter("yfinance"); + fake.set("yfinance:quote:NVDA", { symbol: "NVDA", price: 131 }, "live_quote"); + fake.set("yfinance:symbol:NVDA", { symbol: "NVDA", sector: "Technology", tickerKind: "equity" }, "symbol_meta"); + fake.set( + "yfinance:candles:NVDA:1d", + [ + { ts: "2026-06-27", o: 192, h: 196, l: 191, c: 194.97, v: 1e8, adjClose: 194.9 }, + { ts: "2026-06-28", o: 194, h: 197, l: 193, c: 195, v: 9e7, adjClose: 195 }, + ], + "daily_permanent", + ); + fake.set( + "yfinance:adjustments:NVDA", + [ + { symbol: "NVDA", exDate: "2024-06-10", type: "split" as const, ratio: 10 }, + { symbol: "NVDA", exDate: "2026-02-27", type: "dividend" as const, ratio: 0.01 }, + ], + "daily_permanent", + ); + + const adapters = new Map([["yfinance", fake]]); + const queue = new AdapterQueue({ db, adapters, rateLimitMs: { yfinance: 0 } }); + const cache = createCacheRepository({ db, scheduler: queue }); + queue.cache = cache; + + await cache.subscribe("NVDA", "equity"); + + // adapter_queue should have exactly 4 pending rows. + const pending = db.prepare( + "SELECT key FROM adapter_queue WHERE status='pending'", + ).all() as Array<{ key: string }>; + + assert.equal(pending.length, 4, "expected 4 pending rows after subscribe"); + + const keys = pending.map((r) => r.key).sort(); + assert.deepEqual(keys, [ + "yfinance:adjustments:NVDA", + "yfinance:candles:NVDA:1d", + "yfinance:quote:NVDA", + "yfinance:symbol:NVDA", + ]); +}); + +test("drain processes all 4 keys and writes to cache + DB", async () => { + const db = createDb({ path: ":memory:" }); + initSchema(db); + + const fake = new FakeSourceAdapter("yfinance"); + fake.set("yfinance:quote:NVDA", { symbol: "NVDA", price: 131 }, "live_quote"); + fake.set("yfinance:symbol:NVDA", { symbol: "NVDA", sector: "Technology", tickerKind: "equity" }, "symbol_meta"); + fake.set( + "yfinance:candles:NVDA:1d", + [ + { ts: "2026-06-27", o: 192, h: 196, l: 191, c: 194.97, v: 1e8, adjClose: 194.9 }, + { ts: "2026-06-28", o: 194, h: 197, l: 193, c: 195, v: 9e7, adjClose: 195 }, + ], + "daily_permanent", + ); + fake.set( + "yfinance:adjustments:NVDA", + [ + { symbol: "NVDA", exDate: "2024-06-10", type: "split" as const, ratio: 10 }, + { symbol: "NVDA", exDate: "2026-02-27", type: "dividend" as const, ratio: 0.01 }, + ], + "daily_permanent", + ); + + const adapters = new Map([["yfinance", fake]]); + const queue = new AdapterQueue({ db, adapters, rateLimitMs: { yfinance: 0 } }); + const cache = createCacheRepository({ db, scheduler: queue }); + queue.cache = cache; + + await cache.subscribe("NVDA", "equity"); + await queue.drain(); + + // FakeSourceAdapter should have been called for candles and adjustments at minimum. + assert.ok( + fake.calls.includes("yfinance:candles:NVDA:1d"), + "expected candles fetch to have been invoked", + ); + assert.ok( + fake.calls.includes("yfinance:adjustments:NVDA"), + "expected adjustments fetch to have been invoked", + ); + + // Cache should now hold non-stale entries for candles and adjustments. + const candlesEntry = await cache.get("yfinance:candles:NVDA:1d"); + assert.equal(candlesEntry.isStale, false, "candles should not be stale after drain"); + assert.ok(Array.isArray(candlesEntry.value), "candles value should be an array"); + assert.equal((candlesEntry.value as unknown[]).length, 2, "candles should have 2 rows"); + + const adjEntry = await cache.get("yfinance:adjustments:NVDA"); + assert.equal(adjEntry.isStale, false, "adjustments should not be stale after drain"); + assert.ok(Array.isArray(adjEntry.value), "adjustments value should be an array"); + assert.equal((adjEntry.value as unknown[]).length, 2, "adjustments should have 2 rows"); + + // DB tables should reflect the written data. + const candleRows = db.prepare("SELECT * FROM price_candles WHERE symbol='NVDA'").all(); + assert.equal(candleRows.length, 2, "price_candles table should have 2 rows"); + + const adjRows = db.prepare("SELECT * FROM price_adjustments WHERE symbol='NVDA'").all(); + assert.equal(adjRows.length, 2, "price_adjustments table should have 2 rows"); +}); + +test("re-subscribe increments refcount but does NOT re-queue (no duplicate fetches)", async () => { + const db = createDb({ path: ":memory:" }); + initSchema(db); + + const fake = new FakeSourceAdapter("yfinance"); + fake.set("yfinance:quote:NVDA", { symbol: "NVDA", price: 131 }, "live_quote"); + fake.set("yfinance:symbol:NVDA", { symbol: "NVDA", sector: "Technology", tickerKind: "equity" }, "symbol_meta"); + fake.set( + "yfinance:candles:NVDA:1d", + [ + { ts: "2026-06-27", o: 192, h: 196, l: 191, c: 194.97, v: 1e8, adjClose: 194.9 }, + { ts: "2026-06-28", o: 194, h: 197, l: 193, c: 195, v: 9e7, adjClose: 195 }, + ], + "daily_permanent", + ); + fake.set( + "yfinance:adjustments:NVDA", + [ + { symbol: "NVDA", exDate: "2024-06-10", type: "split" as const, ratio: 10 }, + { symbol: "NVDA", exDate: "2026-02-27", type: "dividend" as const, ratio: 0.01 }, + ], + "daily_permanent", + ); + + const adapters = new Map([["yfinance", fake]]); + const queue = new AdapterQueue({ db, adapters, rateLimitMs: { yfinance: 0 } }); + const cache = createCacheRepository({ db, scheduler: queue }); + queue.cache = cache; + + // First subscribe + drain. + await cache.subscribe("NVDA", "equity"); + await queue.drain(); + + // Capture call count AFTER drain (all 4 fetches completed). + const afterFirstDrain = fake.calls.length; + + // Now re-subscribe — refcount goes to 2, but since NVDA is already first-demand + // and the cache is warm, no new queue entries should appear. + const rowsBefore = db.prepare("SELECT key FROM adapter_queue WHERE status='pending'").all(); + + await cache.subscribe("NVDA", "equity"); + + const rowsAfter = db.prepare("SELECT key FROM adapter_queue WHERE status='pending'").all(); + assert.deepEqual( + rowsAfter.map((r) => r.key).sort(), + rowsBefore.map((r) => r.key).sort(), + "re-subscribe should not add new pending rows", + ); + + // Drain again: no new work, so fake.calls.length must stay the same. + await queue.drain(); + assert.equal( + fake.calls.length, + afterFirstDrain, + "re-subscribe + drain should not trigger additional fetches", + ); + + // DB tables should still have exactly 2 rows each (no duplicates). + const candleRows = db.prepare("SELECT * FROM price_candles WHERE symbol='NVDA'").all(); + assert.equal(candleRows.length, 2, "price_candles should still have 2 rows after re-subscribe"); + + const adjRows = db.prepare("SELECT * FROM price_adjustments WHERE symbol='NVDA'").all(); + assert.equal(adjRows.length, 2, "price_adjustments should still have 2 rows after re-subscribe"); +});