diff --git a/Congress.Trade.xcworkspace/contents.xcworkspacedata b/Congress.Trade.xcworkspace/contents.xcworkspacedata deleted file mode 100644 index fa25973e7..000000000 --- a/Congress.Trade.xcworkspace/contents.xcworkspacedata +++ /dev/null @@ -1,7 +0,0 @@ - - - - - diff --git a/app/.dev.vars.example b/app/.dev.vars.example index 0b8ec0778..219f03cd3 100644 --- a/app/.dev.vars.example +++ b/app/.dev.vars.example @@ -86,11 +86,19 @@ IMPORT_MAX_CLOSES_PER_TICKER="1500" IMPORT_MAX_INSIDER="5000" IMPORT_MAX_SHORT_VOLUME="5000" -# Optional temporary benchmark: polls FMP's House/Senate latest-disclosure -# endpoints and records whether Congress.Trade saw new filings before FMP did. +# Optional disclosure-latency benchmark: polls third-party latest-disclosure +# endpoints and records whether Congress.Trade saw new filings first. +# Supported direct providers: fmp, unusual_whales, quiver. Finnhub and AInvest +# are symbol-scoped and Capitol Trades has no stable server API, so those are +# reported in admin provider status but not automatically probed. # Turn on before waiting for the next new disclosures; default code path is off. +DISCLOSURE_LATENCY_WATCH_ENABLED="" +DISCLOSURE_LATENCY_PROVIDERS="fmp,unusual_whales,quiver" +DISCLOSURE_LATENCY_WATCH_LIMIT="100" FMP_DISCLOSURE_WATCH_ENABLED="" FMP_DISCLOSURE_WATCH_LIMIT="100" +# Optional provider credentials are resolved from .dev.vars or Worker secrets +# when those providers are enabled: Unusual Whales, Quiver, and AInvest. # Cloudflare Access (Zero Trust) sign-in for humans — an ALTERNATIVE/ADDITION to # the bearer token. With these set AND an Access application fronting diff --git a/app/migrations/0021_disclosure_latency_watch.sql b/app/migrations/0021_disclosure_latency_watch.sql index 0b20edd38..1d8891094 100644 --- a/app/migrations/0021_disclosure_latency_watch.sql +++ b/app/migrations/0021_disclosure_latency_watch.sql @@ -1,8 +1,8 @@ -- 0021_disclosure_latency_watch.sql --- Records Congress.Trade-vs-FMP disclosure discovery timing. Candidates are +-- Records Congress.Trade-vs-provider disclosure discovery timing. Candidates are -- created when our watcher first sees a new filing; provider observations are --- populated from FMP latest endpoints so we can tell whether FMP was already --- aware or caught up later. +-- populated from provider latest endpoints so we can tell who was already aware +-- or caught up later. CREATE TABLE IF NOT EXISTS disclosure_latency_candidates ( doc_id TEXT NOT NULL, diff --git a/app/migrations/0023_disclosure_provider_timestamps.sql b/app/migrations/0023_disclosure_provider_timestamps.sql new file mode 100644 index 000000000..32f33fe20 --- /dev/null +++ b/app/migrations/0023_disclosure_provider_timestamps.sql @@ -0,0 +1,6 @@ +-- 0023_disclosure_provider_timestamps.sql +-- Stores provider-supplied publication/upload timestamps separately from the +-- monitor's first-observed timestamp. Not every provider exposes this. + +ALTER TABLE disclosure_latency_candidates ADD COLUMN provider_published_at TEXT; +ALTER TABLE disclosure_provider_observations ADD COLUMN provider_published_at TEXT; diff --git a/app/src/admin/__tests__/disclosureLatency.test.ts b/app/src/admin/__tests__/disclosureLatency.test.ts index da90362fc..e78e7de92 100644 --- a/app/src/admin/__tests__/disclosureLatency.test.ts +++ b/app/src/admin/__tests__/disclosureLatency.test.ts @@ -26,6 +26,7 @@ function fakeDb() { congress_first_seen_at: '2026-06-29T14:00:00.000Z', provider_key: '20012345', provider_first_seen_at: '2026-06-29T14:03:30.000Z', + provider_published_at: '2026-06-29T14:02:00.000Z', match_method: 'doc-token', status: 'matched', attempts: 2, @@ -59,13 +60,36 @@ describe('admin disclosure latency API', () => { ); expect(res.status).toBe(200); - const body = (await res.json()) as { items: Array<{ docId: string; providerDeltaSec: number; status: string }> }; + const body = (await res.json()) as { + items: Array<{ docId: string; providerDeltaSec: number; providerPublishedDeltaSec: number; status: string }>; + }; expect(body.items).toEqual([ expect.objectContaining({ docId: 'H-2026-20012345', providerDeltaSec: 210, + providerPublishedDeltaSec: 120, status: 'matched', }), ]); }); + + it('returns aggregate metrics and a public-safe summary payload', async () => { + const res = await app.request( + '/disclosure-latency/summary', + { headers: { Authorization: 'Bearer admin-secret' } }, + { ADMIN_TOKEN: 'admin-secret', FMP_API_KEY: 'configured', DB: fakeDb() } as never, + ); + + expect(res.status).toBe(200); + const body = (await res.json()) as { + totals: { candidates: number; matched: number; configuredComparableProviders: number }; + providers: Array<{ provider: string; avgMonitorDeltaSec: number | null; avgProviderPublishedDeltaSec: number | null }>; + publicSummary: { providers: Array<{ provider: string }> }; + }; + expect(body.totals).toEqual(expect.objectContaining({ candidates: 1, matched: 1, configuredComparableProviders: 1 })); + expect(body.providers[0]).toEqual( + expect.objectContaining({ provider: 'fmp', avgMonitorDeltaSec: 210, avgProviderPublishedDeltaSec: 120 }), + ); + expect(JSON.stringify(body.publicSummary)).not.toContain('H-2026-20012345'); + }); }); diff --git a/app/src/admin/routes.ts b/app/src/admin/routes.ts index e7c0b5904..cdf375fc3 100644 --- a/app/src/admin/routes.ts +++ b/app/src/admin/routes.ts @@ -74,7 +74,7 @@ import { mergeRefs } from '../enrichment/compute'; import type { SecurityRef } from '../enrichment/types'; import { runPriceRefresh } from '../prices/service'; import { getSecretResolverStatus, refreshSecrets, resolveSecret, resolveSecrets } from '../secrets/infisical'; -import { runFmpDisclosureLatencyProbe } from '../ingestion/fmpDisclosureLatency'; +import { getDisclosureLatencySummary, runDisclosureLatencyProbe } from '../ingestion/fmpDisclosureLatency'; // Optional secrets/vars; not declared on Env (frozen). Read defensively. type EnvWithAdmin = Env & { @@ -1131,12 +1131,22 @@ export function buildAdminRouter(): Hono<{ Bindings: Env }> { return c.json({ ok: true, latencyResetAt }); }); + // --- GET /disclosure-latency/summary ------------------------------------ + // Aggregate provider-race metrics. `publicSummary` intentionally excludes + // filing/member detail so it can be reviewed before any public sharing. + r.get('/disclosure-latency/summary', async (c) => { + return c.json(await getDisclosureLatencySummary(c.env)); + }); + // --- GET /disclosure-latency ------------------------------------------- - // Congress.Trade-vs-FMP race monitor. `providerDeltaSec` is FMP monitor - // first-observed minus Congress.Trade first_seen_at: positive means we observed - // first; negative means FMP was already observed first. + // Congress.Trade-vs-provider race monitor. `providerDeltaSec` is provider + // monitor first-observed minus Congress.Trade first_seen_at: positive means we + // observed first; negative means the provider was already observed first. r.get('/disclosure-latency', async (c) => { const limit = Math.min(Math.max(parseInt(c.req.query('limit') || '50', 10) || 50, 1), 200); + const provider = (c.req.query('provider') || '').trim().toLowerCase(); + const where = provider ? 'WHERE provider = ?' : ''; + const params: SqlParam[] = provider ? [provider, limit] : [limit]; const rows = await optionalAll<{ doc_id: string; provider: string; @@ -1147,6 +1157,7 @@ export function buildAdminRouter(): Hono<{ Bindings: Env }> { congress_first_seen_at: string; provider_key: string | null; provider_first_seen_at: string | null; + provider_published_at: string | null; match_method: string | null; status: string; attempts: number; @@ -1157,13 +1168,14 @@ export function buildAdminRouter(): Hono<{ Bindings: Env }> { }>( c.env, `SELECT doc_id, provider, chamber, source_url, filed_date, filer_name, - congress_first_seen_at, provider_key, provider_first_seen_at, + congress_first_seen_at, provider_key, provider_first_seen_at, provider_published_at, match_method, status, attempts, last_checked_at, error, created_at, updated_at FROM disclosure_latency_candidates + ${where} ORDER BY created_at DESC LIMIT ?`, - [limit], + params, ); const items = rows.map((row) => ({ docId: row.doc_id, @@ -1176,6 +1188,8 @@ export function buildAdminRouter(): Hono<{ Bindings: Env }> { providerKey: row.provider_key, providerFirstSeenAt: row.provider_first_seen_at, providerDeltaSec: deltaSeconds(row.provider_first_seen_at, row.congress_first_seen_at), + providerPublishedAt: row.provider_published_at, + providerPublishedDeltaSec: deltaSeconds(row.provider_published_at, row.congress_first_seen_at), matchMethod: row.match_method, status: row.status, attempts: row.attempts, @@ -1188,10 +1202,15 @@ export function buildAdminRouter(): Hono<{ Bindings: Env }> { }); // --- POST /disclosure-latency/probe ------------------------------------- - // Force a one-off FMP latest probe, useful immediately after new filings land - // or before turning on the continuous cron switch. + // Force a one-off provider latest probe, useful immediately after new filings + // land or before turning on the continuous cron switch. Optional query: + // ?providers=fmp,unusual_whales,quiver r.post('/disclosure-latency/probe', async (c) => { - const result = await runFmpDisclosureLatencyProbe(c.env, new Date(), fetch, { force: true }); + const providers = (c.req.query('providers') || c.req.query('provider') || '') + .split(/[,\s]+/) + .map((part) => part.trim()) + .filter(Boolean); + const result = await runDisclosureLatencyProbe(c.env, new Date(), fetch, { force: true, providers }); return c.json({ ok: result.errors.length === 0, ...result }); }); @@ -2737,6 +2756,7 @@ export function buildAdminRouter(): Hono<{ Bindings: Env }> { congress_first_seen_at TEXT NOT NULL, provider_key TEXT, provider_first_seen_at TEXT, + provider_published_at TEXT, match_method TEXT, status TEXT NOT NULL DEFAULT 'pending', attempts INTEGER NOT NULL DEFAULT 0, @@ -2755,6 +2775,7 @@ export function buildAdminRouter(): Hono<{ Bindings: Env }> { provider_key TEXT NOT NULL, first_observed_at TEXT NOT NULL, last_observed_at TEXT NOT NULL, + provider_published_at TEXT, source_url TEXT, filed_date TEXT, filer_name TEXT, @@ -2772,6 +2793,9 @@ export function buildAdminRouter(): Hono<{ Bindings: Env }> { )`, `CREATE INDEX IF NOT EXISTS idx_stripe_webhook_events_received ON stripe_webhook_events (received_at DESC)`, + // 0023_disclosure_provider_timestamps.sql — provider-side publish/upload timestamp when available. + 'ALTER TABLE disclosure_latency_candidates ADD COLUMN provider_published_at TEXT', + 'ALTER TABLE disclosure_provider_observations ADD COLUMN provider_published_at TEXT', ]; const applied: string[] = []; const skipped: string[] = []; diff --git a/app/src/index.ts b/app/src/index.ts index aca0a8052..9ebe89455 100644 --- a/app/src/index.ts +++ b/app/src/index.ts @@ -36,7 +36,7 @@ import { buildUiRouter } from './ui/routes'; import { maybeRunDailyJobs } from './jobs'; import { maybeRunAgreementAutopublish, handleAgreementCheck } from './extraction/agreement'; import { refreshSecrets } from './secrets/infisical'; -import { runFmpDisclosureLatencyProbe } from './ingestion/fmpDisclosureLatency'; +import { runDisclosureLatencyProbe } from './ingestion/fmpDisclosureLatency'; const app = new Hono<{ Bindings: Env }>(); @@ -161,8 +161,8 @@ export default Sentry.withSentry( await runWatcher(env, new Date()); ctx.waitUntil(refreshSecrets(env).catch((err) => console.warn('infisical secret refresh failed:', (err as Error).message))); ctx.waitUntil( - runFmpDisclosureLatencyProbe(env).catch((err) => - console.warn('fmp disclosure latency probe failed:', (err as Error).message), + runDisclosureLatencyProbe(env).catch((err) => + console.warn('disclosure latency probe failed:', (err as Error).message), ), ); ctx.waitUntil(maybeRunDailyJobs(env)); diff --git a/app/src/ingestion/__tests__/fmpDisclosureLatency.test.ts b/app/src/ingestion/__tests__/fmpDisclosureLatency.test.ts index bdc79e2ae..75d1f86c0 100644 --- a/app/src/ingestion/__tests__/fmpDisclosureLatency.test.ts +++ b/app/src/ingestion/__tests__/fmpDisclosureLatency.test.ts @@ -1,5 +1,11 @@ import { describe, expect, it } from 'vitest'; -import { matchFmpDisclosureCandidate, parseFmpDisclosureRows } from '../fmpDisclosureLatency'; +import { + matchDisclosureCandidate, + matchFmpDisclosureCandidate, + parseFmpDisclosureRows, + parseQuiverDisclosureRows, + parseUnusualWhalesDisclosureRows, +} from '../fmpDisclosureLatency'; describe('parseFmpDisclosureRows', () => { it('extracts a House doc token from PTR PDF URLs', () => { @@ -86,6 +92,65 @@ describe('matchFmpDisclosureCandidate', () => { }, row, ), - ).toEqual({ providerKey: row.providerKey, matchMethod: 'probable-filer-date' }); + ).toEqual({ providerKey: row.providerKey, matchMethod: 'filer-date' }); + }); +}); + +describe('parse third-party disclosure providers', () => { + it('normalizes Unusual Whales recent Congress rows', () => { + const rows = parseUnusualWhalesDisclosureRows({ + data: [ + { + filed_at_date: '2026-06-29', + member_type: 'senate', + name: 'Jane Smith', + politician_id: 'abc', + ticker: 'MSFT', + transaction_date: '2026-06-20', + txn_type: 'Buy', + }, + ], + }); + + expect(rows).toHaveLength(1); + expect(rows[0]).toEqual( + expect.objectContaining({ + provider: 'unusual_whales', + chamber: 'senate', + filedDate: '2026-06-29', + filerName: 'Jane Smith', + providerPublishedAt: null, + }), + ); + expect( + matchDisclosureCandidate( + { doc_id: 'S-hidden', source_url: null, filed_date: '2026-06-29', filer_name: 'Smith, Jane' }, + rows[0], + ), + ).toEqual({ providerKey: rows[0].providerKey, matchMethod: 'filer-date' }); + }); + + it('captures Quiver upload timestamps separately from monitor observation time', () => { + const rows = parseQuiverDisclosureRows('house', [ + { + Representative: 'Jane Smith', + ReportDate: '2026-06-29T00:00:00Z', + Date: '2026-06-20T00:00:00Z', + Ticker: 'MSFT', + Transaction: 'Purchase', + Quiver_Upload_Time: '2026-06-29T14:05:00Z', + }, + ]); + + expect(rows).toHaveLength(1); + expect(rows[0]).toEqual( + expect.objectContaining({ + provider: 'quiver', + chamber: 'house', + filedDate: '2026-06-29', + filerName: 'Jane Smith', + providerPublishedAt: '2026-06-29T14:05:00.000Z', + }), + ); }); }); diff --git a/app/src/ingestion/fmpDisclosureLatency.ts b/app/src/ingestion/fmpDisclosureLatency.ts index 26f705bbf..a721a9c5e 100644 --- a/app/src/ingestion/fmpDisclosureLatency.ts +++ b/app/src/ingestion/fmpDisclosureLatency.ts @@ -2,10 +2,10 @@ * src/ingestion/fmpDisclosureLatency.ts * OWNER: ingestion * - * Measures whether Congress.Trade sees new congressional disclosures before - * FMP's congressional "latest" endpoints. It deliberately records observations - * rather than making the comparison by hand, because FMP may already have a row - * by the time we start looking unless we poll both sides on the same cadence. + * Provider-latency monitor for congressional disclosures. Candidates are + * created when Congress.Trade first sees a new filing; provider observations + * are populated from third-party "latest" endpoints so admins can measure who + * surfaced the disclosure first. */ import type { Env } from '../shared/types'; @@ -16,16 +16,25 @@ import { assertFmpTierOk } from '../shared/fmpStatus'; import type { DiscoveredFiling } from './watcher'; type Chamber = 'house' | 'senate'; +type ProviderId = 'fmp' | 'unusual_whales' | 'quiver' | 'finnhub' | 'ainvest' | 'capitol_trades'; type EnvWithWatch = Env & { + DISCLOSURE_LATENCY_WATCH_ENABLED?: string; + DISCLOSURE_LATENCY_PROVIDERS?: string; + DISCLOSURE_LATENCY_WATCH_LIMIT?: string; FMP_API_KEY?: string; FMP_DISCLOSURE_WATCH_ENABLED?: string; FMP_DISCLOSURE_WATCH_LIMIT?: string; + UNUSUAL_WHALES_API_KEY?: string; + QUIVER_API_KEY?: string; + QUIVER_API_TOKEN?: string; + FINNHUB_API_KEY?: string; + AINVEST_API_KEY?: string; }; interface CandidateRow { doc_id: string; - provider: string; + provider: ProviderId; chamber: Chamber; source_url: string | null; filed_date: string | null; @@ -35,30 +44,53 @@ interface CandidateRow { } interface ProviderObservationRow { - provider: string; + provider: ProviderId; chamber: Chamber; provider_key: string; first_observed_at: string; + provider_published_at: string | null; source_url: string | null; filed_date: string | null; filer_name: string | null; payload: string | null; } -export interface FmpDisclosureRow { +export interface DisclosureProviderRow { + provider: ProviderId; chamber: Chamber; providerKey: string; payload: Record; sourceUrl: string | null; filedDate: string | null; filerName: string | null; + providerPublishedAt: string | null; } +export type FmpDisclosureRow = DisclosureProviderRow; + export interface CandidateMatch { providerKey: string; matchMethod: string; } +export interface DisclosureLatencyProviderStatus { + id: ProviderId; + label: string; + configured: boolean; + requiresMembership: boolean; + supportsDirectLatest: boolean; + timestampKind: 'provider' | 'monitor' | 'none'; + reason?: string; +} + +export interface DisclosureLatencyProviderRun extends DisclosureLatencyProviderStatus { + enabled: boolean; + fetchedRows: number; + pending: number; + matched: number; + errors: string[]; +} + export interface DisclosureLatencyProbeResult { enabled: boolean; reason?: string; @@ -66,28 +98,154 @@ export interface DisclosureLatencyProbeResult { pending: number; matched: number; errors: string[]; + providers: DisclosureLatencyProviderRun[]; +} + +export interface DisclosureLatencyProviderMetrics { + provider: ProviderId; + label: string; + candidates: number; + matched: number; + pending: number; + errored: number; + coveragePct: number; + ctAheadMonitorCount: number; + providerAheadMonitorCount: number; + tieMonitorCount: number; + avgMonitorDeltaSec: number | null; + medianMonitorDeltaSec: number | null; + p90MonitorDeltaSec: number | null; + avgProviderPublishedDeltaSec: number | null; + medianProviderPublishedDeltaSec: number | null; +} + +export interface DisclosureLatencyTotals { + candidates: number; + matched: number; + pending: number; + errored: number; + comparableProviders: number; + configuredComparableProviders: number; +} + +export interface DisclosureLatencySummary { + generatedAt: string; + totals: DisclosureLatencyTotals; + providers: DisclosureLatencyProviderMetrics[]; + providerStatuses: DisclosureLatencyProviderStatus[]; + publicSummary: { + generatedAt: string; + totals: DisclosureLatencyTotals; + providers: DisclosureLatencyProviderMetrics[]; + }; +} + +interface ProviderDefinition { + id: ProviderId; + label: string; + secretNames: string[]; + requiresMembership: boolean; + supportsDirectLatest: boolean; + timestampKind: 'provider' | 'monitor' | 'none'; + reason?: string; + fetchRows?: (apiKey: string, max: number, fetchImpl: typeof fetch) => Promise; } -const PROVIDER = 'fmp'; const DEFAULT_LIMIT = 100; const RECENT_PROVIDER_HOURS = 72; const PAYLOAD_LIMIT = 20_000; +const DIRECT_PROVIDER_IDS: ProviderId[] = ['fmp', 'unusual_whales', 'quiver']; + +const PROVIDERS: ProviderDefinition[] = [ + { + id: 'fmp', + label: 'Financial Modeling Prep', + secretNames: ['FMP_API_KEY'], + requiresMembership: true, + supportsDirectLatest: true, + timestampKind: 'monitor', + reason: 'FMP exposes disclosure/transaction dates, but not a provider first-seen timestamp.', + fetchRows: fetchFmpRows, + }, + { + id: 'unusual_whales', + label: 'Unusual Whales', + secretNames: ['UNUSUAL_WHALES_API_KEY'], + requiresMembership: true, + supportsDirectLatest: true, + timestampKind: 'monitor', + reason: 'Recent Congress trades exposes filed_at_date, but not a provider first-seen timestamp.', + fetchRows: fetchUnusualWhalesRows, + }, + { + id: 'quiver', + label: 'Quiver Quantitative', + secretNames: ['QUIVER_API_KEY', 'QUIVER_API_TOKEN'], + requiresMembership: true, + supportsDirectLatest: true, + timestampKind: 'provider', + reason: 'Quiver V2 rows may include Quiver_Upload_Time; otherwise the monitor first-observed time is used.', + fetchRows: fetchQuiverRows, + }, + { + id: 'finnhub', + label: 'Finnhub', + secretNames: ['FINNHUB_API_KEY'], + requiresMembership: true, + supportsDirectLatest: false, + timestampKind: 'none', + reason: 'Finnhub congressional trading is symbol/date-range scoped, not a global latest-disclosure feed.', + }, + { + id: 'ainvest', + label: 'AInvest', + secretNames: ['AINVEST_API_KEY'], + requiresMembership: true, + supportsDirectLatest: false, + timestampKind: 'none', + reason: 'AInvest congressional trades require a ticker parameter, so they cannot race all new disclosures directly.', + }, + { + id: 'capitol_trades', + label: 'Capitol Trades', + secretNames: [], + requiresMembership: false, + supportsDirectLatest: false, + timestampKind: 'none', + reason: 'No official API found; the public site is protected by a browser checkpoint, so this remains manual/unsupported.', + }, +]; function truthy(v: string | undefined): boolean { return /^(1|true|yes|on)$/i.test((v ?? '').trim()); } function enabled(env: EnvWithWatch): boolean { - return truthy(env.FMP_DISCLOSURE_WATCH_ENABLED); + return truthy(env.DISCLOSURE_LATENCY_WATCH_ENABLED) || truthy(env.FMP_DISCLOSURE_WATCH_ENABLED); } function limit(env: EnvWithWatch): number { - const n = parseInt(env.FMP_DISCLOSURE_WATCH_LIMIT || '', 10); + const n = parseInt(env.DISCLOSURE_LATENCY_WATCH_LIMIT || env.FMP_DISCLOSURE_WATCH_LIMIT || '', 10); return Number.isFinite(n) && n > 0 ? Math.min(n, 500) : DEFAULT_LIMIT; } function storageMissing(err: unknown): boolean { - return /no such table|no column named/i.test(err instanceof Error ? err.message : String(err)); + return /no such table|no column named|no such column/i.test(err instanceof Error ? err.message : String(err)); +} + +function definition(id: ProviderId): ProviderDefinition { + return PROVIDERS.find((p) => p.id === id) ?? PROVIDERS[0]; +} + +function requestedProviderIds(env: EnvWithWatch, opts: { providers?: string[] } = {}): ProviderId[] { + const raw = opts.providers?.length ? opts.providers.join(',') : env.DISCLOSURE_LATENCY_PROVIDERS || ''; + const parsed = raw + .split(/[,\s]+/) + .map((part) => part.trim().toLowerCase()) + .filter(Boolean) as ProviderId[]; + const allowed = new Set(PROVIDERS.map((p) => p.id)); + const ids = parsed.filter((id) => allowed.has(id)); + return ids.length ? Array.from(new Set(ids)) : [...DIRECT_PROVIDER_IDS]; } function normalizeDate(raw: string | null | undefined): string | null { @@ -100,6 +258,13 @@ function normalizeDate(raw: string | null | undefined): string | null { return s.slice(0, 10); } +function normalizeTimestamp(raw: string | null | undefined): string | null { + const s = (raw ?? '').trim(); + if (!s || /^\d{4}-\d{2}-\d{2}$/.test(s)) return null; + const t = Date.parse(s); + return Number.isFinite(t) ? new Date(t).toISOString() : null; +} + function dateVariants(iso: string | null): string[] { if (!iso) return []; const m = /^(\d{4})-(\d{2})-(\d{2})$/.exec(iso); @@ -145,6 +310,7 @@ function fieldString(row: Record, names: string[]): string | nu for (const name of names) { const v = lower.get(name.toLowerCase()); if (typeof v === 'string' && v.trim()) return v.trim(); + if (typeof v === 'number' && Number.isFinite(v)) return String(v); } return null; } @@ -181,7 +347,6 @@ function tokensFromDoc(docId: string, sourceUrl: string | null): string[] { for (const part of docLower.split(/[^a-z0-9]+/)) { if (part.length >= 6) out.add(part); } - // House ids are H-YYYY-DOCID; the trailing doc id is what FMP usually exposes. const house = /^h-\d{4}-(.+)$/i.exec(docId); if (house && house[1].length >= 6) out.add(house[1].toLowerCase()); const senate = /^s-(.+)$/i.exec(docId); @@ -198,20 +363,11 @@ function lastName(name: string | null): string | null { return last && last.length >= 4 ? last.toLowerCase() : null; } -export function matchFmpDisclosureCandidate( - candidate: Pick, - row: FmpDisclosureRow, -): CandidateMatch | null { - const text = rowText(row.payload); - for (const token of tokensFromDoc(candidate.doc_id, candidate.source_url)) { - if (text.includes(token)) return { providerKey: row.providerKey, matchMethod: 'doc-token' }; - } - const filed = normalizeDate(candidate.filed_date); - const last = lastName(candidate.filer_name); - if (filed && last && text.includes(last) && dateVariants(filed).some((d) => text.includes(d))) { - return { providerKey: row.providerKey, matchMethod: 'probable-filer-date' }; - } - return null; +function normalizeChamber(raw: string | null, fallback: Chamber): Chamber { + const s = (raw ?? '').toLowerCase(); + if (s.includes('senate') || s.includes('senator')) return 'senate'; + if (s.includes('house') || s.includes('representative') || s.includes('representatives')) return 'house'; + return fallback; } function extractRows(json: unknown): Record[] { @@ -223,11 +379,21 @@ function extractRows(json: unknown): Record[] { if (Array.isArray(value)) { return value.filter((v): v is Record => !!v && typeof v === 'object' && !Array.isArray(v)); } + if (value && typeof value === 'object' && !Array.isArray(value)) { + const nested = extractRows(value); + if (nested.length) return nested; + } } } return []; } +function rowKeyFromFields(provider: ProviderId, payload: Record, fields: string[]): string { + const parts = fields.map((field) => fieldString(payload, [field]) ?? '').filter(Boolean); + const basis = parts.length ? parts.join('|') : rowText(payload); + return `${provider}:${simpleHash(basis)}`; +} + export function parseFmpDisclosureRows(chamber: Chamber, json: unknown): FmpDisclosureRow[] { return extractRows(json).map((payload) => { const sourceUrl = firstUrl(payload); @@ -235,34 +401,163 @@ export function parseFmpDisclosureRows(chamber: Chamber, json: unknown): FmpDisc const docToken = providerKeyFromUrl(sourceUrl) ?? fieldString(payload, ['docId', 'documentId', 'reportId', 'disclosureId', 'disclosure_id']); const providerKey = docToken ? String(docToken).toLowerCase() : simpleHash(text); return { + provider: 'fmp', chamber, providerKey, payload, sourceUrl, filedDate: normalizeDate(fieldString(payload, ['filedDate', 'filingDate', 'disclosureDate', 'reportedDate'])), filerName: fieldString(payload, ['representative', 'senator', 'filerName', 'name']), + providerPublishedAt: null, }; }); } -async function fetchFmpLatest( - apiKey: string, - chamber: Chamber, - max: number, - fetchImpl: typeof fetch, -): Promise { - const url = - `https://financialmodelingprep.com/stable/${chamber}-latest?page=0&limit=${max}` + - '&apikey=' + - encodeURIComponent(apiKey); +export function parseUnusualWhalesDisclosureRows(json: unknown): DisclosureProviderRow[] { + return extractRows(json).map((payload) => { + const filedDate = normalizeDate(fieldString(payload, ['filed_at_date', 'filingDate', 'filedDate'])); + const filerName = fieldString(payload, ['name', 'reporter']); + return { + provider: 'unusual_whales', + chamber: normalizeChamber(fieldString(payload, ['member_type', 'chamber']), 'house'), + providerKey: rowKeyFromFields('unusual_whales', payload, [ + 'politician_id', + 'filed_at_date', + 'ticker', + 'transaction_date', + 'txn_type', + 'name', + ]), + payload, + sourceUrl: firstUrl(payload), + filedDate, + filerName, + providerPublishedAt: null, + }; + }); +} + +export function parseQuiverDisclosureRows(chamber: Chamber, json: unknown): DisclosureProviderRow[] { + return extractRows(json).map((payload) => { + const filedDate = normalizeDate(fieldString(payload, ['Filed', 'ReportDate', 'report_date', 'filed_date'])); + const filerName = fieldString(payload, ['Representative', 'Senator', 'Name', 'representative', 'senator', 'name']); + return { + provider: 'quiver', + chamber: normalizeChamber(fieldString(payload, ['Chamber', 'House', 'house']), chamber), + providerKey: rowKeyFromFields('quiver', payload, [ + 'BioGuideID', + 'Representative', + 'Senator', + 'Name', + 'Filed', + 'ReportDate', + 'Ticker', + 'TransactionDate', + 'Date', + 'Traded', + 'Transaction', + ]), + payload, + sourceUrl: firstUrl(payload), + filedDate, + filerName, + providerPublishedAt: normalizeTimestamp(fieldString(payload, ['Quiver_Upload_Time'])), + }; + }); +} + +export function matchDisclosureCandidate( + candidate: Pick, + row: DisclosureProviderRow, +): CandidateMatch | null { + const text = rowText(row.payload); + for (const token of tokensFromDoc(candidate.doc_id, candidate.source_url)) { + if (text.includes(token)) return { providerKey: row.providerKey, matchMethod: 'doc-token' }; + } + const filed = normalizeDate(candidate.filed_date); + const candidateLast = lastName(candidate.filer_name); + const rowLast = lastName(row.filerName); + if (filed && candidateLast && rowLast === candidateLast && row.filedDate === filed) { + return { providerKey: row.providerKey, matchMethod: 'filer-date' }; + } + if (filed && candidateLast && text.includes(candidateLast) && dateVariants(filed).some((d) => text.includes(d))) { + return { providerKey: row.providerKey, matchMethod: 'probable-filer-date' }; + } + return null; +} + +export function matchFmpDisclosureCandidate( + candidate: Pick, + row: FmpDisclosureRow, +): CandidateMatch | null { + return matchDisclosureCandidate(candidate, row); +} + +async function fetchJson(url: string, headers: Record, fetchImpl: typeof fetch): Promise { const res = await fetchImpl(url, { - headers: { 'user-agent': 'congress.trade/0.1 (+https://congress.trade)', accept: 'application/json' }, + headers: { 'user-agent': 'congress.trade/0.1 (+https://congress.trade)', accept: 'application/json', ...headers }, }); - if (!res.ok) { - assertFmpTierOk(res.status); - throw new Error(`FMP_${chamber}_LATEST_HTTP_${res.status}`); + if (!res.ok) throw new Error(`HTTP_${res.status}:${url.replace(/[?&](apikey|token)=[^&]+/gi, '$1=[redacted]')}`); + return res.json(); +} + +async function fetchFmpRows(apiKey: string, max: number, fetchImpl: typeof fetch): Promise { + const fetchOne = async (chamber: Chamber) => { + const url = + `https://financialmodelingprep.com/stable/${chamber}-latest?page=0&limit=${max}` + + '&apikey=' + + encodeURIComponent(apiKey); + try { + return parseFmpDisclosureRows(chamber, await fetchJson(url, {}, fetchImpl)); + } catch (err) { + const status = /HTTP_(\d+)/.exec((err as Error).message)?.[1]; + if (status) assertFmpTierOk(Number(status)); + throw err; + } + }; + return (await Promise.all([fetchOne('house'), fetchOne('senate')])).flat(); +} + +async function fetchUnusualWhalesRows(apiKey: string, max: number, fetchImpl: typeof fetch): Promise { + const url = `https://api.unusualwhales.com/api/congress/recent-trades?limit=${Math.min(max, 200)}`; + return parseUnusualWhalesDisclosureRows(await fetchJson(url, { authorization: `Bearer ${apiKey}` }, fetchImpl)); +} + +async function fetchQuiverRows(apiKey: string, _max: number, fetchImpl: typeof fetch): Promise { + const headers = { authorization: `Bearer ${apiKey}` }; + const [house, senate] = await Promise.all([ + fetchJson('https://api.quiverquant.com/beta/live/housetrading?options=true', headers, fetchImpl), + fetchJson('https://api.quiverquant.com/beta/live/senatetrading?options=true', headers, fetchImpl), + ]); + return [...parseQuiverDisclosureRows('house', house), ...parseQuiverDisclosureRows('senate', senate)]; +} + +async function resolveProviderSecret(env: Env, provider: ProviderDefinition): Promise { + for (const name of provider.secretNames) { + const envx = env as unknown as Record; + const value = (await resolveSecret(env, name as keyof Env & string)).value ?? envx[name]; + if (value?.trim()) return value.trim(); } - return parseFmpDisclosureRows(chamber, await res.json()); + return null; +} + +async function providerStatus(env: Env, provider: ProviderDefinition): Promise { + const configured = provider.secretNames.length === 0 || Boolean(await resolveProviderSecret(env, provider)); + return { + id: provider.id, + label: provider.label, + configured, + requiresMembership: provider.requiresMembership, + supportsDirectLatest: provider.supportsDirectLatest, + timestampKind: provider.timestampKind, + reason: provider.reason, + }; +} + +export async function getDisclosureLatencyProviderStatuses(env: Env): Promise { + const statuses: DisclosureLatencyProviderStatus[] = []; + for (const provider of PROVIDERS) statuses.push(await providerStatus(env, provider)); + return statuses; } export async function recordDisclosureLatencyCandidate( @@ -270,51 +565,55 @@ export async function recordDisclosureLatencyCandidate( filing: DiscoveredFiling, nowIso: string, ): Promise { - try { - await run( - env.DB, - `INSERT INTO disclosure_latency_candidates - (doc_id, provider, chamber, source_url, filed_date, filer_name, - congress_first_seen_at, status, attempts, created_at, updated_at) - VALUES (?, ?, ?, ?, ?, ?, ?, 'pending', 0, ?, ?) - ON CONFLICT(doc_id, provider) DO NOTHING`, - [ - filing.docId, - PROVIDER, - filing.chamber, - filing.sourceUrl, - normalizeDate(filing.filedDate), - filing.filerName ?? null, - nowIso, - nowIso, - nowIso, - ], - ); - } catch (err) { - if (!storageMissing(err)) console.warn('disclosure latency candidate write failed:', (err as Error).message); + for (const provider of DIRECT_PROVIDER_IDS) { + try { + await run( + env.DB, + `INSERT INTO disclosure_latency_candidates + (doc_id, provider, chamber, source_url, filed_date, filer_name, + congress_first_seen_at, status, attempts, created_at, updated_at) + VALUES (?, ?, ?, ?, ?, ?, ?, 'pending', 0, ?, ?) + ON CONFLICT(doc_id, provider) DO NOTHING`, + [ + filing.docId, + provider, + filing.chamber, + filing.sourceUrl, + normalizeDate(filing.filedDate), + filing.filerName ?? null, + nowIso, + nowIso, + nowIso, + ], + ); + } catch (err) { + if (!storageMissing(err)) console.warn('disclosure latency candidate write failed:', (err as Error).message); + } } } -async function upsertProviderRows(env: Env, rows: FmpDisclosureRow[], nowIso: string): Promise { +async function upsertProviderRows(env: Env, provider: ProviderId, rows: DisclosureProviderRow[], nowIso: string): Promise { for (const row of rows) { await run( env.DB, `INSERT INTO disclosure_provider_observations (provider, chamber, provider_key, first_observed_at, last_observed_at, - source_url, filed_date, filer_name, payload) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) + provider_published_at, source_url, filed_date, filer_name, payload) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT(provider, chamber, provider_key) DO UPDATE SET last_observed_at=excluded.last_observed_at, + provider_published_at=COALESCE(disclosure_provider_observations.provider_published_at, excluded.provider_published_at), source_url=COALESCE(disclosure_provider_observations.source_url, excluded.source_url), filed_date=COALESCE(disclosure_provider_observations.filed_date, excluded.filed_date), filer_name=COALESCE(disclosure_provider_observations.filer_name, excluded.filer_name), payload=COALESCE(disclosure_provider_observations.payload, excluded.payload)`, [ - PROVIDER, + provider, row.chamber, row.providerKey, nowIso, nowIso, + row.providerPublishedAt, row.sourceUrl, row.filedDate, row.filerName, @@ -324,90 +623,74 @@ async function upsertProviderRows(env: Env, rows: FmpDisclosureRow[], nowIso: st } } -async function alertMatch(env: Env, candidate: CandidateRow, match: ProviderObservationRow): Promise { - const congress = Date.parse(candidate.congress_first_seen_at); - const provider = Date.parse(match.first_observed_at); - const deltaSec = Number.isFinite(congress) && Number.isFinite(provider) ? Math.round((provider - congress) / 1000) : null; +function deltaSeconds(later: string | null, earlier: string | null): number | null { + if (!later || !earlier) return null; + const a = Date.parse(later); + const b = Date.parse(earlier); + return Number.isFinite(a) && Number.isFinite(b) ? Math.round((a - b) / 1000) : null; +} + +async function alertMatch(env: Env, provider: ProviderDefinition, candidate: CandidateRow, match: ProviderObservationRow): Promise { + const deltaSec = deltaSeconds(match.first_observed_at, candidate.congress_first_seen_at); const direction = deltaSec == null ? 'Delta unavailable' : deltaSec > 0 - ? `Congress.Trade observed it ${deltaSec}s before FMP was first observed by this monitor.` + ? `Congress.Trade observed it ${deltaSec}s before ${provider.label} was first observed by this monitor.` : deltaSec < 0 - ? `FMP was already observed ${Math.abs(deltaSec)}s before Congress.Trade first saw it.` - : 'Congress.Trade and FMP were observed in the same second.'; + ? `${provider.label} was already observed ${Math.abs(deltaSec)}s before Congress.Trade first saw it.` + : `Congress.Trade and ${provider.label} were observed in the same second.`; + const published = + match.provider_published_at && deltaSeconds(match.provider_published_at, candidate.congress_first_seen_at) != null + ? `\n${provider.label} provider timestamp: ${match.provider_published_at}` + : ''; await notifyAdmin(env, { - dedupeKey: `disclosure-latency:${PROVIDER}:${candidate.doc_id}`, + dedupeKey: `disclosure-latency:${provider.id}:${candidate.doc_id}`, throttleSec: 30 * 24 * 60 * 60, - subject: 'Congress.Trade vs FMP disclosure latency', + subject: `Congress.Trade vs ${provider.label} disclosure latency`, text: `${direction}\n\n` + `Doc: ${candidate.doc_id}\n` + `Chamber: ${candidate.chamber}\n` + `Congress.Trade first_seen_at: ${candidate.congress_first_seen_at}\n` + - `FMP monitor first_observed_at: ${match.first_observed_at}\n` + - `FMP key: ${match.provider_key}\n` + + `${provider.label} monitor first_observed_at: ${match.first_observed_at}${published}\n` + + `${provider.label} key: ${match.provider_key}\n` + `Source URL: ${candidate.source_url ?? 'n/a'}\n`, }); } -export async function runFmpDisclosureLatencyProbe( - env: Env, - now: Date = new Date(), - fetchImpl: typeof fetch = fetch, - opts: { force?: boolean } = {}, -): Promise { - const envx = env as EnvWithWatch; - if (!opts.force && !enabled(envx)) { - return { enabled: false, reason: 'FMP_DISCLOSURE_WATCH_ENABLED is not true', fetchedRows: 0, pending: 0, matched: 0, errors: [] }; - } - const apiKey = (await resolveSecret(env, 'FMP_API_KEY')).value ?? envx.FMP_API_KEY; - if (!apiKey) return { enabled: false, reason: 'FMP_API_KEY missing', fetchedRows: 0, pending: 0, matched: 0, errors: [] }; - - const nowIso = now.toISOString(); - const max = limit(envx); - const errors: string[] = []; - let fetchedRows = 0; - try { - const rows = ( - await Promise.all([ - fetchFmpLatest(apiKey, 'house', max, fetchImpl), - fetchFmpLatest(apiKey, 'senate', max, fetchImpl), - ]) - ).flat(); - fetchedRows = rows.length; - await upsertProviderRows(env, rows, nowIso); - } catch (err) { - errors.push((err as Error).message); - } - - let candidates: CandidateRow[] = []; - try { - candidates = await all( - env.DB, - `SELECT doc_id, provider, chamber, source_url, filed_date, filer_name, - congress_first_seen_at, attempts - FROM disclosure_latency_candidates - WHERE provider = ? AND status = 'pending' - ORDER BY created_at DESC - LIMIT 100`, - [PROVIDER], - ); - } catch (err) { - if (storageMissing(err)) return { enabled: true, reason: 'latency tables missing; run /api/admin/migrate', fetchedRows, pending: 0, matched: 0, errors }; - throw err; - } - +async function loadProviderRows(env: Env, provider: ProviderId, now: Date): Promise { const cutoff = new Date(now.getTime() - RECENT_PROVIDER_HOURS * 60 * 60 * 1000).toISOString(); - const providerRows = await all( + return all( env.DB, - `SELECT provider, chamber, provider_key, first_observed_at, source_url, filed_date, filer_name, payload + `SELECT provider, chamber, provider_key, first_observed_at, provider_published_at, + source_url, filed_date, filer_name, payload FROM disclosure_provider_observations WHERE provider = ? AND first_observed_at >= ? ORDER BY first_observed_at DESC LIMIT 1000`, - [PROVIDER, cutoff], + [provider, cutoff], ); +} + +async function matchPendingCandidates( + env: Env, + provider: ProviderDefinition, + now: Date, + nowIso: string, + errors: string[], +): Promise<{ pending: number; matched: number }> { + const candidates = await all( + env.DB, + `SELECT doc_id, provider, chamber, source_url, filed_date, filer_name, + congress_first_seen_at, attempts + FROM disclosure_latency_candidates + WHERE provider = ? AND status = 'pending' + ORDER BY created_at DESC + LIMIT 100`, + [provider.id], + ); + const providerRows = await loadProviderRows(env, provider.id, now); let matched = 0; for (const candidate of candidates) { @@ -416,15 +699,17 @@ export async function runFmpDisclosureLatencyProbe( for (const providerRow of providerRows) { if (providerRow.chamber !== candidate.chamber || !providerRow.payload) continue; const payload = JSON.parse(providerRow.payload) as Record; - const parsed: FmpDisclosureRow = { + const parsed: DisclosureProviderRow = { + provider: providerRow.provider, chamber: providerRow.chamber, providerKey: providerRow.provider_key, payload, sourceUrl: providerRow.source_url, filedDate: providerRow.filed_date, filerName: providerRow.filer_name, + providerPublishedAt: providerRow.provider_published_at, }; - const m = matchFmpDisclosureCandidate(candidate, parsed); + const m = matchDisclosureCandidate(candidate, parsed); if (m) { match = providerRow; method = m.matchMethod; @@ -439,6 +724,7 @@ export async function runFmpDisclosureLatencyProbe( SET status = 'matched', provider_key = ?, provider_first_seen_at = ?, + provider_published_at = ?, match_method = ?, payload = ?, attempts = attempts + 1, @@ -449,26 +735,188 @@ export async function runFmpDisclosureLatencyProbe( [ match.provider_key, match.first_observed_at, + match.provider_published_at, method, match.payload, nowIso, nowIso, candidate.doc_id, - PROVIDER, + provider.id, ], ); matched++; - await alertMatch(env, candidate, match); + await alertMatch(env, provider, candidate, match); } else { await run( env.DB, `UPDATE disclosure_latency_candidates SET attempts = attempts + 1, last_checked_at = ?, updated_at = ?, error = ? WHERE doc_id = ? AND provider = ?`, - [nowIso, nowIso, errors[0] ?? null, candidate.doc_id, PROVIDER], + [nowIso, nowIso, errors[0] ?? null, candidate.doc_id, provider.id], ); } } + return { pending: candidates.length, matched }; +} + +async function runProviderProbe( + env: Env, + provider: ProviderDefinition, + now: Date, + fetchImpl: typeof fetch, + max: number, +): Promise { + const base = await providerStatus(env, provider); + const errors: string[] = []; + if (!provider.supportsDirectLatest || !provider.fetchRows) { + return { ...base, enabled: false, fetchedRows: 0, pending: 0, matched: 0, errors }; + } + const apiKey = await resolveProviderSecret(env, provider); + if (!apiKey) { + return { ...base, configured: false, enabled: false, fetchedRows: 0, pending: 0, matched: 0, errors, reason: `${provider.secretNames[0]} missing` }; + } + + const nowIso = now.toISOString(); + let fetchedRows = 0; + try { + const rows = await provider.fetchRows(apiKey, max, fetchImpl); + fetchedRows = rows.length; + await upsertProviderRows(env, provider.id, rows, nowIso); + } catch (err) { + errors.push((err as Error).message); + } + + try { + const matched = await matchPendingCandidates(env, provider, now, nowIso, errors); + return { ...base, configured: true, enabled: true, fetchedRows, pending: matched.pending, matched: matched.matched, errors }; + } catch (err) { + if (storageMissing(err)) { + return { ...base, configured: true, enabled: true, fetchedRows, pending: 0, matched: 0, errors, reason: 'latency tables missing; run /api/admin/migrate' }; + } + throw err; + } +} - return { enabled: true, fetchedRows, pending: candidates.length, matched, errors }; +export async function runDisclosureLatencyProbe( + env: Env, + now: Date = new Date(), + fetchImpl: typeof fetch = fetch, + opts: { force?: boolean; providers?: string[] } = {}, +): Promise { + const envx = env as EnvWithWatch; + if (!opts.force && !enabled(envx)) { + return { + enabled: false, + reason: 'DISCLOSURE_LATENCY_WATCH_ENABLED is not true', + fetchedRows: 0, + pending: 0, + matched: 0, + errors: [], + providers: [], + }; + } + + const runs: DisclosureLatencyProviderRun[] = []; + for (const providerId of requestedProviderIds(envx, opts)) { + runs.push(await runProviderProbe(env, definition(providerId), now, fetchImpl, limit(envx))); + } + return { + enabled: true, + fetchedRows: runs.reduce((sum, r) => sum + r.fetchedRows, 0), + pending: runs.reduce((sum, r) => sum + r.pending, 0), + matched: runs.reduce((sum, r) => sum + r.matched, 0), + errors: runs.flatMap((r) => r.errors.map((err) => `${r.id}: ${err}`)), + providers: runs, + }; +} + +export async function runFmpDisclosureLatencyProbe( + env: Env, + now: Date = new Date(), + fetchImpl: typeof fetch = fetch, + opts: { force?: boolean } = {}, +): Promise { + return runDisclosureLatencyProbe(env, now, fetchImpl, { ...opts, providers: ['fmp'] }); +} + +function median(values: number[]): number | null { + if (!values.length) return null; + const sorted = [...values].sort((a, b) => a - b); + const mid = Math.floor(sorted.length / 2); + return sorted.length % 2 === 0 ? Math.round((sorted[mid - 1] + sorted[mid]) / 2) : sorted[mid]; +} + +function average(values: number[]): number | null { + if (!values.length) return null; + return Math.round(values.reduce((sum, v) => sum + v, 0) / values.length); +} + +function p90(values: number[]): number | null { + if (!values.length) return null; + const sorted = [...values].sort((a, b) => a - b); + return sorted[Math.min(sorted.length - 1, Math.ceil(sorted.length * 0.9) - 1)]; +} + +export async function getDisclosureLatencySummary(env: Env, now: Date = new Date()): Promise { + const rows = await all<{ + provider: ProviderId; + status: string; + congress_first_seen_at: string; + provider_first_seen_at: string | null; + provider_published_at: string | null; + }>( + env.DB, + `SELECT provider, status, congress_first_seen_at, provider_first_seen_at, provider_published_at + FROM disclosure_latency_candidates + ORDER BY created_at DESC + LIMIT 5000`, + ).catch((err) => { + if (storageMissing(err)) return []; + throw err; + }); + + const statuses = await getDisclosureLatencyProviderStatuses(env); + const providers = PROVIDERS.filter((p) => p.supportsDirectLatest).map((provider) => { + const mine = rows.filter((row) => row.provider === provider.id); + const monitorDeltas = mine + .map((row) => deltaSeconds(row.provider_first_seen_at, row.congress_first_seen_at)) + .filter((v): v is number => v != null); + const publishedDeltas = mine + .map((row) => deltaSeconds(row.provider_published_at, row.congress_first_seen_at)) + .filter((v): v is number => v != null); + const matched = mine.filter((row) => row.status === 'matched').length; + return { + provider: provider.id, + label: provider.label, + candidates: mine.length, + matched, + pending: mine.filter((row) => row.status === 'pending').length, + errored: mine.filter((row) => row.status === 'error').length, + coveragePct: mine.length ? Math.round((matched / mine.length) * 1000) / 10 : 0, + ctAheadMonitorCount: monitorDeltas.filter((d) => d > 0).length, + providerAheadMonitorCount: monitorDeltas.filter((d) => d < 0).length, + tieMonitorCount: monitorDeltas.filter((d) => d === 0).length, + avgMonitorDeltaSec: average(monitorDeltas), + medianMonitorDeltaSec: median(monitorDeltas), + p90MonitorDeltaSec: p90(monitorDeltas), + avgProviderPublishedDeltaSec: average(publishedDeltas), + medianProviderPublishedDeltaSec: median(publishedDeltas), + }; + }); + const totals = { + candidates: rows.length, + matched: rows.filter((row) => row.status === 'matched').length, + pending: rows.filter((row) => row.status === 'pending').length, + errored: rows.filter((row) => row.status === 'error').length, + comparableProviders: PROVIDERS.filter((p) => p.supportsDirectLatest).length, + configuredComparableProviders: statuses.filter((p) => p.supportsDirectLatest && p.configured).length, + }; + const generatedAt = now.toISOString(); + return { + generatedAt, + totals, + providers, + providerStatuses: statuses, + publicSummary: { generatedAt, totals, providers }, + }; } diff --git a/app/src/shared/types.ts b/app/src/shared/types.ts index 1acd95c4a..29f106e8f 100644 --- a/app/src/shared/types.ts +++ b/app/src/shared/types.ts @@ -457,10 +457,23 @@ export interface Env { FMP_API_KEY?: string; /** Daily FMP call budget (stringified int); defaults to 230 when unset. */ FMP_DAILY_CALL_CAP?: string; - /** Enables the Congress.Trade-vs-FMP congressional disclosure latency monitor. */ + /** Enables the Congress.Trade-vs-provider congressional disclosure latency monitor. */ + DISCLOSURE_LATENCY_WATCH_ENABLED?: string; + /** Comma-separated provider ids to race: fmp, unusual_whales, quiver. Defaults to direct comparable providers. */ + DISCLOSURE_LATENCY_PROVIDERS?: string; + /** Latest rows to fetch per provider/chamber endpoint when the latency monitor runs. */ + DISCLOSURE_LATENCY_WATCH_LIMIT?: string; + /** Enables the legacy Congress.Trade-vs-FMP monitor switch; kept for backward compatibility. */ FMP_DISCLOSURE_WATCH_ENABLED?: string; - /** Latest rows to fetch per FMP chamber endpoint when the latency monitor runs. */ + /** Legacy FMP-specific latest-row limit; DISCLOSURE_LATENCY_WATCH_LIMIT takes precedence. */ FMP_DISCLOSURE_WATCH_LIMIT?: string; + /** Unusual Whales API key for recent Congress trades. */ + UNUSUAL_WHALES_API_KEY?: string; + /** Quiver API bearer token for live Congress trading endpoints. */ + QUIVER_API_KEY?: string; + QUIVER_API_TOKEN?: string; + /** AInvest key; currently reported as symbol-scoped and not directly comparable for this monitor. */ + AINVEST_API_KEY?: string; /** Which price provider to prefer: 'fmp' or 'massive'. */ PRICE_PROVIDER?: string; /** HMAC key for signing outbound webhook payloads. */ diff --git a/clients/ios/CongressTrade.xcodeproj/project.xcworkspace/contents.xcworkspacedata b/clients/ios/CongressTrade.xcodeproj/project.xcworkspace/contents.xcworkspacedata deleted file mode 100644 index 919434a62..000000000 --- a/clients/ios/CongressTrade.xcodeproj/project.xcworkspace/contents.xcworkspacedata +++ /dev/null @@ -1,7 +0,0 @@ - - - - -