Files
investor-flow/app/server/src/services/analystRatingsService.ts
T
Investor Flow Build 351c104eba fix: analyst ratings timeout + add accession column to institution_filings
- Add 6s server-side timeout to Yahoo Finance quoteSummary (Promise.race) so the backend responds with stale data instead of hanging indefinitely. Stale cache is served on error (router.ts:1895-1898), matching ADR-0009.
- Hoist YahooFinance to a module-level singleton (getYf()) so the crumb/cookie is fetched once, not per call. Matches YFinanceAdapter pattern.
- Add suppressNotices for yahooSurvey to reduce log noise.
- Add accession column to institution_filings (schema.sql + migration in client.ts) and thread it through both secDataFetcher.ts INSERT paths (13D/G and 13F-HR). Fixes recurring `no such column: accession` error in new13daProducer alert tick.
- new13daProducer.ts needs no changes - its SELECT accession query now works.
2026-07-25 11:52:48 -04:00

165 lines
7.2 KiB
TypeScript

import type { DatabaseSync } from 'node:sqlite';
export interface AnalystRating {
firm: string;
action: string | null;
gradeFrom: string | null;
gradeTo: string | null;
ratingDate: string;
targetFrom: number | null;
targetTo: number | null;
}
export interface AnalystConsensus {
strongBuy: number;
buy: number;
hold: number;
sell: number;
strongSell: number;
}
const RATINGS_TTL_MS = 86_400_000; // 24h
const MIN_FETCH_INTERVAL_MS = 5_000; // 5s between fetches per symbol
const RATE_LIMIT_BACKOFF_MS = 120_000; // 2min backoff after rate-limit
const YAHOO_TIMEOUT_MS = 6_000; // 6s — shorter than frontend 8s TRPC timeout
// In-memory per-symbol throttle to prevent stampeding in the absence of queue integration
const lastFetchBySymbol = new Map<string, number>();
const rateLimitUntil = new Map<string, number>();
// Singleton YahooFinance instance (crumb/cookie is fetched once, reused across calls).
let yfInstance: { quoteSummary(...args: unknown[]): Promise<Record<string, unknown>> } | null = null;
async function getYf(): Promise<{ quoteSummary(...args: unknown[]): Promise<Record<string, unknown>> }> {
if (!yfInstance) {
const mod = await import('yahoo-finance2');
yfInstance = new mod.default({ suppressNotices: ['yahooSurvey'] }) as typeof yfInstance;
}
return yfInstance;
}
async function fetchFromYahoo(symbol: string): Promise<{ ratings: AnalystRating[]; consensus: AnalystConsensus } | { error: string }> {
const now = Date.now();
// Check rate-limit backoff first (applies even if MIN_FETCH_INTERVAL would pass)
const rlUntil = rateLimitUntil.get(symbol);
if (rlUntil && rlUntil > now) {
return { error: `Rate-limited for ${symbol} — retry in ${Math.ceil((rlUntil - now) / 1000)}s` };
}
const last = lastFetchBySymbol.get(symbol);
if (last && now - last < MIN_FETCH_INTERVAL_MS) {
return { error: `Throttled — wait ${Math.ceil((MIN_FETCH_INTERVAL_MS - (now - last)) / 1000)}s before retrying ${symbol}` };
}
lastFetchBySymbol.set(symbol, now);
const yf = await getYf();
let raw: Record<string, unknown>;
let timeoutId: ReturnType<typeof setTimeout> | null = null;
try {
const yfPromise = yf.quoteSummary(symbol, { modules: ['upgradeDowngradeHistory', 'recommendationTrend'] }, { validateResult: false });
const timeoutPromise = new Promise<never>((_, reject) => {
timeoutId = setTimeout(() => reject(new Error('Yahoo Finance timed out')), YAHOO_TIMEOUT_MS);
});
raw = await Promise.race([yfPromise, timeoutPromise]) as Record<string, unknown>;
} catch (e) {
if (timeoutId) { clearTimeout(timeoutId); timeoutId = null; }
const msg = (e as Error).message;
if (/too many requests|rate[- ]?limit|429|edge:\s*too many/i.test(msg)) {
rateLimitUntil.set(symbol, Date.now() + RATE_LIMIT_BACKOFF_MS);
return { error: `Edge: Too Many Requests — yfinance rate limit hit for ${symbol}` };
}
const yfErr = e as Record<string, unknown>;
if (yfErr.errors) console.error(`[analystRatings] Schema errors for ${symbol}:`, JSON.stringify(yfErr.errors).slice(0, 500));
if (yfErr.result) console.error(`[analystRatings] Raw result for ${symbol}:`, JSON.stringify(yfErr.result).slice(0, 500));
return { error: msg };
}
const udh = raw.upgradeDowngradeHistory as Record<string, unknown> | undefined;
const hist = (udh?.history ?? []) as Array<Record<string, unknown>>;
const ratings: AnalystRating[] = hist.map((h: Record<string, unknown>) => {
let ratingDate = '';
const rawDate = h.epochGradeDate;
if (rawDate instanceof Date) ratingDate = rawDate.toISOString().slice(0, 10);
else if (typeof rawDate === 'string') ratingDate = rawDate.slice(0, 10);
else if (typeof rawDate === 'number') ratingDate = new Date(rawDate * 1000).toISOString().slice(0, 10);
return {
firm: String(h.firm ?? ''),
action: h.action ? String(h.action) : null,
gradeFrom: h.fromGrade ? String(h.fromGrade) : null,
gradeTo: h.toGrade ? String(h.toGrade) : null,
ratingDate,
targetFrom: typeof h.priorPriceTarget === 'number' ? h.priorPriceTarget : null,
targetTo: typeof h.currentPriceTarget === 'number' ? h.currentPriceTarget : null,
};
}).filter((r) => r.firm && r.ratingDate);
const rt = raw.recommendationTrend as Record<string, unknown> | undefined;
const trend = (rt?.trend ?? []) as Array<Record<string, unknown>>;
const current = trend.find((t: Record<string, unknown>) => t.period === '0m' || t.period === '0q');
const consensus: AnalystConsensus = {
strongBuy: Number(current?.strongBuy ?? 0),
buy: Number(current?.buy ?? 0),
hold: Number(current?.hold ?? 0),
sell: Number(current?.sell ?? 0),
strongSell: Number(current?.strongSell ?? 0),
};
return { ratings, consensus };
}
export function getAnalystRatings(db: DatabaseSync, symbol: string): { ratings: AnalystRating[]; consensus: AnalystConsensus | null; stale: boolean } | null {
const fresh = db.prepare(`
SELECT firm, action, grade_from, grade_to, target_from, target_to, rating_date, fetched_at
FROM analyst_ratings
WHERE symbol = ?
ORDER BY rating_date DESC, firm
`).all(symbol) as Array<{ firm: string; action: string | null; grade_from: string | null; grade_to: string | null; target_from: number | null; target_to: number | null; rating_date: string; fetched_at: string }>;
if (fresh.length === 0) return null;
const latestFetch = fresh[0].fetched_at;
const stale = Date.now() - new Date(latestFetch).getTime() > RATINGS_TTL_MS;
return {
ratings: fresh.map((r) => ({
firm: r.firm,
action: r.action,
gradeFrom: r.grade_from,
gradeTo: r.grade_to,
ratingDate: r.rating_date,
targetFrom: r.target_from,
targetTo: r.target_to,
})),
consensus: (() => {
const row = db.prepare('SELECT strong_buy, buy, hold, sell, strong_sell FROM analyst_consensus WHERE symbol = ?').get(symbol) as { strong_buy: number; buy: number; hold: number; sell: number; strong_sell: number } | undefined;
if (!row) return null;
return { strongBuy: row.strong_buy, buy: row.buy, hold: row.hold, sell: row.sell, strongSell: row.strong_sell };
})(),
stale,
};
}
export async function fetchAndStoreAnalystRatings(db: DatabaseSync, symbol: string): Promise<{ ratings: AnalystRating[]; consensus: AnalystConsensus | null; stale: boolean } | { error: string }> {
const upper = symbol.toUpperCase();
const result = await fetchFromYahoo(upper);
if ('error' in result) return { error: result.error };
const now = new Date().toISOString();
const upsert = db.prepare(`
INSERT OR REPLACE INTO analyst_ratings (symbol, firm, action, grade_from, grade_to, target_from, target_to, rating_date, fetched_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
`);
for (const r of result.ratings) {
upsert.run(upper, r.firm, r.action, r.gradeFrom, r.gradeTo, r.targetFrom, r.targetTo, r.ratingDate, now);
}
db.prepare(`
INSERT OR REPLACE INTO analyst_consensus (symbol, strong_buy, buy, hold, sell, strong_sell, fetched_at)
VALUES (?, ?, ?, ?, ?, ?, ?)
`).run(upper, result.consensus.strongBuy, result.consensus.buy, result.consensus.hold, result.consensus.sell, result.consensus.strongSell, now);
return { ratings: result.ratings, consensus: result.consensus, stale: false };
}