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.
This commit is contained in:
@@ -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<SourceKind, FakeSourceAdapter>([["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<SourceKind, FakeSourceAdapter>([["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<unknown>("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<unknown>("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<SourceKind, FakeSourceAdapter>([["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");
|
||||
});
|
||||
Reference in New Issue
Block a user