feat(sec-lint): lint+backfill system for SEC data gaps (B1-B4)
- SecLintAdapter implements SourceFetch, runs via shared queue/drain loop - Two new SourceKinds: sec-lint-holders, sec-lint-insiders (weekly schedules) - No-op cache handlers so drain->cache.set doesn't throw on lint keys - tRPC admin.queueLint(symbol, kind) — run lint for one symbol, returns LintResult - tRPC admin.queueLintAll(kind) — backfill ALL watched symbols at once - tRPC admin.dataQualityList() — query data_quality rows (filterable by symbol/kind) - InstitutionalDashboard: 'Lint holders' button + status badge in detail panel header - Admin queue page: Data Quality section with per-row status badges, 'Lint all' buttons - DEFAULT_RATE_MS includes 167ms (~6 req/s) for lint kinds matching EDGAR limiter
This commit is contained in:
@@ -0,0 +1,55 @@
|
||||
// Investor Flow — SecLintAdapter (B2). Admin-triggered or scheduled lint+backfill for SEC data.
|
||||
// Implements SourceFetch so the shared queue/drain loop can process it like any other source,
|
||||
// but actual work writes directly to institution_filings/insider_transactions/data_quality tables.
|
||||
// The kv_cache write from drain→cache.set is a harmless no-op side effect (never read).
|
||||
import type { DatabaseSync } from 'node:sqlite';
|
||||
import type { SourceFetch, FetchResult } from './SourceAdapter.ts';
|
||||
import type { CacheKey, TtlClass, Provenance } from '../cache/CacheRepository.ts';
|
||||
import { lintInstitutionalHolders, lintInsiderTransactions, type LintResult } from '../services/secDataFetcher.ts';
|
||||
|
||||
export class SecLintAdapter implements SourceFetch {
|
||||
readonly sourceKind: 'sec-lint-holders' | 'sec-lint-insiders';
|
||||
private db: () => DatabaseSync;
|
||||
|
||||
constructor(db: () => DatabaseSync, sourceKind: 'sec-lint-holders' | 'sec-lint-insiders') {
|
||||
this.db = db;
|
||||
this.sourceKind = sourceKind;
|
||||
}
|
||||
|
||||
async fetchOne(key: CacheKey): Promise<FetchResult> {
|
||||
// key format: sec-lint-holders:holders:{symbol} or sec-lint-insiders:insiders:{symbol}
|
||||
const parts = key.split(':');
|
||||
const symbol = (parts[2] ?? '').toUpperCase();
|
||||
let result: LintResult;
|
||||
if (this.sourceKind === 'sec-lint-holders') {
|
||||
result = await lintInstitutionalHolders(this.db(), symbol);
|
||||
} else {
|
||||
result = await lintInsiderTransactions(this.db(), symbol);
|
||||
}
|
||||
const value: FetchResult['value'] = { lintDone: true, ...result };
|
||||
return {
|
||||
value,
|
||||
ttlClass: 'daily_permanent' as TtlClass,
|
||||
provenance: { fetchedAt: new Date().toISOString(), sourceKind: this.sourceKind },
|
||||
};
|
||||
}
|
||||
|
||||
/** Run lint for ALL watched symbols at once (used by admin "run all" button). */
|
||||
async runAllSymbols(kind: 'sec-lint-holders' | 'sec-lint-insiders', db: DatabaseSync): Promise<LintResult[]> {
|
||||
const symbols = (db.prepare('SELECT symbol FROM symbol_demand WHERE in_demand=1 ORDER BY symbol').all() as Array<{ symbol: string }>).map((r) => r.symbol);
|
||||
const results: LintResult[] = [];
|
||||
for (const sym of symbols) {
|
||||
try {
|
||||
if (kind === 'sec-lint-holders') {
|
||||
results.push(await lintInstitutionalHolders(db, sym));
|
||||
} else {
|
||||
results.push(await lintInsiderTransactions(db, sym));
|
||||
}
|
||||
} catch (e) {
|
||||
const msg = e instanceof Error ? e.message : String(e);
|
||||
results.push({ symbol: sym.toUpperCase(), kind, discoveredCount: 0, storedCount: 0, missingCount: 0, backfilled: 0, stale: 1, status: 'error', detail: { reason: msg } });
|
||||
}
|
||||
}
|
||||
return results;
|
||||
}
|
||||
}
|
||||
+90
-4
@@ -6,7 +6,7 @@
|
||||
import { DatabaseSync } from 'node:sqlite';
|
||||
import { db as defaultDb } from '../db/client.ts';
|
||||
|
||||
export type SourceKind = 'yfinance' | 'sec' | 'reddit' | 'x' | 'macro' | 'llm';
|
||||
export type SourceKind = 'yfinance' | 'sec' | 'sec-fetch' | 'reddit' | 'x' | 'macro' | 'llm' | 'sec-lint-holders' | 'sec-lint-insiders';
|
||||
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 =
|
||||
@@ -139,6 +139,7 @@ const symbolHandler: KindHandler = {
|
||||
isStale(ts, now) { return tsAgeMs(ts, now) > TTL_MS.symbol_meta; },
|
||||
};
|
||||
|
||||
// ----- Options handlers (slice 15) -----
|
||||
const optionsChainHandler: KindHandler = {
|
||||
ttlClass: 'options_snapshot',
|
||||
read(d, id) {
|
||||
@@ -177,9 +178,9 @@ const optionsChainHandler: KindHandler = {
|
||||
const right = r.right === 'put' ? 'put' : 'call';
|
||||
ins.run(
|
||||
symbol, expiry, strike, right,
|
||||
r.bid ?? null, r.ask ?? null, r.impliedVolatility ?? null,
|
||||
r.delta ?? null, r.gamma ?? null, r.theta ?? null, r.vega ?? null,
|
||||
r.openInterest ?? null, r.volume ?? null,
|
||||
numOrNull(r.bid), numOrNull(r.ask), numOrNull(r.impliedVolatility),
|
||||
numOrNull(r.delta), numOrNull(r.gamma), numOrNull(r.theta), numOrNull(r.vega),
|
||||
numOrNull(r.openInterest), numOrNull(r.volume),
|
||||
provenance.fetchedAt
|
||||
);
|
||||
}
|
||||
@@ -205,6 +206,87 @@ const optionsExpiryDatesHandler: KindHandler = {
|
||||
isStale(ts, now) { return tsAgeMs(ts, now) > TTL_MS.intraday; },
|
||||
};
|
||||
|
||||
/** Coerce a possibly-undefined/unknown value to number | null for SQL binding. */
|
||||
function numOrNull(v: unknown): number | null {
|
||||
return typeof v === 'number' ? v : null;
|
||||
}
|
||||
|
||||
const greeksHandler: KindHandler = {
|
||||
ttlClass: 'options_snapshot',
|
||||
read(d, id) {
|
||||
const { symbol, expiry, strike } = parseGreeksId(id);
|
||||
const r = d.prepare(
|
||||
'SELECT delta, gamma, theta, vega, strike, open_interest, iv, ts, type FROM options_chains WHERE symbol=? AND expiry=? AND strike=? LIMIT 1'
|
||||
).get(symbol, expiry, parseFloat(strike ?? '0')) as Record<string, unknown> | undefined;
|
||||
if (!r) return null;
|
||||
return {
|
||||
value: {
|
||||
delta: numOrNull(r.delta),
|
||||
gamma: numOrNull(r.gamma),
|
||||
theta: numOrNull(r.theta),
|
||||
vega: numOrNull(r.vega),
|
||||
strike: typeof r.strike === 'number' ? r.strike : 0,
|
||||
openInterest: numOrNull(r.open_interest),
|
||||
impliedVolatility: numOrNull(r.iv),
|
||||
right: r.type as 'call' | 'put',
|
||||
},
|
||||
stalenessTs: r.ts as string,
|
||||
};
|
||||
},
|
||||
write(d, id, value, provenance) {
|
||||
const { symbol, expiry, strike } = parseGreeksId(id);
|
||||
const v = value as Record<string, unknown>;
|
||||
const right = v.right === 'put' ? 'put' : 'call';
|
||||
d.prepare(
|
||||
'INSERT OR REPLACE INTO options_chains (symbol,expiry,strike,type,bid,ask,iv,delta,gamma,theta,vega,open_interest,volume,ts) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)'
|
||||
).run(
|
||||
symbol, expiry, parseFloat(strike ?? '0'), right,
|
||||
numOrNull(v.bid), numOrNull(v.ask), numOrNull(v.impliedVolatility),
|
||||
numOrNull(v.delta), numOrNull(v.gamma), numOrNull(v.theta), numOrNull(v.vega),
|
||||
numOrNull(v.openInterest), numOrNull(v.volume),
|
||||
provenance.fetchedAt
|
||||
);
|
||||
},
|
||||
isStale(ts, now) { return tsAgeMs(ts, now) > TTL_MS.options_snapshot; },
|
||||
};
|
||||
|
||||
/** Parse 'symbol:expiry:strike' from greeks cache key id. */
|
||||
function parseGreeksId(id: string): { symbol: string; expiry: string; strike: string } {
|
||||
const parts = id.split(':');
|
||||
return { symbol: parts[0] ?? '', expiry: parts[1] ?? '', strike: parts[2] ?? '0' };
|
||||
}
|
||||
|
||||
const fetchHandler: KindHandler = {
|
||||
ttlClass: 'daily_permanent',
|
||||
read(d, id) {
|
||||
const r = d.prepare('SELECT value, observed_at FROM kv_cache WHERE key=?').get(`sec-fetch:${id}`) as { value: string; observed_at: string } | undefined;
|
||||
if (!r) return null;
|
||||
try { return { value: JSON.parse(r.value), stalenessTs: r.observed_at }; } catch { return null; }
|
||||
},
|
||||
write(d, id, value, provenance) {
|
||||
d.prepare('INSERT OR REPLACE INTO kv_cache (key, value, observed_at) VALUES (?,?,?)').run(`sec-fetch:${id}`, JSON.stringify(value), provenance.fetchedAt);
|
||||
},
|
||||
isStale(ts) { return ts === null; }, // never stale once written
|
||||
};
|
||||
|
||||
const lintHoldersHandler: KindHandler = {
|
||||
ttlClass: 'daily_permanent',
|
||||
read() { return null; }, // never read — work happens in DB tables directly
|
||||
write(d, key, _value, provenance) {
|
||||
d.prepare('INSERT OR REPLACE INTO kv_cache (key,value,observed_at) VALUES (?,?,?)').run(`sec-lint-holders:${key}`, JSON.stringify({ ok: true }), provenance.fetchedAt);
|
||||
},
|
||||
isStale() { return false; }, // never stale once written (lint writes are permanent)
|
||||
};
|
||||
|
||||
const lintInsidersHandler: KindHandler = {
|
||||
ttlClass: 'daily_permanent',
|
||||
read() { return null; }, // never read — work happens in DB tables directly
|
||||
write(d, key, _value, provenance) {
|
||||
d.prepare('INSERT OR REPLACE INTO kv_cache (key,value,observed_at) VALUES (?,?,?)').run(`sec-lint-insiders:${key}`, JSON.stringify({ ok: true }), provenance.fetchedAt);
|
||||
},
|
||||
isStale() { return false; }, // never stale once written
|
||||
};
|
||||
|
||||
const HANDLERS = new Map<string, KindHandler>([
|
||||
['quote', quoteHandler],
|
||||
['candles', candlesHandler],
|
||||
@@ -212,6 +294,10 @@ const HANDLERS = new Map<string, KindHandler>([
|
||||
['adjustments', adjustmentsHandler],
|
||||
['chain', optionsChainHandler],
|
||||
['expiry_dates', optionsExpiryDatesHandler],
|
||||
['greeks', greeksHandler],
|
||||
['fetch', fetchHandler],
|
||||
['holders', lintHoldersHandler],
|
||||
['insiders', lintInsidersHandler],
|
||||
]);
|
||||
|
||||
export interface CacheRepository {
|
||||
|
||||
@@ -6,6 +6,7 @@ import { db } from './db/client.ts';
|
||||
import { createCacheRepository, type SourceKind } from './cache/CacheRepository.ts';
|
||||
import { YFinanceAdapter } from './adapters/YFinanceAdapter.ts';
|
||||
import { SecFetchAdapter } from './adapters/SecFetchAdapter.ts';
|
||||
import { SecLintAdapter } from './adapters/SecLintAdapter.ts';
|
||||
import type { SourceFetch } from './adapters/SourceAdapter.ts';
|
||||
import { AdapterQueue } from './queue/AdapterQueue.ts';
|
||||
import { makeCreateContext } from './trpc/context.ts';
|
||||
@@ -13,9 +14,11 @@ import { appRouter } from './trpc/router.ts';
|
||||
|
||||
const PORT = Number(process.env.PORT ?? 3001);
|
||||
const database = db();
|
||||
const adapters: Map<SourceKind, SourceFetch> = new Map([
|
||||
['yfinance', new YFinanceAdapter()],
|
||||
['sec-fetch', new SecFetchAdapter(database)],
|
||||
const adapters = new Map<SourceKind, SourceFetch>([
|
||||
['yfinance' as const, new YFinanceAdapter() as unknown as SourceFetch],
|
||||
['sec-fetch' as const, new SecFetchAdapter(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],
|
||||
]);
|
||||
const queue = new AdapterQueue({ db: database, adapters });
|
||||
const cache = createCacheRepository({ db: database, scheduler: queue });
|
||||
|
||||
@@ -14,7 +14,7 @@ export interface AdapterQueueOptions {
|
||||
rateLimitMs?: Partial<Record<SourceKind, number>>;
|
||||
}
|
||||
|
||||
const DEFAULT_RATE_MS: Record<SourceKind, number> = { yfinance: 1000, sec: 125, 'sec-fetch': 1000, reddit: 1000, x: 3000, macro: 1000, llm: 0 };
|
||||
const DEFAULT_RATE_MS: Record<SourceKind, number> = { yfinance: 1000, sec: 125, 'sec-fetch': 1000, reddit: 1000, x: 3000, macro: 1000, llm: 0, 'sec-lint-holders': 167, 'sec-lint-insiders': 167 };
|
||||
const BACKOFF_MS = [2000, 4000, 8000, 16000, 60000];
|
||||
const MAX_ATTEMPTS = 5;
|
||||
|
||||
@@ -170,6 +170,14 @@ export class AdapterQueue implements CacheScheduler {
|
||||
await this.queue(`yfinance:candles:${sym.symbol}:1d`);
|
||||
await this.queue(`yfinance:symbol:${sym.symbol}`);
|
||||
}
|
||||
} else if (s.source_kind === 'sec-lint-holders') {
|
||||
for (const sym of symbols) {
|
||||
await this.queue(`sec-lint-holders:holders:${sym.symbol}`);
|
||||
}
|
||||
} else if (s.source_kind === 'sec-lint-insiders') {
|
||||
for (const sym of symbols) {
|
||||
await this.queue(`sec-lint-insiders:insiders:${sym.symbol}`);
|
||||
}
|
||||
}
|
||||
const nextEnqueue = new Date(Date.now() + s.interval_ms).toISOString();
|
||||
this._db.prepare("UPDATE queue_schedules SET last_enqueued=?, next_enqueue=? WHERE source_kind=?").run(now, nextEnqueue, s.source_kind);
|
||||
@@ -182,6 +190,8 @@ export class AdapterQueue implements CacheScheduler {
|
||||
const defaults: Array<[string, number]> = [
|
||||
['sec-fetch', 86400000],
|
||||
['yfinance', 300000],
|
||||
['sec-lint-holders', 7 * 86400000],
|
||||
['sec-lint-insiders', 7 * 86400000],
|
||||
];
|
||||
const insert = this._db.prepare("INSERT OR IGNORE INTO queue_schedules (source_kind, interval_ms, last_enqueued, next_enqueue) VALUES (?,?,?,?)");
|
||||
for (const [kind, ms] of defaults) {
|
||||
|
||||
@@ -10,6 +10,7 @@ import { STARTER_WATCHLIST, defaultDrawdownTolerancePct, defaultRiskTolerance, O
|
||||
import type { Quote, PriceCandle, SymbolMeta } from '../cache/CacheRepository.ts';
|
||||
import { emaFromCandles, rsi as rsiFn, relativeVolume } from '../analysis/indicators.ts';
|
||||
import { listUsers, resetPassword, gdprExport, queueHealth, resetQueueBackoff, NotOwnerError, listUserSessions, listAuditLog, queueSecFetch } from '../admin/admin.ts';
|
||||
import type { LintResult } from '../services/secDataFetcher.ts';
|
||||
import { EdgarAdapter } from '../adapters/EdgarAdapter.ts';
|
||||
import { OptionsAdapter, parseOptionChainRows } from '../adapters/OptionsAdapter.ts';
|
||||
import type { OptionChainRow, OptionGreeks } from '../adapters/OptionsAdapter.ts';
|
||||
@@ -328,6 +329,53 @@ const adminRouter = router({
|
||||
ctx.queue.deleteSchedule(input.sourceKind);
|
||||
return { ok: true };
|
||||
}),
|
||||
|
||||
queueLint: adminProcedure
|
||||
.input(z.object({ symbol: z.string().min(1).max(10), kind: z.enum(['sec-lint-holders', 'sec-lint-insiders']) }))
|
||||
.mutation(async ({ ctx, input }) => {
|
||||
const { lintInstitutionalHolders, lintInsiderTransactions } = await import('../services/secDataFetcher.ts');
|
||||
let result: LintResult;
|
||||
if (input.kind === 'sec-lint-holders') {
|
||||
result = await lintInstitutionalHolders(ctx.db, input.symbol);
|
||||
} else {
|
||||
result = await lintInsiderTransactions(ctx.db, input.symbol);
|
||||
}
|
||||
// Also enqueue so the weekly schedule picks it up too.
|
||||
try { ctx.queue.queue(`${input.kind}:${input.kind === 'sec-lint-holders' ? 'holders' : 'insiders'}:${input.symbol}`); } catch { /* non-fatal */ }
|
||||
return result;
|
||||
}),
|
||||
|
||||
dataQualityList: adminProcedure
|
||||
.input(z.object({ symbol: z.string().min(1).max(20).optional(), kind: z.enum(['institution_filings', 'insider_transactions']).nullish() }))
|
||||
.query(({ ctx, input }) => {
|
||||
let sql = 'SELECT symbol, kind, last_checked_at, stored_count, discovered_count, missing_count, stale, status, detail FROM data_quality WHERE 1=1';
|
||||
const params: unknown[] = [];
|
||||
if (input.symbol) { sql += ' AND symbol = ?'; params.push(input.symbol); }
|
||||
if (input.kind) { sql += ' AND kind = ?'; params.push(input.kind); }
|
||||
sql += ' ORDER BY last_checked_at DESC LIMIT 200';
|
||||
return ctx.db.prepare(sql).all(...params) as Array<Record<string, unknown>>;
|
||||
}),
|
||||
|
||||
queueLintAll: adminProcedure
|
||||
.input(z.object({ kind: z.enum(['sec-lint-holders', 'sec-lint-insiders']) }))
|
||||
.mutation(async ({ ctx, input }) => {
|
||||
const { lintInstitutionalHolders, lintInsiderTransactions } = await import('../services/secDataFetcher.ts');
|
||||
const symbols = (ctx.db.prepare('SELECT symbol FROM symbol_demand WHERE in_demand=1 ORDER BY symbol').all() as Array<{ symbol: string }>).map((r) => r.symbol);
|
||||
const results: LintResult[] = [];
|
||||
for (const sym of symbols) {
|
||||
try {
|
||||
if (input.kind === 'sec-lint-holders') {
|
||||
results.push(await lintInstitutionalHolders(ctx.db, sym));
|
||||
} else {
|
||||
results.push(await lintInsiderTransactions(ctx.db, sym));
|
||||
}
|
||||
} catch (e) {
|
||||
const msg = e instanceof Error ? e.message : String(e);
|
||||
results.push({ symbol: sym.toUpperCase(), kind: input.kind === 'sec-lint-holders' ? 'institution_filings' : 'insider_transactions', discoveredCount: 0, storedCount: 0, missingCount: 0, backfilled: 0, stale: 1, status: 'error', detail: { reason: msg } });
|
||||
}
|
||||
}
|
||||
return { total: symbols.length, results };
|
||||
}),
|
||||
});
|
||||
|
||||
|
||||
@@ -646,6 +694,7 @@ interface Form13fHoldingsResult {
|
||||
value: number;
|
||||
sshPrnamt: number;
|
||||
}>;
|
||||
total: number;
|
||||
accession: string;
|
||||
}
|
||||
|
||||
@@ -722,14 +771,23 @@ const edgarRouter = router({
|
||||
}
|
||||
}),
|
||||
|
||||
/** 13F-HR holdings for a CIK + accession. */
|
||||
/** 13F-HR holdings for a CIK + accession (server-side paginated). */
|
||||
form13f_holdings: publicProcedure
|
||||
.input(z.object({ cik: z.string().min(1), accession: z.string().min(1) }))
|
||||
.input(z.object({
|
||||
cik: z.string().min(1),
|
||||
accession: z.string().min(1),
|
||||
limit: z.number().int().positive().optional(),
|
||||
offset: z.number().int().min(0).optional(),
|
||||
}))
|
||||
.query(async ({ ctx, input }) => {
|
||||
const edgar = new EdgarAdapter();
|
||||
try {
|
||||
const result = await edgar.form13f_holdings(input.cik, input.accession);
|
||||
return { holdings: result.value as Form13fHoldingsResult };
|
||||
const result = await edgar.form13f_holdings(input.cik, input.accession, {
|
||||
limit: input.limit,
|
||||
offset: input.offset,
|
||||
});
|
||||
const v = result.value as Form13fHoldingsResult;
|
||||
return { holdings: v.holdings, total: v.total, accession: v.accession };
|
||||
} catch (e) {
|
||||
throw new TRPCError({ code: 'NOT_FOUND', message: (e as Error).message });
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user