feat: dealer flow, mirror portfolio (M21), options convexity, FINRA short interest, alert producers, vendor gate
CI / Test & Type-Check (push) Canceled after 0s

Snapshot of in-progress module work across multiple slices:

- Dealer Flow: dealerExposureEngine, dealerMapService, dealerMapExplain,
  dealerMapIntegrity, dealerMapReplay, dealerStudyEngine, hanStyleLevels
- Mirror Portfolio (M21): fundRepository, captureIngest, mirrorAlertProducers,
  fund holdings strip, live book, position capture ingest
- Options: BSM, NormalizedOptionSurface types, OptionsChainRouter,
  ConvexityGate, option legs panel
- Alert producers: vixLevel, rotation, thesis, unlock, portfolioRisk,
  mirror (fund_capture, fund_13f, mirror_diff)
- FINRA short interest adapter + queue integration
- SEC company tickers adapter + ingest (symbol search index seed)
- Vendor gate (rate-limit-first data plane, ADR-0009)
- CUSIP registry, reverse 13F refresh, stock float service
- LRU cache, portfolio backtest engine
- Frontend: dealer-flow, funds, journal, lab, monitor, plan, portfolio,
  reports, screener, strategies, theses, guided-start, exits, more pages
- Volume profile, workspace profile, visibility-aware poll
- ADRs 0010 (mirror math not advice), 0011 (symbol search index)
- VENDOR_INTEGRATIONS.md, END_USER_TEST.md
- .gitignore: exclude DBs, .DS_Store, local config, agent scratch
This commit is contained in:
Investor Flow Build
2026-08-10 13:36:26 -04:00
parent 04fc11b2fd
commit ac94acf9e3
229 changed files with 32617 additions and 3934 deletions
+158 -1
View File
@@ -90,9 +90,11 @@ export async function downloadAndIngestFinra(
const ingestedAt = new Date().toISOString();
console.log(`[finra] downloading ${url}`);
const resp = await fetch(url, {
const { vendorFetch } = await import('./vendorGate.ts');
const resp = await vendorFetch('finra', url, {
headers: { 'User-Agent': 'InvestorFlow/1.0 (research) node' },
signal: AbortSignal.timeout(30_000),
hostAllowlist: /finra\.org|cdn\.finra\.org|files\.finra\.org/i,
});
if (!resp.ok) throw new Error(`FINRA download failed: ${resp.status} ${resp.statusText}`);
const body = await resp.text();
@@ -153,4 +155,159 @@ export function latestFinraSettlementDate(db: DatabaseSync): string | null {
'SELECT settlement_date FROM finra_short_interest ORDER BY settlement_date DESC LIMIT 1'
).get() as { settlement_date: string } | undefined;
return r?.settlement_date ?? null;
}
/** Backfill FINRA files for the last N calendar days. Skips 404s (weekends/holidays) with 500ms rate-limit. */
export async function backfillFinra(
db: DatabaseSync,
days: number = 30,
baseUrl?: string,
): Promise<{ filesAttempted: number; filesStored: number; firstDate: string | null; lastDate: string | null; errors: string[] }> {
const errors: string[] = [];
let filesAttempted = 0;
let filesStored = 0;
let firstDate: string | null = null;
let lastDate: string | null = null;
const existing = new Set(
(db.prepare('SELECT DISTINCT settlement_date FROM finra_short_interest').all() as Array<{ settlement_date: string }>)
.map((r) => r.settlement_date),
);
for (let i = days; i >= 0; i--) {
const d = new Date(Date.now() - i * 86400000);
const dateStr = d.toISOString().slice(0, 10);
if (existing.has(dateStr)) continue;
try {
await downloadAndIngestFinra(db, dateStr, baseUrl);
filesStored += 1;
if (firstDate === null) firstDate = dateStr;
lastDate = dateStr;
} catch (e) {
const msg = e instanceof Error ? e.message : String(e);
if (msg.includes('404') || msg.includes('40')) {
// Weekend/holiday — expected
} else {
errors.push(`${dateStr}: ${msg}`);
}
}
filesAttempted += 1;
await new Promise((r) => setTimeout(r, 500));
}
return { filesAttempted, filesStored, firstDate, lastDate, errors };
}
// ============================================================================
// FINRA bi-monthly short INTEREST (outstanding positions) — shrt{DATE}.csv
// Separate dataset from daily volume; comparable to NASDAQ/Yahoo.
// ============================================================================
const FINRA_SI_BASE = 'https://cdn.finra.org/equity/otcmarket/biweekly';
function finraSiFilename(settlementDate: string): string {
return `shrt${settlementDate.replace(/-/g, '')}.csv`;
}
function parseFinraSiFile(body: string, settlementDate: string, ingestedAt: string, sourceFile: string): Array<{
symbol: string; settlementDate: string; issueName: string | null; exchangeCode: string | null;
marketClass: string | null; currentShortPosition: number; previousShortPosition: number | null;
avgDailyVolume: number | null; daysToCover: number | null; changePercent: number | null;
changePrevious: number | null; revisionFlag: string | null;
}> {
const lines = body.split(/\r?\n/);
type SiRow = {
symbol: string; settlementDate: string; issueName: string | null; exchangeCode: string | null;
marketClass: string | null; currentShortPosition: number; previousShortPosition: number | null;
avgDailyVolume: number | null; daysToCover: number | null; changePercent: number | null;
changePrevious: number | null; revisionFlag: string | null;
};
const rows: SiRow[] = [];
let headerFound = false;
for (const raw of lines) {
const line = raw.trim();
if (!line || line.startsWith('#')) continue;
if (line.includes('accountingYearMonthNumber|symbolCode')) { headerFound = true; continue; }
if (!headerFound) continue;
const c = line.split('|').map((x) => x.trim());
if (c.length < 6 || !c[1]) continue;
const symbol = c[1].replace(/\/.*$/, '')?.toUpperCase();
if (!symbol) continue;
const cur = parseFloat((c[5] ?? '0').replace(/,/g, ''));
if (Number.isNaN(cur)) continue;
rows.push({
symbol, settlementDate,
issueName: c[2] || null, exchangeCode: c[3] || null, marketClass: c[4] || null,
currentShortPosition: cur,
previousShortPosition: c[6] ? parseFloat(c[6].replace(/,/g, '')) || null : null,
avgDailyVolume: c[8] ? parseFloat(c[8].replace(/,/g, '')) || null : null,
daysToCover: c[9] ? parseFloat(c[9].replace(/,/g, '')) || null : null,
changePercent: c[11] ? parseFloat(c[11].replace(/,/g, '')) || null : null,
changePrevious: c[12] ? parseFloat(c[12].replace(/,/g, '')) || null : null,
revisionFlag: c[10] || null,
});
}
return rows;
}
export async function downloadAndIngestFinraSi(
db: DatabaseSync,
settlementDate: string,
): Promise<{ symbolsStored: number; sourceFile: string }> {
const filename = finraSiFilename(settlementDate);
const url = `${FINRA_SI_BASE}/${filename}`;
const ingestedAt = new Date().toISOString();
console.log(`[finra-si] downloading ${url}`);
const { vendorFetch } = await import('./vendorGate.ts');
const resp = await vendorFetch('finra', url, {
headers: { 'User-Agent': 'Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36' },
signal: AbortSignal.timeout(30_000),
hostAllowlist: /finra\.org|cdn\.finra\.org|files\.finra\.org/i,
});
if (!resp.ok) throw new Error(`FINRA SI download failed: ${resp.status}`);
const body = await resp.text();
if (!body.trim()) throw new Error('FINRA SI file empty');
const rows = parseFinraSiFile(body, settlementDate, ingestedAt, filename);
if (!rows.length) throw new Error('No FINRA SI rows parsed');
db.exec('BEGIN');
let inserted = 0;
for (const r of rows) {
try {
db.exec(`INSERT OR REPLACE INTO finra_short_interest_biweekly (symbol,settlement_date,issue_name,exchange_code,market_class,current_short_position,previous_short_position,avg_daily_volume,days_to_cover,change_percent,change_previous,revision_flag,source_file,ingested_at) VALUES ('${r.symbol.replace(/'/g,"''")}','${r.settlementDate}','${(r.issueName??'').replace(/'/g,"''")}','${(r.exchangeCode??'').replace(/'/g,"''")}','${(r.marketClass??'').replace(/'/g,"''")}',${r.currentShortPosition},${r.previousShortPosition??'NULL'},${r.avgDailyVolume??'NULL'},${r.daysToCover??'NULL'},${r.changePercent??'NULL'},${r.changePrevious??'NULL'},${r.revisionFlag??'NULL'},'${filename.replace(/'/g,"''")}','${ingestedAt}')`);
inserted++;
} catch (re) { console.log(`[finra-si] skip row ${r.symbol} ${(re as Error).message}`); }
}
db.exec('COMMIT');
console.log(`[finra-si] ingested ${inserted}/${rows.length} rows from ${filename}`);
return { symbolsStored: inserted, sourceFile: filename };
}
export function latestFinraSiSettlementDate(db: DatabaseSync): string | null {
const r = db.prepare('SELECT settlement_date FROM finra_short_interest_biweekly ORDER BY settlement_date DESC LIMIT 1').get() as { settlement_date: string } | undefined;
return r?.settlement_date ?? null;
}
/** Backfill FINRA bi-monthly files: probe date windows around 15th and month-end, skip 403s. */
export async function backfillFinraSi(db: DatabaseSync, months: number = 24): Promise<{ filesStored: number; firstDate: string | null; lastDate: string | null; errors: string[] }> {
const errors: string[] = [];
let filesStored = 0;
let firstDate: string | null = null;
let lastDate: string | null = null;
const existing = new Set((db.prepare('SELECT DISTINCT settlement_date FROM finra_short_interest_biweekly').all() as Array<{ settlement_date: string }>).map(r => r.settlement_date));
const now = new Date();
const candidates: string[] = [];
for (let m = 0; m < months; m++) {
const d = new Date(now.getFullYear(), now.getMonth() - m, 1);
for (let day of [13,14,15,16,17]) candidates.push(new Date(d.getFullYear(), d.getMonth(), day).toISOString().slice(0,10));
const lastDay = new Date(d.getFullYear(), d.getMonth()+1, 0).getDate();
for (let day of [28,29,30,31]) if (day <= lastDay) candidates.push(new Date(d.getFullYear(), d.getMonth(), day).toISOString().slice(0,10));
for (let day of [1,2,3]) candidates.push(new Date(d.getFullYear(), d.getMonth()+1, day).toISOString().slice(0,10));
}
const valid = [...new Set(candidates)].filter((x): x is string => !!x).sort().reverse();
for (const dateStr of valid) {
if (existing.has(dateStr)) continue;
try { await downloadAndIngestFinraSi(db, dateStr); filesStored++; if (!firstDate) firstDate = dateStr; lastDate = dateStr; }
catch (e) { const m = e instanceof Error ? e.message : String(e); if (!m.includes('403')) errors.push(`${dateStr}: ${m}`); }
await new Promise(r => setTimeout(r, 500));
}
return { filesStored, firstDate, lastDate, errors };
}
@@ -0,0 +1,127 @@
// Capture materializer tests — X posts → fund_position_records (source='capture').
import { test } from 'node:test';
import assert from 'node:assert/strict';
import { DatabaseSync } from 'node:sqlite';
import { ingestFundCaptures, ingestAllFundCaptures, listTrackedFundHandles } from '../captureIngest.ts';
function makeDb(): DatabaseSync {
const db = new DatabaseSync(':memory:');
db.exec(`
CREATE TABLE tracked_funds (
id TEXT PRIMARY KEY, ci_key TEXT NOT NULL, fund_name TEXT NOT NULL,
manager_name TEXT NOT NULL, x_handle TEXT, paywall_status TEXT NOT NULL DEFAULT 'unknown',
enabled INTEGER NOT NULL DEFAULT 1, created_at TEXT NOT NULL, updated_at TEXT
);
CREATE TABLE fund_position_records (
id TEXT PRIMARY KEY, fund_id TEXT NOT NULL, symbol TEXT NOT NULL,
shares REAL, value_usd REAL, cost_basis REAL, as_of TEXT NOT NULL,
source TEXT NOT NULL, evidence_url TEXT, notes TEXT, created_at TEXT NOT NULL
);
CREATE TABLE x_cookie_posts (
post_id TEXT PRIMARY KEY, source TEXT, author_handle TEXT, cashtag TEXT,
body_text TEXT, posted_at TEXT, engagement TEXT, cached_until TEXT
);
`);
db.prepare(`INSERT INTO tracked_funds (id, ci_key, fund_name, manager_name, x_handle, paywall_status, enabled, created_at)
VALUES ('alpine-fox-capital', '0002096493', 'Alpine Fox Capital LLC', 'Mike Alfred', 'mikealfred', 'paywalled', 1, ?)`)
.run(new Date().toISOString());
return db;
}
function seedPost(db: DatabaseSync, id: string, body: string, at: string): void {
db.prepare(`INSERT INTO x_cookie_posts (post_id, source, author_handle, cashtag, body_text, posted_at)
VALUES (?, 'x', 'mikealfred', NULL, ?, ?)`).run(id, body, at);
}
test('ingestFundCaptures: materializes a total-position capture with source=capture', () => {
const db = makeDb();
seedPost(db, 'p1', '**Real time position update\n\nTook OPEN over 5.7M shares now. Brought average down to $4.42.', '2026-07-30T19:51:12.000Z');
const s = ingestFundCaptures(db, 'alpine-fox-capital');
assert.equal(s.postsScanned, 1);
assert.equal(s.captures, 1);
assert.equal(s.inserted, 1);
const row = db.prepare('SELECT * FROM fund_position_records').get() as any;
assert.equal(row.fund_id, 'alpine-fox-capital');
assert.equal(row.symbol, 'OPEN');
assert.equal(row.shares, 5_700_000);
assert.ok(Math.abs(row.cost_basis - 4.42) < 0.001);
assert.equal(row.source, 'capture');
assert.equal(row.as_of, '2026-07-30');
assert.equal(row.evidence_url, 'https://x.com/mikealfred/status/p1');
});
test('ingestFundCaptures: idempotent — same post converges, no duplicate rows', () => {
const db = makeDb();
seedPost(db, 'p1', '**Real time position update\n\nTook OPEN over 5.7M shares now. Brought average down to $4.42.', '2026-07-30T19:51:12.000Z');
ingestFundCaptures(db, 'alpine-fox-capital');
const again = ingestFundCaptures(db, 'alpine-fox-capital');
assert.equal(again.inserted, 0);
assert.equal(again.refreshed, 1);
assert.equal(db.prepare('SELECT COUNT(*) FROM fund_position_records').get()!['COUNT(*)'], 1);
});
test('ingestFundCaptures: commentary posts and claims are never materialized', () => {
const db = makeDb();
seedPost(db, 'c1', 'Beautiful day in the mountains. Great conversations with subscribers.', '2026-08-04T23:07:20.000Z');
seedPost(db, 'cl1', 'Added 865,000 shares of BKKT at $3.10 average price today.', '2026-08-01T10:00:00.000Z');
const s = ingestFundCaptures(db, 'alpine-fox-capital');
assert.equal(s.skipped, 1); // commentary
assert.equal(s.claims, 1); // claim detected but skipped in v1
assert.equal(db.prepare('SELECT COUNT(*) FROM fund_position_records').get()!['COUNT(*)'], 0);
});
test('ingestFundCaptures: trade narrative with two quantities is refused (no garbage)', () => {
const db = makeDb();
seedPost(db, 'n1', 'I reduced the position from 200,000 shares to 10,000 shares over the last month. Today I added some shares back just above $46.', '2026-08-02T10:00:00.000Z');
const s = ingestFundCaptures(db, 'alpine-fox-capital');
assert.equal(s.skipped, 1);
assert.equal(db.prepare('SELECT COUNT(*) FROM fund_position_records').get()!['COUNT(*)'], 0);
});
test('ingestFundCaptures: disabled fund or missing handle is skipped', () => {
const db = makeDb();
seedPost(db, 'p1', '**Real time position update\n\nTook OPEN over 5.7M shares now. Brought average down to $4.42.', '2026-07-30T19:51:12.000Z');
db.prepare(`UPDATE tracked_funds SET enabled = 0 WHERE id = 'alpine-fox-capital'`).run();
const s = ingestFundCaptures(db, 'alpine-fox-capital');
assert.equal(s.disabled, true);
assert.equal(db.prepare('SELECT COUNT(*) FROM fund_position_records').get()!['COUNT(*)'], 0);
});
test('ingestFundCaptures: multi-name update materializes OPEN total (SLNH claim skipped)', () => {
const db = makeDb();
seedPost(
db,
'p-multi',
`A few quick position updates.
Just filled a 100,000 share order in SLNH at $1.12. Have another for 100,000 shares at $1.11.
OPEN now over 6.45M shares. Have taken it up on the weakness.`,
'2026-08-07T15:45:54.000Z',
);
const s = ingestFundCaptures(db, 'alpine-fox-capital');
assert.ok(s.captures >= 1);
assert.equal(s.inserted, 1);
assert.equal(s.claims, 1); // SLNH fill recorded as claim, not book row
const open = db.prepare(`SELECT * FROM fund_position_records WHERE symbol='OPEN'`).get() as any;
assert.equal(open.shares, 6_450_000);
assert.equal(open.as_of, '2026-08-07');
assert.equal(
db.prepare(`SELECT COUNT(*) AS n FROM fund_position_records WHERE symbol='SLNH'`).get()!['n'],
0,
);
});
test('ingestAllFundCaptures: runs across every enabled tracked fund with a handle', () => {
const db = makeDb();
seedPost(db, 'p1', '**Real time position update\n\nTook OPEN over 5.7M shares now. Brought average down to $4.42.', '2026-07-30T19:51:12.000Z');
const stats = ingestAllFundCaptures(db);
assert.equal(stats.length, 1);
assert.equal(stats[0].captures, 1);
assert.deepEqual(listTrackedFundHandles(db).map((f) => f.id), ['alpine-fox-capital']);
});
@@ -0,0 +1,132 @@
import { test } from 'node:test';
import { strict as assert } from 'node:assert';
import { DatabaseSync } from 'node:sqlite';
import {
MAJOR_13F_FILER_CIKS,
selectMissingMajorFilers,
selectReverseFilerCiks,
} from '../reverse13fRefresh.ts';
function makeDb(): DatabaseSync {
const db = new DatabaseSync(':memory:');
db.exec(`
CREATE TABLE institution_filings (
filer_cik TEXT, filer_name TEXT, filer_sic TEXT, symbol TEXT, form TEXT,
shares REAL, value_usd REAL, reported_quarter TEXT, filed_at TEXT,
fetched_at TEXT, put_call TEXT, accession TEXT
);
CREATE TABLE tracked_funds (
id TEXT PRIMARY KEY, ci_key TEXT, fund_name TEXT, manager_name TEXT,
x_handle TEXT, paywall_status TEXT, enabled INTEGER, created_at TEXT, updated_at TEXT
);
CREATE TABLE kv_cache (key TEXT PRIMARY KEY, value TEXT NOT NULL, observed_at TEXT NOT NULL);
`);
return db;
}
test('selectReverseFilerCiks prefers prior holders then tracked then majors', () => {
const db = makeDb();
db.prepare(
`INSERT INTO institution_filings
(filer_cik, filer_name, symbol, form, shares, value_usd, reported_quarter, filed_at, fetched_at)
VALUES (?, ?, 'IREN', '13F-HR', ?, 1, '2026-Q2', '2026-07-01', '2026-07-01')`,
).run('0001535385', 'Artemis', 739723);
db.prepare(
`INSERT INTO institution_filings
(filer_cik, filer_name, symbol, form, shares, value_usd, reported_quarter, filed_at, fetched_at)
VALUES (?, ?, 'IREN', '13F-HR', ?, 1, '2026-Q2', '2026-07-01', '2026-07-01')`,
).run('0002012155', 'VIMA', 116245);
db.prepare(
`INSERT INTO tracked_funds
(id, ci_key, fund_name, manager_name, paywall_status, enabled, created_at)
VALUES ('alpine', '0002096493', 'Alpine Fox', 'Mike', 'open', 1, '2026-01-01')`,
).run();
const filers = selectReverseFilerCiks(db, 'IREN', 10, 5);
assert.ok(filers.length >= 3);
assert.equal(filers[0].cik, '0001535385');
assert.equal(filers[0].source, 'prior_holder');
assert.ok(filers.some((f) => f.cik === '0002096493' && f.source === 'tracked_fund'));
assert.ok(filers.some((f) => f.source === 'major'));
// no duplicates
const ciks = filers.map((f) => f.cik);
assert.equal(ciks.length, new Set(ciks).size);
});
test('selectReverseFilerCiks skips non-numeric tracked fund keys', () => {
const db = makeDb();
db.prepare(
`INSERT INTO tracked_funds
(id, ci_key, fund_name, manager_name, paywall_status, enabled, created_at)
VALUES ('mr-t', 'mr-t-invests', 'Mr T', 'T', 'open', 1, '2026-01-01')`,
).run();
const filers = selectReverseFilerCiks(db, 'AAPL', 5, 3);
assert.ok(filers.every((f) => /^\d{10}$/.test(f.cik)));
assert.ok(filers.length >= 1);
assert.ok(MAJOR_13F_FILER_CIKS.length >= 10);
});
test('MAJOR_13F includes current BlackRock Inc CIK (BLK), not old Finance CIK', () => {
const br = MAJOR_13F_FILER_CIKS.find((m) => /blackrock/i.test(m.name));
assert.ok(br);
assert.equal(br!.cik, '0002012383');
assert.ok(!MAJOR_13F_FILER_CIKS.some((m) => m.cik === '0001364742'));
});
test('selectMissingMajorFilers excludes majors already in latest quarter', () => {
const db = makeDb();
// Seed latest quarter with BlackRock present, State Street absent
db.prepare(
`INSERT INTO institution_filings
(filer_cik, filer_name, symbol, form, shares, value_usd, reported_quarter, filed_at, fetched_at)
VALUES ('0002012383', 'BlackRock, Inc.', 'IREN', '13F-HR', 1e7, 1, '2026-Q2', '2026-08-07', '2026-08-07')`,
).run();
db.prepare(
`INSERT INTO institution_filings
(filer_cik, filer_name, symbol, form, shares, value_usd, reported_quarter, filed_at, fetched_at)
VALUES ('0001535385', 'Artemis', 'IREN', '13F-HR', 1e5, 1, '2026-Q2', '2026-07-01', '2026-07-01')`,
).run();
const missing = selectMissingMajorFilers(db, 'IREN');
assert.ok(!missing.some((m) => m.cik === '0002012383'), 'BlackRock already in Q2');
assert.ok(missing.some((m) => m.cik === '0000093751'), 'State Street still missing');
assert.ok(missing.every((m) => m.source === 'missing_major'));
});
test('selectMissingMajorFilers does not thrash fresh prior-quarter majors', () => {
const db = makeDb();
// Frontier is Q2 via a small filer; Goldman only has freshly fetched Q1
db.prepare(
`INSERT INTO institution_filings
(filer_cik, filer_name, symbol, form, shares, value_usd, reported_quarter, filed_at, fetched_at)
VALUES ('0001535385', 'Artemis', 'IREN', '13F-HR', 1e5, 1, '2026-Q2', '2026-07-01', '2026-08-07T12:00:00.000Z')`,
).run();
db.prepare(
`INSERT INTO institution_filings
(filer_cik, filer_name, symbol, form, shares, value_usd, reported_quarter, filed_at, fetched_at)
VALUES ('0000886982', 'Goldman', 'IREN', '13F-HR', 1e7, 1, '2026-Q1', '2026-05-15', '2026-08-07T18:00:00.000Z')`,
).run();
const now = Date.parse('2026-08-08T00:00:00.000Z');
const missing = selectMissingMajorFilers(db, 'IREN', { now, staleAfterMs: 10 * 24 * 3600_000 });
assert.ok(!missing.some((m) => m.cik === '0000886982'), 'Goldman Q1 freshly fetched — wait');
assert.ok(missing.some((m) => m.cik === '0002012383'), 'BlackRock never stored');
});
test('selectMissingMajorFilers rechecks stale prior-quarter majors', () => {
const db = makeDb();
db.prepare(
`INSERT INTO institution_filings
(filer_cik, filer_name, symbol, form, shares, value_usd, reported_quarter, filed_at, fetched_at)
VALUES ('0001535385', 'Artemis', 'IREN', '13F-HR', 1e5, 1, '2026-Q2', '2026-07-01', '2026-08-07')`,
).run();
db.prepare(
`INSERT INTO institution_filings
(filer_cik, filer_name, symbol, form, shares, value_usd, reported_quarter, filed_at, fetched_at)
VALUES ('0000886982', 'Goldman', 'IREN', '13F-HR', 1e7, 1, '2026-Q1', '2026-05-15', '2026-07-01T00:00:00.000Z')`,
).run();
const now = Date.parse('2026-08-08T00:00:00.000Z');
const missing = selectMissingMajorFilers(db, 'IREN', { now, staleAfterMs: 10 * 24 * 3600_000 });
assert.ok(missing.some((m) => m.cik === '0000886982'), 'Goldman Q1 stale vs Q2 frontier');
});
@@ -0,0 +1,66 @@
import { test } from 'node:test';
import { strict as assert } from 'node:assert';
import { DatabaseSync } from 'node:sqlite';
import {
curatedCusipForSymbol,
resolveCusipLocal,
seedCuratedCusips,
} from '../cusipRegistry.ts';
/**
* Unit-level tests for CUSIP cache helpers via sec fetch behavior contracts.
* Network resolve is not exercised here.
*/
test('kv_cache can store and read sec:cusip keys', () => {
const db = new DatabaseSync(':memory:');
db.exec(`
CREATE TABLE kv_cache (key TEXT PRIMARY KEY, value TEXT NOT NULL, observed_at TEXT NOT NULL);
`);
const key = 'sec:cusip:IREN';
const cusip = 'Q4982L109';
db.prepare('INSERT OR REPLACE INTO kv_cache (key,value,observed_at) VALUES (?,?,?)').run(
key,
cusip,
new Date().toISOString(),
);
const row = db.prepare('SELECT value FROM kv_cache WHERE key=?').get(key) as { value: string };
assert.equal(row.value, cusip);
});
test('curated map resolves IREN offline', () => {
assert.equal(curatedCusipForSymbol('IREN'), 'Q4982L109');
assert.equal(curatedCusipForSymbol('AAPL'), '037833100');
});
test('resolveCusipLocal seeds cache from curated map', () => {
const db = new DatabaseSync(':memory:');
db.exec(`
CREATE TABLE kv_cache (key TEXT PRIMARY KEY, value TEXT NOT NULL, observed_at TEXT NOT NULL);
`);
const c = resolveCusipLocal(db, 'IREN');
assert.equal(c, 'Q4982L109');
const row = db.prepare('SELECT value FROM kv_cache WHERE key=?').get('sec:cusip:IREN') as {
value: string;
};
assert.equal(row.value, 'Q4982L109');
});
test('seedCuratedCusips is idempotent', () => {
const db = new DatabaseSync(':memory:');
db.exec(`
CREATE TABLE kv_cache (key TEXT PRIMARY KEY, value TEXT NOT NULL, observed_at TEXT NOT NULL);
`);
const n1 = seedCuratedCusips(db);
const n2 = seedCuratedCusips(db);
assert.ok(n1 > 0);
assert.equal(n2, 0);
});
test('isPermanentDataError does not quarantine cusip failures', async () => {
const { isPermanentDataError } = await import('../../queue/sourceRatePolicy.ts');
assert.equal(
isPermanentDataError('sec-fetch:IREN: cusip not resolved - 13F holder refresh blocked'),
false,
);
assert.equal(isPermanentDataError('quote not found for symbol'), true);
});
@@ -0,0 +1,68 @@
import { test, beforeEach } from 'node:test';
import { strict as assert } from 'node:assert';
import {
isSecHttpCoolingDown,
noteSecRateLimitHit,
resetSecHttpStateForTests,
secFetch,
SecRateLimitError,
SEC_MIN_INTERVAL_MS,
} from '../secHttp.ts';
beforeEach(() => {
resetSecHttpStateForTests();
});
test('secFetch refuses non-SEC hosts', async () => {
await assert.rejects(() => secFetch('https://example.com/x'), /non-SEC/);
});
test('secFetch serializes and spaces requests', async () => {
const times: number[] = [];
const original = globalThis.fetch;
globalThis.fetch = (async () => {
times.push(Date.now());
return new Response(JSON.stringify({ ok: true }), { status: 200 });
}) as typeof fetch;
try {
// Use sequential calls to verify spacing between each pair.
await secFetch('https://data.sec.gov/a');
await secFetch('https://data.sec.gov/b');
await secFetch('https://www.sec.gov/c');
assert.equal(times.length, 3);
// Each sequential start must be spaced by ~SEC_MIN_INTERVAL_MS.
assert.ok(times[1] - times[0] >= SEC_MIN_INTERVAL_MS - 20, `gap1=${times[1] - times[0]}ms`);
assert.ok(times[2] - times[1] >= SEC_MIN_INTERVAL_MS - 20, `gap2=${times[2] - times[1]}ms`);
} finally {
globalThis.fetch = original;
}
});
test('429 sets cool-down and throws SecRateLimitError', async () => {
const original = globalThis.fetch;
globalThis.fetch = (async () =>
new Response('Request Rate Threshold Exceeded', { status: 429 })) as typeof fetch;
try {
await assert.rejects(() => secFetch('https://data.sec.gov/x', { retries: 0 }), (e: unknown) => {
assert.ok(e instanceof SecRateLimitError);
return true;
});
assert.equal(isSecHttpCoolingDown(), true);
} finally {
globalThis.fetch = original;
}
});
test('noteSecRateLimitHit escalates cool-down', () => {
const a = noteSecRateLimitHit();
const b = noteSecRateLimitHit();
assert.ok(b >= a);
assert.equal(isSecHttpCoolingDown(), true);
});
test('preflight rejects while cooling', async () => {
noteSecRateLimitHit();
await assert.rejects(() => secFetch('https://data.sec.gov/y', { retries: 0 }), /preflight cool-down|rate limit/i);
});
@@ -0,0 +1,107 @@
import { test } from 'node:test';
import assert from 'node:assert/strict';
import { DatabaseSync } from 'node:sqlite';
import { mergeCompanyTickers } from '../secTickersIngest.ts';
function makeDb(): DatabaseSync {
const db = new DatabaseSync(':memory:');
db.exec(`
CREATE TABLE symbols (
symbol TEXT PRIMARY KEY,
name TEXT, sector TEXT, industry TEXT, exchange TEXT,
ticker_kind TEXT NOT NULL,
peers TEXT, cik TEXT, updated_at TEXT
);
`);
return db;
}
// A small synthetic SEC payload (shape matches company_tickers.json).
const secRows = {
'0': { cik_str: 2096493, ticker: 'IREN', title: 'IREN Limited' },
'1': { cik_str: 320193, ticker: 'AAPL', title: 'Apple Inc.' },
'2': { cik_str: 1652044, ticker: 'NVDA', title: 'NVIDIA CORPORATION' },
};
test('mergeCompanyTickers: inserts new equity rows with zero-padded issuer CIK', () => {
const db = makeDb();
const r = mergeCompanyTickers(db, secRows);
assert.equal(r.created, 3);
assert.equal(r.upserted, 0);
const row = db.prepare('SELECT * FROM symbols WHERE symbol = ?').get('AAPL') as any;
assert.equal(row.cik, '0000320193'); // zero-padded to 10
assert.equal(row.ticker_kind, 'equity');
assert.equal(row.name, 'Apple Inc.');
assert.equal(row.sector, null);
assert.equal(row.exchange, null);
});
test('mergeCompanyTickers: fills cik on existing rows, never downgrades ticker_kind', () => {
const db = makeDb();
// Simulate an existing yfinance-hydrated ETF row WITHOUT a cik.
db.prepare(
"INSERT INTO symbols (symbol, name, sector, ticker_kind, updated_at) VALUES ('IBIT', 'iShares Bitcoin Trust', 'Crypto', 'etf', ?)",
).run(new Date().toISOString());
// SEC payload includes IBIT as a listed issuer.
const r = mergeCompanyTickers(db, {
'0': { cik_str: 1990438, ticker: 'IBIT', title: 'iShares Bitcoin Trust' },
});
assert.equal(r.upserted, 1);
const row = db.prepare('SELECT * FROM symbols WHERE symbol = ?').get('IBIT') as any;
assert.equal(row.cik, '0001990438'); // cik filled by SEC
assert.equal(row.ticker_kind, 'etf'); // NOT downgraded to equity
assert.equal(row.sector, 'Crypto'); // descriptive metadata untouched
});
test('mergeCompanyTickers: keeps richer yfinance name over SEC title on conflict', () => {
const db = makeDb();
db.prepare(
"INSERT INTO symbols (symbol, name, ticker_kind, updated_at) VALUES ('NVDA', 'NVIDIA Corporation (Detailed)', 'equity', ?)",
).run(new Date().toISOString());
const r = mergeCompanyTickers(db, { '0': { cik_str: 1652044, ticker: 'NVDA', title: 'NVIDIA CORPORATION' } });
assert.equal(r.upserted, 1);
const row = db.prepare('SELECT symbol, name, cik FROM symbols WHERE symbol = ?').get('NVDA') as any;
assert.equal(row.name, 'NVIDIA Corporation (Detailed)'); // existing richer name preserved
assert.equal(row.cik, '0001652044'); // cik still filled
});
test('mergeCompanyTickers: rows absent from SEC payload are never purged', () => {
const db = makeDb();
db.prepare(
"INSERT INTO symbols (symbol, name, ticker_kind, updated_at) VALUES ('BTC-USD', 'Bitcoin', 'crypto', ?)",
).run(new Date().toISOString());
mergeCompanyTickers(db, secRows); // BTC-USD not in payload
const row = db.prepare('SELECT symbol FROM symbols WHERE symbol = ?').get('BTC-USD') as any;
assert.ok(row, 'crypto row must survive a SEC-materialization pass');
});
test('mergeCompanyTickers: skips empty tickers and is idempotent', () => {
const db = makeDb();
const withJunk = { ...secRows, '9': { cik_str: 1, ticker: '', title: 'No ticker' } };
const r1 = mergeCompanyTickers(db, withJunk);
assert.equal(r1.created, 3);
assert.equal(r1.skipped, 1);
// Second pass: no new rows, all existing rows counted as upserted.
const r2 = mergeCompanyTickers(db, secRows);
assert.equal(r2.created, 0);
assert.equal(r2.upserted, 3);
});
test('mergeCompanyTickers: is atomic — a failing row rolls back the batch', () => {
const db = makeDb();
// Force a failure mid-batch: duplicate primary key cannot happen here, so instead
// verify a throwing payload leaves zero rows behind.
assert.throws(() => {
mergeCompanyTickers(db, null as unknown as Record<string, never>);
});
const count = (db.prepare('SELECT COUNT(*) AS n FROM symbols').get() as any).n;
assert.equal(count, 0);
});
@@ -0,0 +1,95 @@
import { test } from 'node:test';
import assert from 'node:assert/strict';
import { DatabaseSync } from 'node:sqlite';
import { mergeCompanyTickers } from '../secTickersIngest.ts';
function makeDb(): DatabaseSync {
const db = new DatabaseSync(':memory:');
db.exec(`
CREATE TABLE symbols (
symbol TEXT PRIMARY KEY,
name TEXT, sector TEXT, industry TEXT, exchange TEXT,
ticker_kind TEXT NOT NULL,
peers TEXT, cik TEXT, updated_at TEXT
);
`);
return db;
}
const secRows = {
'0': { cik_str: 1878848, ticker: 'IREN', title: 'IREN Limited' },
'1': { cik_str: 320193, ticker: 'AAPL', title: 'Apple Inc.' },
'2': { cik_str: 1652044, ticker: 'NVDA', title: 'NVIDIA CORPORATION' },
'3': { cik_str: 1631761, ticker: 'YRD', title: 'Yiren Digital Ltd.' },
};
// The search SQL used by the symbols.search tRPC procedure (kept in sync).
function search(db: DatabaseSync, q: string, limit = 10) {
const upper = q.trim().toUpperCase();
const rows = db.prepare(
`SELECT symbol, name, sector, industry, exchange, ticker_kind, cik,
CASE WHEN symbol = ? THEN 0 ELSE 1 END AS rank
FROM symbols
WHERE symbol = ? OR symbol LIKE ?
ORDER BY rank ASC, LENGTH(symbol) ASC
LIMIT ?`,
).all(upper, upper, `${upper}%`, Math.ceil(limit * 2 / 3)) as Array<Record<string, any>>;
let results = rows.map((r) => ({
symbol: r.symbol, name: r.name ?? null, ticker_kind: r.ticker_kind, cik: r.cik ?? null,
}));
if (results.length < limit) {
const seen = new Set(results.map((r) => r.symbol));
const nameRows = db.prepare(
`SELECT symbol, name, sector, industry, exchange, ticker_kind, cik
FROM symbols WHERE name IS NOT NULL AND name != '' AND name LIKE ?
ORDER BY LENGTH(name) ASC LIMIT ?`,
).all(`%${q}%`, limit) as Array<Record<string, any>>;
const fresh = nameRows
.filter((r) => !seen.has(r.symbol as string))
.map((r) => ({ symbol: r.symbol, name: r.name ?? null, ticker_kind: r.ticker_kind, cik: r.cik ?? null }));
results = [...results, ...fresh.slice(0, limit - results.length)];
}
return results;
}
test('search: exact symbol match ranks first', () => {
const db = makeDb();
mergeCompanyTickers(db, secRows);
const r = search(db, 'IREN', 5);
assert.equal(r[0].symbol, 'IREN');
assert.equal(r[0].cik, '0001878848');
assert.equal(r[0].ticker_kind, 'equity');
});
test('search: prefix match returns partial symbols without exact duplication', () => {
const db = makeDb();
mergeCompanyTickers(db, secRows);
const r = search(db, 'ire', 5);
// IREN must appear exactly once.
const iren = r.filter((x) => x.symbol === 'IREN');
assert.equal(iren.length, 1);
assert.equal(r[0].symbol, 'IREN'); // exact-prefix ranks first
});
test('search: name-contains fallback fills beyond symbol prefix matches', () => {
const db = makeDb();
mergeCompanyTickers(db, secRows);
// 'yiren' matches no symbol prefix but matches Yiren Digital Ltd by name.
const r = search(db, 'yiren', 5);
assert.ok(r.some((x) => x.symbol === 'YRD'));
});
test('search: case-insensitive and trims', () => {
const db = makeDb();
mergeCompanyTickers(db, secRows);
const r = search(db, ' aapl ', 5);
assert.equal(r[0]?.symbol, 'AAPL');
});
test('search: empty corpus returns empty results', () => {
const db = makeDb();
const r = search(db, 'zzz', 5);
assert.equal(r.length, 0);
});
@@ -0,0 +1,108 @@
import { test, beforeEach } from 'node:test';
import { strict as assert } from 'node:assert';
import {
familyDrainBudget,
isVendorCoolingDown,
noteVendorRateLimit,
resetVendorGateForTests,
sourceToFamily,
sourcesForFamily,
vendorFetch,
withVendorGate,
VendorRateLimitError,
} from '../vendorGate.ts';
beforeEach(() => {
resetVendorGateForTests();
});
test('sourceToFamily maps queue kinds', () => {
assert.equal(sourceToFamily('yfinance'), 'yfinance');
assert.equal(sourceToFamily('yfinance-quote'), 'yfinance');
assert.equal(sourceToFamily('sec-fetch'), 'sec');
assert.equal(sourceToFamily('sec-sc-fetch'), 'sec');
assert.equal(sourceToFamily('finra-si'), 'finra');
assert.equal(sourceToFamily('fred'), 'fred');
assert.equal(sourceToFamily('nasdaq'), 'nasdaq');
assert.equal(sourceToFamily('x'), 'x');
assert.equal(sourceToFamily('unregistered-vendor-xyz'), null);
});
test('sourcesForFamily includes all yahoo tiers', () => {
const yf = sourcesForFamily('yfinance');
assert.ok(yf.includes('yfinance'));
assert.ok(yf.includes('yfinance-quote'));
});
test('familyDrainBudget keeps Yahoo higher than SEC', () => {
assert.ok(familyDrainBudget('yfinance') >= 2);
assert.equal(familyDrainBudget('sec'), 2);
});
test('withVendorGate serializes concurrent calls', async () => {
const order: number[] = [];
await Promise.all([
withVendorGate('yfinance', async () => {
order.push(1);
await new Promise((r) => setTimeout(r, 50));
order.push(2);
return 'a';
}),
withVendorGate('yfinance', async () => {
order.push(3);
return 'b';
}),
]);
// Second call cannot interleave inside first
assert.deepEqual(order, [1, 2, 3]);
});
test('rate-limit throw notes cool-down', async () => {
await assert.rejects(
() =>
withVendorGate('yfinance', async () => {
throw new Error('Edge: Too Many Requests');
}),
(e: unknown) => e instanceof VendorRateLimitError,
);
assert.equal(isVendorCoolingDown('yfinance'), true);
});
test('preflight blocks while cooling', async () => {
noteVendorRateLimit('fred');
await assert.rejects(
() => withVendorGate('fred', async () => 'ok'),
/preflight cool-down|rate limit/i,
);
});
test('vendorFetch refuses host outside allowlist', async () => {
await assert.rejects(
() =>
vendorFetch('nasdaq', 'https://evil.example.com/x', {
hostAllowlist: /api\.nasdaq\.com/i,
}),
/not allowed/,
);
});
test('vendorFetch paces allowed hosts', async () => {
const times: number[] = [];
const original = globalThis.fetch;
globalThis.fetch = (async () => {
times.push(Date.now());
return new Response('{}', { status: 200 });
}) as typeof fetch;
try {
await vendorFetch('nasdaq', 'https://api.nasdaq.com/a', {
hostAllowlist: /api\.nasdaq\.com/i,
});
await vendorFetch('nasdaq', 'https://api.nasdaq.com/b', {
hostAllowlist: /api\.nasdaq\.com/i,
});
assert.equal(times.length, 2);
assert.ok(times[1] - times[0] >= 700, `expected pacing, gap=${times[1] - times[0]}`);
} finally {
globalThis.fetch = original;
}
});
@@ -0,0 +1,120 @@
/**
* Structural guard: bare `fetch(` must not appear in adapter/service vendor code.
* Future integrations that skip vendorGate fail this test in CI.
*
* Allowlist: vendorGate itself, secHttp (wrapper), tests, LLM (user-configured endpoint).
*/
import { test } from 'node:test';
import { strict as assert } from 'node:assert';
import { readdirSync, readFileSync, statSync } from 'node:fs';
import { join, relative } from 'node:path';
import { fileURLToPath } from 'node:url';
const ROOT = join(fileURLToPath(new URL('../..', import.meta.url))); // app/server/src
/** Paths relative to src/ that may call global fetch (the gate implementations). */
const ALLOW_FETCH_PATHS = [
'services/vendorGate.ts',
'services/secHttp.ts',
'llm/openaiCompatible.ts', // user-configured LLM endpoint; gated via family 'llm' when queued
];
/** Directories scanned for violations. */
const SCAN_DIRS = ['adapters', 'services', 'macro', 'mirror'];
function walk(dir: string, out: string[] = []): string[] {
for (const name of readdirSync(dir)) {
if (name === '__tests__' || name === 'node_modules') continue;
const p = join(dir, name);
const st = statSync(p);
if (st.isDirectory()) walk(p, out);
else if (name.endsWith('.ts') && !name.endsWith('.test.ts')) out.push(p);
}
return out;
}
const BARE_FETCH_RE = /(?<![\w.])fetch\s*\(/g;
test('no bare fetch() outside vendorGate / allowlisted modules', () => {
const violations: string[] = [];
for (const sub of SCAN_DIRS) {
const dir = join(ROOT, sub);
let files: string[];
try {
files = walk(dir);
} catch {
continue;
}
for (const file of files) {
const rel = relative(ROOT, file).replace(/\\/g, '/');
if (ALLOW_FETCH_PATHS.includes(rel)) continue;
const text = readFileSync(file, 'utf8');
// Strip line comments for a coarse check
const stripped = text.replace(/\/\/.*$/gm, '').replace(/\/\*[\s\S]*?\*\//g, '');
if (BARE_FETCH_RE.test(stripped)) {
violations.push(rel);
}
BARE_FETCH_RE.lastIndex = 0;
}
}
assert.deepEqual(
violations,
[],
`Bare fetch() found (use vendorFetch / withVendorGate / secHttp):\n ${violations.join('\n ')}\n` +
`See docs/VENDOR_INTEGRATIONS.md`,
);
});
test('new vendor families can be registered at runtime', async () => {
const {
registerVendorIntegration,
withVendorGate,
sourceToFamily,
resetVendorRegistryForTests,
assertSourcesBound,
} = await import('../vendorGate.ts');
registerVendorIntegration({
family: 'polygon',
sourceKinds: ['polygon', 'polygon-quotes'],
policy: { minIntervalMs: 200, maxInflight: 1, drainJobBudget: 2 },
});
assert.equal(sourceToFamily('polygon'), 'polygon');
assert.equal(sourceToFamily('polygon-quotes'), 'polygon');
assertSourcesBound(['polygon', 'yfinance']);
let ran = false;
await withVendorGate('polygon', async () => {
ran = true;
return 1;
});
assert.equal(ran, true);
resetVendorRegistryForTests();
assert.equal(sourceToFamily('polygon'), null);
});
test('AdapterQueue rejects adapters without a vendor family', async () => {
const { createDb, initSchema } = await import('../../db/client.ts');
const { AdapterQueue } = await import('../../queue/AdapterQueue.ts');
const { FakeSourceAdapter } = await import('../../adapters/SourceAdapter.ts');
const { resetVendorRegistryForTests } = await import('../vendorGate.ts');
resetVendorRegistryForTests();
const db = createDb({ path: ':memory:' });
initSchema(db);
// FakeSourceAdapter claims sourceKind yfinance which is bound — use a fake unbound kind
// by casting (simulates a future SourceKind not yet registered).
const unbound = new FakeSourceAdapter('yfinance' as never);
Object.defineProperty(unbound, 'sourceKind', { value: 'brand-new-vendor' });
assert.throws(
() =>
new AdapterQueue({
db,
adapters: new Map([['brand-new-vendor' as never, unbound as never]]),
}),
/no vendor family|brand-new-vendor/i,
);
resetVendorRegistryForTests();
});
@@ -57,7 +57,10 @@ async function fetchFromYahoo(symbol: string): Promise<{ ratings: AnalystRating[
let raw: Record<string, unknown>;
let timeoutId: ReturnType<typeof setTimeout> | null = null;
try {
const yfPromise = yf.quoteSummary(symbol, { modules: ['upgradeDowngradeHistory', 'recommendationTrend'] }, { validateResult: false });
const { withVendorGate } = await import('./vendorGate.ts');
const yfPromise = withVendorGate('yfinance', () =>
yf.quoteSummary(symbol, { modules: ['upgradeDowngradeHistory', 'recommendationTrend'] }, { validateResult: false }),
);
const timeoutPromise = new Promise<never>((_, reject) => {
timeoutId = setTimeout(() => reject(new Error('Yahoo Finance timed out')), YAHOO_TIMEOUT_MS);
});
@@ -141,10 +144,24 @@ export function getAnalystRatings(db: DatabaseSync, symbol: string): { ratings:
};
}
export async function fetchAndStoreAnalystRatings(db: DatabaseSync, symbol: string): Promise<{ ratings: AnalystRating[]; consensus: AnalystConsensus | null; stale: boolean } | { error: string }> {
export interface FetchAnalystOpts {
/** Called when Yahoo returns a rate-limit signal so the shared queue can cool down. */
onRateLimit?: () => void;
}
export async function fetchAndStoreAnalystRatings(
db: DatabaseSync,
symbol: string,
opts: FetchAnalystOpts = {},
): Promise<{ ratings: AnalystRating[]; consensus: AnalystConsensus | null; stale: boolean } | { error: string }> {
const upper = symbol.toUpperCase();
const result = await fetchFromYahoo(upper);
if ('error' in result) return { error: result.error };
if ('error' in result) {
if (/too many requests|rate[- ]?limit|429|edge:\s*too many/i.test(result.error)) {
try { opts.onRateLimit?.(); } catch { /* ignore */ }
}
return { error: result.error };
}
const now = new Date().toISOString();
const upsert = db.prepare(`
+152
View File
@@ -0,0 +1,152 @@
// Investor Flow — X position-capture materializer (M21 follow-up).
//
// Reads a tracked fund manager's posts from x_cookie_posts (already fetched by
// the `x` timeline job), classifies them with captureParser, and upserts total
// position snapshots into fund_position_records with source='capture'.
//
// Recency rule (mirrorEngine): captures and 13F compete on MAX(as_of) per symbol.
// CLAIMS are explicitly skipped in v1: a claim ("added 865,000 shares") is a
// DELTA, not a total — letting it win recency would corrupt the Live Book.
// Claim folding (delta onto latest total) is future work.
//
// Local-only: no network, no rate limiting — this is a materializer, run from
// the x schedule branch after timeline jobs are enqueued.
import type { DatabaseSync } from 'node:sqlite';
import { extractCaptures } from '../mirror/captureParser.ts';
export interface CaptureIngestStats {
fundId: string;
postsScanned: number;
captures: number; // total-position posts materialized (inserted or refreshed)
inserted: number;
refreshed: number;
claims: number; // delta posts detected — recorded for observability, not materialized
skipped: number; // commentary / unparseable / missing symbol
disabled: boolean;
}
export interface FundWithHandle {
id: string;
x_handle: string | null;
enabled: number;
}
export function listTrackedFundHandles(db: DatabaseSync): FundWithHandle[] {
return db.prepare('SELECT id, x_handle, enabled FROM tracked_funds').all() as FundWithHandle[];
}
interface PostRow {
post_id: string;
body_text: string | null;
posted_at: string;
}
/**
* Materialize captures for one fund. Idempotent: converges on
* (fund_id, symbol, source='capture', evidence_url) — a re-run with the same
* post updates the row instead of duplicating it.
*/
export function ingestFundCaptures(db: DatabaseSync, fundId: string): CaptureIngestStats {
const fund = db.prepare('SELECT id, x_handle, enabled FROM tracked_funds WHERE id = ?').get(fundId) as
| FundWithHandle
| undefined;
const stats: CaptureIngestStats = {
fundId,
postsScanned: 0,
captures: 0,
inserted: 0,
refreshed: 0,
claims: 0,
skipped: 0,
disabled: false,
};
if (!fund || !fund.x_handle) return stats;
if (fund.enabled !== 1) { stats.disabled = true; return stats; }
const posts = db.prepare(
`SELECT post_id, body_text, posted_at FROM x_cookie_posts
WHERE lower(author_handle) = lower(?)
AND body_text IS NOT NULL AND body_text != ''
ORDER BY posted_at ASC`,
).all(fund.x_handle) as PostRow[];
stats.postsScanned = posts.length;
const selectExisting = db.prepare(
`SELECT id FROM fund_position_records WHERE fund_id=? AND symbol=? AND source='capture' AND evidence_url=?`,
);
const updateExisting = db.prepare(
`UPDATE fund_position_records SET shares=?, value_usd=?, cost_basis=?, as_of=?, notes=COALESCE(notes, ?)
WHERE id=?`,
);
const insertNew = db.prepare(
`INSERT INTO fund_position_records (id, fund_id, symbol, shares, value_usd, cost_basis, as_of, source, evidence_url, notes, created_at)
VALUES (?,?,?,?,?,?,?, 'capture', ?, ?, ?)`,
);
for (const p of posts) {
// One post can mention several names (Mike: SLNH fill + OPEN total in one tweet).
const parsedList = extractCaptures(p.body_text ?? '');
if (parsedList.length === 0) {
stats.skipped++;
continue;
}
let anyMaterialized = false;
for (const parsed of parsedList) {
if (parsed.class === 'claim') {
stats.claims++;
continue;
}
if (parsed.class !== 'capture') continue;
if (!parsed.symbol) continue;
const asOf = normalizePostedAtDate(p.posted_at);
const evidenceUrl = `https://x.com/${fund.x_handle}/status/${p.post_id}`;
// Preserve book_reset / pre_reset markers; only stamp instrument notes when empty.
const instrumentNote = parsed.instrument && parsed.instrument !== 'equity'
? parsed.instrument
: null;
const existing = selectExisting.get(fundId, parsed.symbol, evidenceUrl) as { id?: string } | undefined;
if (existing?.id) {
updateExisting.run(
parsed.shares ?? null, parsed.value_usd ?? null, parsed.cost_basis ?? null, asOf,
instrumentNote, existing.id,
);
stats.refreshed++;
} else {
insertNew.run(
crypto.randomUUID(), fundId, parsed.symbol,
parsed.shares ?? null, parsed.value_usd ?? null, parsed.cost_basis ?? null, asOf, evidenceUrl,
instrumentNote, new Date().toISOString(),
);
stats.inserted++;
}
stats.captures++;
anyMaterialized = true;
}
if (!anyMaterialized && parsedList.every((x) => x.class === 'claim')) {
// already counted as claims
} else if (!anyMaterialized) {
stats.skipped++;
}
}
return stats;
}
/** posted_at may be ISO or Twitter "Wed Jul 15 20:43:28 +0000 2026". */
function normalizePostedAtDate(postedAt: string): string {
if (/^\d{4}-\d{2}-\d{2}/.test(postedAt)) return postedAt.slice(0, 10);
const t = Date.parse(postedAt);
if (Number.isFinite(t)) return new Date(t).toISOString().slice(0, 10);
return postedAt.slice(0, 10);
}
/** Materialize captures for every enabled tracked fund with an x_handle. */
export function ingestAllFundCaptures(db: DatabaseSync): CaptureIngestStats[] {
return listTrackedFundHandles(db)
.filter((f) => f.x_handle && f.enabled === 1)
.map((f) => ingestFundCaptures(db, f.id));
}
+133
View File
@@ -0,0 +1,133 @@
// Shared CUSIP ↔ symbol registry for institutional pipeline reliability.
//
// EFTS full-text (efts.sec.gov) is flaky / often 403. Alerts need CUSIP for 13F
// holder refresh and SC XML often carries issuerCusip only after a successful
// download. This module provides offline curated mappings + kv_cache helpers
// so resolve never depends solely on EFTS.
import type { DatabaseSync } from 'node:sqlite';
/** Curated CUSIP → symbol (observed in 13F info tables / SC XML). */
export const CUSIP_TO_SYMBOL: Record<string, string> = {
'05759B305': 'BKKT', // Bakkt, Inc
'09175A206': 'BTM', // BitMine Immersion Technologies, Inc.
'13646K108': 'CP', // Canadian Pacific Kansas City
'17253J106': 'CIFR', // Cipher Digital Inc.
'21036P108': 'STZ', // Constellation Brands, Inc.
Q4982L109: 'IREN', // IREN Limited
'46438F101': 'IBIT', // iShares Bitcoin Trust ETF
'46438R105': 'ETHA', // iShares Ethereum Trust ETF
'670100205': 'NVO', // Novo-Nordisk A/S
'683712103': 'OPEN', // Opendoor Technologies Inc.
'713448108': 'PEP', // PepsiCo, Inc.
'75886F107': 'REGN', // Regeneron Pharmaceuticals, Inc.
'862945300': 'ASST', // Strive, Inc.
'87612E106': 'TGT', // Target Corporation
// High-demand mega-caps (standard 9-char CUSIPs) so pipeline works when EFTS is down
'037833100': 'AAPL',
'023135106': 'AMZN',
'02079K305': 'GOOGL',
'594918104': 'MSFT',
'67066G104': 'NVDA',
'88160R101': 'TSLA',
'30303M102': 'META',
'11135F101': 'AVGO',
'007903107': 'AMD',
'46090E103': 'INTC',
'166764100': 'CVX',
'30231G102': 'XOM',
'46625H100': 'JPM',
'060505104': 'BAC',
'191216100': 'KO',
'742718109': 'PG',
'931142103': 'WMT',
'084670702': 'BRK-B',
'922908769': 'VOO',
'464287200': 'IWM',
'464287655': 'IWF',
'464287499': 'IWD',
'464288281': 'EEM',
'464287465': 'EFA',
'78462F103': 'SPY',
'92189F106': 'GDX',
'922042858': 'VTI',
};
/** Inverted: symbol → CUSIP (first mapping wins). */
export const SYMBOL_TO_CUSIP: Record<string, string> = Object.fromEntries(
Object.entries(CUSIP_TO_SYMBOL).map(([cusip, sym]) => [sym.toUpperCase(), cusip.toUpperCase()]),
);
export function cusipCacheKey(symbol: string): string {
return `sec:cusip:${symbol.toUpperCase()}`;
}
export function isValidCusip(c: string | null | undefined): c is string {
return !!c && /^[0-9A-Z]{8,9}$/i.test(c);
}
/** Offline curated CUSIP for a ticker (no network). */
export function curatedCusipForSymbol(symbol: string): string | null {
const c = SYMBOL_TO_CUSIP[symbol.toUpperCase()];
return isValidCusip(c) ? c.toUpperCase() : null;
}
export function curatedSymbolForCusip(cusip: string): string | null {
const s = CUSIP_TO_SYMBOL[cusip.toUpperCase()];
return s ? s.toUpperCase() : null;
}
export function readCachedCusip(db: DatabaseSync | undefined, symbol: string): string | null {
if (!db) return null;
try {
const row = db.prepare('SELECT value FROM kv_cache WHERE key=?').get(cusipCacheKey(symbol)) as
| { value: string }
| undefined;
const v = row?.value?.trim();
if (isValidCusip(v)) return v.toUpperCase();
} catch {
/* ignore */
}
return null;
}
export function writeCachedCusip(db: DatabaseSync | undefined, symbol: string, cusip: string): void {
if (!db || !isValidCusip(cusip)) return;
try {
db.prepare('INSERT OR REPLACE INTO kv_cache (key, value, observed_at) VALUES (?,?,?)').run(
cusipCacheKey(symbol),
cusip.toUpperCase(),
new Date().toISOString(),
);
} catch {
/* ignore */
}
}
/**
* Resolve CUSIP without network: kv_cache → curated map.
* Also seeds cache when curated hits so subsequent runs are O(1).
*/
export function resolveCusipLocal(db: DatabaseSync | undefined, symbol: string): string | null {
const upper = symbol.toUpperCase();
const cached = readCachedCusip(db, upper);
if (cached) return cached;
const curated = curatedCusipForSymbol(upper);
if (curated) {
writeCachedCusip(db, upper, curated);
return curated;
}
return null;
}
/** Seed kv_cache for every curated symbol (idempotent). */
export function seedCuratedCusips(db: DatabaseSync): number {
let n = 0;
for (const [sym, cusip] of Object.entries(SYMBOL_TO_CUSIP)) {
if (!readCachedCusip(db, sym)) {
writeCachedCusip(db, sym, cusip);
n += 1;
}
}
return n;
}
@@ -129,3 +129,84 @@ export async function sendAlertEmail(
return false;
}
}
/**
* Write an alert into the notification outbox (pending). ID is the alert id so
* re-running a producer never duplicates an entry. Never touches SMTP.
*/
export function enqueueAlertEmail(db: DatabaseSync, alert: Alert): void {
try {
db.prepare(
`INSERT OR IGNORE INTO notification_outbox
(id, user_id, type, severity, title, description, symbol, created_at, status, attempt)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, 'pending', 0)`,
).run(
alert.id,
alert.userId,
alert.type,
alert.severity,
alert.title,
alert.description,
alert.symbol ?? null,
alert.createdAt,
);
} catch {
/* notification_outbox may not exist until migration runs */
}
}
/**
* Drain the outbox, sending one email per pending alert. Returns how many were
* actually sent. Failures stay pending (bounded retry) so SMTP latency never
* blocks the producer loop; this runs on its own timer.
*/
export async function drainEmailOutbox(db: DatabaseSync, limit = 25): Promise<number> {
let sent = 0;
try {
const pending = db.prepare(
`SELECT id FROM notification_outbox
WHERE status = 'pending' AND attempt < 3
ORDER BY created_at ASC
LIMIT ?`,
).all(limit) as Array<{ id: string }>;
if (pending.length === 0) return 0;
const fetchAlert = db.prepare(
`SELECT user_id, type, severity, title, description, symbol, created_at FROM notification_outbox WHERE id = ?`,
);
const markSent = db.prepare(`UPDATE notification_outbox SET status='sent', sent_at=? WHERE id=?`);
const markFailed = db.prepare(
`UPDATE notification_outbox SET attempt = attempt + 1, last_error = ?, status = CASE WHEN attempt + 1 >= 3 THEN 'failed' ELSE 'pending' END WHERE id = ?`,
);
for (const { id } of pending) {
const row = fetchAlert.get(id) as
| { user_id: string; type: string; severity: string; title: string; description: string; symbol: string | null; created_at: string }
| undefined;
if (!row) { markSent.run(new Date().toISOString(), id); continue; }
const alert: Alert = {
id,
userId: row.user_id,
type: row.type as Alert['type'],
severity: row.severity as Alert['severity'],
title: row.title,
description: row.description,
symbol: row.symbol ?? undefined,
createdAt: row.created_at,
acknowledged: false,
dedupKey: '',
payload: {},
};
try {
const ok = await sendAlertEmail(db, alert);
if (ok) { markSent.run(new Date().toISOString(), id); sent += 1; }
else { markFailed.run('no SMTP / rate-limited', id); }
} catch (err) {
markFailed.run(err instanceof Error ? err.message : String(err), id);
}
}
} catch {
/* outbox not present yet */
}
return sent;
}
@@ -0,0 +1,501 @@
// Reverse 13F holder refresh — fund-centric path that does NOT use EFTS.
//
// When efts.sec.gov is 403/blocked, symbol→holders discovery via CUSIP full-text
// search fails. This module instead:
// 1. Collects filer CIKs that already held the symbol (local history)
// 2. Adds tracked funds + a curated set of major 13F managers
// 3. Pulls each filer's latest 13F-HR via data.sec.gov (EdgarAdapter)
// 4. Extracts the target CUSIP and upserts institution_filings
//
// Completeness is lower than full EFTS pagination, but it keeps the alert
// pipeline live without thrashing a dead full-text index.
//
// Even when EFTS is up, CUSIP full-text ranking often misses mega-filers
// (e.g. BlackRock's 50k-line 13F). Call mode "missing_majors" after every
// EFTS pass so Vanguard/BlackRock/Fidelity/etc. still land.
import type { DatabaseSync } from 'node:sqlite';
import { EdgarAdapter, padCik } from '../adapters/EdgarAdapter.ts';
import { writeCachedCusip } from './cusipRegistry.ts';
const edgar = new EdgarAdapter();
/**
* Large / systemically important 13F managers.
* CIKs verified against data.sec.gov submissions (active 13F-HR in 2026).
* Stale entity CIKs (e.g. old BlackRock Finance 1364742) silently miss books.
*/
export const MAJOR_13F_FILER_CIKS: ReadonlyArray<{ cik: string; name: string }> = [
{ cik: '0000102909', name: 'Vanguard Group Inc' },
{ cik: '0002012383', name: 'BlackRock, Inc.' }, // BLK; not old BlackRock Finance 1364742
{ cik: '0000093751', name: 'State Street Corp' },
{ cik: '0000315066', name: 'FMR LLC (Fidelity)' },
{ cik: '0001214717', name: 'Geode Capital Management LLC' },
{ cik: '0000019617', name: 'JPMorgan Chase & Co' },
{ cik: '0000886982', name: 'Goldman Sachs Group Inc' },
{ cik: '0000895421', name: 'Morgan Stanley' },
{ cik: '0000070858', name: 'Bank of America Corp / DE' },
{ cik: '0001423053', name: 'Citadel Advisors LLC' },
{ cik: '0001179392', name: 'Two Sigma Investments LP' },
{ cik: '0001037389', name: 'Renaissance Technologies LLC' },
{ cik: '0001009207', name: 'D. E. Shaw & Co., L.P.' },
{ cik: '0001273087', name: 'Millennium Management LLC' },
{ cik: '0001067983', name: 'Berkshire Hathaway Inc' },
{ cik: '0001422848', name: 'Capital Research Global Investors' },
{ cik: '0000080255', name: 'T. Rowe Price Associates Inc' },
{ cik: '0000914208', name: 'Invesco Ltd.' },
{ cik: '0000073124', name: 'Northern Trust Corp' },
{ cik: '0001374170', name: 'Norges Bank' },
{ cik: '0001446194', name: 'Susquehanna International Group LLP' },
{ cik: '0001595888', name: 'Jane Street Group LLC' },
{ cik: '0001610520', name: 'UBS Group AG' },
{ cik: '0000884546', name: 'Charles Schwab Investment Management Inc' },
{ cik: '0000902219', name: 'Wellington Management Group LLP' },
{ cik: '0001390777', name: 'Bank of New York Mellon Corp' },
{ cik: '0001167557', name: 'AQR Capital Management LLC' },
{ cik: '0001603466', name: 'Point72 Asset Management L.P.' },
];
export interface Reverse13fResult {
filersTried: number;
filersWithMatch: number;
rowsWritten: number;
errors: string[];
path: 'reverse_13f';
mode?: Reverse13fMode;
}
export type Reverse13fMode = 'balanced' | 'missing_majors' | 'majors_first';
function quarterFromPeriod(periodEnding: string, filedAt: string): string {
if (periodEnding && /^\d{4}-\d{2}/.test(periodEnding)) {
const month = parseInt(periodEnding.slice(5, 7), 10);
return `${periodEnding.slice(0, 4)}-Q${Math.ceil(month / 3)}`;
}
if (filedAt) {
const d = new Date(filedAt);
if (!Number.isNaN(d.getTime())) {
return `${d.getFullYear()}-Q${Math.ceil((d.getMonth() + 1) / 3)}`;
}
}
return '';
}
function normalizeCusip(c: string | null | undefined): string {
return (c ?? '').replace(/[^0-9A-Za-z]/g, '').toUpperCase();
}
export type ReverseFiler = { cik: string; name: string | null; source: string };
/** Build ordered unique filer CIK list for reverse refresh of `symbol`. */
export function selectReverseFilerCiks(
db: DatabaseSync,
symbol: string,
maxPrior = 35,
maxMajors = 20,
): ReverseFiler[] {
const upper = symbol.toUpperCase();
const out: ReverseFiler[] = [];
const seen = new Set<string>();
const push = (cik: string, name: string | null, source: string) => {
const p = padCik(cik);
if (!p || p === '0000000000' || seen.has(p)) return;
// Skip non-numeric CIK placeholders (e.g. mr-t-invests)
if (!/^\d{10}$/.test(p)) return;
seen.add(p);
out.push({ cik: p, name, source });
};
// 1. Prior holders of this symbol (largest recent books first)
try {
const prior = db
.prepare(
`SELECT filer_cik AS cik, filer_name AS name, MAX(shares) AS max_shares
FROM institution_filings
WHERE symbol = ? AND form = '13F-HR' AND filer_cik IS NOT NULL
GROUP BY filer_cik
ORDER BY max_shares DESC
LIMIT ?`,
)
.all(upper, maxPrior) as Array<{ cik: string; name: string | null; max_shares: number }>;
for (const r of prior) push(r.cik, r.name, 'prior_holder');
} catch {
/* ignore */
}
// 2. Tracked funds (mirror managers)
try {
const funds = db
.prepare(`SELECT ci_key, fund_name FROM tracked_funds WHERE enabled = 1`)
.all() as Array<{ ci_key: string; fund_name: string }>;
for (const f of funds) push(f.ci_key, f.fund_name, 'tracked_fund');
} catch {
/* ignore */
}
// 3. Major managers (coverage for first-time / thin symbols)
let majors = 0;
for (const m of MAJOR_13F_FILER_CIKS) {
if (majors >= maxMajors) break;
const before = seen.size;
push(m.cik, m.name, 'major');
if (seen.size > before) majors += 1;
}
return out;
}
/** True if a has a strictly newer calendar quarter label than b (YYYY-Qn). */
export function quarterIsAfter(a: string | null | undefined, b: string | null | undefined): boolean {
if (!a) return false;
if (!b) return true;
const pa = /^(\d{4})-Q([1-4])$/.exec(a);
const pb = /^(\d{4})-Q([1-4])$/.exec(b);
if (!pa || !pb) return a > b;
const ya = Number(pa[1]);
const yb = Number(pb[1]);
if (ya !== yb) return ya > yb;
return Number(pa[2]) > Number(pb[2]);
}
/**
* Majors that still need a reverse pull for this symbol.
*
* - Never stored for symbol → missing (e.g. BlackRock never in EFTS hits).
* - Present in the symbol's latest reported quarter → covered.
* - Only older quarter on file → re-check only if last fetch is older than
* ~10 days (Q2 filing season: do not re-download Goldman every job once
* Q1 is stored while peers already show Q2).
*/
export function selectMissingMajorFilers(
db: DatabaseSync,
symbol: string,
opts?: { staleAfterMs?: number; now?: number },
): ReverseFiler[] {
const upper = symbol.toUpperCase();
const staleAfterMs = opts?.staleAfterMs ?? 10 * 24 * 3600_000;
const now = opts?.now ?? Date.now();
let latestQ: string | null = null;
try {
const row = db
.prepare(
`SELECT MAX(reported_quarter) AS q FROM institution_filings
WHERE symbol = ? AND form = '13F-HR'`,
)
.get(upper) as { q: string | null } | undefined;
latestQ = row?.q ?? null;
} catch {
latestQ = null;
}
/** cik → { maxQ, lastFetchedMs } */
const byCik = new Map<string, { maxQ: string | null; lastFetchedMs: number }>();
try {
const rows = db
.prepare(
`SELECT filer_cik AS cik,
MAX(reported_quarter) AS max_q,
MAX(fetched_at) AS last_fetched
FROM institution_filings
WHERE symbol = ? AND form = '13F-HR' AND filer_cik IS NOT NULL
GROUP BY filer_cik`,
)
.all(upper) as Array<{ cik: string; max_q: string | null; last_fetched: string | null }>;
for (const r of rows) {
const p = padCik(r.cik);
if (!p) continue;
const ts = r.last_fetched ? Date.parse(r.last_fetched) : 0;
byCik.set(p, {
maxQ: r.max_q,
lastFetchedMs: Number.isFinite(ts) ? ts : 0,
});
}
} catch {
/* ignore */
}
const missing: ReverseFiler[] = [];
for (const m of MAJOR_13F_FILER_CIKS) {
const p = padCik(m.cik);
if (!p) continue;
const have = byCik.get(p);
if (!have) {
missing.push({ cik: p, name: m.name, source: 'missing_major' });
continue;
}
// Already on the frontier quarter for this symbol
if (latestQ && have.maxQ === latestQ) continue;
// Have some history; only re-poll when stale (new quarter filings may land)
if (have.lastFetchedMs > 0 && now - have.lastFetchedMs < staleAfterMs) continue;
// Never successfully timestamped — treat as missing
if (have.lastFetchedMs <= 0) {
missing.push({ cik: p, name: m.name, source: 'missing_major' });
continue;
}
// Stale relative to frontier: re-check for a newer 13F
if (latestQ && quarterIsAfter(latestQ, have.maxQ)) {
missing.push({ cik: p, name: m.name, source: 'missing_major' });
}
}
return missing;
}
/** Majors first, then prior holders (used when EFTS is down cold). */
export function selectMajorsFirstFilers(
db: DatabaseSync,
symbol: string,
maxMajors = 28,
maxPrior = 10,
): ReverseFiler[] {
const upper = symbol.toUpperCase();
const out: ReverseFiler[] = [];
const seen = new Set<string>();
const push = (cik: string, name: string | null, source: string) => {
const p = padCik(cik);
if (!p || p === '0000000000' || seen.has(p) || !/^\d{10}$/.test(p)) return;
seen.add(p);
out.push({ cik: p, name, source });
};
let majors = 0;
for (const m of MAJOR_13F_FILER_CIKS) {
if (majors >= maxMajors) break;
const before = seen.size;
push(m.cik, m.name, 'major');
if (seen.size > before) majors += 1;
}
try {
const prior = db
.prepare(
`SELECT filer_cik AS cik, filer_name AS name, MAX(shares) AS max_shares
FROM institution_filings
WHERE symbol = ? AND form = '13F-HR' AND filer_cik IS NOT NULL
GROUP BY filer_cik
ORDER BY max_shares DESC
LIMIT ?`,
)
.all(upper, maxPrior) as Array<{ cik: string; name: string | null }>;
for (const r of prior) push(r.cik, r.name, 'prior_holder');
} catch {
/* ignore */
}
return out;
}
/**
* Refresh 13F positions for `symbol` by re-pulling known/major filers' 13F-HR
* history and matching on CUSIP. Safe under ADR-0009 (EdgarAdapter bucket).
*
* Important: 13F is a quarter-end snapshot. Pulling only the latest filing
* leaves majors with a single row and no "vs prior Q" (BlackRock IREN looked
* like a first-time report when Q1–Q4 history existed on EDGAR).
*/
/** Rank 13F filings: newest report period first; prefer non-amendments. */
function rank13fFilings(filings: Array<Record<string, unknown>>): Array<Record<string, unknown>> {
return filings
.map((f) => {
const reportDate = String(f.reportDate ?? f.periodEnding ?? f.period_of_report ?? '');
const fileDate = String(f.filingDate ?? f.fileDate ?? '');
const form = String(f.form ?? '').toUpperCase();
const isAmend = form.includes('/A') ? 0 : 1;
const reportTs = Date.parse(reportDate) || 0;
const fileTs = Date.parse(fileDate) || 0;
return { f, reportTs, fileTs, isAmend, reportDate };
})
.sort((a, b) => b.reportTs - a.reportTs || b.isAmend - a.isAmend || b.fileTs - a.fileTs)
.map((x) => x.f);
}
/** One original (prefer non-/A) filing per report period, newest first, capped. */
function pickRecent13fs(
filings: Array<Record<string, unknown>>,
maxQuarters: number,
): Array<Record<string, unknown>> {
const ranked = rank13fFilings(filings);
const out: Array<Record<string, unknown>> = [];
const seenPeriod = new Set<string>();
for (const f of ranked) {
const reportDate = String(f.reportDate ?? f.periodEnding ?? f.period_of_report ?? '');
const periodKey = reportDate.slice(0, 10) || String(f.accessionNumber ?? f.accession ?? '');
if (!periodKey || seenPeriod.has(periodKey)) continue;
seenPeriod.add(periodKey);
out.push(f);
if (out.length >= maxQuarters) break;
}
return out;
}
export async function refreshHoldersViaReverse13f(
db: DatabaseSync,
symbol: string,
cusip: string,
opts?: { maxFilers?: number; mode?: Reverse13fMode; maxQuartersPerFiler?: number },
): Promise<Reverse13fResult> {
const upper = symbol.toUpperCase();
const targetCusip = normalizeCusip(cusip);
const mode: Reverse13fMode = opts?.mode ?? 'balanced';
if (!targetCusip) {
return {
filersTried: 0,
filersWithMatch: 0,
rowsWritten: 0,
errors: ['empty cusip'],
path: 'reverse_13f',
mode,
};
}
writeCachedCusip(db, upper, targetCusip);
// Keep default modest: each filer is filings_index + N × form13f XML (heavy).
// maxQuartersPerFiler > 1 is required for "vs prior Q" on majors.
const maxFilers = opts?.maxFilers ?? 15;
const maxQuartersPerFiler = Math.max(1, Math.min(opts?.maxQuartersPerFiler ?? 4, 8));
let filers: ReverseFiler[];
if (mode === 'missing_majors') {
filers = selectMissingMajorFilers(db, upper).slice(0, maxFilers);
} else if (mode === 'majors_first') {
filers = selectMajorsFirstFilers(db, upper).slice(0, maxFilers);
} else {
filers = selectReverseFilerCiks(db, upper).slice(0, maxFilers);
}
const now = new Date().toISOString();
let filersWithMatch = 0;
let rowsWritten = 0;
let tried = 0;
const errors: string[] = [];
const insert = db.prepare(`
INSERT INTO institution_filings
(filer_cik, filer_name, filer_sic, symbol, form, shares, value_usd, reported_quarter, filed_at, accession, fetched_at, put_call)
VALUES (?, ?, NULL, ?, '13F-HR', ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(filer_cik, symbol, reported_quarter, form) DO UPDATE SET
shares = excluded.shares,
value_usd = excluded.value_usd,
filer_name = COALESCE(excluded.filer_name, institution_filings.filer_name),
accession = excluded.accession,
put_call = excluded.put_call,
filed_at = excluded.filed_at,
fetched_at = excluded.fetched_at
`);
for (const filer of filers) {
tried += 1;
try {
const indexRes = await edgar.filings_index(filer.cik, {
formTypes: ['13F-HR', '13F-HR/A'],
});
const filings = (indexRes.value ?? []) as Array<Record<string, unknown>>;
const recent = pickRecent13fs(filings, maxQuartersPerFiler);
if (!recent.length) continue;
let matchedAnyQuarter = false;
for (const filing of recent) {
const accession = String(
filing.accessionNumber ?? filing.accession ?? filing.adsh ?? '',
);
if (!accession) continue;
const reportDate = String(
filing.reportDate ?? filing.periodEnding ?? filing.period_of_report ?? '',
);
const fileDate = String(filing.filingDate ?? filing.fileDate ?? now);
const quarter = quarterFromPeriod(reportDate, fileDate);
if (!quarter) continue;
let holdings: Array<{
cusip: string;
issuerName?: string;
value?: number;
sshPrnamt?: number;
putCall?: string;
}>;
try {
const holdRes = await edgar.form13f_holdings(filer.cik, accession);
holdings = ((holdRes.value as { holdings?: Array<Record<string, unknown>> })?.holdings ??
[]) as typeof holdings;
} catch (e) {
const msg = e instanceof Error ? e.message : String(e);
if (/429|403|too many|rate.?limit|access denied|threshold/i.test(msg)) {
errors.push(msg);
return {
filersTried: tried,
filersWithMatch,
rowsWritten,
errors,
path: 'reverse_13f',
mode,
};
}
if (errors.length < 8) errors.push(`${filer.cik}/${accession}: ${msg.slice(0, 100)}`);
continue;
}
// Aggregate matches (managers often split one issuer across many lines)
let shares = 0;
let value = 0;
let putCall = '';
let matched = false;
for (const h of holdings) {
if (normalizeCusip(h.cusip) !== targetCusip) continue;
matched = true;
shares += Number(h.sshPrnamt) || 0;
value += Number(h.value) || 0;
if (h.putCall) putCall = String(h.putCall);
}
if (!matched) continue;
matchedAnyQuarter = true;
if (shares <= 0 && value <= 0) continue;
const fromFiling = String(filing.companyName ?? filing.name ?? '').trim();
const filerName = filer.name || fromFiling || null;
insert.run(
filer.cik,
filerName,
upper,
shares,
value,
quarter,
fileDate,
accession,
now,
putCall || null,
);
rowsWritten += 1;
}
if (matchedAnyQuarter) filersWithMatch += 1;
} catch (e) {
const msg = e instanceof Error ? e.message : String(e);
// Surface rate limits to caller so AdapterQueue can cool the source
if (/429|403|too many|rate.?limit|access denied|threshold/i.test(msg)) {
errors.push(msg);
return {
filersTried: tried,
filersWithMatch,
rowsWritten,
errors,
path: 'reverse_13f',
mode,
};
}
// 503 on a single document is transient — skip filer, keep going
if (errors.length < 8) errors.push(`${filer.cik}: ${msg.slice(0, 120)}`);
}
}
return {
filersTried: tried,
filersWithMatch,
rowsWritten,
errors,
path: 'reverse_13f',
mode,
};
}
File diff suppressed because it is too large Load Diff
+129
View File
@@ -0,0 +1,129 @@
// SEC fair-access HTTP client — thin wrapper over vendorGate('sec').
//
// EVERY outbound call to sec.gov / data.sec.gov / efts.sec.gov must go through
// this module (or vendorGate('sec') directly). User-Agent always identifies
// the operator (SEC fair-access policy).
import {
clearVendorRateLimit,
getVendorCooldownState,
isVendorCoolingDown,
noteVendorRateLimit,
resetVendorGateForTests,
vendorCooldownRemainingMs,
vendorFetch,
VendorRateLimitError,
type VendorFetchOptions,
} from './vendorGate.ts';
export class SecRateLimitError extends VendorRateLimitError {
constructor(status: number, url: string, cooldownMs: number) {
super(
'sec',
`SEC rate limit HTTP ${status} for ${url} - cool down ${Math.round(cooldownMs / 1000)}s, do not thrash`,
cooldownMs,
status,
);
this.name = 'SecRateLimitError';
}
}
const OPERATOR_EMAIL = process.env.SEC_OPERATOR_EMAIL ?? 'research@example.com';
export const SEC_USER_AGENT = `Investor Flow (${OPERATOR_EMAIL})`;
/** Steady-state min gap (overridable via IFLOW_SEC_MIN_INTERVAL_MS). */
export const SEC_MIN_INTERVAL_MS = Number(process.env.IFLOW_SEC_MIN_INTERVAL_MS ?? 350);
const SEC_HOST_RE = /https?:\/\/([^/]*\.)?sec\.gov\b/i;
export function isSecHttpCoolingDown(now = Date.now()): boolean {
return isVendorCoolingDown('sec', now);
}
export function secHttpCooldownRemainingMs(now = Date.now()): number {
return vendorCooldownRemainingMs('sec', now);
}
export function getSecHttpCooldownState(now = Date.now()) {
return getVendorCooldownState('sec', now);
}
export function resetSecHttpStateForTests(): void {
resetVendorGateForTests();
}
export function noteSecRateLimitHit(now = Date.now()): number {
return noteVendorRateLimit('sec', now);
}
export function clearSecRateLimitHits(): void {
clearVendorRateLimit('sec');
}
export type SecFetchOptions = Omit<VendorFetchOptions, 'hostAllowlist'> & {
method?: string;
headers?: Record<string, string>;
accept?: string;
retries?: number;
};
/**
* Rate-limited fetch for SEC hosts. Serialized process-wide via vendorGate('sec').
*/
export async function secFetch(url: string, opts: SecFetchOptions = {}): Promise<Response> {
if (!SEC_HOST_RE.test(url)) {
throw new Error(`secFetch: refusing non-SEC URL ${url}`);
}
try {
return await vendorFetch('sec', url, {
method: opts.method,
headers: {
'User-Agent': SEC_USER_AGENT,
...(opts.headers ?? {}),
},
accept: opts.accept ?? 'application/json',
retries: opts.retries ?? 1,
hostAllowlist: SEC_HOST_RE,
signal: opts.signal,
});
} catch (e) {
if (e instanceof VendorRateLimitError) {
throw new SecRateLimitError(e.status ?? 429, url, e.cooldownMs);
}
throw e;
}
}
export async function secFetchJson<T = unknown>(url: string, opts?: SecFetchOptions): Promise<T> {
const resp = await secFetch(url, { ...opts, accept: opts?.accept ?? 'application/json' });
if (!resp.ok) {
throw new Error(`SEC HTTP ${resp.status} ${resp.statusText} for ${url}`);
}
return (await resp.json()) as T;
}
export async function secFetchText(url: string, opts?: SecFetchOptions): Promise<string> {
const resp = await secFetch(url, {
...opts,
accept: opts?.accept ?? 'application/xml, text/xml, text/html, */*',
});
if (!resp.ok) {
throw new Error(`SEC HTTP ${resp.status} ${resp.statusText} for ${url}`);
}
return resp.text();
}
/** Source kinds that share the SEC rate budget / cool-down family. */
export const SEC_SOURCE_FAMILY = [
'sec',
'sec-fetch',
'sec-sc-fetch',
'sec-tickers',
'sec-lint-holders',
'sec-lint-insiders',
] as const;
export function isSecSourceFamily(source: string): boolean {
return (SEC_SOURCE_FAMILY as readonly string[]).includes(source) || source.startsWith('sec');
}
@@ -0,0 +1,75 @@
// Investor Flow — SEC company_tickers.json materializer (M22 Symbol Search Index).
//
// Pure merge logic: the SEC bulk file is the authoritative source of issuer identity
// (cik, name, exchange) but NOT of classification (ticker_kind) or descriptive metadata
// (sector/industry/peers). Per ADR-0011's Symbol-Metadata Merge Policy:
// - cik: always overwritten (SEC is the only source of issuer CIK)
// - name/exchange: filled only when absent (don't clobber richer yfinance hydration)
// - ticker_kind: never downgraded on conflict (etf/crypto/index preserved — kind gates module applicability)
// - sector/industry/peers: never touched (absent from the SEC file)
// - rows absent from the SEC file are never purged (SEC universe is a subset, not the whole)
import type { DatabaseSync } from 'node:sqlite';
export interface CompanyTickerRow {
cik_str: number;
ticker: string;
title: string;
}
export interface SecTickersIngestResult {
upserted: number; // rows matched + merged (any change incl. cik fill)
created: number; // rows newly inserted
skipped: number; // rows with no ticker (defensive)
}
/** Merge the SEC ticker file into the symbols table. Pure DB logic, no network. */
export function mergeCompanyTickers(
db: DatabaseSync,
tickers: Record<string, CompanyTickerRow> | CompanyTickerRow[],
): SecTickersIngestResult {
const insert = db.prepare(`
INSERT INTO symbols (symbol, name, exchange, ticker_kind, cik, updated_at)
VALUES (?, ?, NULL, 'equity', ?, ?)
ON CONFLICT(symbol) DO UPDATE SET
cik = excluded.cik,
name = CASE WHEN symbols.name IS NULL OR symbols.name = '' THEN excluded.name ELSE symbols.name END,
updated_at = excluded.updated_at
`);
const rows = Array.isArray(tickers)
? tickers
: Object.values(tickers as Record<string, CompanyTickerRow>);
let upserted = 0;
let created = 0;
let skipped = 0;
const now = new Date().toISOString();
db.exec('BEGIN');
try {
for (const r of rows) {
const symbol = (r.ticker ?? '').trim().toUpperCase();
const name = (r.title ?? '').trim();
const cik = String(r.cik_str).padStart(10, '0');
if (!symbol) { skipped++; continue; }
const existing = db.prepare('SELECT symbol, cik FROM symbols WHERE symbol = ?').get(symbol) as
| { symbol: string; cik: string | null }
| undefined;
insert.run(symbol, name, cik, now);
if (existing) {
upserted++;
} else {
created++;
}
}
db.exec('COMMIT');
} catch (e) {
db.exec('ROLLBACK');
throw e;
}
return { upserted, created, skipped };
}
@@ -0,0 +1,175 @@
// Investor Flow — stock float persistence service.
//
// Fetches sharesOutstanding / floatShares from Yahoo Finance (defaultKeyStatistics),
// persists to `stock_float` table, and provides ownership % computation helpers.
//
// Design:
// - Float/outstanding are time-varying market data, not static symbol metadata.
// Separate table so it doesn't pollute the symbols schema and supports history.
// - Co-located with existing short-interest pipeline (same YF module).
// - Rate-limited via vendorGate('yfinance') — never bare HTTP.
import type { DatabaseSync } from 'node:sqlite';
import { withVendorGate } from './vendorGate.ts';
/** TTL for float data before re-fetching (24h default). */
const FLOAT_TTL_MS = 24 * 60 * 60 * 1000;
interface YFKeyStats {
sharesOutstanding?: unknown;
floatShares?: unknown;
}
/** Parse a numeric value from Yahoo's response (handles strings with commas, nulls). */
function parseNum(v: unknown): number | null {
if (v === null || v === undefined) return null;
if (typeof v === 'number') return Number.isFinite(v) ? v : null;
const s = String(v).replace(/[,]/g, '');
const n = parseFloat(s);
return Number.isNaN(n) ? null : n;
}
interface YFInstance {
quoteSummary(symbol: string, opts: { modules: string[] }): Promise<Record<string, unknown>>;
}
/** Lazy-loaded yahoo-finance2 instance. */
let cachedYf: YFInstance | null = null;
async function getYf(): Promise<YFInstance> {
if (!cachedYf) {
const mod = await import('yahoo-finance2');
cachedYf = new mod.default({ suppressNotices: ['yahooSurvey'] }) as unknown as YFInstance;
}
return cachedYf!;
}
/** Fetch float data from Yahoo Finance and persist to stock_float table. */
export async function fetchAndPersistFloat(db: DatabaseSync, symbol: string): Promise<{
sharesOutstanding: number | null;
floatShares: number | null;
fetchedAt: string;
}> {
const yf = await getYf();
const result = await withVendorGate('yfinance', async () => {
return (await yf.quoteSummary(symbol.toUpperCase(), {
modules: ['defaultKeyStatistics'],
})) as unknown as Record<string, unknown>;
});
const stats = (result?.defaultKeyStatistics ?? {}) as YFKeyStats;
const outstanding = parseNum(stats.sharesOutstanding);
const floatShares = parseNum(stats.floatShares);
const fetchedAt = new Date().toISOString();
// Upsert into stock_float — use today's date as the key for latest snapshot.
const today = new Date().toISOString().slice(0, 10);
db.prepare(`
INSERT INTO stock_float (symbol, shares_outstanding, float_shares, as_of)
VALUES (?, ?, ?, ?)
ON CONFLICT(symbol, as_of) DO UPDATE SET
shares_outstanding = excluded.shares_outstanding,
float_shares = excluded.float_shares,
as_of = excluded.as_of
`).run(symbol.toUpperCase(), outstanding, floatShares, fetchedAt);
return { sharesOutstanding: outstanding, floatShares: floatShares, fetchedAt };
}
/** Options for refreshAllStockFloats. */
export interface RefreshFloatOptions {
/** Force refresh all symbols regardless of TTL. Default: false (skip fresh entries). */
forceRefresh?: boolean;
}
/** Refresh float data for all symbols that are stale (>24h) or missing. */
export async function refreshAllStockFloats(db: DatabaseSync, opts: RefreshFloatOptions = {}): Promise<{ refreshed: number; skipped: number; errors: string[] }> {
const today = new Date().toISOString();
const cutoff = opts.forceRefresh ? '1970-01-01T00:00:00.000Z' : new Date(Date.now() - FLOAT_TTL_MS).toISOString();
// Get all unique symbols from fund_position_records (tracked funds) and institution_filings.
const symbols: string[] = [];
const seen = new Set<string>();
const trackedRows = db.prepare(`SELECT DISTINCT symbol FROM fund_position_records WHERE source IN ('13f', 'capture')`).all() as Array<{ symbol: string }>;
for (const r of trackedRows) {
if (!seen.has(r.symbol)) { seen.add(r.symbol); symbols.push(r.symbol); }
}
const filingRows = db.prepare(`SELECT DISTINCT symbol FROM institution_filings WHERE shares IS NOT NULL`).all() as Array<{ symbol: string }>;
for (const r of filingRows) {
if (!seen.has(r.symbol)) { seen.add(r.symbol); symbols.push(r.symbol); }
}
let refreshed = 0;
let skipped = 0;
const errors: string[] = [];
for (const symbol of symbols) {
try {
// Check if we have fresh data.
const existing = db.prepare(
`SELECT as_of FROM stock_float WHERE symbol = ? ORDER BY as_of DESC LIMIT 1`
).get(symbol) as { as_of: string } | undefined;
if (existing && existing.as_of >= cutoff) {
skipped++;
continue;
}
await fetchAndPersistFloat(db, symbol);
refreshed++;
} catch (err) {
errors.push(`${symbol}: ${(err as Error).message}`);
}
}
return { refreshed, skipped, errors };
}
/** Get the latest float snapshot for a symbol. */
export function getLatestFloat(db: DatabaseSync, symbol: string): { sharesOutstanding: number | null; floatShares: number | null; asOf: string } | null {
const row = db.prepare(
`SELECT shares_outstanding, float_shares, as_of FROM stock_float WHERE symbol = ? ORDER BY as_of DESC LIMIT 1`
).get(symbol) as { shares_outstanding: number | null; float_shares: number | null; as_of: string } | undefined;
if (!row) return null;
return {
sharesOutstanding: row.shares_outstanding,
floatShares: row.float_shares,
asOf: row.as_of,
};
}
/** Compute institutional ownership % for a symbol across all tracked funds. */
export function computeOwnershipPercentages(db: DatabaseSync, symbol: string): {
heldShares: number;
outstanding: number | null;
float: number | null;
pctOutstanding: number | null;
pctFloat: number | null;
} {
// Sum all tracked-fund shares for this symbol (latest record per fund).
const totalHeld = db.prepare(`
SELECT COALESCE(SUM(sub.shares), 0) AS held
FROM (
SELECT fpr.shares, MAX(fpr.as_of) AS max_as_of
FROM fund_position_records fpr
JOIN tracked_funds tf ON tf.id = fpr.fund_id AND tf.enabled = 1
WHERE fpr.symbol = ? AND fpr.source IN ('13f', 'capture')
GROUP BY fpr.fund_id
) sub
`).get(symbol) as { held: number } | undefined;
const heldShares = totalHeld?.held ?? 0;
const latestFloat = getLatestFloat(db, symbol);
return {
heldShares,
outstanding: latestFloat?.sharesOutstanding ?? null,
float: latestFloat?.floatShares ?? null,
pctOutstanding: latestFloat?.sharesOutstanding ? heldShares / latestFloat.sharesOutstanding : null,
pctFloat: latestFloat?.floatShares ? heldShares / latestFloat.floatShares : null,
};
}
+487
View File
@@ -0,0 +1,487 @@
// Process-wide vendor rate gate (ADR-0009).
//
// CONTRACT FOR EVERY EXTERNAL INTEGRATION (existing + future):
//
// 1. Register a family before any traffic:
// registerVendorFamily('polygon', { minIntervalMs: 200, maxInflight: 1, drainJobBudget: 2 })
// bindSourceKind('polygon', 'polygon')
// 2. All outbound work goes through:
// withVendorGate('polygon', () => client.get(...))
// vendorFetch('polygon', url, { hostAllowlist: /polygon\.io/i })
// or extend VendorSourceAdapter (auto-gates fetchOne).
// 3. Never use bare fetch()/SDK calls from adapters or services.
// 4. AdapterQueue refuses source_kinds with no family binding at construction.
//
// Guard test: src/services/__tests__/vendorHttpGuard.test.ts fails CI if bare
// fetch sneaks into adapters/services (except this module + allowlisted paths).
import { isRateLimitError, rateLimitCooldownMs } from '../queue/sourceRatePolicy.ts';
/** Open string type so future vendors register without editing a union. */
export type VendorFamily = string;
export class VendorRateLimitError extends Error {
readonly family: VendorFamily;
readonly status: number | null;
readonly cooldownMs: number;
constructor(family: VendorFamily, message: string, cooldownMs: number, status: number | null = null) {
super(message);
this.name = 'VendorRateLimitError';
this.family = family;
this.cooldownMs = cooldownMs;
this.status = status;
}
}
export class UnknownVendorFamilyError extends Error {
constructor(message: string) {
super(message);
this.name = 'UnknownVendorFamilyError';
}
}
interface FamilyState {
cooldownUntil: number;
consecutiveHits: number;
lastEndedAt: number;
/** Waiters blocked waiting for an inflight slot (maxInflight gate). */
waiters: Array<() => void>;
inflight: number;
}
export interface FamilyPolicy {
/** Min ms between completed calls in this family. */
minIntervalMs: number;
/** Max concurrent calls (1 = strict single-flight). */
maxInflight: number;
/** Max jobs of this family per AdapterQueue.drain() cycle. */
drainJobBudget: number;
/** Optional host allowlist regex source for docs / vendorFetch defaults. */
hostPattern?: string;
}
const DEFAULT_POLICY: FamilyPolicy = {
minIntervalMs: 500,
maxInflight: 1,
drainJobBudget: 1,
};
/** Registered family policies (built-ins + future registerVendorFamily calls). */
const POLICIES = new Map<VendorFamily, FamilyPolicy>();
/** Queue source_kind → vendor family. */
const SOURCE_TO_FAMILY = new Map<string, VendorFamily>();
/**
* Register (or replace) a vendor family policy.
* Call this before bindSourceKind / starting traffic for a new integration.
*/
export function registerVendorFamily(family: VendorFamily, policy: Partial<FamilyPolicy> = {}): void {
if (!family || !/^[a-z][a-z0-9_-]*$/i.test(family)) {
throw new UnknownVendorFamilyError(
`registerVendorFamily: invalid family name '${family}' (use [a-z0-9_-]+)`,
);
}
const prev = POLICIES.get(family);
POLICIES.set(family, {
minIntervalMs: policy.minIntervalMs ?? prev?.minIntervalMs ?? DEFAULT_POLICY.minIntervalMs,
maxInflight: policy.maxInflight ?? prev?.maxInflight ?? DEFAULT_POLICY.maxInflight,
drainJobBudget: policy.drainJobBudget ?? prev?.drainJobBudget ?? DEFAULT_POLICY.drainJobBudget,
hostPattern: policy.hostPattern ?? prev?.hostPattern,
});
}
/**
* Bind a queue source_kind to a registered family.
* Required for AdapterQueue to accept adapters of that kind.
*/
export function bindSourceKind(sourceKind: string, family: VendorFamily): void {
if (!POLICIES.has(family)) {
throw new UnknownVendorFamilyError(
`bindSourceKind('${sourceKind}', '${family}'): family not registered. Call registerVendorFamily('${family}', …) first.`,
);
}
SOURCE_TO_FAMILY.set(sourceKind, family);
}
/** Register family + bind one or more source kinds in one call (preferred for new vendors). */
export function registerVendorIntegration(opts: {
family: VendorFamily;
sourceKinds: string[];
policy?: Partial<FamilyPolicy>;
}): void {
registerVendorFamily(opts.family, opts.policy ?? {});
for (const sk of opts.sourceKinds) {
bindSourceKind(sk, opts.family);
}
}
/** True if family was registered (built-in or runtime). */
export function isVendorFamilyRegistered(family: string): boolean {
return POLICIES.has(family);
}
/** True if this source_kind has a family binding. */
export function isSourceBound(sourceKind: string): boolean {
return sourceToFamily(sourceKind) != null;
}
/**
* Assert every source_kind is bound to a family.
* AdapterQueue calls this on construction so unregistered vendors fail fast.
*/
export function assertSourcesBound(sourceKinds: Iterable<string>): void {
const missing: string[] = [];
for (const sk of sourceKinds) {
if (!sourceToFamily(sk)) missing.push(sk);
}
if (missing.length) {
throw new UnknownVendorFamilyError(
`Adapter source_kind(s) have no vendor family: ${missing.join(', ')}. ` +
`Register with registerVendorIntegration({ family, sourceKinds, policy }) ` +
`before constructing AdapterQueue. See docs/VENDOR_INTEGRATIONS.md.`,
);
}
}
function seedBuiltIns(): void {
if (POLICIES.size > 0) return;
registerVendorIntegration({
family: 'yfinance',
sourceKinds: ['yfinance', 'yfinance-quote', 'yfinance-eod', 'yfinance-meta', 'yfinance-holdings'],
policy: { minIntervalMs: 400, maxInflight: 1, drainJobBudget: 3, hostPattern: 'yahoo|finance\\.yahoo' },
});
registerVendorIntegration({
family: 'sec',
sourceKinds: ['sec', 'sec-fetch', 'sec-sc-fetch', 'sec-tickers', 'sec-lint-holders', 'sec-lint-insiders'],
// SEC fair access allows ~10 req/s; concurrent fetches are paced by the
// family min-interval + per-source spacing, so 2 inflight is safe and
// clears long-lived backlogs ~4x faster than strict single-flight.
policy: { minIntervalMs: 350, maxInflight: 2, drainJobBudget: 2, hostPattern: 'sec\\.gov' },
});
registerVendorIntegration({
family: 'fred',
sourceKinds: ['fred', 'macro'],
policy: { minIntervalMs: 500, maxInflight: 1, drainJobBudget: 1, hostPattern: 'stlouisfed\\.org' },
});
registerVendorIntegration({
family: 'finra',
sourceKinds: ['finra-bulk', 'finra-si'],
policy: { minIntervalMs: 1000, maxInflight: 1, drainJobBudget: 1, hostPattern: 'finra\\.org' },
});
registerVendorIntegration({
family: 'nasdaq',
sourceKinds: ['nasdaq'],
policy: { minIntervalMs: 800, maxInflight: 1, drainJobBudget: 1, hostPattern: 'nasdaq\\.com' },
});
registerVendorIntegration({
family: 'reddit',
sourceKinds: ['reddit'],
policy: { minIntervalMs: 1500, maxInflight: 1, drainJobBudget: 1, hostPattern: 'reddit\\.com' },
});
registerVendorIntegration({
family: 'x',
sourceKinds: ['x'],
policy: { minIntervalMs: 3000, maxInflight: 1, drainJobBudget: 1 },
});
registerVendorIntegration({
family: 'llm',
sourceKinds: ['llm'],
policy: { minIntervalMs: 0, maxInflight: 2, drainJobBudget: 2 },
});
}
// Seed on module load so existing adapters work without ceremony.
seedBuiltIns();
/** All queue source kinds belonging to a family (for family-wide cool-down). */
export function sourcesForFamily(family: VendorFamily): string[] {
const out: string[] = [];
for (const [sk, f] of SOURCE_TO_FAMILY) {
if (f === family) out.push(sk);
}
return out;
}
export function sourceToFamily(source: string): VendorFamily | null {
const direct = SOURCE_TO_FAMILY.get(source);
if (direct) return direct;
// Prefix fallback only for registered families (future yfinance-*, sec-*)
if (source.startsWith('yfinance') && POLICIES.has('yfinance')) return 'yfinance';
if (source.startsWith('sec') && POLICIES.has('sec')) return 'sec';
if (source.startsWith('finra') && POLICIES.has('finra')) return 'finra';
return null;
}
/**
* Resolve family or throw. Prefer this for new code paths so silent
* ungated traffic cannot appear for unregistered vendors.
*/
export function requireSourceFamily(source: string): VendorFamily {
const f = sourceToFamily(source);
if (!f) {
throw new UnknownVendorFamilyError(
`No vendor family for source_kind '${source}'. ` +
`Call registerVendorIntegration({ family, sourceKinds: ['${source}'], policy }) first.`,
);
}
return f;
}
export function familyPolicy(family: VendorFamily): FamilyPolicy {
const p = POLICIES.get(family);
if (!p) {
throw new UnknownVendorFamilyError(
`Unknown vendor family '${family}'. Call registerVendorFamily('${family}', …) first.`,
);
}
return p;
}
export function familyDrainBudget(family: VendorFamily): number {
return familyPolicy(family).drainJobBudget;
}
export function listRegisteredFamilies(): VendorFamily[] {
return [...POLICIES.keys()].sort();
}
const states = new Map<VendorFamily, FamilyState>();
function state(family: VendorFamily): FamilyState {
let s = states.get(family);
if (!s) {
s = {
cooldownUntil: 0,
consecutiveHits: 0,
lastEndedAt: 0,
waiters: [],
inflight: 0,
};
states.set(family, s);
}
return s;
}
export function isVendorCoolingDown(family: VendorFamily, now = Date.now()): boolean {
return now < state(family).cooldownUntil;
}
export function vendorCooldownRemainingMs(family: VendorFamily, now = Date.now()): number {
return Math.max(0, state(family).cooldownUntil - now);
}
export function getVendorCooldownState(
family: VendorFamily,
now = Date.now(),
): {
family: VendorFamily;
active: boolean;
remainingMs: number;
consecutiveHits: number;
until: string | null;
} {
const s = state(family);
const remainingMs = Math.max(0, s.cooldownUntil - now);
return {
family,
active: remainingMs > 0,
remainingMs,
consecutiveHits: s.consecutiveHits,
until: remainingMs > 0 ? new Date(s.cooldownUntil).toISOString() : null,
};
}
/** Record a rate-limit hit; returns cool-down duration applied (ms). */
export function noteVendorRateLimit(family: VendorFamily, now = Date.now()): number {
// Ensure family exists so ad-hoc note doesn't create zombie state
if (!POLICIES.has(family)) {
registerVendorFamily(family, DEFAULT_POLICY);
}
const s = state(family);
s.consecutiveHits += 1;
const ms = rateLimitCooldownMs(s.consecutiveHits);
s.cooldownUntil = Math.max(s.cooldownUntil, now + ms);
return ms;
}
export function clearVendorRateLimit(family: VendorFamily): void {
const s = state(family);
s.consecutiveHits = 0;
s.cooldownUntil = 0;
}
/**
* Test helper: wipe runtime state (cool-downs, chains).
* Does NOT unregister built-in families (re-seed if maps were cleared).
*/
export function resetVendorGateForTests(): void {
states.clear();
// If a test called registerVendorFamily for ad-hoc names, leave policies;
// always re-seed built-ins so bindings exist.
if (!POLICIES.has('yfinance')) {
POLICIES.clear();
SOURCE_TO_FAMILY.clear();
seedBuiltIns();
}
}
/** Full reset including custom registrations (unit tests only). */
export function resetVendorRegistryForTests(): void {
states.clear();
POLICIES.clear();
SOURCE_TO_FAMILY.clear();
seedBuiltIns();
}
function envMinInterval(family: VendorFamily): number {
const key = `IFLOW_${family.toUpperCase().replace(/[^A-Z0-9]/g, '_')}_MIN_INTERVAL_MS`;
const raw = process.env[key];
if (raw && Number.isFinite(Number(raw))) return Math.max(0, Number(raw));
return familyPolicy(family).minIntervalMs;
}
async function waitTurn(family: VendorFamily): Promise<void> {
const s = state(family);
const now = Date.now();
if (now < s.cooldownUntil) {
const remaining = s.cooldownUntil - now;
throw new VendorRateLimitError(
family,
`${family} rate limit preflight cool-down ${Math.round(remaining / 1000)}s - do not thrash`,
remaining,
);
}
const minGap = envMinInterval(family);
const sinceLast = now - s.lastEndedAt;
if (s.lastEndedAt > 0 && sinceLast < minGap) {
await new Promise((r) => setTimeout(r, minGap - sinceLast));
}
}
/**
* Run an async vendor call under the family gate (pace + inflight-cap +
* cool-down). Up to `maxInflight` calls run concurrently; starts are spaced by
* `minIntervalMs`. On thrown errors that look like rate limits, notes cool-down
* then rethrows.
*/
export async function withVendorGate<T>(family: VendorFamily, fn: () => Promise<T>): Promise<T> {
if (!POLICIES.has(family)) {
throw new UnknownVendorFamilyError(
`withVendorGate('${family}'): family not registered. Call registerVendorFamily first.`,
);
}
const s = state(family);
const policy = familyPolicy(family);
await acquireSlot(s, policy.maxInflight);
try {
await waitTurn(family);
const result = await fn();
s.consecutiveHits = 0;
return result;
} catch (e) {
if (e instanceof VendorRateLimitError) throw e;
const msg = e instanceof Error ? e.message : String(e);
if (isRateLimitError(msg) || /edge:\s*too many|429|throttl/i.test(msg)) {
const ms = noteVendorRateLimit(family);
throw new VendorRateLimitError(
family,
`${family} rate limit: ${msg} - cool down ${Math.round(ms / 1000)}s, do not thrash`,
ms,
);
}
throw e;
} finally {
s.inflight = Math.max(0, s.inflight - 1);
s.lastEndedAt = Date.now();
const next = s.waiters.shift();
if (next) next();
}
}
/**
* Block until fewer than `maxInflight` calls are in flight for the family.
* The released waiter takes its slot synchronously inside `next()` (called by
* the releaser's finally), so inflight never dips below the cap between
* release and resume.
*/
function acquireSlot(s: FamilyState, maxInflight: number): Promise<void> {
if (s.inflight < maxInflight) {
s.inflight += 1;
return Promise.resolve();
}
return new Promise<void>((resolve) => {
s.waiters.push(() => {
s.inflight += 1;
resolve();
});
});
}
export interface VendorFetchOptions {
method?: string;
headers?: Record<string, string>;
accept?: string;
retries?: number;
signal?: AbortSignal;
/** Required for safety — callers must declare allowed hosts. */
hostAllowlist?: RegExp;
}
function isRateLimitStatus(status: number): boolean {
return status === 429 || status === 503 || status === 403;
}
/**
* Rate-limited fetch for a vendor family.
* Prefer hostAllowlist so callers cannot mis-route.
*/
export async function vendorFetch(
family: VendorFamily,
url: string,
opts: VendorFetchOptions = {},
): Promise<Response> {
if (opts.hostAllowlist && !opts.hostAllowlist.test(url)) {
throw new Error(`vendorFetch(${family}): URL not allowed: ${url}`);
}
const retries = opts.retries ?? 0;
return withVendorGate(family, async () => {
for (let attempt = 0; attempt <= retries; attempt++) {
const headers: Record<string, string> = {
...(opts.accept ? { Accept: opts.accept } : {}),
...(opts.headers ?? {}),
};
const resp = await fetch(url, {
method: opts.method ?? 'GET',
headers,
signal: opts.signal,
});
if (resp.ok || resp.status === 304) return resp;
if (isRateLimitStatus(resp.status)) {
if (attempt < retries) {
await new Promise((r) => setTimeout(r, 1000 * (attempt + 1)));
continue;
}
const ms = noteVendorRateLimit(family);
throw new VendorRateLimitError(
family,
`${family} rate limit HTTP ${resp.status} for ${url} - cool down ${Math.round(ms / 1000)}s, do not thrash`,
ms,
resp.status,
);
}
return resp;
}
const ms = noteVendorRateLimit(family);
throw new VendorRateLimitError(family, `${family} rate limit for ${url}`, ms);
});
}