From 2eb6e6ca1273f8d8dc16071a23ea258557844fcd Mon Sep 17 00:00:00 2001
From: Jay Wedgeworth <12656028+jaywedgeworth22@users.noreply.github.com>
Date: Tue, 30 Jun 2026 15:24:45 -0500
Subject: [PATCH] feat: expand disclosure latency providers
---
.../contents.xcworkspacedata | 7 -
app/.dev.vars.example | 12 +-
.../0021_disclosure_latency_watch.sql | 6 +-
.../0023_disclosure_provider_timestamps.sql | 6 +
.../admin/__tests__/disclosureLatency.test.ts | 26 +-
app/src/admin/routes.ts | 42 +-
app/src/index.ts | 6 +-
.../__tests__/fmpDisclosureLatency.test.ts | 69 +-
app/src/ingestion/fmpDisclosureLatency.ts | 716 ++++++++++++++----
app/src/shared/types.ts | 17 +-
.../contents.xcworkspacedata | 7 -
11 files changed, 744 insertions(+), 170 deletions(-)
delete mode 100644 Congress.Trade.xcworkspace/contents.xcworkspacedata
create mode 100644 app/migrations/0023_disclosure_provider_timestamps.sql
delete mode 100644 clients/ios/CongressTrade.xcodeproj/project.xcworkspace/contents.xcworkspacedata
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 @@
-
-
-
-
-