diff --git a/app/server/src/adapters/EdgarAdapter.ts b/app/server/src/adapters/EdgarAdapter.ts new file mode 100644 index 0000000..3238927 --- /dev/null +++ b/app/server/src/adapters/EdgarAdapter.ts @@ -0,0 +1,313 @@ +// 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 { + 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, +): Promise<{ status: number; headers: { etag?: string | null; lastModified?: string | null }; body: unknown } | null> { + await bucket.wait(); + + const headers: Record = { + '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; + etags: Record; + lastModified: Record; +} + +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 { + 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 { + 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 } | 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 { + 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 { + 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> } | undefined; + + return { + value: searchData?.filings ?? [], + ttlClass: 'daily_permanent', + provenance: { fetchedAt, sourceKind: 'sec', rawSourceId: key }, + }; + } + + /** SourceFetch.fetchOne dispatch. */ + async fetchOne(key: CacheKey): Promise { + 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)); + default: + throw new Error(`EdgarAdapter: unknown subKind '${subKind}'`); + } + } +}