feat: Phase 3 FINRA bulk adapter + Phase 4 three-way merge with discrepancy flagging
Phase 3 — FINRA bulk short-interest ingest: - Add finra_short_interest table to schema (per-symbol, per-settlement-date, per-exchange, with short/exempt/total volume, avg daily vol, days to cover) - Create FinraIngestService: downloads FINRA consolidated pipe-delimited file from configurable base URL, parses Market|Symbol|ShortVolume|ShortExemptVolume|TotalVolume, stores in finra_short_interest table - Create FinraBulkAdapter: SourceFetch that calls downloadAndIngestFinra, registers as finra-bulk source kind with finraShortinterest cache handler - finraShortinterest handler reads latest settlement row per symbol from finra_short_interest table (no per-symbol kv_cache write; data is bulk-ingested) - Register in index.ts adapter map + HANDLERS + del case Phase 4 — three-way merge with discrepancy detection: - shortInterest tRPC procedure now reads all 3 caches (yfinance, nasdaq, finra-bulk) in parallel - Reconciliation hierarchy: FINRA (shares short) > NASDAQ > Yahoo - daysToCover: NASDAQ (specific) > FINRA (computed) > Yahoo (short ratio fallback) - settlementDate: FINRA > NASDAQ > Yahoo - Discrepancy detection: compares sharesShort across available sources, flags >10% difference with discrepancyPct + discrepancyBetween - Updated ShortInterestPanel: FINRA source badge, discrepancy warning banner, three-source disclaimer - Updated trpc.ts client type for new shape
This commit is contained in:
@@ -0,0 +1,138 @@
|
||||
import type { DatabaseSync } from 'node:sqlite';
|
||||
|
||||
const FINRA_BASE_URL = process.env.FINRA_BASE_URL ?? 'https://www.finra.org/sites/default/files';
|
||||
|
||||
/** Format a FINRA consolidated-short-interest filename: CAshvol{YYYYMMDD}.txt */
|
||||
function finraFilename(settlementDate: string): string {
|
||||
const d = settlementDate.replace(/-/g, '');
|
||||
const ym = settlementDate.slice(0, 7).replace(/-/, '-');
|
||||
return `${ym}/CAshvol${d}.txt`;
|
||||
}
|
||||
|
||||
/** Parse a FINRA consolidated-short-interest file body (pipe-delimited).
|
||||
* Expected columns: Market|Symbol|ShortVolume|ShortExemptVolume|TotalVolume
|
||||
* Returns per-symbol rows aggregated across all exchanges. */
|
||||
function parseFinraFile(
|
||||
body: string,
|
||||
settlementDate: string,
|
||||
ingestedAt: string,
|
||||
sourceFile: string,
|
||||
): Array<{
|
||||
symbol: string;
|
||||
exchange: string;
|
||||
shortVolume: number;
|
||||
shortExempt: number;
|
||||
totalVolume: number;
|
||||
}> {
|
||||
const lines = body.split(/\r?\n/);
|
||||
const rows: Array<{
|
||||
symbol: string;
|
||||
exchange: string;
|
||||
shortVolume: number;
|
||||
shortExempt: number;
|
||||
totalVolume: number;
|
||||
}> = [];
|
||||
let headerFound = false;
|
||||
|
||||
for (const raw of lines) {
|
||||
const line = raw.trim();
|
||||
if (!line || line.startsWith('#')) continue;
|
||||
if (line.startsWith('Date Range') || line.startsWith('Period')) continue;
|
||||
if (line.includes('Market|Symbol|')) { headerFound = true; continue; }
|
||||
if (!headerFound) continue;
|
||||
|
||||
const cols = line.split('|').map((c) => c.trim());
|
||||
if (cols.length < 4) continue;
|
||||
|
||||
const market = cols[0];
|
||||
const symbol = cols[1];
|
||||
const shortVolume = parseFloat(cols[2]?.replace(/,/g, ''));
|
||||
const shortExempt = cols[3] ? parseFloat(cols[3].replace(/,/g, '')) : 0;
|
||||
const totalVolume = cols[4] ? parseFloat(cols[4].replace(/,/g, '')) : shortVolume + shortExempt;
|
||||
|
||||
if (!symbol || Number.isNaN(shortVolume)) continue;
|
||||
|
||||
rows.push({
|
||||
symbol: symbol.toUpperCase(),
|
||||
exchange: market.toUpperCase(),
|
||||
shortVolume,
|
||||
shortExempt,
|
||||
totalVolume,
|
||||
});
|
||||
}
|
||||
return rows;
|
||||
}
|
||||
|
||||
/** Download and ingest a FINRA consolidated-short-interest file.
|
||||
* Returns count of symbols stored. */
|
||||
export async function downloadAndIngestFinra(
|
||||
db: DatabaseSync,
|
||||
settlementDate: string,
|
||||
): Promise<{ symbolsStored: number; sourceFile: string; exchanges: string[] }> {
|
||||
const filename = finraFilename(settlementDate);
|
||||
const url = `${FINRA_BASE_URL}/${filename}`;
|
||||
const ingestedAt = new Date().toISOString();
|
||||
|
||||
console.log(`[finra] downloading ${url}`);
|
||||
const resp = await fetch(url, {
|
||||
headers: { 'User-Agent': 'InvestorFlow/1.0 (research) node' },
|
||||
signal: AbortSignal.timeout(30_000),
|
||||
});
|
||||
if (!resp.ok) throw new Error(`FINRA download failed: ${resp.status} ${resp.statusText}`);
|
||||
const body = await resp.text();
|
||||
if (!body.trim()) throw new Error('FINRA file is empty');
|
||||
|
||||
const rows = parseFinraFile(body, settlementDate, ingestedAt, filename);
|
||||
if (!rows.length) throw new Error('No FINRA short interest rows parsed');
|
||||
|
||||
const exchanges = [...new Set(rows.map((r) => r.exchange))];
|
||||
const exchangeMap: Record<string, string> = {};
|
||||
exchanges.forEach((e) => { exchangeMap[e] = e; });
|
||||
|
||||
const upsert = db.prepare(
|
||||
`INSERT OR REPLACE INTO finra_short_interest
|
||||
(symbol, settlement_date, exchange, short_volume, short_exempt, total_volume, avg_daily_vol, days_to_cover, source_file, ingested_at)
|
||||
VALUES (?, ?, ?, ?, ?, ?, NULL, NULL, ?, ?)`
|
||||
);
|
||||
|
||||
const tx = db.transaction(() => {
|
||||
for (const r of rows) {
|
||||
upsert.run(
|
||||
r.symbol,
|
||||
settlementDate,
|
||||
r.exchange,
|
||||
r.shortVolume,
|
||||
r.shortExempt,
|
||||
r.totalVolume,
|
||||
filename,
|
||||
ingestedAt,
|
||||
);
|
||||
}
|
||||
});
|
||||
tx();
|
||||
|
||||
console.log(`[finra] ingested ${rows.length} symbols from ${filename} (exchanges: ${exchanges.join(', ')})`);
|
||||
return { symbolsStored: rows.length, sourceFile: filename, exchanges };
|
||||
}
|
||||
|
||||
/** Compute days-to-cover for finra_short_interest rows that have avg_daily_vol set.
|
||||
* Called after avg_daily_vol is populated from external volume data. */
|
||||
export function computeDaysToCover(db: DatabaseSync): number {
|
||||
const r = db.exec(
|
||||
`UPDATE finra_short_interest
|
||||
SET days_to_cover = CASE
|
||||
WHEN avg_daily_vol IS NOT NULL AND avg_daily_vol > 0 THEN short_volume / avg_daily_vol
|
||||
ELSE NULL
|
||||
END
|
||||
WHERE days_to_cover IS NULL AND avg_daily_vol IS NOT NULL`
|
||||
);
|
||||
return r.changes;
|
||||
}
|
||||
|
||||
/** Get the latest settlement date available in the finra_short_interest table. */
|
||||
export function latestFinraSettlementDate(db: DatabaseSync): string | null {
|
||||
const r = db.prepare(
|
||||
'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;
|
||||
}
|
||||
Reference in New Issue
Block a user