""" Analyst & News Cache Builder ============================ Runs nightly via GitHub Actions. Fetches analyst consensus data and news sentiment for every US-listed stock and upserts results into the Supabase analyst_cache table. Phase 1 — Analyst data via yf.Ticker.info (~50 min for 6,500 tickers) Phase 2 — News sentiment via Google News RSS (~6 min, 20 workers) Each phase writes to Supabase incrementally every UPSERT_EVERY records so a partial run is never fully lost. """ import html import os import random import re import time import urllib.parse import xml.etree.ElementTree as ET from concurrent.futures import ThreadPoolExecutor, as_completed from datetime import datetime, timezone import requests import yfinance as yf try: from vaderSentiment.vaderSentiment import SentimentIntensityAnalyzer _VADER = SentimentIntensityAnalyzer() except ImportError: _VADER = None # --------------------------------------------------------------------------- # Supabase credentials — injected as GitHub Actions secrets # --------------------------------------------------------------------------- _SUPABASE_URL = os.environ["SUPABASE_URL"] _SUPABASE_KEY = os.environ["SUPABASE_SERVICE_KEY"] _SUPABASE_HEADERS = { "apikey": _SUPABASE_KEY, "Authorization": f"Bearer {_SUPABASE_KEY}", "Content-Type": "application/json", "Prefer": "resolution=merge-duplicates", } # --------------------------------------------------------------------------- # Constants # --------------------------------------------------------------------------- _TICKER_RE = re.compile(r'^[A-Z]{1,5}(-[A-Z]{1,2})?$') _HTTP_HEADERS = {"User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64)"} INFO_BATCH = 20 INFO_PAUSE = 4.0 MAX_WORKERS = 5 NEWS_WORKERS = 20 UPSERT_EVERY = 500 # --------------------------------------------------------------------------- # Ticker fetching (NASDAQ Trader files — same source as the screener) # --------------------------------------------------------------------------- def _clean_ticker(symbol: str) -> str | None: t = symbol.strip().upper().replace(".", "-") return t if _TICKER_RE.match(t) else None def fetch_tickers() -> list[str]: """Return all US common-stock tickers from NASDAQ Trader listing files.""" tickers = [] base = "https://www.nasdaqtrader.com/dynamic/SymDir/" # NASDAQ-listed stocks try: resp = requests.get(base + "nasdaqlisted.txt", headers=_HTTP_HEADERS, timeout=20) resp.raise_for_status() for line in resp.text.splitlines()[1:]: if line.startswith("File Creation"): break parts = line.split("|") if len(parts) < 7: continue if parts[3].strip() == "Y" or parts[6].strip() == "Y": # test/ETF continue t = _clean_ticker(parts[0].strip()) if t: tickers.append(t) print(f" NASDAQ listed : {sum(1 for _ in tickers)} tickers") except Exception as e: print(f" [WARN] nasdaqlisted.txt: {e}") # NYSE / AMEX listed stocks try: resp = requests.get(base + "otherlisted.txt", headers=_HTTP_HEADERS, timeout=20) resp.raise_for_status() count = 0 for line in resp.text.splitlines()[1:]: if line.startswith("File Creation"): break parts = line.split("|") if len(parts) < 7: continue if parts[2].strip() not in {"N", "A"}: # NYSE=N, AMEX=A only continue if parts[4].strip() == "Y" or parts[6].strip() == "Y": # ETF/test continue t = _clean_ticker(parts[0].strip()) if t: tickers.append(t) count += 1 print(f" NYSE/AMEX : {count} tickers") except Exception as e: print(f" [WARN] otherlisted.txt: {e}") combined = list(dict.fromkeys(tickers)) print(f" Total : {len(combined)} unique tickers\n") return combined # --------------------------------------------------------------------------- # Phase 1 — Analyst data # --------------------------------------------------------------------------- def _fetch_analyst(ticker: str) -> dict | None: """Fetch analyst consensus fields for one ticker via yf.Ticker.info.""" time.sleep(random.uniform(0.0, 1.0)) try: info = yf.Ticker(ticker).info if not info or not isinstance(info, dict) or len(info) < 3: return None try: rec_mean = float(info["recommendationMean"]) if info.get("recommendationMean") is not None else None except (TypeError, ValueError): rec_mean = None analyst_norm = (5.0 - rec_mean) / 4.0 if rec_mean is not None and 1 <= rec_mean <= 5 else None try: target = float(info["targetMeanPrice"]) if info.get("targetMeanPrice") is not None else None except (TypeError, ValueError): target = None price = info.get("currentPrice") or info.get("regularMarketPrice") try: price = float(price) if price is not None else None except (TypeError, ValueError): price = None analyst_upside = (target / price - 1.0) if target and price and price > 0 else None return { "ticker": ticker, "pe_forward": info.get("forwardPE"), "eps_forward": info.get("forwardEps"), "analyst_norm": analyst_norm, "analyst_upside": analyst_upside, "analyst_count": info.get("numberOfAnalystOpinions"), "analyst_target": target, "recommendation": info.get("recommendationKey") or "N/A", "sector": info.get("sector") or None, "industry": info.get("industry") or None, "updated_at": datetime.now(timezone.utc).isoformat(), } except Exception: return None def run_analyst_phase(tickers: list[str]) -> None: print(f"Phase 1: Analyst data for {len(tickers)} tickers ...") total = len(tickers) done = 0 pending: list[dict] = [] start = time.time() for batch_start in range(0, total, INFO_BATCH): batch = tickers[batch_start: batch_start + INFO_BATCH] with ThreadPoolExecutor(max_workers=MAX_WORKERS) as ex: futures = {ex.submit(_fetch_analyst, t): t for t in batch} for fut in as_completed(futures): done += 1 result = fut.result() if result: pending.append(result) elapsed = time.time() - start eta = (elapsed / done) * (total - done) if done else 0 print( f"\r {done}/{total} ({done / total * 100:.0f}%)" f" cached={len(pending)} ETA={eta:.0f}s ", end="", flush=True, ) if len(pending) >= UPSERT_EVERY: print() _upsert(pending) pending.clear() if batch_start + INFO_BATCH < total: time.sleep(INFO_PAUSE) if pending: print() _upsert(pending) print(f"\n Phase 1 complete in {time.time() - start:.0f}s\n") # --------------------------------------------------------------------------- # Phase 2 — News sentiment (Google News RSS + VADER) # --------------------------------------------------------------------------- def _fetch_news_sentiment(ticker: str) -> dict | None: """Fetch Google News RSS headlines and compute VADER compound score.""" if _VADER is None: return None try: query = urllib.parse.quote(f"{ticker} stock") rss_url = ( f"https://news.google.com/rss/search" f"?q={query}&hl=en-US&gl=US&ceid=US:en" ) resp = requests.get(rss_url, headers=_HTTP_HEADERS, timeout=8) resp.raise_for_status() root = ET.fromstring(resp.content) scores = [] for item in root.findall(".//item")[:15]: title_el = item.find("title") if title_el is not None and title_el.text: title = html.unescape( re.sub(r"<[^>]+>", " ", title_el.text) ).strip() scores.append(_VADER.polarity_scores(title)["compound"]) sentiment = sum(scores) / len(scores) if scores else None return { "ticker": ticker, "news_sentiment": sentiment, "updated_at": datetime.now(timezone.utc).isoformat(), } except Exception: return None def run_news_phase(tickers: list[str]) -> None: if _VADER is None: print("Phase 2: Skipped — vaderSentiment not installed\n") return print(f"Phase 2: News sentiment for {len(tickers)} tickers ...") total = len(tickers) done = 0 pending: list[dict] = [] start = time.time() with ThreadPoolExecutor(max_workers=NEWS_WORKERS) as ex: futures = {ex.submit(_fetch_news_sentiment, t): t for t in tickers} for fut in as_completed(futures): done += 1 result = fut.result() if result: pending.append(result) if done % 200 == 0: elapsed = time.time() - start eta = (elapsed / done) * (total - done) if done else 0 print( f"\r {done}/{total} ({done / total * 100:.0f}%)" f" ETA={eta:.0f}s ", end="", flush=True, ) if len(pending) >= UPSERT_EVERY: _upsert(pending) pending.clear() if pending: _upsert(pending) print(f"\n Phase 2 complete in {time.time() - start:.0f}s\n") # --------------------------------------------------------------------------- # Supabase upsert # --------------------------------------------------------------------------- def _upsert(rows: list[dict]) -> None: """Upsert a batch of rows into analyst_cache (merge on primary key).""" if not rows: return try: resp = requests.post( f"{_SUPABASE_URL}/rest/v1/analyst_cache", headers=_SUPABASE_HEADERS, json=rows, timeout=30, ) resp.raise_for_status() print(f" Upserted {len(rows)} rows to Supabase") except Exception as e: print(f" [WARN] Upsert failed: {e}") # --------------------------------------------------------------------------- # Entry point # --------------------------------------------------------------------------- def main() -> None: print("=" * 55) print(" ANALYST + NEWS CACHE BUILDER") print("=" * 55 + "\n") tickers = fetch_tickers() if not tickers: print("ERROR: No tickers fetched — aborting.") raise SystemExit(1) run_analyst_phase(tickers) run_news_phase(tickers) print("Cache build complete.") if __name__ == "__main__": main()