Files
investor-flow/app/server/src/adapters/EdgarAdapter.ts
T

663 lines
24 KiB
TypeScript
Raw Normal View History

// Investor Flow — EdgarAdapter (DESIGN.md §3a Module 2 + §5 sec policy).
//
// Concrete SourceFetch for SEC EDGAR public API endpoints.
// Uses a shared `edgarFetch` helper with token-bucket rate limiting (8 req/s,
// min 125ms between calls) and ETag / If-Modified-Since conditional requests
// so immutable data (filings index, company facts) is revalidated only when
// the server indicates a change.
import type { CacheKey, TtlClass, Provenance, SourceKind } from '../cache/CacheRepository.ts';
import { parseCacheKey } from '../cache/CacheRepository.ts';
import type { SourceFetch, FetchResult } from './SourceAdapter.ts';
// ---------------------------------------------------------------------------
// Rate-limit + fetch helpers (module-private)
// ---------------------------------------------------------------------------
const OPERATOR_EMAIL = process.env.SEC_OPERATOR_EMAIL ?? 'research@example.com';
const UA = `Investor Flow (${OPERATOR_EMAIL})`;
/**
* Token-bucket rate limiter: 8 tokens, refilled at 8/sec, min 125ms between
* calls. Guarantees we stay under EDGAR's 10 req/s rule with headroom.
*/
class TokenBucket {
private tokens = 8;
private lastDrain = Date.now();
async wait(): Promise<void> {
const now = Date.now();
const elapsed = (now - this.lastDrain) / 1000;
// Refill tokens up to maxRate (8).
this.tokens = Math.min(8, this.tokens + elapsed * 8);
this.lastDrain = now;
if (this.tokens < 1) {
const waitMs = Math.ceil(((1 - this.tokens) / 8) * 1000);
await new Promise((r) => setTimeout(r, waitMs));
this.tokens = 0;
this.lastDrain = Date.now();
} else {
this.tokens -= 1;
}
}
}
const bucket = new TokenBucket();
/**
* Fetch a URL with EDGAR-compliant headers, rate limiting, and ETag caching.
* Returns null on 304 (caller should return cached row).
*/
async function edgarFetch(
url: string,
extraHeaders?: Record<string, string>,
): Promise<{ status: number; headers: { etag?: string | null; lastModified?: string | null }; body: unknown } | null> {
await bucket.wait();
const headers: Record<string, string> = {
'User-Agent': UA,
Accept: 'application/json',
...extraHeaders,
};
const resp = await fetch(url, { method: 'GET', headers });
const etag = resp.headers.get('etag');
const lastModified = resp.headers.get('last-modified');
if (resp.status === 304) {
return { status: 304, headers: { etag, lastModified }, body: null };
}
if (resp.status >= 400) {
throw new Error(`EDGAR ${resp.status} ${resp.statusText} for ${url}`);
}
const body = (await resp.json()) as unknown;
return { status: resp.status, headers: { etag, lastModified }, body };
}
// ---------------------------------------------------------------------------
// In-process cache (module-level, shared across adapter instances).
// Used for ETag / Last-Modified revalidation between calls within a process.
// ---------------------------------------------------------------------------
interface EdgarCacheStore {
cache: Record<string, unknown>;
etags: Record<string, string>;
lastModified: Record<string, string>;
}
const store: EdgarCacheStore = { cache: {}, etags: {}, lastModified: {} };
function readStore(): EdgarCacheStore { return store; }
function writeStore(s: EdgarCacheStore): void { Object.assign(store, s); }
// ---------------------------------------------------------------------------
// Cache key builders (per EDGAR endpoint)
// ---------------------------------------------------------------------------
/** Pad a CIK string to 10 digits (EDGAR requires zero-padded). */
export function padCik(cik: string): string {
const digits = cik.replace(/\D/g, '');
return digits.padStart(10, '0');
}
function filingsIndexKey(cik: string, formTypes?: string[], dateRange?: { from?: string; to?: string }): CacheKey {
const parts = [`sec:filings_index:${padCik(cik)}`];
if (formTypes && formTypes.length > 0) parts.push(`types:${formTypes.sort().join(',')}`);
if (dateRange) {
if (dateRange.from) parts.push(`from:${dateRange.from}`);
if (dateRange.to) parts.push(`to:${dateRange.to}`);
}
return parts.join(':');
}
function companyFactsKey(cik: string): CacheKey {
return `sec:company_facts:${padCik(cik)}`;
}
function filerMetaKey(cik: string): CacheKey {
return `sec:filer_meta:${padCik(cik)}`;
}
function searchIndexKey(query: string): CacheKey {
return `sec:search_index:${encodeURIComponent(query).toLowerCase()}`;
}
// ---------------------------------------------------------------------------
// Public adapter
// ---------------------------------------------------------------------------
export class EdgarAdapter implements SourceFetch {
readonly sourceKind = 'sec' as const;
/**
* Recent filings for a CIK, filtered by form type + date range.
* GET https://data.sec.gov/submissions/CIK{padded10}.json
*/
async filings_index(
cik: string,
opts?: { formTypes?: string[]; dateRange?: { from?: string; to?: string } },
): Promise<FetchResult> {
const key = filingsIndexKey(cik, opts?.formTypes, opts?.dateRange);
const padded = padCik(cik);
const url = `https://data.sec.gov/submissions/CIK${padded}.json`;
const fetchedAt = new Date().toISOString();
const s = readStore();
const cachedEtag = s.etags[key] ?? null;
const cachedLms = s.lastModified[key] ?? null;
const resp = await edgarFetch(url, {
...(cachedEtag ? { 'If-None-Match': cachedEtag } : {}),
...(cachedLms ? { 'If-Modified-Since': cachedLms } : {}),
});
if (resp?.status === 304) {
const cached = s.cache[key];
if (cached) return { value: cached, ttlClass: 'daily_permanent', provenance: { fetchedAt, sourceKind: 'sec', rawSourceId: key } };
throw new Error(`EDGAR 304 but no cached data for ${key}`);
}
const rawData = resp?.body as { name?: string; filings?: { recent?: Array<{ form?: string; dateReporter?: string; accessionNumber?: string; accessionNormalization?: string; reportDate?: string; reportFile?: string; primaryDocument?: string }> } } | undefined;
const recentFilings = rawData?.filings?.recent;
if (!recentFilings) {
throw new Error(`EDGAR filings_index: no recent filings for CIK${padded}`);
}
// Filter by form type.
let filings = recentFilings;
if (opts?.formTypes && opts.formTypes.length > 0) {
const allowed = new Set(opts.formTypes.map((f) => f.toUpperCase()));
filings = filings.filter((f) => allowed.has((f.form ?? '').toUpperCase()));
}
// Filter by date range.
if (opts?.dateRange) {
const from = opts.dateRange.from ? new Date(opts.dateRange.from).getTime() : null;
const to = opts.dateRange.to ? new Date(opts.dateRange.to).getTime() : null;
filings = filings.filter((f) => {
const ts = new Date(f.dateReporter ?? f.accessionNormalization ?? '').getTime();
if (Number.isNaN(ts)) return false;
if (from !== null && ts < from) return false;
if (to !== null && ts > to) return false;
return true;
});
}
// Persist ETag + Last-Modified for future revalidation.
const newStore: EdgarCacheStore = { cache: { ...s.cache }, etags: { ...s.etags }, lastModified: { ...s.lastModified } };
newStore.cache[key] = filings;
if (resp?.headers?.etag) newStore.etags[key] = resp.headers.etag;
if (resp?.headers?.lastModified) newStore.lastModified[key] = resp.headers.lastModified;
writeStore(newStore);
return { value: filings, ttlClass: 'daily_permanent', provenance: { fetchedAt, sourceKind: 'sec', rawSourceId: key } };
}
/**
* Company facts (XBRL-derived financials) for a CIK.
* GET https://data.sec.gov/api/xbrl/companyfacts/CIK{padded10}.json
* Cached forever; revalidates on ETag change.
*/
async company_facts(cik: string): Promise<FetchResult> {
const key = companyFactsKey(cik);
const padded = padCik(cik);
const url = `https://data.sec.gov/api/xbrl/companyfacts/CIK${padded}.json`;
const fetchedAt = new Date().toISOString();
const s = readStore();
const cachedEtag = s.etags[key] ?? null;
const cachedLms = s.lastModified[key] ?? null;
const resp = await edgarFetch(url, {
...(cachedEtag ? { 'If-None-Match': cachedEtag } : {}),
...(cachedLms ? { 'If-Modified-Since': cachedLms } : {}),
});
if (resp?.status === 304) {
const cached = s.cache[key];
if (cached) return { value: cached, ttlClass: 'daily_permanent', provenance: { fetchedAt, sourceKind: 'sec', rawSourceId: key } };
throw new Error(`EDGAR 304 but no cached data for ${key}`);
}
const factsData = resp?.body as { entityName?: string; facts?: Record<string, unknown> } | undefined;
if (!factsData) {
throw new Error(`EDGAR company_facts: no data for CIK${padded}`);
}
const newStore: EdgarCacheStore = { cache: { ...s.cache }, etags: { ...s.etags }, lastModified: { ...s.lastModified } };
newStore.cache[key] = factsData;
if (resp?.headers?.etag) newStore.etags[key] = resp.headers.etag;
if (resp?.headers?.lastModified) newStore.lastModified[key] = resp.headers.lastModified;
writeStore(newStore);
return { value: factsData, ttlClass: 'daily_permanent', provenance: { fetchedAt, sourceKind: 'sec', rawSourceId: key } };
}
/**
* Filer metadata (CIK + SIC + name). Fetched once, cached forever.
* Reuses company_facts under the hood (it includes entityName + SIC).
*/
async filer_cik_meta(cik: string): Promise<FetchResult> {
const key = filerMetaKey(cik);
const factsResult = await this.company_facts(cik);
const facts = factsResult.value as { entityName?: string; sic?: string | number } | undefined;
return {
value: {
cik: padCik(cik),
name: facts?.entityName ?? null,
sic: typeof facts?.sic === 'string' ? facts.sic : facts?.sic != null ? String(facts.sic) : null,
},
ttlClass: 'daily_permanent',
provenance: { fetchedAt: new Date().toISOString(), sourceKind: 'sec', rawSourceId: key },
};
}
/**
* Full-text search across EDGAR filings (EFTS search-index).
* GET https://efts.sec.gov/LATEST/search-index?q=...
*/
async full_text_search(query: string): Promise<FetchResult> {
const key = searchIndexKey(query);
const encoded = encodeURIComponent(query);
const url = `https://efts.sec.gov/LATEST/search-index?q=${encoded}`;
const fetchedAt = new Date().toISOString();
const s = readStore();
const cachedEtag = s.etags[key] ?? null;
const cachedLms = s.lastModified[key] ?? null;
const resp = await edgarFetch(url, {
...(cachedEtag ? { 'If-None-Match': cachedEtag } : {}),
...(cachedLms ? { 'If-Modified-Since': cachedLms } : {}),
});
if (resp?.status === 304) {
const cached = s.cache[key];
if (cached) return { value: cached, ttlClass: 'daily_permanent', provenance: { fetchedAt, sourceKind: 'sec', rawSourceId: key } };
throw new Error(`EDGAR 304 but no cached data for ${key}`);
}
const searchData = resp?.body as { filings?: Array<Record<string, unknown>> } | undefined;
return {
value: searchData?.filings ?? [],
ttlClass: 'daily_permanent',
provenance: { fetchedAt, sourceKind: 'sec', rawSourceId: key },
};
}
// -----------------------------------------------------------------------
// 13F-HR Holdings
// -----------------------------------------------------------------------
/**
* Parse a 13F-HR filing's holdings table.
*
* Steps:
* 1. Fetch the filing index JSON at
* https://www.sec.gov/Archives/edgar/data/{paddedCik}/{accessionNoDashes}/index.json
* and locate the primary .txt or .xml document.
* 2. Fetch that document (XML/HTML) and extract the holdings table.
* 3. Return an array of { cusip, issuerName, value, sshPrnamt }.
*
* Uses `edgarFetch` for the index (JSON) and a secondary XML-aware fetch
* for the primary document. Rate-limiting is handled by the shared bucket.
*/
async form13f_holdings(
cik: string,
accession: string,
): Promise<FetchResult> {
const padded = padCik(cik);
const accessionNoDashes = accession.replace(/-/g, '');
// --- Step 1: filing index (JSON) ---------------------------------------
const indexUrl = `https://www.sec.gov/Archives/edgar/data/${padded}/${accessionNoDashes}/index.json`;
const indexResp = await edgarFetch(indexUrl);
if (!indexResp || indexResp.status >= 400) {
throw new Error(`EDGAR 13F index: ${indexResp?.status ?? 'unknown'} for ${padded}/${accessionNoDashes}`);
}
const indexBody = indexResp.body as {
fileDate?: string;
documents?: Array<{ name?: string; type?: string; size?: string | number; path?: string }>;
partialSubmissionIndicator?: unknown;
};
if (!indexBody?.documents || indexBody.documents.length === 0) {
throw new Error(`EDGAR 13F: no documents in index for ${padded}/${accessionNoDashes}`);
}
// Pick the primary document (usually the .txt or .xml filing).
const primaryDoc = indexBody.documents.find(
(d) => d.type === '13F' || d.type === '13F-infoTable'
) ?? indexBody.documents[0];
const primaryName = primaryDoc.name;
if (!primaryName) {
throw new Error(`EDGAR 13F: primary document has no name for ${padded}/${accessionNoDashes}`);
}
// --- Step 2: fetch the primary document (XML/HTML) ---------------------
const docUrl = `https://www.sec.gov/Archives/edgar/data/${padded}/${accessionNoDashes}/${primaryName}`;
const docText = await this.edgarXmlFetch(docUrl);
// --- Step 3: parse holdings table --------------------------------------
const holdings = parse13fHoldings(docText);
return {
value: { holdings, accession: `${padded}/${accessionNoDashes}` },
ttlClass: 'daily_permanent' as const,
provenance: {
fetchedAt: new Date().toISOString(),
sourceKind: 'sec' as const,
rawSourceId: `sec:13f:${padded}:${accessionNoDashes}`,
},
};
}
// -----------------------------------------------------------------------
// Form 4 Transactions
// -----------------------------------------------------------------------
/**
* Parse a Form 4 filing's transaction table.
*
* Steps:
* 1. Construct the XML URL:
* https://www.sec.gov/Archives/edgar/data/{paddedCik}/{accessionNoDashes}/{filename}.xml
* 2. Fetch the XML document.
* 3. Extract transaction rows (reporter, relationship, securityTitle,
* transactionDate, transactionCode, shares, price).
* 4. Return the transactions array.
*
* Uses `edgarFetch` for the index (JSON) and a secondary XML-aware fetch
* for the primary document. Rate-limiting is handled by the shared bucket.
*/
async form4_tx(
cik: string,
accession: string,
): Promise<FetchResult> {
const padded = padCik(cik);
const accessionNoDashes = accession.replace(/-/g, '');
// --- Step 1: filing index (JSON) ---------------------------------------
const indexUrl = `https://www.sec.gov/Archives/edgar/data/${padded}/${accessionNoDashes}/index.json`;
const indexResp = await edgarFetch(indexUrl);
if (!indexResp || indexResp.status >= 400) {
throw new Error(`EDGAR Form 4 index: ${indexResp?.status ?? 'unknown'} for ${padded}/${accessionNoDashes}`);
}
const indexBody = indexResp.body as {
fileDate?: string;
documents?: Array<{ name?: string; type?: string; size?: string | number; path?: string }>;
partialSubmissionIndicator?: unknown;
};
if (!indexBody?.documents || indexBody.documents.length === 0) {
throw new Error(`EDGAR Form 4: no documents in index for ${padded}/${accessionNoDashes}`);
}
// Form 4 is typically filed as a single XML. Pick the .xml document.
const xmlDoc = indexBody.documents.find(
(d) => d.name?.toLowerCase().endsWith('.xml')
) ?? indexBody.documents[0];
const xmlName = xmlDoc.name;
if (!xmlName) {
throw new Error(`EDGAR Form 4: document has no name for ${padded}/${accessionNoDashes}`);
}
// --- Step 2: fetch the XML document ------------------------------------
const xmlUrl = `https://www.sec.gov/Archives/edgar/data/${padded}/${accessionNoDashes}/${xmlName}.xml`;
const xmlText = await this.edgarXmlFetch(xmlUrl);
// --- Step 3: parse transactions ----------------------------------------
const transactions = parseForm4Transactions(xmlText);
return {
value: { transactions, accession: `${padded}/${accessionNoDashes}` },
ttlClass: 'daily_permanent' as const,
provenance: {
fetchedAt: new Date().toISOString(),
sourceKind: 'sec' as const,
rawSourceId: `sec:form4:${padded}:${accessionNoDashes}`,
},
};
}
// -----------------------------------------------------------------------
// Internal helpers: XML fetch + parsing
// -----------------------------------------------------------------------
/**
* Fetch a URL that returns XML/HTML (not JSON), using the rate limiter.
* EDGAR returns XML for primary filing documents; we strip the body here
* rather than calling `resp.json()`.
*/
async edgarXmlFetch(url: string): Promise<string> {
await bucket.wait();
const resp = await fetch(url, {
method: 'GET',
headers: {
'User-Agent': UA,
Accept: 'application/xml, text/xml, */*',
},
});
if (resp.status >= 400) {
throw new Error(`EDGAR XML ${resp.status} ${resp.statusText} for ${url}`);
}
const text = await resp.text();
return text;
}
/** SourceFetch.fetchOne dispatch. */
async fetchOne(key: CacheKey): Promise<FetchResult> {
const { id } = parseCacheKey(key);
const segments = id.split(':');
const subKind = segments[0];
const rest = segments.slice(1).join(':');
switch (subKind) {
case 'filings_index':
return this.filings_index(rest.split(':')[0]);
case 'company_facts':
return this.company_facts(rest);
case 'filer_meta':
return this.filer_cik_meta(rest);
case 'search_index':
return this.full_text_search(decodeURIComponent(rest));
case '13f':
return this.form13f_holdings(rest.split(':')[0], rest.split(':')[1]);
case 'form4':
return this.form4_tx(rest.split(':')[0], rest.split(':')[1]);
default:
throw new Error(`EdgarAdapter: unknown subKind '${subKind}'`);
}
}
}
// ---------------------------------------------------------------------------
// Module-level parsing helpers (used by EdgarAdapter class methods above).
// ---------------------------------------------------------------------------
/**
* Extract holdings from a 13F-HR filing document (XML or HTML).
*
* 13F filings use a table structure with columns:
* CUSIP | Name of Issuer and Title of Class | Value | sshPrnamt (shares)
*
* This is a thin regex-based parser that extracts the key fields.
*/
function parse13fHoldings(text: string): Array<{
cusip: string;
issuerName: string;
value: number;
sshPrnamt: number;
}> {
const holdings: Array<{ cusip: string; issuerName: string; value: number; sshPrnamt: number }> = [];
// Try XML-style <table> parsing first (13F filings are often XML).
const xmlMatch = text.match(/<table[^>]*>([\s\S]*?)<\/table>/gi);
if (xmlMatch) {
for (const tbl of xmlMatch) {
// Extract rows that look like data rows (not header).
const rowMatches = tbl.match(/<tr[^>]*>([\s\S]*?)<\/tr>/gi);
if (!rowMatches) continue;
for (const row of rowMatches) {
// Skip header rows.
if (/Class\s+of\s+Issuer|CUSIP\s+No/i.test(row)) continue;
// Extract CUSIP from <td> or <th>.
const cusipMatch = row.match(/<t[dh][^>]*>(\d{9,12})<\/t[dh]>/i);
if (!cusipMatch) continue;
// Extract issuer name (usually the second cell).
const cells = row.match(/<t[dh][^>]*>([\s\S]*?)<\/t[dh]>/gi);
if (!cells || cells.length < 4) continue;
const issuerRaw = cells[1].replace(/<[^>]+>/g, '').trim();
if (!issuerRaw) continue;
// Extract value (numeric, third cell).
const valueMatch = cells[2].match(/([\d,.]+)/);
const value = valueMatch ? parseFloat(valueMatch[1].replace(/,/g, '')) : 0;
// Extract shares (sshPrnamt, fourth cell).
const sharesMatch = cells[3].match(/([\d,.]+)/);
const sshPrnamt = sharesMatch ? parseFloat(sharesMatch[1].replace(/,/g, '')) : 0;
holdings.push({ cusip: cusipMatch[1], issuerName: issuerRaw, value, sshPrnamt });
}
}
}
// Fallback: if no XML holdings found, try regex on the raw text.
if (holdings.length === 0) {
const lines = text.split(/\r?\n/);
let currentCusip: string | null = null;
let currentIssuer = '';
let currentValue = 0;
let currentShares = 0;
for (const line of lines) {
// Match CUSIP line: typically 9-12 digits.
const cusipLine = line.match(/\b(\d{9,12})\b/);
if (cusipLine && !/Class|CUSIP|Name|Issuer/i.test(line)) {
// Save previous holding if we have one.
if (currentCusip) {
holdings.push({ cusip: currentCusip, issuerName: currentIssuer, value: currentValue, sshPrnamt: currentShares });
}
currentCusip = cusipLine[1];
currentIssuer = '';
currentValue = 0;
currentShares = 0;
} else if (currentCusip) {
// Accumulate issuer name from non-numeric lines.
if (!/\d/.test(line)) {
currentIssuer += line.trim() + ' ';
} else {
// Numeric data: value and shares.
const nums = line.match(/([\d,.]+)/g);
if (nums) {
currentValue = nums[0] ? parseFloat(nums[0].replace(/,/g, '')) : currentValue;
if (nums.length > 1) {
currentShares = parseFloat(nums[1].replace(/,/g, ''));
}
}
}
}
}
// Push last holding.
if (currentCusip) {
holdings.push({ cusip: currentCusip, issuerName: currentIssuer.trim(), value: currentValue, sshPrnamt: currentShares });
}
}
return holdings;
}
/**
* Extract transactions from a Form 4 XML filing.
*
* Form 4 uses an XML structure with <infotable> elements containing:
* <reportOwner>, <securityTitle>, <transactionDate>,
* <transactionCode>, <shares>, <pricePerShare>
*
* This is a thin parser that extracts the key fields.
*/
function parseForm4Transactions(text: string): Array<{
reporter: string;
relationship: string;
securityTitle: string;
transactionDate: string;
transactionCode: string;
shares: number;
price: number;
}> {
const transactions: Array<{
reporter: string;
relationship: string;
securityTitle: string;
transactionDate: string;
transactionCode: string;
shares: number;
price: number;
}> = [];
// Extract infotable blocks.
const infotables = text.match(/<infotable[\s\S]*?<\/infotable>/gi);
if (!infotables) return transactions;
for (const table of infotables) {
// Extract reporter name.
const reporterMatch = table.match(/<reporterCik>\s*(\d+[^<]*)/);
const reporter = reporterMatch ? reporterMatch[1].trim() : '';
// Extract relationship (direct/indirect).
const relMatch = table.match(/<directOrIndirectOwner>\s*([ADI])/i);
const relationship = relMatch ? (relMatch[1] === 'D' ? 'Direct' : 'Indirect') : '';
// Extract security title.
const titleMatch = table.match(/<securityTitle>\s*([\s\S]*?)<\/securityTitle>/i);
const securityTitle = titleMatch ? titleMatch[1].trim() : '';
// Extract transaction date.
const dateMatch = table.match(/<transactionDate>\s*(\d{4}-\d{2}-\d{2})/);
const transactionDate = dateMatch ? dateMatch[1] : '';
// Extract transaction code.
const codeMatch = table.match(/<transactionCode>\s*([A-HV])/i);
const transactionCode = codeMatch ? codeMatch[1].toUpperCase() : '';
// Extract shares (non-decimal, integer).
const sharesMatch = table.match(/<nonDerivativeShares>\s*([\d]+)/);
const shares = sharesMatch ? parseInt(sharesMatch[1], 10) : 0;
// Extract price per share.
const priceMatch = table.match(/<priceOrStrike\s*Price>\s*([\d,.]+)/);
const price = priceMatch ? parseFloat(priceMatch[1].replace(/,/g, '')) : 0;
transactions.push({
reporter,
relationship,
securityTitle,
transactionDate,
transactionCode,
shares,
price,
});
}
return transactions;
}