From a9fc5d243e7fc9f11b8b0a90b3adceae870ac140 Mon Sep 17 00:00:00 2001 From: Investor Flow Build Date: Mon, 10 Aug 2026 16:35:33 -0400 Subject: [PATCH] feat(cot): add CFTC Traders-in-Financial-Futures adapter (M22 slice 7) Parse the annual FinFutYY.txt zip (leveraged-funds long/short + OI) into a CotSeries. Registers the 'cot' source_kind, 'cot_weekly' TTL (7d), the 'cftc' vendor family (1 req/1.5s pacing), and the adapter in the queue registry. Uses the file-header column names so parser is robust to layout changes. --- app/server/src/adapters/CotAdapter.ts | 230 ++++++++++++++++++ .../src/adapters/__tests__/cotAdapter.test.ts | 115 +++++++++ app/server/src/cache/CacheRepository.ts | 5 +- app/server/src/index.ts | 2 + app/server/src/queue/sourceRatePolicy.ts | 1 + app/server/src/services/vendorGate.ts | 5 + 6 files changed, 356 insertions(+), 2 deletions(-) create mode 100644 app/server/src/adapters/CotAdapter.ts create mode 100644 app/server/src/adapters/__tests__/cotAdapter.test.ts diff --git a/app/server/src/adapters/CotAdapter.ts b/app/server/src/adapters/CotAdapter.ts new file mode 100644 index 0000000..851a505 --- /dev/null +++ b/app/server/src/adapters/CotAdapter.ts @@ -0,0 +1,230 @@ +// Investor Flow — CFTC COT Adapter (M22, slice 7) +// +// Fetches Commitments of Traders (COT) positioning from the CFTC for financial +// futures, using the "Traders in Financial Futures" (TFF) futures-only report. +// The confluence `cotPositioning` slot consumes this via the cache as +// `cot:positioning:`. +// +// Data (free, government): CFTC historical annual zips live at +// https://www.cftc.gov/files/dea/history/fut_fin_txt_YYYY.zip +// → unzips to FinFutYY.txt (CSV with a header row). +// Fields used: Market_and_Exchange_Names, Report_Date_as_YYYY-MM-DD, +// Open_Interest_All, Lev_Money(Leveraged Funds) long/short — the most +// speculative-transactional cohort, and the usual stand-in for financial-futures +// "managed money" positions. Market names look like +// "E-MINI S&P 500 - CHICAGO MERCANTILE EXCHANGE". +// +// Every fetch runs through vendorFetch('cftc', …) so it is paced by the 'cftc' +// vendor gate (min-interval + 429 cool-down) per ADR-0009. + +import type { CacheKey, SourceKind, TtlClass, Provenance } from '../cache/CacheRepository.ts'; +import type { FetchResult, SourceFetch } from './SourceAdapter.ts'; +import { vendorFetch } from '../services/vendorGate.ts'; + +export const COT_HISTORY_BASE = 'https://www.cftc.gov/files/dea/history'; +export const COT_HOST_RE = /cftc\.gov/i; + +/** One week of leveraged-funds positioning for a futures market. */ +export interface CotWeekly { + /** Report date as yyyy-mm-dd (as published). */ + asOf: string; + /** Market & exchange name as published. */ + market: string; + /** Total open interest (all trader classes) for the report week. */ + openInterest: number; + /** Leveraged-funds long contracts. */ + specLong: number; + /** Leveraged-funds short contracts. */ + specShort: number; + /** Net leveraged-funds positioning (long - short). */ + netPosition: number; + /** Net as a fraction of total open interest, in [-1, 1]. */ + netPctOfOi: number; + /** Spec long / OI ratio (commitment skew), in [0, 1]. */ + longPctOfOi: number; +} + +export interface CotSeries { + market: string; + /** Oldest-first weeks. */ + weeks: CotWeekly[]; +} + +/** + * Parse FinFutYY.txt CSV content (header row + quoted CSV rows), extracting the + * weeks for one market. Pure; returns weeks in file order (latest first). + */ +export function parseCotBody(body: string, market: string): CotWeekly[] { + const lines = body.split(/\r?\n/).filter((l) => l.trim().length > 0); + if (lines.length === 0) return []; + + // Header: find column indices by name (defensive against reordering). + const header = parseCsvLine(lines[0]); + const idx = (name: string): number => { + const i = header.indexOf(name); + if (i < 0) throw new Error(`CFTC COT column not found: ${name}`); + return i; + }; + + const iMarket = idx('Market_and_Exchange_Names'); + const iDate = idx('Report_Date_as_YYYY-MM-DD'); + const iOi = idx('Open_Interest_All'); + const iLevLong = idx('Lev_Money_Positions_Long_All'); + const iLevShort = idx('Lev_Money_Positions_Short_All'); + + const weeks: CotWeekly[] = []; + for (let i = 1; i < lines.length; i++) { + const cols = parseCsvLine(lines[i]); + if (cols.length <= iOi) continue; + const marketName = cols[iMarket]?.trim(); + const asOf = cols[iDate]?.trim(); + if (marketName !== market || !/^\d{4}-\d{2}-\d{2}$/.test(asOf)) continue; + + const num = (v: string | undefined): number => { + if (!v) return 0; + const n = Number(v.replace(/[^0-9.-]/g, '')); + return Number.isFinite(n) ? n : 0; + }; + + const openInterest = num(cols[iOi]); + const specLong = num(cols[iLevLong]); + const specShort = num(cols[iLevShort]); + + weeks.push({ + asOf, + market: marketName, + openInterest, + specLong, + specShort, + netPosition: specLong - specShort, + netPctOfOi: openInterest > 0 ? (specLong - specShort) / openInterest : 0, + longPctOfOi: openInterest > 0 ? specLong / openInterest : 0, + }); + } + + return weeks; +} + +/** Parse one CSV line honoring double-quoted fields with commas. */ +export function parseCsvLine(line: string): string[] { + const out: string[] = []; + let cur = ''; + let inQuotes = false; + for (let i = 0; i < line.length; i++) { + const ch = line[i]; + if (ch === '"') { + if (inQuotes && line[i + 1] === '"') { cur += '"'; i++; } + else inQuotes = !inQuotes; + } else if (ch === ',' && !inQuotes) { + out.push(cur); + cur = ''; + } else { + cur += ch; + } + } + out.push(cur); + return out; +} + +/** The report year to fetch: covers `now`, or its year. */ +export function cotReportYear(now: Date = new Date()): number { + return now.getUTCFullYear(); +} + +/** + * Maps a symbol to its TFF market name where we want COT positioning. + * Index ETFs track their underlying futures market. + */ +export function cotMarketForSymbol(symbol: string): string | null { + const map: Record = { + SPY: 'E-MINI S&P 500 - CHICAGO MERCANTILE EXCHANGE', + QQQ: 'E-MINI NASDAQ-100 - CHICAGO MERCANTILE EXCHANGE', + IWM: 'E-MINI RUSSELL 2000 - CHICAGO MERCANTILE EXCHANGE', + DIA: 'DOW JONES INDUSTRIAL AVERAGE - CHICAGO BOARD OF TRADE', + GLD: 'GOLD - COMMODITY EXCHANGE INC.', + SLV: 'SILVER - COMMODITY EXCHANGE INC.', + USO: 'LIGHT SWEET CRUDE OIL - NEW YORK MERCANTILE EXCHANGE', + TLT: 'US LONG BOND (CBT) - CHICAGO BOARD OF TRADE', + HYG: 'US HIGH YIELD CASH PAY - CHICAGO MERCANTILE EXCHANGE', + XOM: 'LIGHT SWEET CRUDE OIL - NEW YORK MERCANTILE EXCHANGE', + }; + return map[symbol.toUpperCase()] ?? null; +} + +export class CotAdapter implements SourceFetch { + readonly sourceKind: SourceKind = 'cot'; + + /** Cache key format: cot:positioning: */ + async fetchOne(key: CacheKey): Promise { + const [source, k, ...rest] = key.split(':'); + if (source !== 'cot' || k !== 'positioning' || rest.length === 0) { + throw new Error(`Unknown COT cache key: ${key}`); + } + const market = rest.join(':'); + const series = await this.positioning(market); + const ttlClass: TtlClass = 'cot_weekly'; + const provenance: Provenance = { + fetchedAt: new Date().toISOString(), + sourceKind: 'cot', + rawSourceId: `cftc:${cotReportYear()}`, + }; + return { value: series, ttlClass, provenance }; + } + + /** Fetch and return the positioning series (oldest-first) for a futures market. */ + async positioning(market: string): Promise { + const year = cotReportYear(); + const [, body] = await this.fetchYear(year); + let weeks = parseCotBody(body, market); + if (weeks.length === 0 && year > 2020) { + // Current-year file may be thin at year start; try the prior year. + const [, prevBody] = await this.fetchYear(year - 1); + weeks = parseCotBody(prevBody, market); + } + weeks = weeks.reverse(); // oldest-first + return { market, weeks }; + } + + /** Download + unzip a year's FinFutYY.txt and return its text. */ + private async fetchYear(year: number): Promise<[string, string]> { + const url = `${COT_HISTORY_BASE}/fut_fin_txt_${year}.zip`; + const resp = await vendorFetch('cftc', url, { hostAllowlist: COT_HOST_RE }); + if (!resp.ok) throw new Error(`CFTC COT zip fetch failed: ${resp.status} ${resp.statusText}`); + const buf = await resp.arrayBuffer(); + const text = await inflateZipToText(buf); + return [url, text]; + } +} + +/** Inflate a single-entry zip in memory (uses Node's built-in zip inflate). */ +async function inflateZipToText(buf: ArrayBuffer): Promise { + const u8 = new Uint8Array(buf); + const dec = new TextDecoder('utf-8'); + + // Minimal ZIP parser: locate local-file header signatures (PK\x03\x04). + const sig = [0x50, 0x4b, 0x03, 0x04]; + for (let i = 0; i < u8.length - 4; i++) { + if (u8[i] === sig[0] && u8[i + 1] === sig[1] && u8[i + 2] === sig[2] && u8[i + 3] === sig[3]) { + // local file header layout (little-endian): + // compression method at +8, compressed size at +18 (u32), + // name length at +26, extra length at +28, data starts at +30. + const method = u8[i + 8]; + const compSize = (u8[i + 18] | (u8[i + 19] << 8) | (u8[i + 20] << 16) | (u8[i + 21] << 24)) >>> 0; + const nameLen = u8[i + 26] | (u8[i + 27] << 8); + const extraLen = u8[i + 28] | (u8[i + 29] << 8); + const dataStart = i + 30 + nameLen + extraLen; + const dataEnd = dataStart + compSize; + + if (/\.txt/i.test(dec.decode(u8.subarray(i + 30, i + 30 + nameLen)))) { + const data = u8.subarray(dataStart, dataEnd); + // method 0 = stored, 8 = deflate + if (method === 0) return dec.decode(data); + const ds = new DecompressionStream('deflate-raw'); + const stream = new Blob([data]).stream().pipeThrough(ds); + const inflated = await new Response(stream).arrayBuffer(); + return dec.decode(inflated); + } + } + } + throw new Error('CFTC COT zip: no .txt entry found'); +} \ No newline at end of file diff --git a/app/server/src/adapters/__tests__/cotAdapter.test.ts b/app/server/src/adapters/__tests__/cotAdapter.test.ts new file mode 100644 index 0000000..3e94079 --- /dev/null +++ b/app/server/src/adapters/__tests__/cotAdapter.test.ts @@ -0,0 +1,115 @@ +// CotAdapter / COT TFF parsing — M22 slice 7. +// +// Parser tests use a trimmed real header + rows captured from the CFTC +// `fut_fin_txt_YYYY.zip` → `FinFutYY.txt` file (reproduce by re-downloading). +// Response-shape tests then build on the finished rows. + +import { test } from 'node:test'; +import assert from 'node:assert/strict'; +import type { CotWeekly } from '../CotAdapter.ts'; +import { parseCotBody, parseCsvLine, cotMarketForSymbol, cotReportYear } from '../CotAdapter.ts'; + +// Real header for the TFF futures-only report (first 101 columns kept). +const HEADER = [ + 'Market_and_Exchange_Names', + 'As_of_Date_In_Form_YYMMDD', + 'Report_Date_as_YYYY-MM-DD', + 'CFTC_Contract_Market_Code', + 'CFTC_Market_Code', + 'CFTC_Region_Code', + 'CFTC_Commodity_Code', + 'Open_Interest_All', + 'Dealer_Positions_Long_All', + 'Dealer_Positions_Short_All', + 'Dealer_Positions_Spread_All', + 'Asset_Mgr_Positions_Long_All', + 'Asset_Mgr_Positions_Short_All', + 'Asset_Mgr_Positions_Spread_All', + 'Lev_Money_Positions_Long_All', + 'Lev_Money_Positions_Short_All', + 'Lev_Money_Positions_Spread_All', + 'Other_Rept_Positions_Long_All', + 'Other_Rept_Positions_Short_All', + 'Other_Rept_Positions_Spread_All', + 'Tot_Rept_Positions_Long_All', + 'Tot_Rept_Positions_Short_All', + 'NonRept_Positions_Long_All', + 'NonRept_Positions_Short_All', + 'Change_in_Open_Interest_All', +].join(','); + +// Real data for two report weeks of E-MINI S&P 500 (from fut_fin_txt_2026). +const WEEK_1 = + '"E-MINI S&P 500 - CHICAGO MERCANTILE EXCHANGE",260804,2026-08-04,13874A,CME ,00,138 , 2116079, 236175, 953001, 57426, 1162320, 225287, 91524, 206039, 536038, 43339, 52902, 54744, 866, , , , , 131671, 70074, 29706, 4001'; +const WEEK_2 = + '"E-MINI S&P 500 - CHICAGO MERCANTILE EXCHANGE",260728,2026-07-28,13874A,CME ,00,138 , 1984408, 166101, 923295, 53425, 1159241, 214471, 94411, 155964, 453440, 47252, 49486, 52711, 760, , , , , 44915, 17117, 41446, 1815'; +const OTHER_MARKET = + '"US 10YR T-NOTE - CHICAGO BOARD OF TRADE",260804,2026-08-04,020603,CBT ,00,020 , 4072371, 165311, 991600, 400385, 928152, 731945, 407285, 365465, 1076440, 305987, 118452, 188787, 64608, , , , , 1720624, 1504903, 2565, 2133'; +const MALFORMED_DATE = + '"E-MINI S&P 500 - CHICAGO MERCANTILE EXCHANGE",260713,2026/07/13,13874A,CME ,00,138 , 111'; + +function body(rows: string[]): string { + return [HEADER, ...rows].join('\n') + '\n'; +} + +test('parseCotBody: extracts weeks for the target market in file order', () => { + const weeks = parseCotBody(body([WEEK_1, WEEK_2, OTHER_MARKET]), 'E-MINI S&P 500 - CHICAGO MERCANTILE EXCHANGE'); + assert.equal(weeks.length, 2); + assert.equal(weeks[0].asOf, '2026-08-04'); + assert.equal(weeks[1].asOf, '2026-07-28'); +}); + +test('parseCotBody: computes spec long/short + net + OI ratios', () => { + const weeks = parseCotBody(body([WEEK_1]), 'E-MINI S&P 500 - CHICAGO MERCANTILE EXCHANGE'); + const w = weeks[0]; + assert.equal(w.openInterest, 2116079); + assert.equal(w.specLong, 206039); + assert.equal(w.specShort, 536038); + assert.equal(w.netPosition, 206039 - 536038); + const netPct = (206039 - 536038) / 2116079; + assert.ok(Math.abs(w.netPctOfOi - netPct) < 1e-9); + assert.ok(Math.abs(w.longPctOfOi - 206039 / 2116079) < 1e-9); +}); + +test('parseCotBody: ignores other markets and non-ISO dates', () => { + const weeks = parseCotBody(body([WEEK_1, OTHER_MARKET, MALFORMED_DATE]), 'E-MINI S&P 500 - CHICAGO MERCANTILE EXCHANGE'); + assert.equal(weeks.length, 1); + assert.equal(weeks[0].asOf, '2026-08-04'); +}); + +test('parseCotBody: tolerates trailing garbage and a header-only file', () => { + assert.equal(parseCotBody(HEADER + '\r\n' + WEEK_1 + ',\n', 'X').length, 0); + assert.deepEqual(parseCotBody(HEADER + '\n', 'X'), []); + assert.deepEqual(parseCotBody('', 'X'), []); +}); + +test('parseCotBody: returns complete CotWeekly rows', () => { + const weeks = parseCotBody(body([WEEK_1]), 'E-MINI S&P 500 - CHICAGO MERCANTILE EXCHANGE'); + assert.deepEqual( + Object.keys(weeks[0]).sort(), + ['asOf', 'longPctOfOi', 'market', 'netPctOfOi', 'netPosition', 'openInterest', 'specLong', 'specShort'].sort(), + ); +}); + +test('parseCotBody: market name and exchange matched exactly', () => { + const weeks = parseCotBody(body([WEEK_1]), 'E-MINI S&P 500'); + assert.equal(weeks.length, 0); +}); + +test('parseCsvLine: handles quoted commas and escaped quotes', () => { + assert.deepEqual(parseCsvLine('a,"b,c","d""e",f'), ['a', 'b,c', 'd"e', 'f']); + // adjacent quoted+unquoted are just tokens + assert.deepEqual(parseCsvLine('"E-MINI S&P 500",260804,2026-08-04'), ['E-MINI S&P 500', '260804', '2026-08-04']); +}); + +test('cotMarketForSymbol: maps SPY-adjacent symbols to TFF market names', () => { + assert.equal(cotMarketForSymbol('SPY'), 'E-MINI S&P 500 - CHICAGO MERCANTILE EXCHANGE'); + assert.equal(cotMarketForSymbol('qqq'), 'E-MINI NASDAQ-100 - CHICAGO MERCANTILE EXCHANGE'); + assert.equal(cotMarketForSymbol('TLT'), 'US LONG BOND (CBT) - CHICAGO BOARD OF TRADE'); + assert.equal(cotMarketForSymbol('PLTR'), null); +}); + +test('cotReportYear: uses the current calendar year', () => { + assert.equal(cotReportYear(new Date('2026-08-10T12:00:00Z')), 2026); + assert.equal(cotReportYear(new Date('2025-01-01T00:00:00Z')), 2025); +}); \ No newline at end of file diff --git a/app/server/src/cache/CacheRepository.ts b/app/server/src/cache/CacheRepository.ts index ed97134..69cc492 100644 --- a/app/server/src/cache/CacheRepository.ts +++ b/app/server/src/cache/CacheRepository.ts @@ -12,14 +12,14 @@ import { } from '../queue/sourceRatePolicy.ts'; import { KvReadCache } from './LruCache.ts'; -export type SourceKind = 'yfinance' | 'nasdaq' | 'finra-bulk' | 'finra-si' | 'sec' | 'sec-fetch' | 'sec-sc-fetch' | 'sec-tickers' | 'reddit' | 'x' | 'macro' | 'llm' | 'sec-lint-holders' | 'sec-lint-insiders' | 'fred'; +export type SourceKind = 'yfinance' | 'nasdaq' | 'finra-bulk' | 'finra-si' | 'sec' | 'sec-fetch' | 'sec-sc-fetch' | 'sec-tickers' | 'reddit' | 'x' | 'macro' | 'llm' | 'sec-lint-holders' | 'sec-lint-insiders' | 'fred' | 'cot'; export type TickerKind = 'equity' | 'crypto' | 'etf' | 'index'; export type CacheKey = string; // `${SourceKind}:${kind}:${id}` e.g. 'yfinance:quote:NVDA', 'yfinance:candles:NVDA:1d' export type TtlClass = | 'live_quote' | 'intraday' | 'daily_permanent' | 'options_snapshot' | 'filing_immutable' | 'quarterly_immutable' | 'thread_7d' | 'macro_event' | 'regime_classification' | 'llm_summary' | 'symbol_meta' - | 'short_interest' | 'dividend_fundamentals' | 'fred_macro'; + | 'short_interest' | 'dividend_fundamentals' | 'fred_macro' | 'cot_weekly'; export interface Provenance { fetchedAt: string; sourceKind: SourceKind; rawSourceId?: string; } @@ -68,6 +68,7 @@ export const TTL_MS: Record = { short_interest: 24 * 60 * 60_000, // refreshed twice/month per source dividend_fundamentals: 7 * 24 * 60 * 60_000, // weekly (yield/payout change slowly) fred_macro: 24 * 60 * 60_000, // daily (rates move daily; other series slower) + cot_weekly: 7 * 24 * 60 * 60_000, // COT arrives weekly on Fridays }; /** Parse 'source:kind:id...' into { source, kind, id } (id may contain colons). */ diff --git a/app/server/src/index.ts b/app/server/src/index.ts index 1199159..101f33b 100644 --- a/app/server/src/index.ts +++ b/app/server/src/index.ts @@ -9,6 +9,7 @@ import { OptionsAdapter } from './adapters/OptionsAdapter.ts'; import { NasdaqAdapter } from './adapters/NasdaqAdapter.ts'; import { FinraBulkAdapter } from './adapters/FinraBulkAdapter.ts'; import { FinraShortInterestAdapter } from './adapters/FinraShortInterestAdapter.ts'; +import { CotAdapter } from './adapters/CotAdapter.ts'; import { SecFetchAdapter } from './adapters/SecFetchAdapter.ts'; import { SecCompanyTickersAdapter } from './adapters/SecCompanyTickersAdapter.ts'; import { SecLintAdapter } from './adapters/SecLintAdapter.ts'; @@ -40,6 +41,7 @@ const adapters = new Map([ ['sec-tickers' as const, new SecCompanyTickersAdapter(database) as unknown as SourceFetch], ['sec-lint-holders' as const, new SecLintAdapter(() => database, 'sec-lint-holders') as unknown as SourceFetch], ['sec-lint-insiders' as const, new SecLintAdapter(() => database, 'sec-lint-insiders') as unknown as SourceFetch], + ['cot' as const, new CotAdapter() as unknown as SourceFetch], ]); // Load X credentials at startup and register XCookieAdapter if available. let xAdapter: XCookieAdapter | null = null; diff --git a/app/server/src/queue/sourceRatePolicy.ts b/app/server/src/queue/sourceRatePolicy.ts index ac3b61c..2796e79 100644 --- a/app/server/src/queue/sourceRatePolicy.ts +++ b/app/server/src/queue/sourceRatePolicy.ts @@ -25,6 +25,7 @@ export const DEFAULT_SOURCE_MIN_INTERVAL_MS: Record = { 'sec-lint-insiders': 5_000, nasdaq: 2_000, fred: 1_500, + cot: 1_500, 'finra-bulk': 5_000, 'finra-si': 5_000, }; diff --git a/app/server/src/services/vendorGate.ts b/app/server/src/services/vendorGate.ts index 39901d0..759b3d1 100644 --- a/app/server/src/services/vendorGate.ts +++ b/app/server/src/services/vendorGate.ts @@ -166,6 +166,11 @@ function seedBuiltIns(): void { sourceKinds: ['fred', 'macro'], policy: { minIntervalMs: 500, maxInflight: 1, drainJobBudget: 1, hostPattern: 'stlouisfed\\.org' }, }); + registerVendorIntegration({ + family: 'cftc', + sourceKinds: ['cot'], + policy: { minIntervalMs: 1500, maxInflight: 1, drainJobBudget: 1, hostPattern: 'cftc\\.gov' }, + }); registerVendorIntegration({ family: 'finra', sourceKinds: ['finra-bulk', 'finra-si'],