Isolate Yahoo drain from SEC backlog and protect the guest account.
CI / Test (push) Canceled after 0s
CI / Build and push (push) Canceled after 0s

Yahoo and SEC drain on separate workers so hung EDGAR jobs cannot freeze watchlist prices. Queue health splits yahooHealthy from institutional backfill, watchlist snapshots only enqueue missing quotes, and first subscribe seeds quote/candles/symbol. Hide the anonymous system user from Admin so it cannot be deleted.
This commit is contained in:
Investor Flow Build
2026-09-09 09:55:21 -04:00
parent 7649eaf399
commit 7426103d73
16 changed files with 583 additions and 188 deletions
@@ -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<void> {
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<unknown>("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<unknown>("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");
});
+26 -1
View File
@@ -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.
});
+19 -3
View File
@@ -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
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[];
.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) {
+33 -20
View File
@@ -627,6 +627,8 @@ const HANDLERS = new Map<string, KindHandler>([
export interface CacheRepository {
get<T>(key: CacheKey): Promise<CacheEntry<T>>;
/** Read cache without scheduling a refresh (watchlist snapshots). */
peek<T>(key: CacheKey): Promise<CacheEntry<T>>;
set<T>(key: CacheKey, value: T, ttlClass: TtlClass, provenance: Provenance): Promise<void>;
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<void>;
demandSet(): Promise<string[]>;
getMany<T>(keys: CacheKey[]): Promise<Array<{ key: CacheKey; value: T | null; isStale: boolean; fetchedAt: string | null }>>;
getMany<T>(keys: CacheKey[], opts?: { refresh?: boolean }): Promise<Array<{ key: CacheKey; value: T | null; isStale: boolean; fetchedAt: string | null }>>;
/** Delete a cache entry by key (or, for wildcard keys ending in `:*`, all matching entries). */
del(key: CacheKey): Promise<void>;
/** 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<T>(key: CacheKey): Promise<CacheEntry<T>> {
private readEntry<T>(key: CacheKey): CacheEntry<T> {
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<T>(key: CacheKey): Promise<CacheEntry<T>> {
return this.readEntry<T>(key);
}
async get<T>(key: CacheKey): Promise<CacheEntry<T>> {
const e = this.readEntry<T>(key);
if (e.isStale) {
try { await this._scheduler.queue(key); } catch { /* background refresh; never block readers */ }
}
return e;
}
async set<T>(key: CacheKey, value: T, ttlClass: TtlClass, provenance: Provenance): Promise<void> {
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<void> {
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<T>(keys: CacheKey[]): Promise<Array<{ key: CacheKey; value: T | null; isStale: boolean; fetchedAt: string | null }>> {
async getMany<T>(keys: CacheKey[], opts?: { refresh?: boolean }): Promise<Array<{ key: CacheKey; value: T | null; isStale: boolean; fetchedAt: string | null }>> {
const refresh = opts?.refresh !== false;
return Promise.all(keys.map(async (key) => {
const e = await this.get<T>(key);
const e = refresh ? await this.get<T>(key) : await this.peek<T>(key);
return { key, value: e.value, isStale: e.isStale, fetchedAt: e.provenance?.fetchedAt ?? null };
}));
}
+54 -1
View File
@@ -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<void> { this.queued.push(key); } reset() { this.queued = []; } }
class FakeScheduler {
queued: CacheKey[] = [];
prioritized: CacheKey[] = [];
async queue(key: CacheKey): Promise<void> { this.queued.push(key); }
async prioritize(key: CacheKey): Promise<void> { 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<Quote>('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<Quote>(['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<Quote>('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');
+21 -9
View File
@@ -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)`);
});
+138 -36
View File
@@ -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<SourceKind, Promise<unknown>>();
/** Header / page-view quotes jump ahead of the rest of the Yahoo pile. */
private readonly _priorityKeys = new Set<string>();
/** 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<VendorFamily>();
constructor(opts: AdapterQueueOptions) {
this._db = opts.db;
@@ -443,18 +451,59 @@ export class AdapterQueue implements CacheScheduler {
}
async drain(): Promise<void> {
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<VendorFamily>();
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<void> {
/** 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<void> {
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<Record<VendorFamily, number>> = {};
// 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');
try {
const FETCH_TIMEOUT_MS = fetchTimeoutMs(key);
let timeoutHandle: ReturnType<typeof setTimeout> | undefined;
try {
const fetchPromise = adapter.fetchOne(key);
const timeoutPromise = new Promise<never>((_, reject) =>
setTimeout(() => reject(new Error(`fetchOne timeout after ${FETCH_TIMEOUT_MS}ms for ${key}`)), FETCH_TIMEOUT_MS),
const timeoutPromise = new Promise<never>((_, 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,14 +1052,19 @@ this._db.prepare("UPDATE adapter_queue SET status='backoff', last_attempt=?, ret
const d = this._db;
if (s.source_kind === 'sec-fetch') {
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)
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');
} else if (s.source_kind === 'yfinance' || s.source_kind === 'yfinance-quote') {
@@ -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<number> {
// 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,
};
@@ -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<void>((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');
+12
View File
@@ -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;
+12 -6
View File
@@ -421,9 +421,15 @@ const marketRouter = router({
`yfinance:symbol:${symbol}`,
]);
// Batch read all at once
const entries = await ctx.cache.getMany<unknown>(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<unknown>(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.' }); }
+13 -3
View File
@@ -31,6 +31,8 @@ interface QueueStatus {
stoppedSources?: string[];
pendingByKind?: Record<string, number>;
demandSize?: number;
yahooHealthy?: boolean;
secHealthy?: boolean;
dataPlaneHealthy?: boolean;
dataPlaneNotes?: string[];
}
@@ -585,11 +587,19 @@ export default function QueuePage() {
</span>
</div>
{status && status.dataPlaneHealthy === false && (status.dataPlaneNotes?.length ?? 0) > 0 && (
{status && status.yahooHealthy === false && (
<div className="mx-4 mt-3 flex items-center gap-2 rounded border border-red-500/30 bg-red-950/20 px-3 py-1.5">
<span className="h-1.5 w-1.5 rounded-full bg-danger animate-pulse" />
<span className="text-xs text-danger font-medium">
Yahoo: {(status.dataPlaneNotes ?? []).filter((n) => !/institutional backfill/i.test(n)).join(" · ") || "quote lane unhealthy"}
</span>
</div>
)}
{status && status.yahooHealthy !== false && status.secHealthy === false && (
<div className="mx-4 mt-3 flex items-center gap-2 rounded border border-amber-500/30 bg-amber-950/20 px-3 py-1.5">
<span className="h-1.5 w-1.5 rounded-full bg-[#fbbf24] animate-pulse" />
<span className="h-1.5 w-1.5 rounded-full bg-[#fbbf24]" />
<span className="text-xs text-[#fbbf24] font-medium">
Data plane: {status.dataPlaneNotes!.join(" · ")}
Institutional backfill: {(status.dataPlaneNotes ?? []).filter((n) => /institutional backfill/i.test(n)).join(" · ") || "SEC jobs queued"}
</span>
</div>
)}
+11 -2
View File
@@ -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
</button>
{open && (
<div className="absolute right-0 mt-1 w-48 rounded-md border border-line-strong bg-surface-raised shadow-xl z-10 py-1">
{!isPending && (
{isSystemAccount && (
<span className="block px-4 py-2 text-xs text-fg-faint">Guest book owner. Not a login.</span>
)}
{!isPending && !isSystemAccount && (
<>
<button onClick={() => { setOpen(false); setResetOpen(true); }} className="block w-full text-left px-4 py-2 text-xs text-fg hover:bg-surface-overlay transition">Reset Password</button>
<button onClick={() => { setOpen(false); setModulesOpen(true); }} className="block w-full text-left px-4 py-2 text-xs text-fg hover:bg-surface-overlay transition">Manage Modules</button>
@@ -323,7 +330,9 @@ function UserActions({ user, onUpdated }: { user: UserRow; onUpdated: () => void
) : (
<button onClick={() => { setOpen(false); setDisableOpen(true); }} disabled={loading} className="block w-full text-left px-4 py-2 text-xs text-warn hover:bg-surface-overlay transition disabled:opacity-50">Disable</button>
)}
{isValidId && (
<button onClick={() => { setOpen(false); setDeleteOpen(true); }} disabled={loading} className="block w-full text-left px-4 py-2 text-xs text-error hover:bg-surface-overlay transition disabled:opacity-50">Delete</button>
)}
</>
)}
{isPending && (
+20 -1
View File
@@ -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) => {
+51 -45
View File
@@ -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<Set<string>> => {
const got = new Set<string>();
const fetchQuotes = useCallback(async () => {
if (entries.length === 0) return;
const updates = new Map<string, MiniQuote>();
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);
}
}
} catch (e) {
console.log(`[watchlist] batched snapshot fail:`, (e as Error).message.slice(0, 50));
}
if (!cancelled && updates.size > 0) {
if (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; };
}, [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 (
<aside className={containerClass}>
@@ -377,7 +383,7 @@ export function WatchlistSidebar({ compact = false, collapsed = false, onToggle
{/* Portfolio snapshot — aligned rows, one clear CTA */}
<div className="shrink-0 px-3 pt-3 pb-2">
<div className="rounded-md border border-line bg-surface-sunken/80 overflow-hidden">
<div className="flex items-center justify-between gap-2 px-3 py-2 border-b border-line">
<div className="flex items-center justify-between gap-2 px-3.5 py-2 border-b border-line">
<p className="text-[10px] font-semibold uppercase tracking-[0.12em] text-fg-muted">
Portfolio
</p>
@@ -386,25 +392,25 @@ export function WatchlistSidebar({ compact = false, collapsed = false, onToggle
</span>
</div>
<dl className="px-3 py-2.5 space-y-1.5">
<div className="grid grid-cols-[1fr_auto] items-baseline gap-x-3 gap-y-0">
<dt className="text-[11px] text-fg-muted">Cost basis</dt>
<dd className="text-[13px] font-semibold tabular-nums text-fg text-right">
<dl className="px-3.5 py-2.5 space-y-1.5">
<div className="flex items-baseline gap-3 min-w-0">
<dt className="w-[6.75rem] shrink-0 text-[11px] text-fg-muted">Cost basis</dt>
<dd className="m-0 min-w-0 text-[13px] font-semibold tabular-nums text-fg whitespace-nowrap">
${fmtUsd(Math.round(costBasis))}
</dd>
</div>
<div className="grid grid-cols-[1fr_auto] items-baseline gap-x-3">
<dt className="text-[11px] text-fg-muted">Market value</dt>
<dd className="text-[13px] font-semibold tabular-nums text-fg text-right">
<div className="flex items-baseline gap-3 min-w-0">
<dt className="w-[6.75rem] shrink-0 text-[11px] text-fg-muted">Market value</dt>
<dd className="m-0 min-w-0 text-[13px] font-semibold tabular-nums text-fg whitespace-nowrap">
{marketDisplay}
</dd>
</div>
<div className="grid grid-cols-[1fr_auto] items-baseline gap-x-3">
<dt className="text-[11px] text-fg-muted">Unrealized</dt>
<dd className={`text-[13px] font-semibold tabular-nums text-right ${pnlTone}`}>
{pnlDisplay}
<div className="flex items-baseline gap-3 min-w-0">
<dt className="w-[6.75rem] shrink-0 text-[11px] text-fg-muted">Unrealized</dt>
<dd className={`m-0 min-w-0 flex flex-wrap items-baseline gap-x-1.5 text-[13px] font-semibold tabular-nums ${pnlTone}`}>
<span className="whitespace-nowrap">{pnlDisplay}</span>
{unrealizedPnl != null && unrealizedPnlPct != null && (
<span className="ml-1.5 text-[11px] font-medium opacity-90">
<span className="text-[11px] font-medium opacity-90 whitespace-nowrap">
({unrealizedPnlPct >= 0 ? "+" : "−"}
{Math.abs(unrealizedPnlPct).toFixed(1)}%)
</span>
@@ -419,7 +425,7 @@ export function WatchlistSidebar({ compact = false, collapsed = false, onToggle
</p>
)}
<div className="px-3 py-2 border-t border-line">
<div className="px-3.5 py-2 border-t border-line">
<a
href="/portfolio"
className="block w-full text-center rounded-md border border-line bg-surface px-2 py-1.5 text-[11px] font-medium text-fg hover:border-accent/40 hover:text-accent transition-colors"
+2
View File
@@ -1826,6 +1826,8 @@ export interface QueueStatus {
stoppedSources?: string[];
pendingByKind?: Record<string, number>;
demandSize?: number;
yahooHealthy?: boolean;
secHealthy?: boolean;
dataPlaneHealthy?: boolean;
dataPlaneNotes?: string[];
}
+1 -1
View File
@@ -204,7 +204,7 @@ Admin default: every catalog type is seeded ON for `is_admin=1` (idempotent).
| Capability | Status |
|------------|--------|
| Users, sessions, reset password, enable/disable/delete | **Live** |
| Users, sessions, reset password, enable/disable/delete | **Live** | System sentinel `anonymous@investor-flow.local` (guest book owner, id `anonymous`) is hidden from the list and cannot be deleted. |
| Per-user module assignment | **Live** |
| Pending approval queue | **Live** |
| Adapter queue health, global pause/resume, **per-source pause/stop/start**, retry, schedules, error stacks | **Live** | Cooldown column ticks 429 pauses; blank is not-paused. Unhealthy Yahoo shows backlog/notes. T0/focused quotes drain first. |