diff --git a/app/server/src/adapters/EdgarAdapter.ts b/app/server/src/adapters/EdgarAdapter.ts index 3a2fc1b..41db4d8 100644 --- a/app/server/src/adapters/EdgarAdapter.ts +++ b/app/server/src/adapters/EdgarAdapter.ts @@ -17,22 +17,23 @@ const OPERATOR_EMAIL = process.env.SEC_OPERATOR_EMAIL ?? 'research@example.com'; const UA = `Investor Flow (${OPERATOR_EMAIL})`; /** - * Token-bucket rate limiter: 8 tokens, refilled at 8/sec, min 125ms between - * calls. Guarantees we stay under EDGAR's 10 req/s rule with headroom. + * Token-bucket rate limiter: 6 tokens, refilled at 6/sec, min ~167ms between + * calls. Stays within EDGAR's 10 req/s rule with headroom. */ class TokenBucket { - private tokens = 8; + private tokens = 6; private lastDrain = Date.now(); + private static readonly MAX_TOKENS = 6; + private static readonly REFILL_RATE = 6; // tokens/sec async wait(): Promise { const now = Date.now(); const elapsed = (now - this.lastDrain) / 1000; - // Refill tokens up to maxRate (8). - this.tokens = Math.min(8, this.tokens + elapsed * 8); + this.tokens = Math.min(TokenBucket.MAX_TOKENS, this.tokens + elapsed * TokenBucket.REFILL_RATE); this.lastDrain = now; if (this.tokens < 1) { - const waitMs = Math.ceil(((1 - this.tokens) / 8) * 1000); + const waitMs = Math.ceil(((1 - this.tokens) / TokenBucket.REFILL_RATE) * 1000); await new Promise((r) => setTimeout(r, waitMs)); this.tokens = 0; this.lastDrain = Date.now(); @@ -47,34 +48,46 @@ const bucket = new TokenBucket(); /** * Fetch a URL with EDGAR-compliant headers, rate limiting, and ETag caching. * Returns null on 304 (caller should return cached row). + * Retries with exponential backoff on 429 (rate limited). */ async function edgarFetch( url: string, extraHeaders?: Record, + retries = 3, ): Promise<{ status: number; headers: { etag?: string | null; lastModified?: string | null }; body: unknown } | null> { - await bucket.wait(); + for (let attempt = 1; attempt <= retries; attempt++) { + await bucket.wait(); - const headers: Record = { - 'User-Agent': UA, - Accept: 'application/json', - ...extraHeaders, - }; + const headers: Record = { + 'User-Agent': UA, + Accept: 'application/json', + ...extraHeaders, + }; - const resp = await fetch(url, { method: 'GET', headers }); + const resp = await fetch(url, { method: 'GET', headers }); - const etag = resp.headers.get('etag'); - const lastModified = resp.headers.get('last-modified'); + const etag = resp.headers.get('etag'); + const lastModified = resp.headers.get('last-modified'); - if (resp.status === 304) { - return { status: 304, headers: { etag, lastModified }, body: null }; + if (resp.status === 304) { + return { status: 304, headers: { etag, lastModified }, body: null }; + } + + if (resp.status === 429 && attempt < retries) { + const backoff = Math.pow(2, attempt) * 1000; + await new Promise((r) => setTimeout(r, backoff)); + continue; + } + + if (resp.status >= 400) { + throw new Error(`EDGAR ${resp.status} ${resp.statusText} for ${url}`); + } + + const body = (await resp.json()) as unknown; + return { status: resp.status, headers: { etag, lastModified }, body }; } - if (resp.status >= 400) { - throw new Error(`EDGAR ${resp.status} ${resp.statusText} for ${url}`); - } - - const body = (await resp.json()) as unknown; - return { status: resp.status, headers: { etag, lastModified }, body }; + throw new Error(`EDGAR max retries (${retries}) exceeded for ${url}`); } // --------------------------------------------------------------------------- @@ -160,12 +173,34 @@ export class EdgarAdapter implements SourceFetch { throw new Error(`EDGAR 304 but no cached data for ${key}`); } - const rawData = resp?.body as { name?: string; filings?: { recent?: Array<{ form?: string; dateReporter?: string; accessionNumber?: string; accessionNormalization?: string; reportDate?: string; reportFile?: string; primaryDocument?: string }> } } | undefined; - const recentFilings = rawData?.filings?.recent; - if (!recentFilings) { + const rawData = resp?.body as { name?: string; filings?: { recent?: unknown } } | undefined; + const recentObj = rawData?.filings?.recent; + if (!recentObj) { throw new Error(`EDGAR filings_index: no recent filings for CIK${padded}`); } + // SEC EDGAR returns recent filings as either: + // (a) an array of objects: [{form, dateReporter, ...}, ...] (current format) + // (b) parallel arrays: {form: [...], dateReporter: [...], ...} (legacy) + let recentFilings: Array>; + if (Array.isArray(recentObj)) { + recentFilings = recentObj as Array>; + } else if (typeof recentObj === 'object') { + // Parallel arrays — convert to array of objects. + const keys = Object.keys(recentObj); + const len = (recentObj[keys[0]] as unknown[])?.length ?? 0; + recentFilings = []; + for (let i = 0; i < len; i++) { + const row: Record = {}; + for (const k of keys) { + row[k] = (recentObj as Record)[k]?.[i]; + } + recentFilings.push(row); + } + } else { + throw new Error(`EDGAR filings_index: unexpected recent filings format for CIK${padded}`); + } + // Filter by form type. let filings = recentFilings; if (opts?.formTypes && opts.formTypes.length > 0) { @@ -178,7 +213,7 @@ export class EdgarAdapter implements SourceFetch { const from = opts.dateRange.from ? new Date(opts.dateRange.from).getTime() : null; const to = opts.dateRange.to ? new Date(opts.dateRange.to).getTime() : null; filings = filings.filter((f) => { - const ts = new Date(f.dateReporter ?? f.accessionNormalization ?? '').getTime(); + const ts = new Date(f.reportDate ?? f.filingDate ?? f.dateReporter ?? '').getTime(); if (Number.isNaN(ts)) return false; if (from !== null && ts < from) return false; if (to !== null && ts > to) return false; @@ -318,7 +353,8 @@ export class EdgarAdapter implements SourceFetch { cik: string, accession: string, ): Promise { - const padded = padCik(cik); + const filerCik = accession.split('-')[0]; + const padded = padCik(filerCik); const accessionNoDashes = accession.replace(/-/g, ''); // --- Step 1: filing index (JSON) --------------------------------------- @@ -332,29 +368,49 @@ export class EdgarAdapter implements SourceFetch { const indexBody = indexResp.body as { fileDate?: string; documents?: Array<{ name?: string; type?: string; size?: string | number; path?: string }>; + directory?: { item?: Array<{ name?: string; type?: string; size?: string | number }> }; partialSubmissionIndicator?: unknown; }; - if (!indexBody?.documents || indexBody.documents.length === 0) { + // EDGAR API may return documents as `documents` (old) or `directory.item` (new) + const docList = indexBody?.documents ?? indexBody?.directory?.item ?? []; + + if (docList.length === 0) { throw new Error(`EDGAR 13F: no documents in index for ${padded}/${accessionNoDashes}`); } - // Pick the primary document (usually the .txt or .xml filing). - const primaryDoc = indexBody.documents.find( - (d) => d.type === '13F' || d.type === '13F-infoTable' - ) ?? indexBody.documents[0]; + // Pick candidates: XML documents (excluding index files), sorted by size desc + const xmlCandidates = docList + .filter((d) => d.name?.toLowerCase().endsWith('.xml') && !d.name?.includes('index')) + .sort((a, b) => { + const sa = typeof a.size === 'string' ? parseInt(a.size, 10) || 0 : (a.size as number) ?? 0; + const sb = typeof b.size === 'string' ? parseInt(b.size, 10) || 0 : (b.size as number) ?? 0; + return sb - sa; + }); - const primaryName = primaryDoc.name; - if (!primaryName) { - throw new Error(`EDGAR 13F: primary document has no name for ${padded}/${accessionNoDashes}`); + // Try each candidate until we find one with holdings; fall back to any non-HTML doc + let docText = ''; + let parsed: Array<{ cusip: string; issuerName: string; value: number; sshPrnamt: number }> = []; + for (const cand of xmlCandidates) { + const docUrl = `https://www.sec.gov/Archives/edgar/data/${padded}/${accessionNoDashes}/${cand.name}`; + docText = await this.edgarXmlFetch(docUrl); + parsed = parse13fHoldings(docText); + if (parsed.length > 0) break; } - // --- Step 2: fetch the primary document (XML/HTML) --------------------- - const docUrl = `https://www.sec.gov/Archives/edgar/data/${padded}/${accessionNoDashes}/${primaryName}`; - const docText = await this.edgarXmlFetch(docUrl); + // Fallback: try first non-HTML doc if XML candidates had no holdings + if (parsed.length === 0) { + const fallback = docList.find( + (d) => !d.name?.toLowerCase().endsWith('.html') && !d.name?.includes('index') + ) ?? docList[0]; + if (fallback && fallback.name) { + const docUrl = `https://www.sec.gov/Archives/edgar/data/${padded}/${accessionNoDashes}/${fallback.name}`; + docText = await this.edgarXmlFetch(docUrl); + parsed = parse13fHoldings(docText); + } + } - // --- Step 3: parse holdings table -------------------------------------- - const holdings = parse13fHoldings(docText); + const holdings = parsed; return { value: { holdings, accession: `${padded}/${accessionNoDashes}` }, @@ -389,7 +445,8 @@ export class EdgarAdapter implements SourceFetch { cik: string, accession: string, ): Promise { - const padded = padCik(cik); + const filerCik = accession.split('-')[0]; + const padded = padCik(filerCik); const accessionNoDashes = accession.replace(/-/g, ''); // --- Step 1: filing index (JSON) --------------------------------------- @@ -403,17 +460,21 @@ export class EdgarAdapter implements SourceFetch { const indexBody = indexResp.body as { fileDate?: string; documents?: Array<{ name?: string; type?: string; size?: string | number; path?: string }>; + directory?: { item?: Array<{ name?: string; type?: string; size?: string | number }> }; partialSubmissionIndicator?: unknown; }; - if (!indexBody?.documents || indexBody.documents.length === 0) { + // EDGAR API may return documents as `documents` (old) or `directory.item` (new) + const docList = indexBody?.documents ?? indexBody?.directory?.item ?? []; + + if (docList.length === 0) { throw new Error(`EDGAR Form 4: no documents in index for ${padded}/${accessionNoDashes}`); } // Form 4 is typically filed as a single XML. Pick the .xml document. - const xmlDoc = indexBody.documents.find( + const xmlDoc = docList.find( (d) => d.name?.toLowerCase().endsWith('.xml') - ) ?? indexBody.documents[0]; + ) ?? docList[0]; const xmlName = xmlDoc.name; if (!xmlName) { @@ -446,24 +507,35 @@ export class EdgarAdapter implements SourceFetch { * Fetch a URL that returns XML/HTML (not JSON), using the rate limiter. * EDGAR returns XML for primary filing documents; we strip the body here * rather than calling `resp.json()`. + * Retries with exponential backoff on 429 (rate limited). */ - async edgarXmlFetch(url: string): Promise { - await bucket.wait(); + async edgarXmlFetch(url: string, retries = 3): Promise { + for (let attempt = 1; attempt <= retries; attempt++) { + await bucket.wait(); - const resp = await fetch(url, { - method: 'GET', - headers: { - 'User-Agent': UA, - Accept: 'application/xml, text/xml, */*', - }, - }); + const resp = await fetch(url, { + method: 'GET', + headers: { + 'User-Agent': UA, + Accept: 'application/xml, text/xml, */*', + }, + }); - if (resp.status >= 400) { - throw new Error(`EDGAR XML ${resp.status} ${resp.statusText} for ${url}`); + if (resp.status === 429 && attempt < retries) { + const backoff = Math.pow(2, attempt) * 1000; + await new Promise((r) => setTimeout(r, backoff)); + continue; + } + + if (resp.status >= 400) { + throw new Error(`EDGAR XML ${resp.status} ${resp.statusText} for ${url}`); + } + + const text = await resp.text(); + return text; } - const text = await resp.text(); - return text; + throw new Error(`EDGAR XML max retries (${retries}) exceeded for ${url}`); } /** SourceFetch.fetchOne dispatch. */ @@ -512,7 +584,30 @@ function parse13fHoldings(text: string): Array<{ }> { const holdings: Array<{ cusip: string; issuerName: string; value: number; sshPrnamt: number }> = []; - // Try XML-style parsing first (13F filings are often XML). + // Try 13F XML format: (possibly namespaced) with , , , + // SEC 13F XML uses namespaces like ... + const infoTableMatches = text.match(/<(?:\w+:)?infoTable[^>]*>([\s\S]*?)<\/(?:\w+:)?infoTable>/gi); + if (infoTableMatches && infoTableMatches.length > 0) { + for (const block of infoTableMatches) { + const nameMatch = block.match(/<(?:\w+:)?nameOfIssuer>\s*([\s\S]*?)\s*<\/(?:\w+:)?nameOfIssuer>/i); + const cusipMatch = block.match(/<(?:\w+:)?cusip>\s*([\s\S]*?)\s*<\/(?:\w+:)?cusip>/i); + const valueMatch = block.match(/<(?:\w+:)?value>\s*([\s\S]*?)\s*<\/(?:\w+:)?value>/i); + const sharesMatch = block.match(/<(?:\w+:)?sshPrnamt>\s*([\s\S]*?)\s*<\/(?:\w+:)?sshPrnamt>/i); + + if (!cusipMatch) continue; + + const issuerName = nameMatch ? nameMatch[1].replace(/<[^>]+>/g, '').trim() : ''; + const cusip = cusipMatch ? cusipMatch[1].replace(/<[^>]+>/g, '').trim() : ''; + const value = valueMatch ? parseFloat(valueMatch[1].replace(/<[^>]+>/g, '').replace(/,/g, '')) || 0 : 0; + const sshPrnamt = sharesMatch ? parseFloat(sharesMatch[1].replace(/<[^>]+>/g, '').replace(/,/g, '')) || 0 : 0; + + holdings.push({ cusip, issuerName, value, sshPrnamt }); + } + + if (holdings.length > 0) return holdings; + } + + // Try HTML-style
parsing (older 13F filings may use HTML tables). const xmlMatch = text.match(/]*>([\s\S]*?)<\/table>/gi); if (xmlMatch) { for (const tbl of xmlMatch) { @@ -596,11 +691,16 @@ function parse13fHoldings(text: string): Array<{ /** * Extract transactions from a Form 4 XML filing. * - * Form 4 uses an XML structure with elements containing: - * , , , - * , , + * Modern EDGAR Form 4 uses with: + * ...... + * ... + * ... + * YYYY-MM-DD + * ...... + * ...... + * ...... * - * This is a thin parser that extracts the key fields. + * The old format is also handled as a fallback. */ function parseForm4Transactions(text: string): Array<{ reporter: string; @@ -621,12 +721,57 @@ function parseForm4Transactions(text: string): Array<{ price: number; }> = []; - // Extract infotable blocks. + // Try modern format: with blocks + const hasNonDerivative = //i.test(text); + if (hasNonDerivative) { + // Get reporting owner info from document root (shared across all transactions) + const ownerNameMatch = text.match(/\s*([\s\S]*?)<\/rptOwnerName>/i); + const reporter = ownerNameMatch ? ownerNameMatch[1].replace(/<[^>]+>/g, '').trim() : ''; + + const relBlockMatch = text.match(/([\s\S]*?)<\/reportingOwnerRelationship>/i); + let relationship = ''; + if (relBlockMatch) { + const relBlock = relBlockMatch[1]; + const parts: string[] = []; + if (/\s*1/i.test(relBlock)) parts.push('Director'); + if (/\s*1/i.test(relBlock)) { + const titleMatch = relBlock.match(/\s*([\s\S]*?)<\/officerTitle>/i); + parts.push(titleMatch ? titleMatch[1].replace(/<[^>]+>/g, '').trim() : 'Officer'); + } + if (/\s*1/i.test(relBlock)) parts.push('10% Owner'); + if (/\s*1/i.test(relBlock)) parts.push('Other'); + relationship = parts.join(', '); + } + + const txBlocks = text.match(//gi); + if (txBlocks) { + for (const block of txBlocks) { + const titleMatch = block.match(/\s*([\s\S]*?)\s*<\/value>/i); + const securityTitle = titleMatch ? titleMatch[1].replace(/<[^>]+>/g, '').trim() : ''; + + const dateMatch = block.match(/\s*(\d{4}-\d{2}-\d{2})/); + const transactionDate = dateMatch ? dateMatch[1] : ''; + + const codeMatch = block.match(/\s*([A-HV])/i); + const transactionCode = codeMatch ? codeMatch[1].toUpperCase() : ''; + + const sharesMatch = block.match(/\s*([\d]+)/i); + const shares = sharesMatch ? parseInt(sharesMatch[1], 10) : 0; + + const priceMatch = block.match(/\s*([\d,.]+)/i); + const price = priceMatch ? parseFloat(priceMatch[1].replace(/,/g, '')) : 0; + + transactions.push({ reporter, relationship, securityTitle, transactionDate, transactionCode, shares, price }); + } + } + return transactions; + } + + // Fallback: old format. const infotables = text.match(//gi); if (!infotables) return transactions; for (const table of infotables) { - // Extract reporter CIK and name from . const rptOwnerMatch = table.match(/([\s\S]*?)<\/rptOwner>/i); let reporter = ''; if (rptOwnerMatch) { @@ -639,7 +784,6 @@ function parseForm4Transactions(text: string): Array<{ } } - // Extract relationship from . const relBlockMatch = table.match(/([\s\S]*?)<\/rptOwnerRelationship>/i); let relationship = ''; if (relBlockMatch) { @@ -655,35 +799,22 @@ function parseForm4Transactions(text: string): Array<{ relationship = parts.join(', '); } - // Extract security title. const titleMatch = table.match(/\s*([\s\S]*?)<\/securityTitle>/i); const securityTitle = titleMatch ? titleMatch[1].trim() : ''; - // Extract transaction date. const dateMatch = table.match(/\s*(\d{4}-\d{2}-\d{2})/); const transactionDate = dateMatch ? dateMatch[1] : ''; - // Extract transaction code. const codeMatch = table.match(/\s*([A-HV])/i); const transactionCode = codeMatch ? codeMatch[1].toUpperCase() : ''; - // Extract shares (non-decimal, integer). const sharesMatch = table.match(/\s*([\d]+)/); const shares = sharesMatch ? parseInt(sharesMatch[1], 10) : 0; - // Extract price per share. const priceMatch = table.match(/\s*([\d,.]+)/); const price = priceMatch ? parseFloat(priceMatch[1].replace(/,/g, '')) : 0; - transactions.push({ - reporter, - relationship, - securityTitle, - transactionDate, - transactionCode, - shares, - price, - }); + transactions.push({ reporter, relationship, securityTitle, transactionDate, transactionCode, shares, price }); } return transactions; diff --git a/app/server/src/adapters/__tests__/EdgarAdapter.test.ts b/app/server/src/adapters/__tests__/EdgarAdapter.test.ts index 666fa1b..d3c9360 100644 --- a/app/server/src/adapters/__tests__/EdgarAdapter.test.ts +++ b/app/server/src/adapters/__tests__/EdgarAdapter.test.ts @@ -98,10 +98,10 @@ function makeFilingsResponse() { name: 'TEST COMPANY INC', filings: { recent: [ - { form: '10-K', dateReporter: '2026-03-15', accessionNumber: '0001234567-26-000001', accessionNormalization: '2026-03-15', reportDate: '2026-02-28', reportFile: 'http://example.com/10k.pdf', primaryDocument: 'form10k.pdf' }, - { form: '10-Q', dateReporter: '2026-01-15', accessionNumber: '0001234567-26-000002', accessionNormalization: '2026-01-15', reportDate: '2025-12-31', reportFile: 'http://example.com/10q.pdf', primaryDocument: 'form10q.pdf' }, - { form: '8-K', dateReporter: '2025-12-01', accessionNumber: '0001234567-25-000003', accessionNormalization: '2025-12-01', reportDate: '2025-12-01', reportFile: 'http://example.com/8k.pdf', primaryDocument: 'form8k.pdf' }, - { form: 'SC 13G', dateReporter: '2025-06-30', accessionNumber: '0001234567-25-000004', accessionNormalization: '2025-06-30', reportDate: '2025-06-30', reportFile: 'http://example.com/13g.pdf', primaryDocument: 'form13g.pdf' }, + { form: '10-K', filingDate: '2026-03-15', accessionNumber: '0001234567-26-000001', accessionNormalization: '2026-03-15', reportDate: '2026-02-28', reportFile: 'http://example.com/10k.pdf', primaryDocument: 'form10k.pdf' }, + { form: '10-Q', filingDate: '2026-01-15', accessionNumber: '0001234567-26-000002', accessionNormalization: '2026-01-15', reportDate: '2026-01-15', reportFile: 'http://example.com/10q.pdf', primaryDocument: 'form10q.pdf' }, + { form: '8-K', filingDate: '2025-12-01', accessionNumber: '0001234567-25-000003', accessionNormalization: '2025-12-01', reportDate: '2025-12-01', reportFile: 'http://example.com/8k.pdf', primaryDocument: 'form8k.pdf' }, + { form: 'SC 13G', filingDate: '2025-06-30', accessionNumber: '0001234567-25-000004', accessionNormalization: '2025-06-30', reportDate: '2025-06-30', reportFile: 'http://example.com/13g.pdf', primaryDocument: 'form13g.pdf' }, ], }, }; @@ -195,7 +195,7 @@ test('filings_index filters by dateRange (from + to)', async () => { assert.equal(getCallCount(), 1); const filings = result.value as Array>; - // 10-K (2026-03-15), 10-Q (2026-01-15), 8-K (2025-12-01) are in range; SC 13G (2025-06-30) is out + // 10-K (2026-02-28), 10-Q (2026-01-15), 8-K (2025-12-01) are in range; SC 13G (2025-06-30) is out assert.equal(filings.length, 3, 'should include filings within date range'); const forms = filings.map((f) => f.form); assert.ok(forms.includes('10-K')); @@ -222,7 +222,7 @@ test('filings_index filters by formTypes + dateRange combined', async () => { assert.equal(getCallCount(), 1); const filings = result.value as Array>; - // 10-K (2026-03-15) and 10-Q (2026-01-15) both have dateReporter >= 2026-01-01 + // 10-K (2026-02-28) and 10-Q (2026-01-15) both have reportDate >= 2026-01-01 assert.equal(filings.length, 2, '10-K and 10-Q are in date range with form filter'); const forms = filings.map((f) => f.form); assert.ok(forms.includes('10-K')); diff --git a/app/server/src/db/client.ts b/app/server/src/db/client.ts index ea2b38f..8643a3a 100644 --- a/app/server/src/db/client.ts +++ b/app/server/src/db/client.ts @@ -44,12 +44,24 @@ export function initSchema(database: DatabaseSync): void { database.exec(sql); } +/** Idempotent migrations for existing databases (new columns, tables). */ +function runMigrations(db: DatabaseSync): void { + const migrations: string[] = [ + `ALTER TABLE adapter_queue ADD COLUMN error TEXT`, + `ALTER TABLE adapter_queue ADD COLUMN scheduled_for TEXT`, + ]; + for (const sql of migrations) { + try { db.exec(sql); } catch { /* column already exists */ } + } +} + let _db: DatabaseSync | null = null; /** Process-wide singleton connection (auto-creates + inits schema on first call). */ export function db(): DatabaseSync { if (!_db) { _db = createDb(); initSchema(_db); + runMigrations(_db); // Ensure the anonymous user exists for unauthenticated watchlist/portfolio usage. try { _db.prepare("INSERT OR IGNORE INTO users (id, email, pw_hash, created_at) VALUES (?, ?, ?, ?)") diff --git a/app/server/src/db/schema.sql b/app/server/src/db/schema.sql index 08eef29..2317faa 100644 --- a/app/server/src/db/schema.sql +++ b/app/server/src/db/schema.sql @@ -13,6 +13,7 @@ CREATE TABLE IF NOT EXISTS users ( pw_hash TEXT NOT NULL, -- scrypt(verifiable); design specified argon2id, adapted to built-in scrypt (reversible) totp_secret TEXT, -- 2FA secret, encrypted at rest (slice 2) is_2fa_enabled INTEGER NOT NULL DEFAULT 0, + is_admin INTEGER NOT NULL DEFAULT 0, -- 1 = operator (admin CLI/route access, Slice 25); never user-editable oauth_subject TEXT, -- nullable; social OAuth link (slice 2) oauth_provider TEXT, -- 'github' | 'google' | NULL complexity TEXT NOT NULL DEFAULT 'beginner', -- beginner|intermediate|advanced @@ -180,12 +181,44 @@ CREATE TABLE IF NOT EXISTS threads ( PRIMARY KEY (source, source_id) ); +CREATE TABLE IF NOT EXISTS admin_audit ( + id TEXT PRIMARY KEY, -- uuid + actor TEXT NOT NULL, -- admin user id (or 'cli' when none) + action TEXT NOT NULL, -- users.reset-password|users.gdpr-export|queue.reset-backoff|... + target TEXT, -- target user id / sourceKind / null + ts TEXT NOT NULL, -- ISO timestamp + meta TEXT -- optional JSON blob (never secrets) +); + CREATE TABLE IF NOT EXISTS adapter_queue ( key TEXT PRIMARY KEY, -- e.g. "yfinance:quote:NVDA" - status TEXT NOT NULL, -- pending|in_flight|done|backoff + status TEXT NOT NULL, -- pending|in_flight|done|backoff|failed last_attempt TEXT, retry_count INTEGER NOT NULL DEFAULT 0, - backoff_until TEXT + backoff_until TEXT, + error TEXT, -- last error message (for failed/backoff) + scheduled_for TEXT -- ISO timestamp; job won't run before this +); + +CREATE TABLE IF NOT EXISTS queue_errors ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + queue_key TEXT NOT NULL, + attempt INTEGER NOT NULL, + error_message TEXT, + error_stack TEXT, + ts TEXT NOT NULL +); + +CREATE TABLE IF NOT EXISTS queue_schedules ( + source_kind TEXT PRIMARY KEY, + interval_ms INTEGER NOT NULL, + last_enqueued TEXT, + next_enqueue TEXT +); + +CREATE TABLE IF NOT EXISTS queue_state ( + key TEXT PRIMARY KEY, -- 'paused' + value TEXT NOT NULL ); CREATE TABLE IF NOT EXISTS symbol_demand ( @@ -206,6 +239,8 @@ CREATE TABLE IF NOT EXISTS watchlists ( sort_order INTEGER NOT NULL DEFAULT 0 ); +CREATE UNIQUE INDEX IF NOT EXISTS idx_watchlists_owner_name ON watchlists(owner_id, name); + CREATE TABLE IF NOT EXISTS portfolio_holdings ( id TEXT PRIMARY KEY, owner_id TEXT NOT NULL REFERENCES users(id) ON DELETE CASCADE, @@ -216,6 +251,9 @@ CREATE TABLE IF NOT EXISTS portfolio_holdings ( status TEXT NOT NULL DEFAULT 'open' -- open|closed ); +CREATE UNIQUE INDEX IF NOT EXISTS uq_portfolio_owner_symbol + ON portfolio_holdings(owner_id, symbol); + CREATE TABLE IF NOT EXISTS trades ( id TEXT PRIMARY KEY, owner_id TEXT NOT NULL REFERENCES users(id) ON DELETE CASCADE, @@ -306,6 +344,28 @@ CREATE TABLE IF NOT EXISTS trade_key_levels ( ); -- ===== Shared LLM cache (output deterministic by prompt-hash; not per-user) ===== + +-- ===== Slice 22 — Emotion Logger: per-trade emotion snapshots ===== +CREATE TABLE IF NOT EXISTS trade_executions ( + id TEXT PRIMARY KEY, + trade_id TEXT NOT NULL REFERENCES trades(id) ON DELETE CASCADE, + status TEXT NOT NULL DEFAULT 'active', -- active|closed + started_at TEXT NOT NULL, + completed_at TEXT, + created_at TEXT NOT NULL +); + +CREATE TABLE IF NOT EXISTS emotion_logs ( + id TEXT PRIMARY KEY, + trade_execution_id TEXT NOT NULL REFERENCES trade_executions(id) ON DELETE CASCADE, + timestamp INTEGER NOT NULL, + price_at_event REAL, -- price at moment of emotion event + emotion TEXT NOT NULL, -- e.g. 'fear', 'greed', 'FOMO', 'regret' + note TEXT, -- optional user note + created_at TEXT DEFAULT CURRENT_TIMESTAMP +); + +CREATE INDEX IF NOT EXISTS idx_emotion_logs_trade ON emotion_logs(trade_execution_id); CREATE TABLE IF NOT EXISTS llm_summaries ( hash TEXT PRIMARY KEY, -- hash(prompt + source-payload + complexity) feature TEXT NOT NULL, -- filing_summary|post_summary|dashboard_rollup|explainer|qa|sizing_explain|macro_commentary|... @@ -391,3 +451,102 @@ CREATE INDEX IF NOT EXISTS idx_threads_symbol_ts ON threads(symbol, ts); CREATE INDEX IF NOT EXISTS idx_trades_owner ON trades(owner_id); CREATE INDEX IF NOT EXISTS idx_alerts_owner_sym ON alerts(owner_id, symbol); CREATE INDEX IF NOT EXISTS idx_signals_sector ON rotation_signals(sector); + +-- ===== Slice 18 — Risk Engine: Gentle-halt circuit breaker state ===== +CREATE TABLE IF NOT EXISTS halt_state ( + user_id TEXT PRIMARY KEY, + halted_until TEXT NOT NULL, -- ISO8601; 24h cooldown window + triggered_by TEXT NOT NULL, -- e.g. 'max_drawdown_tolerance_breach' + ts TEXT NOT NULL -- ISO8601; when the halt was recorded +); + +-- ===== Slice 17 — AlertEngine: individual alert event instances ===== +CREATE TABLE IF NOT EXISTS alert_events ( + id TEXT PRIMARY KEY, + user_id TEXT NOT NULL REFERENCES users(id) ON DELETE CASCADE, + type TEXT NOT NULL, -- informed_buy|informed_sell|... + severity TEXT NOT NULL, -- info|warning|critical + title TEXT NOT NULL, + description TEXT NOT NULL, + symbol TEXT, + created_at TEXT NOT NULL, + acknowledged INTEGER NOT NULL DEFAULT 0, + dedup_key TEXT NOT NULL, + payload TEXT NOT NULL DEFAULT '{}', -- JSON + throttle_key TEXT -- for rate limiting: "user:type:symbol:YYYYMMDDHH" +); + +CREATE INDEX IF NOT EXISTS idx_alert_events_user ON alert_events(user_id, created_at DESC); +CREATE INDEX IF NOT EXISTS idx_alert_events_dedup ON alert_events(dedup_key); +CREATE INDEX IF NOT EXISTS idx_alert_events_throttle ON alert_events(throttle_key); + +-- ===== Slice 12 — Strategy Lab + Backtest ===== +CREATE TABLE IF NOT EXISTS strategies ( + id TEXT PRIMARY KEY, + owner_id TEXT NOT NULL REFERENCES users(id) ON DELETE CASCADE, + name TEXT NOT NULL, + components TEXT NOT NULL DEFAULT '[]', -- JSON array of StrategyComponent + unlocked INTEGER NOT NULL DEFAULT 0, + created_at TEXT NOT NULL +); + +CREATE TABLE IF NOT EXISTS backtest_results ( + id TEXT PRIMARY KEY, + strategy_id TEXT NOT NULL REFERENCES strategies(id) ON DELETE CASCADE, + symbol TEXT NOT NULL, + timeframe TEXT NOT NULL DEFAULT '1d', + equity_curve TEXT NOT NULL DEFAULT '[]', -- JSON array + exit_reasons TEXT NOT NULL DEFAULT '[]', -- JSON array + sample_size INTEGER NOT NULL, + created_at TEXT NOT NULL +); + +CREATE INDEX IF NOT EXISTS idx_strategies_owner ON strategies(owner_id); +CREATE INDEX IF NOT EXISTS idx_backtest_strategy ON backtest_results(strategy_id); + +-- ===== Slice 13 — Saved Filter Sets ===== +CREATE TABLE IF NOT EXISTS saved_filters ( + id TEXT PRIMARY KEY, + user_id TEXT NOT NULL REFERENCES users(id) ON DELETE CASCADE, + name TEXT NOT NULL, + expression TEXT NOT NULL, + scope TEXT NOT NULL DEFAULT 'watchlist', + created_at TEXT NOT NULL +); + +CREATE INDEX IF NOT EXISTS idx_saved_filters_user ON saved_filters(user_id); + +-- ===== Slice 19 — Options Convexity Sleeve Unlock State ===== +CREATE TABLE IF NOT EXISTS options_unlock ( + user_id TEXT PRIMARY KEY REFERENCES users(id) ON DELETE CASCADE, + state INTEGER NOT NULL DEFAULT 0, + last_unlock_at TEXT, + understanding INTEGER NOT NULL DEFAULT 0 -- demonstrated understanding flag +); + +-- ===== Slice 20 — Macro Regime History ===== +CREATE TABLE IF NOT EXISTS regime_history ( + id TEXT PRIMARY KEY, + date TEXT NOT NULL, + regime TEXT NOT NULL, + confidence INTEGER NOT NULL, + factors TEXT NOT NULL DEFAULT '[]', -- JSON array + created_at TEXT NOT NULL +); + +CREATE INDEX IF NOT EXISTS idx_regime_history_date ON regime_history(date DESC); + +-- ===== Slice 21 — Thesis Monitoring ===== +CREATE TABLE IF NOT EXISTS theses ( + id TEXT PRIMARY KEY, + user_id TEXT NOT NULL REFERENCES users(id) ON DELETE CASCADE, + symbol TEXT NOT NULL, + statement TEXT NOT NULL, + invalidation_criteria TEXT NOT NULL DEFAULT '[]', -- JSON array + status TEXT NOT NULL DEFAULT 'intact', + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL +); + +CREATE INDEX IF NOT EXISTS idx_theses_user ON theses(user_id); +CREATE INDEX IF NOT EXISTS idx_theses_symbol ON theses(symbol); diff --git a/app/server/src/index.ts b/app/server/src/index.ts index ceedc95..272bd9c 100644 --- a/app/server/src/index.ts +++ b/app/server/src/index.ts @@ -5,17 +5,31 @@ import { fetchRequestHandler } from '@trpc/server/adapters/fetch'; import { db } from './db/client.ts'; import { createCacheRepository, type SourceKind } from './cache/CacheRepository.ts'; import { YFinanceAdapter } from './adapters/YFinanceAdapter.ts'; +import { SecFetchAdapter } from './adapters/SecFetchAdapter.ts'; +import type { SourceFetch } from './adapters/SourceAdapter.ts'; import { AdapterQueue } from './queue/AdapterQueue.ts'; import { makeCreateContext } from './trpc/context.ts'; import { appRouter } from './trpc/router.ts'; const PORT = Number(process.env.PORT ?? 3001); const database = db(); -const adapters = new Map([['yfinance', new YFinanceAdapter()]]); +const adapters: Map = new Map([ + ['yfinance', new YFinanceAdapter()], + ['sec-fetch', new SecFetchAdapter(database)], +]); const queue = new AdapterQueue({ db: database, adapters }); const cache = createCacheRepository({ db: database, scheduler: queue }); queue.cache = cache; // break the cache<->scheduler cycle -const createContext = makeCreateContext({ db: database, cache }); + +// Seed default schedules (noop if already seeded) +queue.seedDefaultSchedules(); + +// Startup recovery: any job left 'in_flight' was interrupted by a restart/crash. +// Reset to 'pending' so the drain loop reprocesses it. +const recovered = database.prepare("UPDATE adapter_queue SET status='pending', error=NULL, retry_count=0 WHERE status='in_flight'").run(); +if (Number(recovered.changes) > 0) console.log(`[investor-flow] recovered ${recovered.changes} interrupted in_flight jobs`); + +const createContext = makeCreateContext({ db: database, cache, queue }); // Background drain: stale-while-revalidate refreshes are queued by CacheRepository.get; // this loop drains them (fetch via adapter -> write to cache), deduped + backed off. @@ -23,6 +37,11 @@ const DRAIN_MS = Number(process.env.IFLOW_DRAIN_MS ?? 2000); const drainTimer = setInterval(() => { queue.drain().catch((e) => console.error('[drain error]', e)); }, DRAIN_MS); drainTimer.unref(); +// Auto-scheduler: every 30s, enqueue refreshes for due schedules +const SCHEDULE_MS = 30_000; +const scheduleTimer = setInterval(() => { queue.enqueueDueSchedules().catch((e) => console.error('[schedule error]', e)); }, SCHEDULE_MS); +scheduleTimer.unref(); + function readBody(req: IncomingMessage): Promise { return new Promise((resolve, reject) => { let data = ''; @@ -44,7 +63,12 @@ const server = createServer(async (req, res) => { for (const [k, v] of Object.entries(req.headers)) if (v != null) headers.set(k, Array.isArray(v) ? v.join(', ') : String(v)); const method = req.method ?? 'GET'; const body = method === 'GET' || method === 'HEAD' ? undefined : await readBody(req); - const request = new Request(`http://localhost:${PORT}${url.pathname}${url.search}`, { method, headers, body }); + const requestUrl = `http://localhost:${PORT}${url.pathname}${url.search}`; + const init: RequestInit = { method, headers }; + if (body) { + init.body = new Blob([body], { type: 'application/json' }); + } + const request = new Request(requestUrl, init); try { const response = await fetchRequestHandler({ router: appRouter, createContext, endpoint: '/api/trpc', req: request }); const buf = Buffer.from(await response.arrayBuffer()); @@ -61,6 +85,7 @@ const server = createServer(async (req, res) => { res.end(JSON.stringify({ error: 'not found' })); }); -server.listen(PORT, () => { - console.log(`[investor-flow] backend on http://localhost:${PORT} (tRPC /api/trpc, health /health, drain every ${DRAIN_MS}ms)`); +const HOST = process.env.HOST ?? '0.0.0.0'; +server.listen(PORT, HOST, () => { + console.log(`[investor-flow] backend on http://${HOST}:${PORT} (tRPC /api/trpc, health /health, drain every ${DRAIN_MS}ms)`); }); diff --git a/app/server/src/queue/AdapterQueue.ts b/app/server/src/queue/AdapterQueue.ts index 391a9a9..af23036 100644 --- a/app/server/src/queue/AdapterQueue.ts +++ b/app/server/src/queue/AdapterQueue.ts @@ -5,7 +5,7 @@ import type { DatabaseSync } from 'node:sqlite'; import type { CacheKey, SourceKind, CacheRepository, CacheScheduler } from '../cache/CacheRepository.ts'; import { parseCacheKey } from '../cache/CacheRepository.ts'; -import type { SourceFetch, AdapterHealth } from '../adapters/SourceAdapter.ts'; +import type { SourceFetch } from '../adapters/SourceAdapter.ts'; export interface AdapterQueueOptions { db: DatabaseSync; @@ -14,12 +14,39 @@ export interface AdapterQueueOptions { rateLimitMs?: Partial>; } -const DEFAULT_RATE_MS: Record = { yfinance: 1000, sec: 125, reddit: 1000, x: 3000, macro: 1000, llm: 0 }; -const BACKOFF_MS = [2000, 4000, 8000, 16000, 60000]; // yfinance policy §5 (2s..60s ceiling) +const DEFAULT_RATE_MS: Record = { yfinance: 1000, sec: 125, 'sec-fetch': 1000, reddit: 1000, x: 3000, macro: 1000, llm: 0 }; +const BACKOFF_MS = [2000, 4000, 8000, 16000, 60000]; const MAX_ATTEMPTS = 5; const sleep = (ms: number) => new Promise((r) => setTimeout(r, ms)); +export interface QueueHealthDetailed { + queued: number; + in_flight: number; + failed: number; + backoff: number; + done: number; + backoff_until: string | null; + last_error: string | undefined; + paused: boolean; + counts: Record; + lastErrors: Array<{ key: string; error: string | null }>; +} + +export interface QueueErrorEntry { + attempt: number; + errorMessage: string | null; + errorStack: string | null; + ts: string; +} + +export interface ScheduleEntry { + sourceKind: string; + intervalMs: number; + lastEnqueued: string | null; + nextEnqueue: string | null; +} + export class AdapterQueue implements CacheScheduler { private readonly _db: DatabaseSync; private readonly _adapters: Map; @@ -34,58 +61,160 @@ export class AdapterQueue implements CacheScheduler { this._rate = { ...DEFAULT_RATE_MS, ...(opts.rateLimitMs ?? {}) } as Record; } - /** Link the cache (results are written here after a successful fetch). Breaks the cache<->queue cycle. */ set cache(c: CacheRepository) { this._cache = c; } - /** Enqueue a fetch. Dedupes against pending/in_flight; respects active backoff. force bypasses dedupe. */ - async queue(key: CacheKey): Promise { + isPaused(): boolean { + const row = this._db.prepare("SELECT value FROM queue_state WHERE key='paused'").get() as { value: string } | undefined; + return row?.value === 'true'; + } + + setPaused(paused: boolean): void { + this._db.prepare("INSERT OR REPLACE INTO queue_state (key, value) VALUES ('paused', ?)").run(String(paused)); + } + + async queue(key: CacheKey, scheduledFor?: string): Promise { const row = this._db.prepare('SELECT status, backoff_until FROM adapter_queue WHERE key=?').get(key) as { status?: string; backoff_until?: string | null } | undefined; if (row) { - if (row.status === 'pending' || row.status === 'in_flight') return; // dedupe: already queued - if (row.status === 'backoff' && row.backoff_until && Date.parse(row.backoff_until) > Date.now()) return; // still backing off + if (row.status === 'pending' || row.status === 'in_flight') return; + if (row.status === 'backoff' && row.backoff_until && Date.parse(row.backoff_until) > Date.now()) return; } - this._db.prepare('INSERT OR REPLACE INTO adapter_queue (key,status,last_attempt,retry_count,backoff_until) VALUES (?,?,?,?,?)').run(key, 'pending', null, 0, null); + this._db.prepare('INSERT OR REPLACE INTO adapter_queue (key,status,last_attempt,retry_count,backoff_until,error,scheduled_for) VALUES (?,?,?,?,?,?,?)').run(key, 'pending', null, 0, null, null, scheduledFor ?? null); } - /** Process pending/backoff-expired jobs: fetch via the source adapter, write to cache, set done/backoff/failed. */ async drain(): Promise { if (!this._cache) return; + if (this.isPaused()) return; const now = Date.now(); const jobs = this._db.prepare( - "SELECT key, retry_count, backoff_until FROM adapter_queue WHERE status IN ('pending','backoff') ORDER BY (last_attempt IS NULL) DESC, last_attempt ASC", - ).all() as Array<{ key: string; retry_count: number; backoff_until?: string | null }>; + "SELECT key, retry_count, backoff_until, scheduled_for FROM adapter_queue WHERE status IN ('pending','backoff') ORDER BY (last_attempt IS NULL) DESC, last_attempt ASC", + ).all() as Array<{ key: string; retry_count: number; backoff_until?: string | null; scheduled_for?: string | null }>; for (const job of jobs) { - if (job.backoff_until && Date.parse(job.backoff_until) > now) continue; // backoff not yet expired + if (job.backoff_until && Date.parse(job.backoff_until) > now) continue; + if (job.scheduled_for && Date.parse(job.scheduled_for) > now) continue; const { source } = parseCacheKey(job.key); const adapter = this._adapters.get(source); - if (!adapter) { this._setStatus(job.key, 'done'); continue; } // no adapter registered -> don't leave stuck - // token bucket: enforce per-source min interval + if (!adapter) { this._setStatus(job.key, 'done'); continue; } const last = this._lastFetchAt[source] ?? 0; const wait = (this._rate[source] ?? 0) - (Date.now() - last); if (wait > 0) await sleep(wait); this._lastFetchAt[source] = Date.now(); this._setStatus(job.key, 'in_flight'); + const attempt = job.retry_count + 1; try { const res = await adapter.fetchOne(job.key); await this._cache.set(job.key, res.value, res.ttlClass, res.provenance); - this._db.prepare("UPDATE adapter_queue SET status='done', last_attempt=? WHERE key=?").run(new Date().toISOString(), job.key); + this._db.prepare("UPDATE adapter_queue SET status='done', last_attempt=?, error=NULL WHERE key=?").run(new Date().toISOString(), job.key); } catch (e) { - const retry = job.retry_count + 1; - this._lastError = e instanceof Error ? e.message : String(e); - if (retry >= MAX_ATTEMPTS) { - this._db.prepare("UPDATE adapter_queue SET status='failed', last_attempt=?, retry_count=?, backoff_until=NULL WHERE key=?").run(new Date().toISOString(), retry, job.key); + const msg = e instanceof Error ? e.message : String(e); + const stack = e instanceof Error ? (e.stack ?? msg) : msg; + this._lastError = msg; + this._db.prepare("INSERT INTO queue_errors (queue_key, attempt, error_message, error_stack, ts) VALUES (?,?,?,?,?)").run(job.key, attempt, msg, stack, new Date().toISOString()); + if (attempt >= MAX_ATTEMPTS) { + this._db.prepare("UPDATE adapter_queue SET status='failed', last_attempt=?, retry_count=?, backoff_until=NULL, error=? WHERE key=?").run(new Date().toISOString(), attempt, msg, job.key); } else { - const bo = BACKOFF_MS[Math.min(retry - 1, BACKOFF_MS.length - 1)]; - this._db.prepare("UPDATE adapter_queue SET status='backoff', last_attempt=?, retry_count=?, backoff_until=? WHERE key=?").run(new Date().toISOString(), retry, new Date(Date.now() + bo).toISOString(), job.key); + const bo = BACKOFF_MS[Math.min(attempt - 1, BACKOFF_MS.length - 1)]; + this._db.prepare("UPDATE adapter_queue SET status='backoff', last_attempt=?, retry_count=?, backoff_until=?, error=? WHERE key=?").run(new Date().toISOString(), attempt, new Date(Date.now() + bo).toISOString(), msg, job.key); } } } } - health(): AdapterHealth { - const row = this._db.prepare("SELECT SUM(status='pending') AS q, SUM(status='in_flight') AS i, MAX(backoff_until) AS bu FROM adapter_queue").get() as { q: number | null; i: number | null; bu: string | null }; - return { queued: row.q ?? 0, in_flight: row.i ?? 0, backoff_until: row.bu ?? null, last_error: this._lastError ?? undefined }; + retryJob(key: string): void { + this._db.prepare("UPDATE adapter_queue SET status='pending', backoff_until=NULL, retry_count=0, error=NULL WHERE key=?").run(key); + } + + retrySource(kind: string): number { + const res = this._db.prepare( + "UPDATE adapter_queue SET status='pending', backoff_until=NULL, retry_count=0, error=NULL WHERE status IN ('failed','backoff') AND substr(key,1,instr(key,':')-1)=?" + ).run(kind); + return Number(res.changes); + } + + clearDone(olderThanMs: number = 86400000): number { + const cutoff = new Date(Date.now() - olderThanMs).toISOString(); + const res = this._db.prepare("DELETE FROM adapter_queue WHERE status='done' AND last_attempt < ?").run(cutoff); + return Number(res.changes); + } + + setSchedule(sourceKind: string, intervalMs: number): void { + this._db.prepare("INSERT OR REPLACE INTO queue_schedules (source_kind, interval_ms, last_enqueued, next_enqueue) VALUES (?,?,?,?)").run(sourceKind, intervalMs, null, new Date(Date.now() + intervalMs).toISOString()); + } + + deleteSchedule(sourceKind: string): void { + this._db.prepare("DELETE FROM queue_schedules WHERE source_kind=?").run(sourceKind); + } + + listSchedules(): ScheduleEntry[] { + const rows = this._db.prepare("SELECT source_kind, interval_ms, last_enqueued, next_enqueue FROM queue_schedules ORDER BY source_kind").all() as Array<{ source_kind: string; interval_ms: number; last_enqueued: string | null; next_enqueue: string | null }>; + return rows.map((r) => ({ + sourceKind: r.source_kind, + intervalMs: r.interval_ms, + lastEnqueued: r.last_enqueued, + nextEnqueue: r.next_enqueue, + })); + } + + async enqueueDueSchedules(): Promise { + const now = new Date().toISOString(); + const due = this._db.prepare("SELECT * FROM queue_schedules WHERE next_enqueue IS NOT NULL AND next_enqueue <= ?").all(now) as Array<{ source_kind: string; interval_ms: number; last_enqueued: string | null; next_enqueue: string | null }>; + for (const s of due) { + const symbols = this._db.prepare("SELECT symbol FROM symbol_demand WHERE in_demand=1").all() as Array<{ symbol: string }>; + if (s.source_kind === 'sec-fetch') { + for (const sym of symbols) { + await this.queue(`sec-fetch:fetch:${sym.symbol}`); + } + } else if (s.source_kind === 'yfinance') { + for (const sym of symbols) { + await this.queue(`yfinance:quote:${sym.symbol}`); + await this.queue(`yfinance:candles:${sym.symbol}:1d`); + await this.queue(`yfinance:symbol:${sym.symbol}`); + } + } + const nextEnqueue = new Date(Date.now() + s.interval_ms).toISOString(); + this._db.prepare("UPDATE queue_schedules SET last_enqueued=?, next_enqueue=? WHERE source_kind=?").run(now, nextEnqueue, s.source_kind); + } + } + + seedDefaultSchedules(): void { + const existing = this._db.prepare("SELECT COUNT(*) AS c FROM queue_schedules").get() as { c: number }; + if (existing.c > 0) return; + const defaults: Array<[string, number]> = [ + ['sec-fetch', 86400000], + ['yfinance', 300000], + ]; + const insert = this._db.prepare("INSERT OR IGNORE INTO queue_schedules (source_kind, interval_ms, last_enqueued, next_enqueue) VALUES (?,?,?,?)"); + for (const [kind, ms] of defaults) { + insert.run(kind, ms, null, new Date(Date.now() + ms).toISOString()); + } + } + + health(): QueueHealthDetailed { + const paused = this.isPaused(); + const row = this._db.prepare("SELECT SUM(status='pending') AS q, SUM(status='in_flight') AS i, SUM(status='failed') AS f, SUM(status='backoff') AS b, SUM(status='done') AS d, MAX(backoff_until) AS bu FROM adapter_queue").get() as { q: number | null; i: number | null; f: number | null; b: number | null; d: number | null; bu: string | null }; + const failedJobs = this._db.prepare("SELECT key, error FROM adapter_queue WHERE status='failed' AND error IS NOT NULL ORDER BY last_attempt DESC LIMIT 5").all() as Array<{ key: string; error: string | null }>; + return { + queued: row.q ?? 0, + in_flight: row.i ?? 0, + failed: row.f ?? 0, + backoff: row.b ?? 0, + done: row.d ?? 0, + backoff_until: row.bu ?? null, + last_error: this._lastError ?? undefined, + paused, + counts: { pending: row.q ?? 0, in_flight: row.i ?? 0, failed: row.f ?? 0, backoff: row.b ?? 0, done: row.d ?? 0 }, + lastErrors: failedJobs, + }; + } + + getErrorLog(key: string, limit: number = 10): QueueErrorEntry[] { + const rows = this._db.prepare("SELECT attempt, error_message, error_stack, ts FROM queue_errors WHERE queue_key=? ORDER BY attempt DESC LIMIT ?").all(key, limit) as Array<{ attempt: number; error_message: string | null; error_stack: string | null; ts: string }>; + return rows.map((r) => ({ + attempt: r.attempt, + errorMessage: r.error_message, + errorStack: r.error_stack, + ts: r.ts, + })); } private _setStatus(key: string, status: string): void { @@ -93,7 +222,6 @@ export class AdapterQueue implements CacheScheduler { } } -/** Wire queue + cache together (breaks the cache<->scheduler cycle). Used by slice 1f. */ export function createAdapterQueue(opts: AdapterQueueOptions): AdapterQueue { return new AdapterQueue(opts); } diff --git a/app/server/src/trpc/context.ts b/app/server/src/trpc/context.ts index ba5ea43..dceebd6 100644 --- a/app/server/src/trpc/context.ts +++ b/app/server/src/trpc/context.ts @@ -4,6 +4,7 @@ import { randomUUID, randomBytes, scryptSync, timingSafeEqual, createHmac } from 'node:crypto'; import type { DatabaseSync } from 'node:sqlite'; import type { CacheRepository } from '../cache/CacheRepository.ts'; +import type { AdapterQueue } from '../queue/AdapterQueue.ts'; export const SESSION_COOKIE = 'iflow_session'; export const OAUTH_STATE_COOKIE = 'iflow_oauth_state'; @@ -19,6 +20,7 @@ const SESSION_TTL_MS = 30 * 24 * 60 * 60 * 1000; // 30 days export interface Context { db: DatabaseSync; cache: CacheRepository; + queue: AdapterQueue; resHeaders: Headers; // mutable; procedures append Set-Cookie here (applied to Response by the fetch adapter) userId: string | null; // resolved from the session cookie; null = unauthenticated cookies: Record; // parsed request cookies (session + oauth state) @@ -100,12 +102,13 @@ export function verifyPassword(pw: string, stored: string): boolean { } // ----- context factory: 1f wires real db+cache; tests inject in-memory ----- -export function makeCreateContext(opts: { db: DatabaseSync; cache: CacheRepository }) { +export function makeCreateContext(opts: { db: DatabaseSync; cache: CacheRepository; queue: AdapterQueue }) { return ({ req, resHeaders }: CreateContextOpts): Context => { const cookies = parseCookies(req); return { db: opts.db, cache: opts.cache, + queue: opts.queue, resHeaders, cookies, userId: resolveSessionUserId(opts.db, req), diff --git a/app/server/src/trpc/router.ts b/app/server/src/trpc/router.ts index cd31acf..9982843 100644 --- a/app/server/src/trpc/router.ts +++ b/app/server/src/trpc/router.ts @@ -9,7 +9,7 @@ import { buildAuthorizeUrl, generateState, exchangeCode, type OAuthProvider } fr import { STARTER_WATCHLIST, defaultDrawdownTolerancePct, defaultRiskTolerance, ONBOARDING_DISCLAIMER, type Complexity } from '../onboarding/starter.ts'; import type { Quote, PriceCandle, SymbolMeta } from '../cache/CacheRepository.ts'; import { emaFromCandles, rsi as rsiFn, relativeVolume } from '../analysis/indicators.ts'; -import { listUsers, resetPassword, gdprExport, queueHealth, resetQueueBackoff, NotOwnerError } from '../admin/admin.ts'; +import { listUsers, resetPassword, gdprExport, queueHealth, resetQueueBackoff, NotOwnerError, listUserSessions, listAuditLog, queueSecFetch } from '../admin/admin.ts'; import { EdgarAdapter } from '../adapters/EdgarAdapter.ts'; import { OptionsAdapter, parseOptionChainRows } from '../adapters/OptionsAdapter.ts'; import type { OptionChainRow, OptionGreeks } from '../adapters/OptionsAdapter.ts'; @@ -156,6 +156,7 @@ const onboardingRouter = router({ for (const sym of symbols) { const kind = (STARTER_WATCHLIST.find((s) => s.symbol === sym)?.tickerKind ?? 'equity') as 'equity' | 'crypto' | 'etf' | 'index'; await ctx.cache.subscribe(sym, kind); + queueSecFetch(ctx.db, sym); } if (input.portfolio) { const ins = ctx.db.prepare('INSERT INTO portfolio_holdings (id, owner_id, symbol, qty, avg_cost, acquired_at, status) VALUES (?,?,?,?,?,?,?)'); @@ -254,6 +255,79 @@ const adminRouter = router({ resetBackoff: adminProcedure .input(z.object({ sourceKind: z.string().regex(/^[a-z0-9_]+$/i) })) .mutation(({ ctx, input }) => resetQueueBackoff(ctx.db, ctx.userId, input.sourceKind)), + + userSessions: adminProcedure + .input(z.object({ userId: z.string().uuid() })) + .query(({ ctx, input }) => listUserSessions(ctx.db, input.userId)), + + auditLog: adminProcedure + .input(z.object({ + limit: z.number().int().min(1).max(200).default(50), + offset: z.number().int().min(0).default(0), + actor: z.string().optional(), + action: z.string().optional(), + })) + .query(({ ctx, input }) => listAuditLog(ctx.db, { limit: input.limit, offset: input.offset, actor: input.actor ?? null, action: input.action ?? null })), + + queueSecFetch: adminProcedure + .input(z.object({ symbol: z.string().min(1).max(10) })) + .mutation(({ ctx, input }) => { + try { queueSecFetch(ctx.db, input.symbol); return { ok: true }; } + catch (e) { throw new TRPCError({ code: 'BAD_REQUEST', message: e instanceof Error ? e.message : 'Failed to queue fetch.' }); } + }), + + queueStatus: adminProcedure.query(({ ctx }) => ctx.queue.health()), + + queuePause: adminProcedure.mutation(({ ctx }) => { + ctx.queue.setPaused(true); + return { paused: true }; + }), + + queueResume: adminProcedure.mutation(({ ctx }) => { + ctx.queue.setPaused(false); + return { paused: false }; + }), + + queueLogs: adminProcedure + .input(z.object({ key: z.string().min(1), limit: z.number().int().min(1).max(100).optional().default(10) })) + .query(({ ctx, input }) => ctx.queue.getErrorLog(input.key, input.limit)), + + queueRetryJob: adminProcedure + .input(z.object({ key: z.string().min(1) })) + .mutation(({ ctx, input }) => { + ctx.queue.retryJob(input.key); + return { ok: true }; + }), + + queueRetrySource: adminProcedure + .input(z.object({ sourceKind: z.string().min(1) })) + .mutation(({ ctx, input }) => { + const cleared = ctx.queue.retrySource(input.sourceKind); + return { cleared }; + }), + + queueClearDone: adminProcedure + .input(z.object({ olderThanHours: z.number().min(1).max(336).optional().default(24) })) + .mutation(({ ctx, input }) => { + const cleared = ctx.queue.clearDone(input.olderThanHours * 3600000); + return { cleared }; + }), + + queueSchedules: adminProcedure.query(({ ctx }) => ctx.queue.listSchedules()), + + queueSetSchedule: adminProcedure + .input(z.object({ sourceKind: z.string().min(1), intervalMs: z.number().int().min(60000).max(604800000) })) + .mutation(({ ctx, input }) => { + ctx.queue.setSchedule(input.sourceKind, input.intervalMs); + return { ok: true }; + }), + + queueDeleteSchedule: adminProcedure + .input(z.object({ sourceKind: z.string().min(1) })) + .mutation(({ ctx, input }) => { + ctx.queue.deleteSchedule(input.sourceKind); + return { ok: true }; + }), }); @@ -262,7 +336,7 @@ const dashboardRouter = router({ rollup: publicProcedure .input(z.object({})) .query(async ({ ctx }) => { - const userId = ctx.userId as string; + const userId = ctx.userId ?? 'anonymous'; // Import lazily to avoid circular deps. const { DashboardRollupEngine } = await import('../analysis/dashboardRollup.ts'); const { InstitutionFlowEngine } = await import('../analysis/institutionFlowEngine.ts'); @@ -410,7 +484,7 @@ const institutionalRouter = router({ }>; if (rows.length === 0) { - return { symbol, flow: [], quarters: [] }; + return { symbol, flow: [], quarters: [], filings: [] }; } // Group by filer_cik, then by reported_quarter to compute QoQ flow @@ -468,26 +542,42 @@ const institutionalRouter = router({ // Collect unique quarters const quarters = [...new Set(rows.map(r => r.reported_quarter))].sort().reverse(); - return { symbol, flow: allFlow, quarters }; + // Return individual filing rows for month-granularity grouping + const filings = rows.map(r => ({ + filerCik: r.filer_cik, + filerName: r.filer_name, + shares: r.shares ?? 0, + valueUsd: r.value_usd ?? 0, + reportedQuarter: r.reported_quarter, + filedAt: r.filed_at, + form: r.form, + })); + + return { symbol, flow: allFlow, quarters, filings }; }), /** Get insider activity stream for a symbol (M5 quarterly price strip). */ insiderStream: publicProcedure .input(z.object({ symbol: z.string().min(1), - limit: z.number().int().min(1).max(100).optional().default(50), + limit: z.number().int().min(1).max(500).optional().default(200), })) .query(async ({ ctx, input }) => { const symbol = input.symbol.toUpperCase(); + // Go back 5 years so the monthly chart has real depth. + const fiveYearsAgo = new Date(); + fiveYearsAgo.setFullYear(fiveYearsAgo.getFullYear() - 5); + const since = fiveYearsAgo.toISOString().slice(0, 10); + const transactions = ctx.db.prepare(` SELECT form4_id, symbol, insider_name, insider_role, tx_date, tx_code, tx_type, shares, price, is_10b5_1, classification, filed_at FROM insider_transactions - WHERE symbol = ? + WHERE symbol = ? AND tx_date >= ? ORDER BY tx_date DESC LIMIT ? - `).all(symbol, input.limit) as Array<{ + `).all(symbol, since, input.limit) as Array<{ form4_id: string; symbol: string; insider_name: string; @@ -509,11 +599,11 @@ const institutionalRouter = router({ netShares += shares * direction; return { - insiderName: tx.insider_name, - role: tx.insider_role, - txDate: tx.tx_date, - txCode: tx.tx_code, - txType: tx.tx_type, + reporter: tx.insider_name, + relationship: tx.insider_role, + transactionDate: tx.tx_date, + transactionCode: tx.tx_code, + transactionType: tx.tx_type, shares, price: tx.price, is10b5: tx.is_10b5_1 === 1, @@ -678,14 +768,8 @@ const watchlistRouter = router({ const { addSymbol } = await import('../db/watchlistRepository.ts'); const added = addSymbol(ctx.db, userId, input.symbol); - // Fire-and-forget: fetch SEC data (13F + Form 4) for this symbol in background. - // Don't await — let the response return immediately. if (added) { - import('../services/secDataFetcher.ts').then(({ fetchAndStoreSecData }) => { - fetchAndStoreSecData(ctx.db, input.symbol).catch((e) => { - console.error(`[sec fetch] failed for ${input.symbol}:`, e); - }); - }).catch(() => { /* sec fetcher not available */ }); + queueSecFetch(ctx.db, input.symbol); } return { added }; diff --git a/app/src/app/admin/queue/page.tsx b/app/src/app/admin/queue/page.tsx new file mode 100644 index 0000000..fb52d7f --- /dev/null +++ b/app/src/app/admin/queue/page.tsx @@ -0,0 +1,614 @@ +"use client"; +import { useEffect, useState, useCallback } from "react"; +import { AdminLayout } from "@/components/AdminLayout"; +import { api } from "@/lib/trpc"; + +interface QueueItem { + source: string; + status: string; + count: number; + last_attempt: string | null; + retry_count: number; + backoff_until: string | null; +} + +interface QueueStatus { + paused: boolean; + queued: number; + in_flight: number; + failed: number; + backoff: number; + done: number; + backoff_until: string | null; + last_error: string | undefined; + counts: Record; + lastErrors: Array<{ key: string; error: string | null }>; +} + +interface ErrorLogEntry { + attempt: number; + errorMessage: string | null; + errorStack: string | null; + ts: string; +} + +interface ScheduleEntry { + sourceKind: string; + intervalMs: number; + lastEnqueued: string | null; + nextEnqueue: string | null; +} + +function Spinner() { + return ( +
+
+
+ ); +} + +function Banner({ type, message, onDismiss }: { type: "success" | "error"; message: string; onDismiss?: () => void }) { + const styles = type === "success" + ? "border-green-500/30 bg-green-950/20 text-up" + : "border-red-500/30 bg-red-950/20 text-danger"; + return ( +
+ {message} + {onDismiss && ( + + )} +
+ ); +} + +function formatTimestamp(ts: string | null): string { + if (!ts) return "Never"; + try { + const d = new Date(ts); + return d.toLocaleString("en-US", { month: "short", day: "numeric", hour: "2-digit", minute: "2-digit" }); + } catch { + return ts.slice(0, 16); + } +} + +function msToHuman(ms: number): string { + if (ms < 60000) return `${Math.round(ms / 1000)}s`; + if (ms < 3600000) return `${Math.round(ms / 60000)}m`; + if (ms < 86400000) return `${Math.round(ms / 3600000)}h`; + return `${Math.round(ms / 86400000)}d`; +} + +export default function QueuePage() { + const [queue, setQueue] = useState([]); + const [status, setStatus] = useState(null); + const [schedules, setSchedules] = useState([]); + const [loading, setLoading] = useState(true); + const [error, setError] = useState(null); + const [success, setSuccess] = useState(null); + const [logs, setLogs] = useState>({}); + const [expandedLog, setExpandedLog] = useState(null); + + const [symbol, setSymbol] = useState(""); + const [triggering, setTriggering] = useState(false); + + const [resettingSource, setResettingSource] = useState(null); + const [resetResult, setResetResult] = useState>({}); + + const [newScheduleKind, setNewScheduleKind] = useState("sec-fetch"); + const [newScheduleInterval, setNewScheduleInterval] = useState("86400000"); + const [addingSchedule, setAddingSchedule] = useState(false); + + const [pausing, setPausing] = useState(false); + + const loadData = useCallback(async () => { + try { + const [queueData, statusData, scheduleData] = await Promise.all([ + api.admin.queueHealth(), + api.admin.queueStatus(), + api.admin.queueSchedules(), + ]); + setQueue(queueData); + setStatus(statusData); + setSchedules(scheduleData); + setError(null); + } catch (e) { + setError(e instanceof Error ? e.message : "Failed to load queue data"); + } finally { + setLoading(false); + } + }, []); + + useEffect(() => { + loadData(); + const interval = setInterval(loadData, 5000); + return () => { clearInterval(interval); }; + }, [loadData]); + + const handlePause = async () => { + setPausing(true); + try { + if (status?.paused) { + await api.admin.queueResume(); + setSuccess("Queue resumed"); + } else { + await api.admin.queuePause(); + setSuccess("Queue paused"); + } + await loadData(); + } catch (e) { + setError(e instanceof Error ? e.message : "Failed to toggle pause"); + } finally { + setPausing(false); + } + }; + + const handleClearDone = async () => { + try { + const result = await api.admin.queueClearDone(24); + setSuccess(`Cleared ${result.cleared} done jobs`); + await loadData(); + } catch (e) { + setError(e instanceof Error ? e.message : "Failed to clear done jobs"); + } + }; + + const handleQueueFetch = async (e: React.FormEvent) => { + e.preventDefault(); + if (!symbol.trim()) return; + setTriggering(true); + setError(null); + setSuccess(null); + try { + await api.admin.queueSecFetch(symbol.trim().toUpperCase()); + setSuccess(`Queued SEC fetch for ${symbol.toUpperCase()}`); + setSymbol(""); + await loadData(); + } catch (e) { + setError(e instanceof Error ? e.message : "Failed to queue fetch"); + } finally { + setTriggering(false); + } + }; + + const handleResetBackoff = async (sourceKind: string) => { + setResettingSource(sourceKind); + try { + const result = await api.admin.queueRetrySource(sourceKind); + setResetResult((prev) => ({ ...prev, [sourceKind]: result.cleared })); + await loadData(); + } catch (e) { + setError(e instanceof Error ? e.message : "Failed to reset backoff"); + } finally { + setResettingSource(null); + } + }; + + const handleRetryJob = async (key: string) => { + try { + await api.admin.queueRetryJob(key); + setSuccess(`Retried job: ${key}`); + await loadData(); + } catch (e) { + setError(e instanceof Error ? e.message : "Failed to retry job"); + } + }; + + const toggleLog = async (key: string) => { + if (expandedLog === key) { + setExpandedLog(null); + return; + } + setExpandedLog(key); + if (!logs[key]) { + try { + const entries = await api.admin.queueLogs(key, 10); + setLogs((prev) => ({ ...prev, [key]: entries })); + } catch { + // ignore + } + } + }; + + const handleAddSchedule = async (e: React.FormEvent) => { + e.preventDefault(); + setAddingSchedule(true); + try { + await api.admin.queueSetSchedule(newScheduleKind, Number(newScheduleInterval)); + setSuccess(`Schedule added for ${newScheduleKind}`); + setNewScheduleKind("sec-fetch"); + setNewScheduleInterval("86400000"); + await loadData(); + } catch (e) { + setError(e instanceof Error ? e.message : "Failed to add schedule"); + } finally { + setAddingSchedule(false); + } + }; + + const handleDeleteSchedule = async (sourceKind: string) => { + try { + await api.admin.queueDeleteSchedule(sourceKind); + setSuccess(`Schedule deleted for ${sourceKind}`); + await loadData(); + } catch (e) { + setError(e instanceof Error ? e.message : "Failed to delete schedule"); + } + }; + + type SourceAgg = { + pending: number; + inFlight: number; + done: number; + backoff: number; + failed: number; + lastAttempt: string | null; + hasBackoff: boolean; + }; + + const sources: Record = {}; + for (const item of queue) { + if (!sources[item.source]) { + sources[item.source] = { pending: 0, inFlight: 0, done: 0, backoff: 0, failed: 0, lastAttempt: null, hasBackoff: false }; + } + const agg = sources[item.source]; + if (item.status === "backoff") { agg.backoff += item.count; agg.hasBackoff = true; } + else if (item.status === "failed") { agg.failed += item.count; } + else if (item.status === "in_flight") { agg.inFlight += item.count; } + else if (item.status === "pending") { agg.pending += item.count; } + else if (item.status === "done") { agg.done += item.count; } + if (item.last_attempt) { + if (!agg.lastAttempt || new Date(item.last_attempt) > new Date(agg.lastAttempt)) { + agg.lastAttempt = item.last_attempt; + } + } + } + + const sortedSources = Object.entries(sources).sort(([a], [b]) => a.localeCompare(b)); + + if (loading) { + return ; + } + + return ( + +
+
+

Queue Management

+

Monitor queue health, schedules, and manually trigger data fetches

+
+ + {error && setError(null)} />} + {success && setSuccess(null)} />} + {Object.entries(resetResult).map(([src, cleared]) => ( + setResetResult({})} /> + ))} + + {/* Control Bar */} +
+
+
+ + +
+
+ + + {status?.paused ? "Paused" : "Active"} + + + {status?.counts.pending ?? 0}p / {status?.counts.in_flight ?? 0}i / {status?.counts.failed ?? 0}f / {status?.counts.backoff ?? 0}b + +
+
+
+ + {/* Queue Health Overview */} +
+
+

+ + Queue Health Overview +

+ {sortedSources.length} sources +
+ + {sortedSources.length === 0 ? ( +

No queue data available

+ ) : ( + <> +
+ + Pending + + + In-flight + + + Done + + + Backoff + + + Failed + +
+ + {sortedSources.some(([, v]) => v.backoff > 0) && ( +
+ + Backoff active — jobs stalled and awaiting retry +
+ )} + +
+ + + + + + + + + + + + + + {sortedSources.map(([source, agg]) => ( + + + + + + + + + + + ))} + +
SourcePendingIn-flightDoneBackoffFailedLast AttemptActions
{source} 0 ? "text-accent" : "text-fg-faint"}`}>{agg.pending} 0 ? "text-up" : "text-fg-faint"}`}>{agg.inFlight}{agg.done}{agg.backoff} 0 ? "text-danger" : "text-fg-faint"}`}>{agg.failed}{formatTimestamp(agg.lastAttempt)} +
+ {agg.hasBackoff && ( + + )} + {agg.failed > 0 && ( + + )} +
+
+ + {sortedSources.some(([, v]) => v.failed > 0) && ( +
+ + Failed jobs detected — expand below for details +
+ )} + + )} + + + {/* Failed Jobs Log */} + {status && status.lastErrors.length > 0 && ( +
+
+

+ + Failed Jobs +

+ {status.lastErrors.length} recent +
+
+ {status.lastErrors.map((job) => ( +
+ + {expandedLog === job.key && ( +
+
+ +
+
+ {logs[job.key]?.length ? ( + logs[job.key].map((entry, i) => ( +
+
+ Attempt #{entry.attempt} — {formatTimestamp(entry.ts)} +
+
+ {entry.errorStack ?? entry.errorMessage ?? "Unknown error"} +
+
+ )) + ) : ( +
+ + Loading logs... +
+ )} +
+
+ )} +
+ ))} +
+
+ )} + + {/* Schedule Configuration */} +
+
+

+ + Schedule Configuration +

+ {schedules.length} schedules +
+ + {schedules.length > 0 && ( + + + + + + + + + + + + {schedules.map((s) => ( + + + + + + + + ))} + +
SourceIntervalLast EnqueuedNext EnqueueActions
{s.sourceKind}{msToHuman(s.intervalMs)}{formatTimestamp(s.lastEnqueued)}{formatTimestamp(s.nextEnqueue)} + +
+ )} + + {/* Add Schedule Form */} +
+
+
+ + +
+
+ + setNewScheduleInterval(e.target.value)} + min={60000} + max={604800000} + className="w-full rounded border border-line-strong bg-surface px-3 py-2 text-xs font-mono text-fg focus:outline-none focus:border-accent/50 transition-colors" + disabled={addingSchedule} + /> +
+ +
+
+
+ + {/* Manual Trigger */} +
+
+

+ + Manual Trigger +

+ Queue a fresh SEC data fetch for any symbol +
+
+
+ + setSymbol(e.target.value.toUpperCase().slice(0, 10))} + placeholder="NVDA" + maxLength={10} + className="w-full rounded border border-line-strong bg-surface px-3 py-2 text-sm font-mono text-fg placeholder:text-fg-faint focus:outline-none focus:border-accent/50 transition-colors" + disabled={triggering} + /> +
+ +
+
+ + + ); +} diff --git a/app/src/lib/trpc.ts b/app/src/lib/trpc.ts index d9cdb78..dc60b65 100644 --- a/app/src/lib/trpc.ts +++ b/app/src/lib/trpc.ts @@ -210,12 +210,12 @@ export const api = { removeHolding: (symbol: string) => trpcMutate<{ removed: boolean }>("portfolio.removeHolding", { symbol }), }, dashboard: { - rollup: () => trpcQuery[]>("dashboard.rollup"), + rollup: () => trpcQuery[]>("dashboard.rollup", {}), }, institutional: { flow: (symbol: string) => trpcQuery<{ symbol: string; flow: Array>; quarters: string[] }>("institutional.flow", { symbol }), - insiderStream: (symbol: string, limit = 50) => + insiderStream: (symbol: string, limit = 200) => trpcQuery<{ symbol: string; events: Array>; netShares: number; count: number }>("institutional.insiderStream", { symbol, limit }), }, emotionLogger: { @@ -230,6 +230,26 @@ export const api = { unackedCount: () => trpcQuery<{ count: number }>("alerts.unackedCount"), acknowledgeAll: () => trpcMutate<{ ok: boolean }>("alerts.acknowledgeAll"), }, + admin: { + usersList: () => trpcQuery>("admin.usersList"), + userSessions: (userId: string) => trpcQuery>("admin.userSessions", { userId }), + resetPassword: (email: string, tempPassword: string) => trpcMutate<{ userId: string }>("admin.resetPassword", { email, tempPassword }), + queueHealth: () => trpcQuery>("admin.queueHealth"), + resetBackoff: (sourceKind: string) => trpcMutate<{ cleared: number }>("admin.resetBackoff", { sourceKind }), + queueSecFetch: (symbol: string) => trpcMutate<{ ok: boolean }>("admin.queueSecFetch", { symbol }), + queueStatus: () => trpcQuery("admin.queueStatus"), + queuePause: () => trpcMutate<{ paused: boolean }>("admin.queuePause", undefined), + queueResume: () => trpcMutate<{ paused: boolean }>("admin.queueResume", undefined), + queueClearDone: (olderThanHours: number) => trpcMutate<{ cleared: number }>("admin.queueClearDone", { olderThanHours }), + queueSchedules: () => trpcQuery("admin.queueSchedules"), + queueSetSchedule: (sourceKind: string, intervalMs: number) => trpcMutate<{ ok: boolean }>("admin.queueSetSchedule", { sourceKind, intervalMs }), + queueDeleteSchedule: (sourceKind: string) => trpcMutate<{ ok: boolean }>("admin.queueDeleteSchedule", { sourceKind }), + queueLogs: (key: string, limit: number) => trpcQuery("admin.queueLogs", { key, limit }), + queueRetryJob: (key: string) => trpcMutate<{ ok: boolean }>("admin.queueRetryJob", { key }), + queueRetrySource: (sourceKind: string) => trpcMutate<{ cleared: number }>("admin.queueRetrySource", { sourceKind }), + auditLog: (limit = 50, offset = 0, actor?: string, action?: string) => + trpcQuery>("admin.auditLog", { limit, offset, ...(actor ? { actor } : {}), ...(action ? { action } : {}) }), + }, }; export interface WatchlistEntry { @@ -244,3 +264,31 @@ export interface PortfolioHolding { avg_cost: number; added_at: string; } + +export interface QueueStatus { + paused: boolean; + queued: number; + in_flight: number; + failed: number; + backoff: number; + done: number; + backoff_until: string | null; + last_error: string | undefined; + counts: Record; + lastErrors: Array<{ key: string; error: string | null }>; +} + +export interface ScheduleEntry { + sourceKind: string; + intervalMs: number; + lastEnqueued: string | null; + nextEnqueue: string | null; +} + +export interface ErrorLogEntry { + attempt: number; + errorMessage: string | null; + errorStack: string | null; + ts: string; +} +