slice 6a EdgarAdapter (ornith-35): filings_index/company_facts/filer_cik_meta/full_text_search + edgarFetch (UA, 8 req/s, ETag/If-Modified-Since, 304 no-op)
Cross-review by qwopus35b pending (after tests 6b).
This commit is contained in:
@@ -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<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 },
|
||||
};
|
||||
}
|
||||
|
||||
/** 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));
|
||||
default:
|
||||
throw new Error(`EdgarAdapter: unknown subKind '${subKind}'`);
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user