diff --git a/app/server/src/adapters/__tests__/backfill.test.ts b/app/server/src/adapters/__tests__/backfill.test.ts index 8993e87..4ccabab 100644 --- a/app/server/src/adapters/__tests__/backfill.test.ts +++ b/app/server/src/adapters/__tests__/backfill.test.ts @@ -1,7 +1,7 @@ // Investor Flow — slice-4 backfill integration test. // // Verifies the full subscribe → queue → drain → cache-write path end-to-end: -// 1. CacheRepository.subscribe('NVDA','equity') enqueues 4 keys into the AdapterQueue. +// 1. CacheRepository.subscribe('NVDA','equity') enqueues quote + candles + symbol. // 2. AdapterQueue.drain() processes them; FakeSourceAdapter serves canned responses. // 3. Re-subscribing increments refcount but does NOT re-queue (no duplicate fetches). // @@ -17,7 +17,15 @@ import { AdapterQueue } from "../../queue/AdapterQueue.ts"; import { FakeSourceAdapter } from "../SourceAdapter.ts"; import { createDb, initSchema } from "../../db/client.ts"; -test("subscribe enqueues 4 pending rows for NVDA/equity", async () => { +async function drainUntilIdle(queue: AdapterQueue, db: DatabaseSync, max = 8): Promise { + for (let i = 0; i < max; i++) { + const pending = db.prepare("SELECT COUNT(*) AS c FROM adapter_queue WHERE status IN ('pending','backoff')").get() as { c: number }; + if (pending.c === 0) return; + await queue.drain(); + } +} + +test("subscribe enqueues quote + candles + symbol for NVDA/equity", async () => { const db = createDb({ path: ":memory:" }); initSchema(db); @@ -53,20 +61,17 @@ test("subscribe enqueues 4 pending rows for NVDA/equity", async () => { "SELECT key FROM adapter_queue WHERE status='pending'", ).all() as Array<{ key: string }>; - assert.equal(pending.length, 6, "expected 6 pending rows after subscribe"); + assert.equal(pending.length, 3, "expected quote + candles + symbol on first subscribe"); const keys = pending.map((r) => r.key).sort(); assert.deepEqual(keys, [ - "nasdaq:nasdaqShortinterest:NVDA", - "yfinance:adjustments:NVDA", "yfinance:candles:NVDA:1d", "yfinance:quote:NVDA", - "yfinance:shortinterest:NVDA", "yfinance:symbol:NVDA", ]); }); -test("drain processes all 4 keys and writes to cache + DB", async () => { +test("drain processes first-seed keys and writes to cache + DB", async () => { const db = createDb({ path: ":memory:" }); initSchema(db); @@ -96,35 +101,20 @@ test("drain processes all 4 keys and writes to cache + DB", async () => { queue.cache = cache; await cache.subscribe("NVDA", "equity"); - await queue.drain(); + await drainUntilIdle(queue, db); - // FakeSourceAdapter should have been called for candles and adjustments at minimum. - assert.ok( - fake.calls.includes("yfinance:candles:NVDA:1d"), - "expected candles fetch to have been invoked", - ); - assert.ok( - fake.calls.includes("yfinance:adjustments:NVDA"), - "expected adjustments fetch to have been invoked", - ); + assert.ok(fake.calls.includes("yfinance:quote:NVDA"), "expected quote fetch"); + assert.ok(fake.calls.includes("yfinance:candles:NVDA:1d"), "expected candles fetch"); + assert.ok(fake.calls.includes("yfinance:symbol:NVDA"), "expected symbol meta fetch"); + assert.ok(!fake.calls.includes("yfinance:adjustments:NVDA"), "adjustments deferred off first seed"); - // Cache should now hold non-stale entries for candles and adjustments. const candlesEntry = await cache.get("yfinance:candles:NVDA:1d"); assert.equal(candlesEntry.isStale, false, "candles should not be stale after drain"); assert.ok(Array.isArray(candlesEntry.value), "candles value should be an array"); assert.equal((candlesEntry.value as unknown[]).length, 2, "candles should have 2 rows"); - const adjEntry = await cache.get("yfinance:adjustments:NVDA"); - assert.equal(adjEntry.isStale, false, "adjustments should not be stale after drain"); - assert.ok(Array.isArray(adjEntry.value), "adjustments value should be an array"); - assert.equal((adjEntry.value as unknown[]).length, 2, "adjustments should have 2 rows"); - - // DB tables should reflect the written data. const candleRows = db.prepare("SELECT * FROM price_candles WHERE symbol='NVDA'").all(); assert.equal(candleRows.length, 2, "price_candles table should have 2 rows"); - - const adjRows = db.prepare("SELECT * FROM price_adjustments WHERE symbol='NVDA'").all(); - assert.equal(adjRows.length, 2, "price_adjustments table should have 2 rows"); }); test("re-subscribe increments refcount but does NOT re-queue (no duplicate fetches)", async () => { @@ -162,7 +152,7 @@ test("re-subscribe increments refcount but does NOT re-queue (no duplicate fetch // First subscribe + drain. await cache.subscribe("NVDA", "equity"); - await queue.drain(); + await drainUntilIdle(queue, db); // Capture call count AFTER drain (all 4 fetches completed). const afterFirstDrain = fake.calls.length; @@ -189,10 +179,6 @@ test("re-subscribe increments refcount but does NOT re-queue (no duplicate fetch "re-subscribe + drain should not trigger additional fetches", ); - // DB tables should still have exactly 2 rows each (no duplicates). const candleRows = db.prepare("SELECT * FROM price_candles WHERE symbol='NVDA'").all(); assert.equal(candleRows.length, 2, "price_candles should still have 2 rows after re-subscribe"); - - const adjRows = db.prepare("SELECT * FROM price_adjustments WHERE symbol='NVDA'").all(); - assert.equal(adjRows.length, 2, "price_adjustments should still have 2 rows after re-subscribe"); }); diff --git a/app/server/src/admin/__tests__/admin.test.ts b/app/server/src/admin/__tests__/admin.test.ts index 2ea37f4..fb684ac 100644 --- a/app/server/src/admin/__tests__/admin.test.ts +++ b/app/server/src/admin/__tests__/admin.test.ts @@ -8,7 +8,7 @@ import { fileURLToPath } from 'node:url'; import { listUsers, resetPassword, gdprExport, queueHealth, resetQueueBackoff, - recordAudit, requireAdminOrOwner, NotOwnerError, + recordAudit, requireAdminOrOwner, NotOwnerError, deleteUser, isSystemAccount, } from '../admin.ts'; import { hashPassword } from '../../trpc/context.ts'; @@ -33,6 +33,31 @@ test('listUsers returns users without pw_hash', () => { assert.ok(users.some((u) => u.is_admin === 1)); }); +test('listUsers omits the anonymous system account', () => { + const db = freshDb(); + db.prepare("INSERT INTO users (id,email,pw_hash,created_at) VALUES (?,?,?,?)") + .run('anonymous', 'anonymous@investor-flow.local', '', '2026-01-01T00:00:00Z'); + const users = listUsers(db); + assert.equal(users.length, 2); + assert.ok(!users.some((u) => u.id === 'anonymous' || u.email.startsWith('anonymous@'))); +}); + +test('deleteUser refuses the anonymous system account', () => { + const db = freshDb(); + db.prepare("INSERT INTO users (id,email,pw_hash,created_at) VALUES (?,?,?,?)") + .run('anonymous', 'anonymous@investor-flow.local', '', '2026-01-01T00:00:00Z'); + assert.throws(() => deleteUser(db, 'admin_1', 'anonymous'), /System account cannot be deleted/); + const still = db.prepare("SELECT id FROM users WHERE id='anonymous'").get() as { id: string }; + assert.equal(still.id, 'anonymous'); +}); + +test('isSystemAccount matches sentinel ids and emails', () => { + assert.equal(isSystemAccount('anonymous', 'x@y.z'), true); + assert.equal(isSystemAccount('system', null), true); + assert.equal(isSystemAccount('abc', 'anonymous@investor-flow.local'), true); + assert.equal(isSystemAccount('admin_1', 'admin@example.com'), false); +}); + test('resetPassword: admin actor updates pw_hash and audits', { skip: true }, async () => { // Hashing tested indirectly; logic verified in resetPassword-audits test. }); diff --git a/app/server/src/admin/admin.ts b/app/server/src/admin/admin.ts index c33a607..fb73390 100644 --- a/app/server/src/admin/admin.ts +++ b/app/server/src/admin/admin.ts @@ -71,13 +71,23 @@ export function requireAdminOrOwner( if (!row || !row.is_admin) throw new NotOwnerError('Actor is not admin and does not own the target.'); } -/** list users — id, email, complexity, created_at, is_admin, status, modules. Never returns pw_hash. */ +/** Guest / sentinel rows used as owner_id when nobody is signed in. Not product accounts. */ +export function isSystemAccount(id: string, email?: string | null): boolean { + if (id === 'anonymous' || id === 'system') return true; + const e = (email ?? '').toLowerCase(); + return e.startsWith('anonymous@'); +} + +/** list users — id, email, complexity, created_at, is_admin, status, modules. Never returns pw_hash. + * System sentinels (anonymous guest book) are omitted. */ export function listUsers(db: DatabaseSync): UserRecord[] { - return db - .prepare( - 'SELECT id, email, complexity, created_at, is_admin, status, modules FROM users ORDER BY created_at ASC', - ) - .all() as unknown as UserRecord[]; + return ( + db + .prepare( + 'SELECT id, email, complexity, created_at, is_admin, status, modules FROM users ORDER BY created_at ASC', + ) + .all() as unknown as UserRecord[] + ).filter((u) => !isSystemAccount(u.id, u.email)); } const VALID_MODULES = ['research', 'execution', 'analytics', 'settings']; @@ -102,6 +112,8 @@ export function setUserModules( * Refuses to disable the actor's own account or the last remaining admin. */ export function disableUser(db: DatabaseSync, actorId: string | null, targetUserId: string): { ok: boolean } { if (actorId === targetUserId) throw new NotOwnerError('Cannot disable your own account.'); + const target = db.prepare('SELECT email FROM users WHERE id=?').get(targetUserId) as { email: string } | undefined; + if (target && isSystemAccount(targetUserId, target.email)) throw new Error('System account cannot be disabled.'); const row = db.prepare('SELECT is_admin FROM users WHERE id=?').get(targetUserId) as { is_admin: number } | undefined; if (!row) throw new Error('admin: target user not found'); if (row.is_admin) { @@ -127,6 +139,10 @@ export function enableUser(db: DatabaseSync, actorId: string | null, targetUserI * own account or the last remaining admin. */ export function deleteUser(db: DatabaseSync, actorId: string | null, targetUserId: string): { ok: boolean } { if (actorId === targetUserId) throw new NotOwnerError('Cannot delete your own account.'); + const target = db.prepare('SELECT email FROM users WHERE id=?').get(targetUserId) as { email: string } | undefined; + if (target && isSystemAccount(targetUserId, target.email)) { + throw new Error('System account cannot be deleted. It owns guest watchlists when nobody is signed in.'); + } const row = db.prepare('SELECT is_admin FROM users WHERE id=?').get(targetUserId) as { is_admin: number } | undefined; if (!row) throw new Error('admin: target user not found'); if (row.is_admin) { diff --git a/app/server/src/cache/CacheRepository.ts b/app/server/src/cache/CacheRepository.ts index 6a0b0ef..aaf087e 100644 --- a/app/server/src/cache/CacheRepository.ts +++ b/app/server/src/cache/CacheRepository.ts @@ -627,6 +627,8 @@ const HANDLERS = new Map([ export interface CacheRepository { get(key: CacheKey): Promise>; + /** Read cache without scheduling a refresh (watchlist snapshots). */ + peek(key: CacheKey): Promise>; set(key: CacheKey, value: T, ttlClass: TtlClass, provenance: Provenance): Promise; stale(key: CacheKey): boolean; /** @@ -645,7 +647,7 @@ export interface CacheRepository { /** Permanent system pin (rotation universe, SPY, VIX) — survives unsubscribe. */ pinSystemSymbol(symbol: string, tickerKind: TickerKind): Promise; demandSet(): Promise; - getMany(keys: CacheKey[]): Promise>; + getMany(keys: CacheKey[], opts?: { refresh?: boolean }): Promise>; /** Delete a cache entry by key (or, for wildcard keys ending in `:*`, all matching entries). */ del(key: CacheKey): Promise; /** Underlying DB for schedule TTL checks (queue only). */ @@ -665,28 +667,34 @@ export class CacheRepositoryImpl implements CacheRepository { if (!h) throw new Error(`unknown cache kind: ${kind}`); return h; } - async get(key: CacheKey): Promise> { + private readEntry(key: CacheKey): CacheEntry { const { source, kind, id } = parseCacheKey(key); const h = this.handler(kind); const row = h.read(this._db, id); const now = Date.now(); let stale = h.isStale(row ? row.stalenessTs : null, now, id); - // Incomplete symbol meta (null name) is always treated as stale for SWR re-fetch. if (kind === 'symbol' && row) { const meta = row.value as SymbolMeta; if (!meta?.name && tsAgeMs(row.stalenessTs, now) > SYMBOL_META_INCOMPLETE_TTL_MS) { stale = true; } } - if (stale) { - try { await this._scheduler.queue(key); } catch { /* background refresh; never block readers */ } - } return { value: (row ? row.value : null) as T | null, provenance: row ? { fetchedAt: row.stalenessTs, sourceKind: source } : null, isStale: stale, }; } + async peek(key: CacheKey): Promise> { + return this.readEntry(key); + } + async get(key: CacheKey): Promise> { + const e = this.readEntry(key); + if (e.isStale) { + try { await this._scheduler.queue(key); } catch { /* background refresh; never block readers */ } + } + return e; + } async set(key: CacheKey, value: T, ttlClass: TtlClass, provenance: Provenance): Promise { const { kind, id } = parseCacheKey(key); const h = this.handler(kind); @@ -736,25 +744,29 @@ export class CacheRepositoryImpl implements CacheRepository { async subscribe(symbol: string, tickerKind: TickerKind): Promise { const d = this._db; - d.prepare('INSERT OR IGNORE INTO symbol_demand (symbol,refcount,ticker_kind,in_demand,last_refreshed_at) VALUES (?,?,?,?,?)').run(symbol, 0, tickerKind, 1, null); + // User watch/portfolio demand is watched (T2). MIN() keeps T0 holdings / T1 alerts. + d.prepare( + 'INSERT OR IGNORE INTO symbol_demand (symbol,refcount,ticker_kind,in_demand,last_refreshed_at,tier) VALUES (?,?,?,?,?,?)', + ).run(symbol, 0, tickerKind, 1, null, 2); const prev = d.prepare('SELECT refcount FROM symbol_demand WHERE symbol=?').get(symbol) as { refcount: number } | undefined; const before = prev?.refcount ?? 0; - d.prepare('UPDATE symbol_demand SET refcount = refcount + 1, in_demand = 1, ticker_kind=? WHERE symbol=?').run(tickerKind, symbol); + d.prepare( + 'UPDATE symbol_demand SET refcount = refcount + 1, in_demand = 1, ticker_kind=?, tier = MIN(COALESCE(tier, 3), 2) WHERE symbol=?', + ).run(tickerKind, symbol); if (before === 0) { - // First user demand: full seed once (not every schedule tick). - for (const k of [ - `yfinance:quote:${symbol}`, - `yfinance:symbol:${symbol}`, - `yfinance:candles:${symbol}:1d`, - `yfinance:adjustments:${symbol}`, - `yfinance:shortinterest:${symbol}`, - `nasdaq:nasdaqShortinterest:${symbol}`, - ]) { + // First user demand: jump the quote to the front so the watchlist mark + // is not parked behind a T2/T3 backlog (or skipped as default T3). + const quoteKey = `yfinance:quote:${symbol}`; + try { + if (this._scheduler.prioritize) await this._scheduler.prioritize(quoteKey); + else await this._scheduler.queue(quoteKey); + } catch { /* ignore */ } + for (const k of [`yfinance:candles:${symbol}:1d`, `yfinance:symbol:${symbol}`]) { try { await this._scheduler.queue(k); } catch { /* ignore */ } } } else { // Subsequent demand: re-check staleness and queue missing kinds (no refcount change). - await this.queueIfNeeded(symbol); + await this.queueIfNeeded(symbol, { prioritizeQuote: true }); } } @@ -797,9 +809,10 @@ export class CacheRepositoryImpl implements CacheRepository { 'SELECT symbol FROM symbol_demand WHERE in_demand = 1 OR COALESCE(system_pin, 0) = 1 ORDER BY symbol', ).all() as Array<{ symbol: string }>).map((r) => r.symbol); } - async getMany(keys: CacheKey[]): Promise> { + async getMany(keys: CacheKey[], opts?: { refresh?: boolean }): Promise> { + const refresh = opts?.refresh !== false; return Promise.all(keys.map(async (key) => { - const e = await this.get(key); + const e = refresh ? await this.get(key) : await this.peek(key); return { key, value: e.value, isStale: e.isStale, fetchedAt: e.provenance?.fetchedAt ?? null }; })); } diff --git a/app/server/src/cache/__tests__/CacheRepository.test.ts b/app/server/src/cache/__tests__/CacheRepository.test.ts index 19482f0..be2d970 100644 --- a/app/server/src/cache/__tests__/CacheRepository.test.ts +++ b/app/server/src/cache/__tests__/CacheRepository.test.ts @@ -6,7 +6,13 @@ import { type CacheScheduler, type CacheKey, type Quote, type PriceCandle, type SymbolMeta, } from '../CacheRepository.ts'; -class FakeScheduler { queued: CacheKey[] = []; async queue(key: CacheKey): Promise { this.queued.push(key); } reset() { this.queued = []; } } +class FakeScheduler { + queued: CacheKey[] = []; + prioritized: CacheKey[] = []; + async queue(key: CacheKey): Promise { this.queued.push(key); } + async prioritize(key: CacheKey): Promise { this.prioritized.push(key); this.queued.push(key); } + reset() { this.queued = []; this.prioritized = []; } +} function setup() { const db = createDb({ path: ':memory:' }); @@ -40,6 +46,27 @@ test('get returns stale value (stale-while-revalidate) AND schedules refresh', a assert.ok(scheduler.queued.includes('yfinance:quote:NVDA')); }); +test('peek does not schedule a refresh', async () => { + const { repo, scheduler } = setup(); + const e = await repo.peek('yfinance:quote:NVDA'); + assert.equal(e.value, null); + assert.equal(e.isStale, true); + assert.equal(scheduler.queued.length, 0); +}); + +test('getMany refresh:false does not enqueue a 2-minute-old T2 quote', async () => { + const { repo, scheduler } = setup(); + await repo.subscribe('RIVN', 'equity'); + scheduler.reset(); + await repo.set('yfinance:quote:RIVN', { symbol: 'RIVN', price: 16 } as Quote, 'live_quote', { + fetchedAt: iso(-2 * 60_000), sourceKind: 'yfinance', + }); + scheduler.reset(); + const res = await repo.getMany(['yfinance:quote:RIVN'], { refresh: false }); + assert.equal(res[0].value?.price, 16); + assert.equal(scheduler.queued.length, 0); +}); + test('get on never-cached key returns null + schedules refresh', async () => { const { repo, scheduler } = setup(); const e = await repo.get('yfinance:quote:NVDA'); @@ -55,7 +82,9 @@ test('subscribe bumps refcount and schedules full seed on FIRST demand', async ( assert.equal(r.refcount, 1); assert.equal(r.in_demand, 1); assert.ok(scheduler.queued.includes('yfinance:quote:NVDA')); + assert.ok(scheduler.queued.includes('yfinance:candles:NVDA:1d')); assert.ok(scheduler.queued.includes('yfinance:symbol:NVDA')); + assert.ok(!scheduler.queued.includes('yfinance:adjustments:NVDA')); scheduler.reset(); // Second subscribe: TTL-aware requeue of missing/stale kinds only (no adjustments seed). await repo.subscribe('NVDA', 'equity'); @@ -171,6 +200,30 @@ test('fresh 5m candles are not stale; old observed_at is', async () => { assert.equal(scheduler.queued.length, 0); }); +test('subscribe marks the symbol watched (tier 2) and prioritizes the quote', async () => { + const { repo, scheduler, db } = setup(); + await repo.subscribe('RIVN', 'equity'); + const r = db.prepare('SELECT tier, refcount, in_demand FROM symbol_demand WHERE symbol=?').get('RIVN') as { + tier: number; refcount: number; in_demand: number; + }; + assert.equal(r.tier, 2); + assert.equal(r.refcount, 1); + assert.equal(r.in_demand, 1); + assert.equal(scheduler.prioritized[0], 'yfinance:quote:RIVN'); + assert.ok(scheduler.queued.includes('yfinance:quote:RIVN')); +}); + +test('subscribe does not raise a portfolio (T0) symbol to watched', async () => { + const { repo, db } = setup(); + db.prepare( + 'INSERT INTO symbol_demand (symbol,refcount,ticker_kind,in_demand,last_refreshed_at,system_pin,tier) VALUES (?,?,?,?,?,?,?)', + ).run('AAPL', 1, 'equity', 1, null, 0, 0); + await repo.subscribe('AAPL', 'equity'); + const r = db.prepare('SELECT tier, refcount FROM symbol_demand WHERE symbol=?').get('AAPL') as { tier: number; refcount: number }; + assert.equal(r.tier, 0); + assert.equal(r.refcount, 2); +}); + test('subscribe does not enqueue 1m or 5m', async () => { const { repo, scheduler } = setup(); await repo.subscribe('NVDA', 'equity'); diff --git a/app/server/src/index.ts b/app/server/src/index.ts index d19c135..0f004f8 100644 --- a/app/server/src/index.ts +++ b/app/server/src/index.ts @@ -19,6 +19,7 @@ import type { SourceFetch } from './adapters/SourceAdapter.ts'; import { composeYFinanceWithOptions } from './options/OptionsChainRouter.ts'; import cryptoMod from './lib/crypto.ts'; import { AdapterQueue } from './queue/AdapterQueue.ts'; +import { listRegisteredFamilies } from './services/vendorGate.ts'; import { seedCuratedCusips } from './services/cusipRegistry.ts'; import { seedAdminDefaultAlertSubscriptions } from './db/alertSubscriptionRepository.ts'; import { seedConfluence } from './confluence/confluenceSeed.ts'; @@ -126,18 +127,29 @@ queue.pinSystemUniverse().then(() => { console.log('[investor-flow] system pins (rotation universe) ready'); }).catch((e) => console.error('[investor-flow] pinSystemUniverse failed', e)); -// 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`); +// Startup recovery: Yahoo quotes resume immediately; SEC/lint wait with jitter +// so a restart cannot stampede unofficial Yahoo and 429 the quote lane. +const recovered = queue.recoverInterruptedJobs(); +if (recovered.yahoo + recovered.other > 0) { + console.log(`[investor-flow] recovered in_flight jobs yahoo=${recovered.yahoo} other=${recovered.other} (SEC delayed)`); +} const createContext = makeCreateContext({ db: database, cache, queue, xAdapter }); -// Background drain: stale-while-revalidate refreshes are queued by CacheRepository.get; -// this loop drains them (fetch via adapter -> write to cache), deduped + backed off. +// Per-family drain: Yahoo quotes never wait on 10-minute SEC jobs. 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(); +const SEC_DRAIN_MS = Number(process.env.IFLOW_SEC_DRAIN_MS ?? 5000); +const drainYahooTimer = setInterval(() => { queue.drainFamily('yfinance').catch((e) => console.error('[drain yfinance]', e)); }, DRAIN_MS); +drainYahooTimer.unref(); +const drainSecTimer = setInterval(() => { queue.drainFamily('sec').catch((e) => console.error('[drain sec]', e)); }, SEC_DRAIN_MS); +drainSecTimer.unref(); +const drainOtherTimer = setInterval(() => { + for (const family of listRegisteredFamilies()) { + if (family === 'yfinance' || family === 'sec') continue; + queue.drainFamily(family).catch((e) => console.error(`[drain ${family}]`, e)); + } +}, DRAIN_MS); +drainOtherTimer.unref(); // Auto-scheduler: every 30s, enqueue refreshes for due schedules const SCHEDULE_MS = 30_000; @@ -351,5 +363,5 @@ const server = createServer(async (req, res) => { 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)`); + console.log(`[investor-flow] backend on http://${HOST}:${PORT} (tRPC /api/trpc, health /health, yahoo drain ${DRAIN_MS}ms, sec drain ${SEC_DRAIN_MS}ms)`); }); diff --git a/app/server/src/queue/AdapterQueue.ts b/app/server/src/queue/AdapterQueue.ts index 75ef693..e132919 100644 --- a/app/server/src/queue/AdapterQueue.ts +++ b/app/server/src/queue/AdapterQueue.ts @@ -22,6 +22,10 @@ import type { SourceFetch } from '../adapters/SourceAdapter.ts'; import { DEFAULT_SOURCE_MIN_INTERVAL_MS, DEMAND_SET_SOFT_CAP, + WATCHED_QUOTE_SCHEDULE_CAP, + SEC_PENDING_CAP, + SEC_HEAL_INTERVAL_MS, + MAX_SEC_HEAL_PER_CYCLE, DRAIN_KIND_BUDGET, fetchTimeoutMs, isPermanentDataError, @@ -88,6 +92,10 @@ export interface QueueHealthDetailed { systemPinCount: number; quarantinedCount: number; candleLagMs: number | null; + /** Yahoo quote lane is the terminal SLO (watchlist/header). */ + yahooHealthy: boolean; + /** EDGAR/13F lane. A large backlog is backfill, not "restart the app". */ + secHealthy: boolean; dataPlaneHealthy: boolean; dataPlaneNotes: string[]; } @@ -126,8 +134,8 @@ export class AdapterQueue implements CacheScheduler { private readonly _sourceChains = new Map>(); /** Header / page-view quotes jump ahead of the rest of the Yahoo pile. */ private readonly _priorityKeys = new Set(); - /** Overlapping setInterval drains must not select a second batch mid-fetch. */ - private _drainBusy = false; + /** Per-family select lock. Yahoo can claim while SEC fetches. */ + private readonly _familyBusy = new Set(); constructor(opts: AdapterQueueOptions) { this._db = opts.db; @@ -443,18 +451,59 @@ export class AdapterQueue implements CacheScheduler { } async drain(): Promise { - if (!this._cache) return; - if (this.isPaused()) return; - if (this._drainBusy) return; - this._drainBusy = true; - try { - await this.drainOnce(); - } finally { - this._drainBusy = false; + const families = new Set(); + for (const sk of this._adapters.keys()) { + const f = sourceToFamily(String(sk)); + if (f) families.add(f); } + await Promise.all([...families].map((f) => this.drainFamily(f))); } - private async drainOnce(): Promise { + /** Restart recovery: Yahoo quotes resume now; slow families wait with jitter. */ + recoverInterruptedJobs(): { yahoo: number; other: number } { + const yf = this._db.prepare( + "UPDATE adapter_queue SET status='pending', error=NULL, retry_count=0, scheduled_for=NULL WHERE status='in_flight' AND key LIKE 'yfinance:%'", + ).run(); + const stuck = this._db.prepare( + "SELECT key FROM adapter_queue WHERE status='in_flight' AND key NOT LIKE 'yfinance:%'", + ).all() as Array<{ key: string }>; + const upd = this._db.prepare( + "UPDATE adapter_queue SET status='pending', error=NULL, retry_count=0, scheduled_for=? WHERE key=? AND status='in_flight'", + ); + for (const row of stuck) { + const delayMs = 15_000 + Math.floor(Math.random() * 120_000); + upd.run(new Date(Date.now() + delayMs).toISOString(), row.key); + } + return { yahoo: Number(yf.changes), other: stuck.length }; + } + + async drainFamily(family: VendorFamily): Promise { + if (!this._cache) return; + if (this.isPaused()) return; + if (this._familyBusy.has(family)) return; + this._familyBusy.add(family); + let specs: Array<{ + key: string; + source: SourceKind; + family: VendorFamily | null; + sym: string; + attempt: number; + }> = []; + try { + specs = this.claimDrainBatch(family); + } finally { + this._familyBusy.delete(family); + } + await Promise.all(specs.map((spec) => this.withSourceChain(spec.source, () => this.fetchSpec(spec)))); + } + + private claimDrainBatch(family?: VendorFamily): Array<{ + key: string; + source: SourceKind; + family: VendorFamily | null; + sym: string; + attempt: number; + }> { const now = Date.now(); // Recover hung workers: jobs left in_flight after crash/hang never complete. @@ -500,10 +549,20 @@ export class AdapterQueue implements CacheScheduler { if (key.startsWith('yfinance:adjustments:')) return 8; return 9; }; - const jobs = this._db.prepare( + let jobs = this._db.prepare( `SELECT key, retry_count, backoff_until, scheduled_for, last_attempt FROM adapter_queue WHERE status IN ('pending','backoff')`, ).all() as Array<{ key: string; retry_count: number; backoff_until?: string | null; scheduled_for?: string | null; last_attempt?: string | null }>; + if (family) { + jobs = jobs.filter((j) => { + try { + const { source } = parseCacheKey(j.key); + return sourceToFamily(source) === family; + } catch { + return false; + } + }); + } // Tier-aware sort: jumped header quotes first, then portfolio (T0), T1, T2, T3. // Within each tier: quote > candles > symbol > topHoldings > adjustments > other. @@ -573,6 +632,23 @@ export class AdapterQueue implements CacheScheduler { // fresh subscribe needs (quote+symbol+candles+adjustments+…) always fit // because spare drain capacity is filled back in pass 2. const familyJobsThisDrain: Partial> = {}; + // Count jobs already in flight from a previous drain so overlapping + // drains (Yahoo while SEC hangs) cannot exceed the family budget. + const inflightRows = this._db.prepare( + "SELECT key FROM adapter_queue WHERE status='in_flight'", + ).all() as Array<{ key: string }>; + for (const r of inflightRows) { + try { + const { source } = parseCacheKey(r.key); + const family = sourceToFamily(source); + if (family) familyJobsThisDrain[family] = (familyJobsThisDrain[family] ?? 0) + 1; + } catch { /* malformed key */ } + } + const nowIso = new Date(now).toISOString(); + const claimStmt = this._db.prepare( + "UPDATE adapter_queue SET status='in_flight', last_attempt=? WHERE key=? AND status IN ('pending','backoff')", + ); + const claim = (key: string): boolean => Number(claimStmt.run(nowIso, key).changes) === 1; // Reserve selected jobs, then run them concurrently (one at a time per // source so per-source pacing/cool-downs stay intact across parallel paths). const specs: Array<{ @@ -655,6 +731,7 @@ export class AdapterQueue implements CacheScheduler { } } + if (!claim(job.key)) continue; specs.push({ key: job.key, source, family, sym, attempt: job.retry_count + 1 }); this._priorityKeys.delete(job.key); processed += 1; @@ -672,6 +749,8 @@ export class AdapterQueue implements CacheScheduler { if (this.isSourceCoolingDown(source, Date.now())) continue; if (this.isSourceControlled(source)) continue; if (family && this.isVendorFamilyCoolingDown(family, Date.now())) continue; + if (family && (familyJobsThisDrain[family] ?? 0) >= familyDrainBudget(family) + && !isMinuteCandle(job.key) && !isOptionsSurface(job.key)) continue; const kindBudget = DRAIN_KIND_BUDGET[kind] ?? DRAIN_KIND_BUDGET._default; if ((kindUsed[kind] ?? 0) >= kindBudget) continue; const tier = tierMap.get(sym) ?? 99; @@ -683,16 +762,15 @@ export class AdapterQueue implements CacheScheduler { // Backlog pressure: skip T3 (background) symbols entirely so portfolio/alert symbols drain first. if ((hotQuotePending || backlogPressure) && source === 'yfinance' && tier >= 3) continue; if (backlogPressure && sym && tier >= 3) continue; + if (!claim(job.key)) continue; specs.push({ key: job.key, source, family, sym, attempt: job.retry_count + 1 }); this._priorityKeys.delete(job.key); processed += 1; kindUsed[kind] = (kindUsed[kind] ?? 0) + 1; + if (family) familyJobsThisDrain[family] = (familyJobsThisDrain[family] ?? 0) + 1; } - // Execute selected jobs concurrently. Same-source jobs are serialized by a - // per-source chain so per-source min-interval pacing and cool-downs apply - // exactly as in the old single-threaded loop. - await Promise.all(specs.map((spec) => this.withSourceChain(spec.source, () => this.fetchSpec(spec)))); + return specs; } /** Serialize per-source work (pacing + fetch) so parallel jobs never overlap a source. */ @@ -749,12 +827,16 @@ export class AdapterQueue implements CacheScheduler { this._lastFetchAt[source] = Date.now(); this._setStatus(key, 'in_flight'); + const FETCH_TIMEOUT_MS = fetchTimeoutMs(key); + let timeoutHandle: ReturnType | undefined; try { - const FETCH_TIMEOUT_MS = fetchTimeoutMs(key); const fetchPromise = adapter.fetchOne(key); - const timeoutPromise = new Promise((_, reject) => - setTimeout(() => reject(new Error(`fetchOne timeout after ${FETCH_TIMEOUT_MS}ms for ${key}`)), FETCH_TIMEOUT_MS), - ); + const timeoutPromise = new Promise((_, reject) => { + timeoutHandle = setTimeout( + () => reject(new Error(`fetchOne timeout after ${FETCH_TIMEOUT_MS}ms for ${key}`)), + FETCH_TIMEOUT_MS, + ); + }); const res = await Promise.race([fetchPromise, timeoutPromise]); try { await this._cache?.set(key, res.value, res.ttlClass, res.provenance); } catch { /* adapter may persist directly */ } this.clearSourceCooldown(source); @@ -826,6 +908,8 @@ export class AdapterQueue implements CacheScheduler { const bo = jobBackoffMs(attempt); 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, key); } + } finally { + if (timeoutHandle) clearTimeout(timeoutHandle); } } @@ -897,11 +981,20 @@ this._db.prepare("UPDATE adapter_queue SET status='backoff', last_attempt=?, ret } /** Symbols at a specific tier — used by per-tier schedule branches. */ - private demandSymbolsByTier(tier: number): string[] { + private demandSymbolsByTier(tier: number, cap?: number): string[] { const rows = this._db.prepare( 'SELECT symbol FROM symbol_demand WHERE tier = ? AND (in_demand = 1 OR COALESCE(system_pin, 0) = 1) ORDER BY symbol', ).all(tier) as Array<{ symbol: string }>; - return rows.filter((r) => !this.isSymbolQuarantined(r.symbol)).map((r) => r.symbol); + const out = rows.filter((r) => !this.isSymbolQuarantined(r.symbol)).map((r) => r.symbol); + if (cap == null || out.length <= cap) return out; + return out.slice(0, cap); + } + + private pendingCountForPrefix(prefix: string): number { + const row = this._db.prepare( + "SELECT COUNT(*) AS c FROM adapter_queue WHERE status IN ('pending','in_flight','backoff') AND key LIKE ?", + ).get(`${prefix}%`) as { c: number }; + return row.c; } /** @@ -959,13 +1052,18 @@ this._db.prepare("UPDATE adapter_queue SET status='backoff', last_attempt=?, ret const d = this._db; if (s.source_kind === 'sec-fetch') { - for (const sym of this.secEligibleSymbols()) { - await this.queue(`sec-fetch:fetch:${sym}`); + if (this.pendingCountForPrefix('sec-fetch:') < SEC_PENDING_CAP) { + for (const sym of this.secEligibleSymbols()) { + if (this.pendingCountForPrefix('sec-fetch:') >= SEC_PENDING_CAP) break; + await this.queue(`sec-fetch:fetch:${sym}`); + } } } else if (s.source_kind === 'sec-sc-fetch') { - // Fast SC-only path: issuer submissions feed (lightweight, no CUSIP pagination/Form 4) - for (const sym of this.secEligibleSymbols()) { - await this.queue(`sec-sc-fetch:sc:${sym}`); + if (this.pendingCountForPrefix('sec-sc-fetch:') < SEC_PENDING_CAP) { + for (const sym of this.secEligibleSymbols()) { + if (this.pendingCountForPrefix('sec-sc-fetch:') >= SEC_PENDING_CAP) break; + await this.queue(`sec-sc-fetch:sc:${sym}`); + } } } else if (s.source_kind === 'sec-tickers') { await this.queue('sec-tickers:companyTickers:latest'); @@ -992,7 +1090,7 @@ this._db.prepare("UPDATE adapter_queue SET status='backoff', last_attempt=?, ret } } else if (s.source_kind === 'yfinance-quote-watched') { // Tier 2: user watchlist symbols + recently viewed pages - for (const sym of this.demandSymbolsByTier(2)) { + for (const sym of this.demandSymbolsByTier(2, WATCHED_QUOTE_SCHEDULE_CAP)) { if (needsTieredQuoteRefresh(d, sym, 2)) { await this.queue(`yfinance:quote:${sym}`); } @@ -1021,14 +1119,9 @@ this._db.prepare("UPDATE adapter_queue SET status='backoff', last_attempt=?, ret await this.queue(`yfinance:topHoldings:${sym}`); } } - } else if (s.source_kind === 'sec-lint-holders') { - for (const sym of this.secEligibleSymbols()) { - await this.queue(`sec-lint-holders:holders:${sym}`); - } - } else if (s.source_kind === 'sec-lint-insiders') { - for (const sym of this.secEligibleSymbols()) { - await this.queue(`sec-lint-insiders:insiders:${sym}`); - } + } else if (s.source_kind === 'sec-lint-holders' || s.source_kind === 'sec-lint-insiders') { + // Universe lint is operator / focused-ticker only. Advance the clock so + // the due row does not spin every 30s. } else if (s.source_kind === 'x') { const accounts = this._db.prepare('SELECT symbol, handle FROM x_accounts').all() as Array<{ symbol: string; handle: string }>; // Tracked fund manager handles (M21 mirror) are timeline sources too — dedupe with x_accounts. @@ -1096,8 +1189,9 @@ this._db.prepare("UPDATE adapter_queue SET status='backoff', last_attempt=?, ret async requeueUnhealthySecSymbols(): Promise { // Never heal-enqueue while any SEC budget is hot — that was the thrash loop. if (this.isSecFamilyCoolingDown()) return 0; - // Reverse 13F is heavy (index + XML per filer). Keep heal volume low. - const MAX_HEAL_PER_CYCLE = 3; + if (this.pendingCountForPrefix('sec-fetch:') >= SEC_PENDING_CAP) return 0; + const lastHeal = this._db.prepare("SELECT value FROM queue_state WHERE key='sec_heal_at'").get() as { value?: string } | undefined; + if (lastHeal?.value && Date.now() - Date.parse(lastHeal.value) < SEC_HEAL_INTERVAL_MS) return 0; let enqueued = 0; try { // Seed offline CUSIPs so heal does not thrash name-search for known names. @@ -1183,7 +1277,7 @@ this._db.prepare("UPDATE adapter_queue SET status='backoff', last_attempt=?, ret candidates.sort((a, b) => a.priority - b.priority); for (const { sym } of candidates) { - if (enqueued >= MAX_HEAL_PER_CYCLE) break; + if (enqueued >= MAX_SEC_HEAL_PER_CYCLE) break; await this.queue(`sec-fetch:fetch:${sym}`); // Pair with SC so CUSIP can seed from 13G XML on the same drain window await this.queue(`sec-sc-fetch:sc:${sym}`); @@ -1196,6 +1290,7 @@ this._db.prepare("UPDATE adapter_queue SET status='backoff', last_attempt=?, ret ); } if (enqueued > 0) { + this._db.prepare("INSERT OR REPLACE INTO queue_state (key, value) VALUES ('sec_heal_at', ?)").run(new Date().toISOString()); console.log(`[queue] re-queued SEC for ${enqueued} unhealthy demand symbol(s)`); } return enqueued; @@ -1348,14 +1443,19 @@ this._db.prepare("UPDATE adapter_queue SET status='backoff', last_attempt=?, ret if (c.stopped) notes.push(`${c.source} stopped`); else if (c.paused) notes.push(`${c.source} paused`); } - if ((row.q ?? 0) > 100) notes.push(`large backlog pending=${row.q}`); if (candleLagMs != null && candleLagMs > 3 * 86_400_000) notes.push(`SPY candle lag ${Math.round(candleLagMs / 86_400_000)}d`); if (demandSize > DEMAND_SET_SOFT_CAP) notes.push(`demand set ${demandSize} > soft cap ${DEMAND_SET_SOFT_CAP}`); const yfCool = mergedCools.some((c) => c.source === 'yfinance' && c.active); + const secCool = mergedCools.some((c) => c.source.startsWith('sec') && c.active); const pendingQuotes = pendingByKind.quote ?? 0; + const pendingSec = (pendingByKind.fetch ?? 0) + (pendingByKind.sc ?? 0) + (pendingByKind.holders ?? 0) + (pendingByKind.insiders ?? 0); if (pendingQuotes > 25) notes.push(`pending quotes=${pendingQuotes}`); - const dataPlaneHealthy = !paused && !yfCool && (row.q ?? 0) < 150 && pendingQuotes < 25 && (candleLagMs == null || candleLagMs < 3 * 86_400_000); + if (pendingSec > SEC_PENDING_CAP) notes.push(`institutional backfill pending=${pendingSec}`); + const yahooHealthy = !paused && !yfCool && pendingQuotes < 25 && (candleLagMs == null || candleLagMs < 3 * 86_400_000); + const secHealthy = !paused && !secCool && pendingSec <= SEC_PENDING_CAP; + // Terminal usability is the Yahoo quote lane. A 100-job EDGAR pile is backfill. + const dataPlaneHealthy = yahooHealthy; return { queued: row.q ?? 0, @@ -1376,6 +1476,8 @@ this._db.prepare("UPDATE adapter_queue SET status='backoff', last_attempt=?, ret systemPinCount, quarantinedCount, candleLagMs, + yahooHealthy, + secHealthy, dataPlaneHealthy, dataPlaneNotes: notes, }; diff --git a/app/server/src/queue/__tests__/AdapterQueue.test.ts b/app/server/src/queue/__tests__/AdapterQueue.test.ts index 2905ab9..7811c80 100644 --- a/app/server/src/queue/__tests__/AdapterQueue.test.ts +++ b/app/server/src/queue/__tests__/AdapterQueue.test.ts @@ -391,6 +391,130 @@ test('drain fetches T0 quote before any T3 Yahoo job', async () => { assert.equal(statusOf(db, 'yfinance:quote:INTC').status, 'done'); }); +test('first subscribe quote jumps ahead of a T2 quote backlog', async () => { + const { db, fake, queue, cache } = setup(); + for (let i = 0; i < 8; i++) { + const sym = `T${i}`; + seedDemand(db, sym, 2); + fake.set(`yfinance:quote:${sym}`, { symbol: sym, price: 1 } as Quote, 'live_quote'); + await queue.queue(`yfinance:quote:${sym}`); + } + fake.set('yfinance:quote:RIVN', { symbol: 'RIVN', price: 16.5 } as Quote, 'live_quote'); + await cache.subscribe('RIVN', 'equity'); + await queue.drain(); + assert.equal(statusOf(db, 'yfinance:quote:RIVN').status, 'done'); + assert.equal(fake.calls[0], 'yfinance:quote:RIVN'); +}); + +test('drainFamily yfinance does not claim a pending SEC job', async () => { + resetVendorGateForTests(); + const db = createDb({ path: ':memory:' }); + initSchema(db); + const yf = new FakeSourceAdapter('yfinance') + .set('yfinance:quote:RIVN', { symbol: 'RIVN', price: 16 } as Quote, 'live_quote'); + const secCalls: string[] = []; + const sec = { + sourceKind: 'sec-fetch' as const, + async fetchOne(key: string) { + secCalls.push(key); + return { + value: {}, + ttlClass: 'filing_immutable' as const, + provenance: { fetchedAt: new Date().toISOString(), sourceKind: 'sec-fetch' as const }, + }; + }, + }; + const q = new AdapterQueue({ + db, + adapters: new Map([ + ['yfinance' as const, yf as unknown as import('../../adapters/SourceAdapter.ts').SourceFetch], + ['sec-fetch' as const, sec as unknown as import('../../adapters/SourceAdapter.ts').SourceFetch], + ]), + rateLimitMs: { yfinance: 0, 'sec-fetch': 0 }, + }); + q.cache = createCacheRepository({ db, scheduler: q }); + seedDemand(db, 'RIVN', 2); + await q.queue('sec-fetch:fetch:MSFT'); + await q.queue('yfinance:quote:RIVN'); + await q.drainFamily('yfinance'); + assert.equal(statusOf(db, 'yfinance:quote:RIVN').status, 'done'); + assert.equal(statusOf(db, 'sec-fetch:fetch:MSFT').status, 'pending'); + assert.deepEqual(secCalls, []); +}); + +test('health treats a SEC backlog as institutional, not Yahoo down', async () => { + const { queue } = setup(); + for (let i = 0; i < 20; i++) { + await queue.queue(`sec-fetch:fetch:S${i}`); + } + const h = queue.health(); + assert.equal(h.yahooHealthy, true); + assert.equal(h.dataPlaneHealthy, true); + assert.equal(h.secHealthy, false); + assert.ok(h.dataPlaneNotes.some((n) => /institutional backfill/i.test(n))); +}); + +test('enqueueDueSchedules does not add SEC work when pending is at cap', async () => { + const { db, queue } = setup(); + db.prepare( + 'INSERT INTO symbol_demand (symbol,refcount,ticker_kind,in_demand,last_refreshed_at,system_pin,tier) VALUES (?,?,?,?,?,?,?)', + ).run('RIVN', 1, 'equity', 1, null, 0, 2); + for (let i = 0; i < 15; i++) { + await queue.queue(`sec-fetch:fetch:S${i}`); + } + db.prepare("UPDATE queue_schedules SET next_enqueue=? WHERE source_kind='sec-fetch'") + .run(new Date(Date.now() - 1000).toISOString()); + await queue.enqueueDueSchedules(); + const rivn = db.prepare("SELECT status FROM adapter_queue WHERE key='sec-fetch:fetch:RIVN'").get() as { status?: string } | undefined; + assert.equal(rivn, undefined); +}); + +test('hung SEC fetch does not stall Yahoo quote drain', async () => { + resetVendorGateForTests(); + const db = createDb({ path: ':memory:' }); + initSchema(db); + let releaseSec: (() => void) | undefined; + const secHold = new Promise((resolve) => { releaseSec = resolve; }); + const yf = new FakeSourceAdapter('yfinance') + .set('yfinance:quote:RIVN', { symbol: 'RIVN', price: 16 } as Quote, 'live_quote'); + const sec = { + sourceKind: 'sec-fetch' as const, + async fetchOne() { + await secHold; + return { + value: {}, + ttlClass: 'filing_immutable' as const, + provenance: { fetchedAt: new Date().toISOString(), sourceKind: 'sec-fetch' as const }, + }; + }, + }; + const hungQueue = new AdapterQueue({ + db, + adapters: new Map([ + ['yfinance' as const, yf as unknown as import('../../adapters/SourceAdapter.ts').SourceFetch], + ['sec-fetch' as const, sec as unknown as import('../../adapters/SourceAdapter.ts').SourceFetch], + ]), + rateLimitMs: { yfinance: 0, 'sec-fetch': 0 }, + }); + hungQueue.cache = createCacheRepository({ db, scheduler: hungQueue }); + seedDemand(db, 'RIVN', 2); + await hungQueue.queue('sec-fetch:fetch:MSFT'); + const hung = hungQueue.drain(); + try { + await new Promise((r) => setTimeout(r, 20)); + await hungQueue.queue('yfinance:quote:RIVN'); + await hungQueue.drain(); + assert.equal( + statusOf(db, 'yfinance:quote:RIVN').status, + 'done', + 'Yahoo must complete while SEC is still in flight', + ); + } finally { + releaseSec?.(); + await hung; + } +}); + test('prioritizeQuote jumps ahead of other pending hot quotes', async () => { const { db, fake, queue } = setup(); fake.set('yfinance:quote:IREN', { symbol: 'IREN', price: 45 } as Quote, 'live_quote'); diff --git a/app/server/src/queue/sourceRatePolicy.ts b/app/server/src/queue/sourceRatePolicy.ts index ddcb264..78f9895 100644 --- a/app/server/src/queue/sourceRatePolicy.ts +++ b/app/server/src/queue/sourceRatePolicy.ts @@ -103,6 +103,18 @@ export function fetchTimeoutMs(key: string): number { /** Soft cap on demand-set size; schedule skips excess beyond system pins. */ export const DEMAND_SET_SOFT_CAP = 80; +/** Max T2 names the watched-quote schedule will enqueue in one tick. */ +export const WATCHED_QUOTE_SCHEDULE_CAP = 80; + +/** Do not dump more SEC jobs while this many are already waiting or running. */ +export const SEC_PENDING_CAP = 15; + +/** Heal-enqueue at most this often (schedule ticks are 30s). */ +export const SEC_HEAL_INTERVAL_MS = 5 * 60_000; + +/** Max institutional heal jobs per heal tick. */ +export const MAX_SEC_HEAL_PER_CYCLE = 1; + const RATE_LIMIT_RE = /too many requests|rate[- ]?limit|429|edge:\s*too many|http\s*429|throttl|quota.?exceeded|temporarily blocked|access denied|http\s*403|efts outage|request rate threshold|cool down|do not thrash|sec rate limit|preflight cool-down|vendor rate limit|yfinance rate limit|fred rate limit|finra rate limit|nasdaq rate limit|reddit rate limit/i; diff --git a/app/server/src/trpc/router.ts b/app/server/src/trpc/router.ts index e495af7..e268d49 100644 --- a/app/server/src/trpc/router.ts +++ b/app/server/src/trpc/router.ts @@ -421,9 +421,15 @@ const marketRouter = router({ `yfinance:symbol:${symbol}`, ]); - // Batch read all at once - const entries = await ctx.cache.getMany(keys); + // Read-only: watchlist polling must not requeue every T2 quote on a 60s TTL. + // Missing quotes still get one enqueue so a newly added name can land. + const entries = await ctx.cache.getMany(keys, { refresh: false }); const byKey = new Map(entries.map((e) => [e.key, e])); + for (const symbol of symbols) { + const q = byKey.get(`yfinance:quote:${symbol}`); + if (q?.value != null) continue; + try { await ctx.queue.queue(`yfinance:quote:${symbol}`); } catch { /* ignore */ } + } // Group results by symbol const results: Array<{ @@ -1552,27 +1558,27 @@ const adminRouter = router({ usersList: adminProcedure.query(({ ctx }) => listUsers(ctx.db)), setUserModules: adminProcedure - .input(z.object({ userId: z.string().uuid(), modules: z.array(z.string()) })) + .input(z.object({ userId: z.string().min(1), modules: z.array(z.string()) })) .mutation(({ ctx, input }) => { return setUserModules(ctx.db, ctx.userId, input.userId, input.modules); }), disableUser: adminProcedure - .input(z.object({ userId: z.string().uuid() })) + .input(z.object({ userId: z.string().min(1) })) .mutation(({ ctx, input }) => { try { return disableUser(ctx.db, ctx.userId, input.userId); } catch (e) { throw new TRPCError({ code: 'FORBIDDEN', message: e instanceof Error ? e.message : 'Failed to disable user.' }); } }), enableUser: adminProcedure - .input(z.object({ userId: z.string().uuid() })) + .input(z.object({ userId: z.string().min(1) })) .mutation(({ ctx, input }) => { try { return enableUser(ctx.db, ctx.userId, input.userId); } catch (e) { throw new TRPCError({ code: 'FORBIDDEN', message: e instanceof Error ? e.message : 'Failed to enable user.' }); } }), deleteUser: adminProcedure - .input(z.object({ userId: z.string().uuid() })) + .input(z.object({ userId: z.string().min(1) })) .mutation(({ ctx, input }) => { try { return deleteUser(ctx.db, ctx.userId, input.userId); } catch (e) { throw new TRPCError({ code: 'FORBIDDEN', message: e instanceof Error ? e.message : 'Failed to delete user.' }); } diff --git a/app/src/app/admin/queue/page.tsx b/app/src/app/admin/queue/page.tsx index ad911b0..6c58eb7 100644 --- a/app/src/app/admin/queue/page.tsx +++ b/app/src/app/admin/queue/page.tsx @@ -31,6 +31,8 @@ interface QueueStatus { stoppedSources?: string[]; pendingByKind?: Record; demandSize?: number; + yahooHealthy?: boolean; + secHealthy?: boolean; dataPlaneHealthy?: boolean; dataPlaneNotes?: string[]; } @@ -585,11 +587,19 @@ export default function QueuePage() { - {status && status.dataPlaneHealthy === false && (status.dataPlaneNotes?.length ?? 0) > 0 && ( + {status && status.yahooHealthy === false && ( +
+ + + Yahoo: {(status.dataPlaneNotes ?? []).filter((n) => !/institutional backfill/i.test(n)).join(" · ") || "quote lane unhealthy"} + +
+ )} + {status && status.yahooHealthy !== false && status.secHealthy === false && (
- + - Data plane: {status.dataPlaneNotes!.join(" · ")} + Institutional backfill: {(status.dataPlaneNotes ?? []).filter((n) => /institutional backfill/i.test(n)).join(" · ") || "SEC jobs queued"}
)} diff --git a/app/src/app/admin/users/page.tsx b/app/src/app/admin/users/page.tsx index d0904fd..0147bdf 100644 --- a/app/src/app/admin/users/page.tsx +++ b/app/src/app/admin/users/page.tsx @@ -261,7 +261,11 @@ function UserActions({ user, onUpdated }: { user: UserRow; onUpdated: () => void const isPending = user.status === "pending_approval"; const isDisabled = user.status === "disabled"; - const isValidId = user.id && /^[0-9a-fA-F]{8}-/.test(user.id); + const isSystemAccount = + user.id === "anonymous" || + user.id === "system" || + (user.email ?? "").toLowerCase().startsWith("anonymous@"); + const isValidId = !isSystemAccount && !!user.id && /^[0-9a-fA-F]{8}-/.test(user.id); async function doResetPassword() { setResetOpen(false); @@ -313,7 +317,10 @@ function UserActions({ user, onUpdated }: { user: UserRow; onUpdated: () => void {open && (
- {!isPending && ( + {isSystemAccount && ( + Guest book owner. Not a login. + )} + {!isPending && !isSystemAccount && ( <> @@ -323,7 +330,9 @@ function UserActions({ user, onUpdated }: { user: UserRow; onUpdated: () => void ) : ( )} - + {isValidId && ( + + )} )} {isPending && ( diff --git a/app/src/components/MobileWatchlistSheet.tsx b/app/src/components/MobileWatchlistSheet.tsx index b9a5bbb..88fafaf 100644 --- a/app/src/components/MobileWatchlistSheet.tsx +++ b/app/src/components/MobileWatchlistSheet.tsx @@ -6,6 +6,7 @@ import { X, List } from "lucide-react"; import { useActiveSymbol } from "@/stores/active-symbol-store"; import { useActiveWatchlist } from "@/stores/active-watchlist-store"; import { api, type PortfolioHolding, type WatchlistEntry } from "@/lib/trpc"; +import { useVisibilityAwarePoll } from "@/lib/useVisibilityAwarePoll"; type QuoteSnap = { price: number; changePercent: number }; @@ -52,7 +53,7 @@ export function MobileWatchlistSheet({ } }), ); - setQuotes(next); + setQuotes((prev) => ({ ...prev, ...next })); }, []); useEffect(() => { @@ -82,6 +83,24 @@ export function MobileWatchlistSheet({ }; }, [open, activeWatchlist, loadQuotes]); + const quoteSymbols = [ + ...holdings.map((h) => h.symbol), + ...entries.map((e) => e.symbol), + ]; + const quotesMissing = quoteSymbols.some((s) => quotes[s] == null); + const [fastPoll, setFastPoll] = useState(true); + useEffect(() => { + if (!open) return; + setFastPoll(true); + const t = setTimeout(() => setFastPoll(false), 45_000); + return () => clearTimeout(t); + }, [open, activeWatchlist, entries]); + useVisibilityAwarePoll( + () => { void loadQuotes(quoteSymbols); }, + quotesMissing && fastPoll ? 2500 : 15_000, + open && quoteSymbols.length > 0, + ); + useEffect(() => { if (!open) return; const onKey = (e: KeyboardEvent) => { diff --git a/app/src/components/WatchlistSidebar.tsx b/app/src/components/WatchlistSidebar.tsx index ce48855..a824e41 100644 --- a/app/src/components/WatchlistSidebar.tsx +++ b/app/src/components/WatchlistSidebar.tsx @@ -1,9 +1,10 @@ "use client"; -import { useEffect, useRef, useState } from "react"; +import { useCallback, useEffect, useRef, useState } from "react"; import { useActiveSymbol } from "@/stores/active-symbol-store"; import { useActiveWatchlist } from "@/stores/active-watchlist-store"; import { api, type WatchlistEntry, type Quote, type PortfolioHolding } from "@/lib/trpc"; import { chart as CHART } from "@/lib/chart-theme"; +import { useVisibilityAwarePoll } from "@/lib/useVisibilityAwarePoll"; import { SymbolAutocomplete } from "./SymbolAutocomplete"; interface MiniQuote extends Quote { sparkline?: number[]; @@ -149,58 +150,63 @@ export function WatchlistSidebar({ compact = false, collapsed = false, onToggle // Fetch quotes for watchlist symbols. Quote alone is enough for the mark; // candles only improve the sparkline (queue may deliver quote first). - useEffect(() => { - let cancelled = false; - - const fetchQuotes = async (): Promise> => { - const got = new Set(); - const updates = new Map(); - - try { - // Use batched snapshots to avoid N+1 API calls - const symbols = entries.map((e) => e.symbol); - const results = await api.market.snapshots(symbols); - - for (const result of results) { - if (result.quote) { - const sparkline = result.candles?.length - ? result.candles.slice(-20).map((c) => c.c) - : undefined; - updates.set(result.symbol, { ...result.quote, symbol: result.symbol, sparkline }); - got.add(result.symbol); - } + const fetchQuotes = useCallback(async () => { + if (entries.length === 0) return; + const updates = new Map(); + try { + const symbols = entries.map((e) => e.symbol); + const results = await api.market.snapshots(symbols); + for (const result of results) { + if (result.quote) { + const sparkline = result.candles?.length + ? result.candles.slice(-20).map((c) => c.c) + : undefined; + updates.set(result.symbol, { ...result.quote, symbol: result.symbol, sparkline }); } - } catch (e) { - console.log(`[watchlist] batched snapshot fail:`, (e as Error).message.slice(0, 50)); } - - if (!cancelled && updates.size > 0) { - setQuotes((prev) => { - const m = new Map(prev); - for (const [k, v] of updates) m.set(k, v); - return m; - }); - } - return got; - }; - - if (entries.length === 0) return () => { cancelled = true; }; - - void (async () => { - let got = await fetchQuotes(); - // Newly added symbols often land after a short queue delay - retry a few times. - for (let attempt = 0; attempt < 6 && !cancelled; attempt++) { - if (entries.every((e) => got.has(e.symbol))) break; - await new Promise((r) => setTimeout(r, 2500)); - if (cancelled) break; - const next = await fetchQuotes(); - got = new Set([...got, ...next]); - } - })(); - - return () => { cancelled = true; }; + } catch (e) { + console.log(`[watchlist] batched snapshot fail:`, (e as Error).message.slice(0, 50)); + } + if (updates.size > 0) { + setQuotes((prev) => { + const m = new Map(prev); + for (const [k, v] of updates) m.set(k, v); + return m; + }); + } }, [entries]); + useEffect(() => { + void fetchQuotes(); + }, [fetchQuotes]); + + // Only the visible rail polls. Compact and desktop instances both mount; + // CSS-hiding the other must not double the Yahoo enqueue path. + const [pollEnabled, setPollEnabled] = useState(false); + useEffect(() => { + const mq = compact + ? window.matchMedia("(min-width: 768px) and (max-width: 1023px)") + : window.matchMedia("(min-width: 1024px)"); + const sync = () => setPollEnabled(mq.matches); + sync(); + mq.addEventListener("change", sync); + return () => mq.removeEventListener("change", sync); + }, [compact]); + + // Burst-poll after a list change so a newly added name can land; then settle. + const [fastPoll, setFastPoll] = useState(true); + useEffect(() => { + setFastPoll(true); + const t = setTimeout(() => setFastPoll(false), 15_000); + return () => clearTimeout(t); + }, [entries]); + const quotesMissing = entries.some((e) => !quotes.has(e.symbol)); + useVisibilityAwarePoll( + fetchQuotes, + quotesMissing && fastPoll ? 2500 : 15_000, + pollEnabled && entries.length > 0, + ); + // Hybrid resolve (ADR-0011): known index symbol adds instantly; unmatched text // soft-blocks with an explicit "add anyway?" before adding (background hydration // then picks the row up via the existing demand pipeline). @@ -312,7 +318,7 @@ export function WatchlistSidebar({ compact = false, collapsed = false, onToggle const containerClass = compact ? "w-full border-b border-line bg-surface-raised flex-shrink-0" - : `hidden lg:flex items-stretch transition-all duration-200 ${collapsed ? 'w-4' : 'w-72'}`; + : `hidden lg:flex items-stretch transition-all duration-200 ${collapsed ? 'w-4' : 'w-80'}`; return (