"""Helper functions for stock analysis endpoints. Extracted from fastapi_app.py to break up the monolithic router. All functions here are async-friendly (run blocking I/O in thread pools). """ import asyncio import html import json import logging import os import re import threading import time as _time from concurrent.futures import ThreadPoolExecutor from datetime import datetime, timedelta from pathlib import Path import duckdb import httpx import pandas as pd from core.shared.cache import cache_get, cache_set from core.shared.cache import cache_key as _cache_key_fn from core.shared.deps import ( data_manager, defeatbeta, finviz, smallcaplab, stockanalysis, ) from core.shared.file_io import atomic_write_json from core.splits_provider import adjust_for_splits, splits_provider # Pool global de threads para providers — evita crear/destruir threads por request _PROVIDER_EXECUTOR = ThreadPoolExecutor(max_workers=4, thread_name_prefix="provider") logger = logging.getLogger(__name__) # Re-export cache_key for use in routes cache_key = _cache_key_fn # ── NASDAQ non-compliant list cache (6h) ── NON_COMPLIANT_CACHE_FILE = Path("data/cache/nasdaq_non_compliant.json") NON_COMPLIANT_TTL = 21600 # 6 hours NON_COMPLIANT_CACHE: dict[str, dict] | None = None # Serializes the lazy NASDAQ fetch+write so two concurrent loaders # don't both hit the network and race on the cache file. _NON_COMPLIANT_LOCK = threading.Lock() def load_non_compliant_cache() -> dict[str, dict]: """Load/refresh the NASDAQ non-compliant list.""" global NON_COMPLIANT_CACHE now = _time.time() if NON_COMPLIANT_CACHE is not None: return NON_COMPLIANT_CACHE try: if NON_COMPLIANT_CACHE_FILE.exists(): mtime = os.path.getmtime(NON_COMPLIANT_CACHE_FILE) if (now - mtime) < NON_COMPLIANT_TTL: with open(NON_COMPLIANT_CACHE_FILE) as f: data = json.load(f) NON_COMPLIANT_CACHE = data return data except Exception: pass try: url = "https://api.nasdaq.com/api/quote/list-type-extended/listing?queryString=deficient" headers = { "origin": "https://www.nasdaq.com", "referer": "https://www.nasdaq.com/", "user-agent": "Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36", } resp = httpx.get(url, headers=headers, timeout=15) resp.raise_for_status() payload = resp.json() rows = payload.get("data", {}).get("noncomplaintCompanyList", {}).get("rows", []) lookup: dict[str, dict] = {} for issuer in rows: issuer_name = issuer.get("IssuerName", "") for company_entry in issuer.get("companies", []): info = { "issuer_name": issuer_name, "deficiency": company_entry.get("Deficiency", ""), "market": company_entry.get("Market", ""), "notification_date": company_entry.get("NotificationDate", ""), } for sym in company_entry.get("AffectedIssues", []): if sym: lookup[sym.upper()] = info with _NON_COMPLIANT_LOCK: # Re-check: a concurrent loader may have already written it. if NON_COMPLIANT_CACHE is None: NON_COMPLIANT_CACHE_FILE.parent.mkdir(parents=True, exist_ok=True) atomic_write_json(NON_COMPLIANT_CACHE_FILE, lookup) NON_COMPLIANT_CACHE = lookup logger.info(f"Loaded NASDAQ non-compliant list: {len(lookup)} securities") return lookup except Exception as e: logger.warning(f"Failed to fetch NASDAQ non-compliant list: {e}") NON_COMPLIANT_CACHE = {} return {} def non_compliant_info(symbol: str) -> dict | None: """Return non-compliant info for symbol, or None if compliant/unknown.""" lookup = load_non_compliant_cache() return lookup.get(symbol.upper()) # ── Executive title filters ── EXECUTIVE_TITLE_FILTERS = [ "chairman", "chairwoman", "board of director", "board member", "independent director", "lead independent", "managing director", "ceo", "chief executive", "cfo", "chief financial", "coo", "chief operating", "cto", "chief technology", "cio", "chief information", "cmo", "chief marketing", "chro", "chief human", "chief scientific", "chief medical", "chief compliance", "chief accounting", "chief legal", "chief admin", "cso", "chief strategy", "chief revenue", "chief growth", "president", "co-president", "vice president", "svp", "senior vice president", "evp", "executive vice president", "vp", "founder", "secretary", "treasurer", "general counsel", "managing member", ] def title_matches(title: str) -> bool: tl = title.lower().strip() return any(kw in tl for kw in EXECUTIVE_TITLE_FILTERS) def extract_board(company_officers: list[dict]) -> list[dict]: if not company_officers: return [] seen: set[str] = set() result: list[dict] = [] for o in company_officers: raw_name = o.get("name", "") raw_title = o.get("title", "") if not raw_name: continue title = raw_title.strip() if raw_title else "" if title and not title_matches(title): continue name = raw_name.strip() if name.lower() in seen: continue seen.add(name.lower()) role = title if title_matches(title) else raw_title.strip() result.append({"role": role or "", "name": name}) return result # ── SEC dilution risk ── DILUTION_FORM_TYPES = { "S-3", "S-3/A", "S-1", "S-1/A", "424B1", "424B2", "424B3", "424B4", "424B5", "D", "D/A", "S-8", "S-8/A", } WARRANT_KEYWORDS = ["warrant", "warrants", "Warrant", "Warrants"] ATM_KEYWORDS = [ "at the market", "atm offering", "ATM", "equity distribution", "sales agreement", "controlled equity", ] def fetch_sec_dilution_risk(symbol: str) -> dict: try: import httpx cik = None try: with duckdb.connect(":memory:") as con: meta_path = data_manager.metadata_db_path con.execute(f"ATTACH '{meta_path}' AS meta (TYPE DUCKDB)") res = con.execute("SELECT cik FROM meta.tickers WHERE symbol = ?", [symbol.upper()]).fetchone() if res and res[0]: cik = str(res[0]).strip().zfill(10) except Exception: pass if not cik: try: headers = { "User-Agent": "Stock Scanner (email@example.com)", "Accept": "application/json", } resp = httpx.get( "https://www.sec.gov/files/company_tickers.json", headers=headers, timeout=10, ) resp.raise_for_status() tickers = resp.json() for entry in tickers.values(): if entry.get("ticker", "").upper() == symbol.upper(): cik = str(entry["cik_str"]).zfill(10) break except Exception: pass if not cik: return {} headers = { "User-Agent": "Stock Scanner (email@example.com)", "Accept": "application/json", } url = f"https://data.sec.gov/submissions/CIK{cik}.json" resp = httpx.get(url, headers=headers, timeout=15) resp.raise_for_status() data = resp.json() filings_data = [] has_warrants = False has_atm = False has_s3 = False recent = data.get("filings", {}).get("recent", {}) forms = recent.get("form", []) dates = recent.get("filingDate", []) descs = recent.get("description", []) for i, form in enumerate(forms): if form in DILUTION_FORM_TYPES: date = dates[i] if i < len(dates) else "" desc = (descs[i] if i < len(descs) else "") or "" clean_desc = html.unescape(re.sub(r"<[^>]+>", "", desc)) if desc else "" filings_data.append( { "form": form, "date": date, "description": clean_desc[:300] if clean_desc else "", } ) if form.startswith("S-3"): has_s3 = True desc = (descs[i] if i < len(descs) else "") or "" if desc: desc_lower = desc.lower() for kw in WARRANT_KEYWORDS: if kw.lower() in desc_lower: has_warrants = True break for kw in ATM_KEYWORDS: if kw.lower() in desc_lower: has_atm = True break return { "filings": filings_data[:10], "has_warrants": has_warrants, "has_atm": has_atm, "active_shelf": has_s3, } except Exception as e: logger.debug(f"SEC dilution risk failed for {symbol}: {e}") return {} # ── SEC resale-unlock (EFFECT) extraction ── RESALE_REG_TYPES = {"S-1", "S-1/A", "S-3", "S-3/A", "F-1", "F-3", "F-1/A", "F-3/A"} EFFECT_FORM = "EFFECT" EFFECT_WINDOW_DAYS = 21 # unlock within N days = distribution-danger window LOCKUP_DEFAULT_DAYS = 180 # typical placement/insider lock-up length LOCKUP_WINDOW_DAYS = 30 # |days to est. expiry| within N = volatility window def fetch_sec_unlocks(symbol: str) -> dict | None: """Return resale-registration unlock info, or None. Detects resale registration statements (S-1/S-3/F-1/F-3 + /A) and their EFFECT filing (filingDate == effectiveness/unlock date). Mirrors the CIK + submissions plumbing already in fetch_sec_dilution_risk. """ try: from datetime import datetime as _dt from datetime import timedelta import httpx cik = None try: with duckdb.connect(":memory:") as con: meta_path = data_manager.metadata_db_path con.execute(f"ATTACH '{meta_path}' AS meta (TYPE DUCKDB)") res = con.execute("SELECT cik FROM meta.tickers WHERE symbol = ?", [symbol.upper()]).fetchone() if res and res[0]: cik = str(res[0]).strip().zfill(10) except Exception: pass if not cik: try: headers = { "User-Agent": "Stock Scanner (email@example.com)", "Accept": "application/json", } resp = httpx.get( "https://www.sec.gov/files/company_tickers.json", headers=headers, timeout=10, ) resp.raise_for_status() tickers = resp.json() for entry in tickers.values(): if entry.get("ticker", "").upper() == symbol.upper(): cik = str(entry["cik_str"]).zfill(10) break except Exception: pass if not cik: return None headers = { "User-Agent": "Stock Scanner (email@example.com)", "Accept": "application/json", } url = f"https://data.sec.gov/submissions/CIK{cik}.json" resp = httpx.get(url, headers=headers, timeout=15) resp.raise_for_status() data = resp.json() recent = data.get("filings", {}).get("recent", {}) forms = recent.get("form", []) dates = recent.get("filingDate", []) descs = recent.get("description", []) resale_forms = [] effect_dates = [] has_lockup = False lockup_days = None for i, form in enumerate(forms): d = dates[i] if i < len(dates) else "" if form in RESALE_REG_TYPES: resale_forms.append({"form": form, "date": d}) if form == EFFECT_FORM and d: effect_dates.append(d) # lock-up detection from filing descriptions (prose heuristic) desc = (descs[i] if i < len(descs) else "") or "" if desc: dl = desc.lower() if "lock" in dl: has_lockup = True if lockup_days is None: m = re.search(r"(\d{1,3})\s*-?\s*day", dl) if m: lockup_days = int(m.group(1)) if not resale_forms and not effect_dates: return None last_effect = max(effect_dates) if effect_dates else None days_since = None if last_effect: try: # datetime.strptime returns datetime; .date() -> date; now().date() - date -> days eff_date = _dt.strptime(last_effect, "%Y-%m-%d").date() days_since = (_dt.now().date() - eff_date).days except Exception: pass resale_dates = [rf["date"] for rf in resale_forms if rf["date"]] # pending = a resale registration filed with no effective EFFECT yet has_pending = bool(resale_forms) and (last_effect is None or any(rd > last_effect for rd in resale_dates)) in_window = days_since is not None and 0 <= days_since <= EFFECT_WINDOW_DAYS # Lock-up estimate: anchored on the resale-reg EFFECT, length from # filing prose if detectable else typical 180d. Clearly an ESTIMATE. lockup_est_end = None days_to_lockup_est = None in_lockup_est_window = False if last_effect: try: base = _dt.strptime(last_effect, "%Y-%m-%d").date() end = base + timedelta(days=lockup_days or LOCKUP_DEFAULT_DAYS) lockup_est_end = end.isoformat() days_to_lockup_est = (_dt.now().date() - end).days in_lockup_est_window = abs(days_to_lockup_est) <= LOCKUP_WINDOW_DAYS except Exception: pass return { "has_pending_registration": has_pending, "last_effect_date": last_effect, "days_since_effect": days_since, "resale_forms": [rf["form"] for rf in resale_forms], "in_unlock_window": in_window, "has_lockup": has_lockup, "lockup_est_end_date": lockup_est_end, "days_to_lockup_est": days_to_lockup_est, "in_lockup_est_window": in_lockup_est_window, } except Exception as e: logger.debug(f"SEC unlock check failed for {symbol}: {e}") return None def fetch_press_release(symbol: str) -> list[dict]: try: import yfinance as yf ticker = yf.Ticker(symbol.upper()) news = ticker.get_news(count=1, tab="press releases") if not news: return [] item = news[0].get("content", {}) title = item.get("title", "") pub_date = item.get("pubDate", "") click_url = item.get("clickThroughUrl", {}) url = click_url.get("url", "") if isinstance(click_url, dict) else "" if title: return [{"headline": title, "url": url, "created_at": pub_date}] return [] except Exception: return [] def parse_finviz_financial_number(val: str | None) -> float | None: if not val or val == "-" or val == "": return None try: from core.finviz_provider import safe_finviz_number return safe_finviz_number(val) except Exception: pass try: val_str = str(val).upper().strip() if "B" in val_str: return float(val_str.replace("B", "")) * 1_000_000_000 if "M" in val_str: return float(val_str.replace("M", "")) * 1_000_000 if "K" in val_str: return float(val_str.replace("K", "")) * 1_000 if "T" in val_str: return float(val_str.replace("T", "")) * 1_000_000_000_000 return float(val_str.replace("%", "")) except (ValueError, TypeError): return None def parse_earnings_date(raw: str) -> str | None: if not raw or raw == "-": return None clean = raw.split()[0:2] if len(clean) < 2: return None try: dt = datetime.strptime(f"{clean[0]} {clean[1]}", "%b %d") now = datetime.now() # Check both current and next year to handle year boundaries correctly. # e.g., on Dec 30, "Jan 5" should map to next year's Jan 5, not this year's. candidates = [ dt.replace(year=now.year), dt.replace(year=now.year + 1), ] # Pick the closest future date future = [c for c in candidates if c > now] if not future: return None return min(future).strftime("%Y-%m-%d") except ValueError: return None def compute_sma( daily_df: pd.DataFrame, period: int, symbol: str | None = None, data_manager=None, ) -> float | None: """Compute SMA over the last `period` closes from split-adjusted daily data. Returns None when daily_df is empty or has fewer than `period` rows. SQL fallback removed — raw unadjusted data gives wrong results for reverse-split tickers. """ if daily_df is not None and not daily_df.empty: closes = pd.to_numeric(daily_df.get("close", pd.Series(dtype=float)), errors="coerce").dropna() if len(closes) >= period: return float(closes.tail(period).mean()) return None def get_52w_high(daily_df: pd.DataFrame) -> float | None: if daily_df.empty: return None highs = pd.to_numeric(daily_df.get("high", pd.Series()), errors="coerce").dropna() if len(highs) < 252: return float(highs.max()) if not highs.empty else None return float(highs.tail(252).max()) def get_corporate_actions_dict(symbol: str) -> dict: try: df = data_manager.get_corporate_actions(symbol, days=365) if df.empty: return {} result: dict = {} for _, row in df.iterrows(): ca_type = row.get("ca_sub_type", "") if not ca_type: continue entry = {"ex_date": str(row.get("ex_date", ""))} ca_ratio = row.get("ca_ratio", "") ca_amount = row.get("ca_amount", "") if ca_ratio and pd.notna(ca_ratio) and ca_ratio != "": entry["ca_ratio"] = str(ca_ratio) if ca_amount and pd.notna(ca_amount): entry["ca_amount"] = float(ca_amount) result.setdefault(ca_type, []).append(entry) return result except Exception as e: logger.warning(f"Corporate actions fetch failed for {symbol}: {e}") return {} def compute_resistance_levels( daily_df: pd.DataFrame, price_now: float, symbol: str | None = None, data_manager=None, ) -> list[dict]: """Compute overhead resistance levels from split-adjusted daily data. Returns empty list when daily_df is empty or has insufficient data. SQL fallback removed — raw unadjusted data gives wrong results for reverse-split tickers. """ if daily_df is None or daily_df.empty: return [] df = daily_df.copy() df["date"] = pd.to_datetime(df["timestamp"]).dt.date window_start = datetime.now().date() - timedelta(days=63) window = df[df["date"] >= window_start] levels = [] if len(window) >= 10: window["vol_avg"] = window["volume"].rolling(20, min_periods=5).mean() above = window[window["high"] > price_now] spiked = above[above["volume"] > 2 * above["vol_avg"]] if not above.empty else above candidate = spiked.iloc[-1] if not spiked.empty else (above.iloc[-1] if not above.empty else None) if candidate is not None: vol_ratio = ( round(float(candidate["volume"]) / float(candidate["vol_avg"]), 1) if pd.notna(candidate["vol_avg"]) and candidate["vol_avg"] > 0 else None ) levels.append( { "level": "OH", "label": "Overhead", "price": float(candidate["high"]), "volume": int(candidate["volume"]), "vol_ratio": vol_ratio, "date": candidate["date"].isoformat(), "source": "daily_adjusted", } ) levels.sort(key=lambda x: x["price"]) return levels def compute_support_levels( daily_df: pd.DataFrame, price_now: float, symbol: str | None = None, data_manager=None, ) -> list[dict]: """Compute support levels below current price from split-adjusted daily data. Returns empty list when daily_df is empty or has insufficient data. SQL fallback removed — raw unadjusted data gives wrong results for reverse-split tickers. """ if daily_df is None or daily_df.empty: return [] df = daily_df.copy() df["date"] = pd.to_datetime(df["timestamp"]).dt.date window_start = datetime.now().date() - timedelta(days=63) window = df[df["date"] >= window_start] levels = [] if len(window) >= 10: window["vol_avg"] = window["volume"].rolling(20, min_periods=5).mean() below = window[window["low"] < price_now] spiked = below[below["volume"] > 2 * below["vol_avg"]] if not below.empty else below candidate = spiked.iloc[-1] if not spiked.empty else (below.iloc[-1] if not below.empty else None) if candidate is not None: vol_ratio = ( round(float(candidate["volume"]) / float(candidate["vol_avg"]), 1) if pd.notna(candidate["vol_avg"]) and candidate["vol_avg"] > 0 else None ) levels.append( { "level": "SL", "label": "Support", "price": float(candidate["low"]), "volume": int(candidate["volume"]), "vol_ratio": vol_ratio, "date": candidate["date"].isoformat(), "source": "daily_adjusted", } ) levels.sort(key=lambda x: x["price"]) return levels def fetch_duckdb_data(symbol: str) -> dict: """Fetch all DuckDB data for a symbol in one shot (gap/run history).""" daily_df = data_manager.get_all_daily_data(symbol) gap_history: list = [] run_history: list = [] if not daily_df.empty: try: daily_df = daily_df.copy() prev_close = daily_df["close"].shift(1) daily_df["gap_pct"] = (daily_df["open"] - prev_close) / prev_close.replace(0, float("nan")) daily_df["run_pct"] = (daily_df["close"] - daily_df["open"]) / daily_df["open"].replace(0, float("nan")) daily_df["change_pct"] = (daily_df["close"] - prev_close) / prev_close.replace(0, float("nan")) except Exception as e: logger.warning(f"Column computation failed for {symbol}: {e}") try: mask = (daily_df["gap_pct"] >= 0.20) & (daily_df["volume"] > 1_000_000) gh = daily_df[mask].copy() if not gh.empty: gh = gh.sort_values("timestamp", ascending=False).head(10) gh["date"] = gh["timestamp"].dt.strftime("%Y-%m-%d") gh["gap_pct"] = (gh["gap_pct"] * 100).round(1) gh["vol_M"] = (gh["volume"] / 1_000_000).round(1) gh["return_pct"] = (gh["change_pct"] * 100).round(1) gh["close_dir"] = gh["close"].gt(gh["open"]).map({True: "UP", False: "DOWN"}) gap_history = gh[["date", "gap_pct", "vol_M", "return_pct", "close_dir"]].to_dict(orient="records") except Exception as e: logger.warning(f"Gap history failed for {symbol}: {e}") try: mask = (daily_df["run_pct"] >= 0.50) & (daily_df["volume"] > 1_000_000) rh = daily_df[mask].copy() if not rh.empty: rh = rh.sort_values("timestamp", ascending=False).head(10) rh["date"] = rh["timestamp"].dt.strftime("%Y-%m-%d") rh["run_pct"] = (rh["run_pct"] * 100).round(1) rh["vol_M"] = (rh["volume"] / 1_000_000).round(1) rh["return_pct"] = (rh["change_pct"] * 100).round(1) rh["close_dir"] = rh["close"].gt(rh["open"]).map({True: "UP", False: "DOWN"}) run_history = rh[["date", "run_pct", "vol_M", "return_pct", "close_dir"]].to_dict(orient="records") except Exception as e: logger.warning(f"Run history failed for {symbol}: {e}") return {"daily_df": daily_df, "gap_history": gap_history, "run_history": run_history} def build_stock_result( symbol: str, finviz_overview: dict, corporate_actions: dict, next_earnings_date: str | None, daily_df: pd.DataFrame, gap_history: list, run_history: list, finviz_catalyst: dict | None = None, employees: int | None = None, company_description: str = "", float_finviz: int = 0, company_officers: list | None = None, nasdaq_non_compliant: dict | None = None, sec_unlocks: dict | None = None, financials: dict | None = None, dilution_risk: dict | None = None, stockanalysis_data: dict | None = None, smallcaplab_data: dict | None = None, press_releases: list | None = None, ) -> dict: from core.chinese_profile import score_chinese_profile from core.theme_inference import infer_themes finviz_catalyst = finviz_catalyst or {"symbol": symbol, "has_catalyst": False} company_officers = company_officers or finviz_overview.get("company_officers", []) reverse_splits = corporate_actions.get("reverse_split", []) if isinstance(corporate_actions, dict) else [] last_reverse_split = reverse_splits[0] if reverse_splits else None forward_splits = corporate_actions.get("forward_split", []) if isinstance(corporate_actions, dict) else [] last_forward_split = forward_splits[0] if forward_splits else None price_now = finviz_overview.get("price", 0) or 0 resistance_levels = compute_resistance_levels(daily_df, price_now, symbol, data_manager) for key, label in [("52w_high", "52W High"), ("sma_200", "SMA 200d")]: val = finviz_overview.get(key) if val and float(val) > price_now: resistance_levels.append( { "level": key, "label": label, "price": float(val), "source": "daily_adjusted", } ) resistance_levels.sort(key=lambda x: x["price"]) oh_level = next((r for r in resistance_levels if r.get("level") == "OH"), None) overhead_resistance = oh_level or (resistance_levels[0] if resistance_levels else None) if overhead_resistance: overhead_resistance = { "price": overhead_resistance["price"], "volume": overhead_resistance.get("volume"), "date": overhead_resistance.get("date"), } support_levels = compute_support_levels(daily_df, price_now, symbol, data_manager) sma_200_val = finviz_overview.get("sma_200") if sma_200_val and float(sma_200_val) < price_now: support_levels.append( { "level": "sma_200", "label": "SMA 200d", "price": float(sma_200_val), "source": "daily_adjusted", } ) support_levels.sort(key=lambda x: x["price"]) nearest_support = support_levels[-1] if support_levels else None if nearest_support: nearest_support = { "price": nearest_support["price"], "volume": nearest_support.get("volume"), "date": nearest_support.get("date"), } volume_today = finviz_overview.get("volume", 0) or 0 avg_vol = finviz_overview.get("avg_volume", 0) or 0 rel_vol = round(volume_today / avg_vol, 1) if avg_vol else None prev_close = finviz_overview.get("previous_close", 0) or 0 gap_today = round((price_now - prev_close) / prev_close * 100, 1) if prev_close else None inst_own = finviz_overview.get("inst_own") company_data = { "country": finviz_overview.get("country", ""), "city": finviz_overview.get("city", ""), "phone": finviz_overview.get("phone", ""), "business_description": company_description or finviz_overview.get("business_description", ""), "officers": company_officers, } chinese_profile = score_chinese_profile(company_data) themes = infer_themes(company_data["business_description"]) return { "catalysts": { "finviz": finviz_catalyst, "next_earnings_date": next_earnings_date, }, "company": { "name": finviz_overview.get("company_name"), "sector": finviz_overview.get("sector"), "industry": finviz_overview.get("industry"), "country": finviz_overview.get("country"), "city": finviz_overview.get("city"), "phone": finviz_overview.get("phone"), "employees": finviz_overview.get("employees") or employees, "business_description": company_description or finviz_overview.get("business_description", ""), "chinese_profile": chinese_profile, "themes": themes, "nasdaq_non_compliant": nasdaq_non_compliant, "sec_unlocks": sec_unlocks, }, "metrics": { "price": finviz_overview.get("price"), "volume": finviz_overview.get("volume"), "avg_volume": finviz_overview.get("avg_volume"), "market_cap": finviz_overview.get("market_cap"), "float_shares": finviz_overview.get("float_shares") or 0, "float_shares_finviz": float_finviz or 0, "rel_vol": rel_vol, "gap_pct": gap_today, "gap_history": gap_history, "run_history": run_history, "employees": employees, "inst_own": inst_own, "overhead_resistance": overhead_resistance, "resistance_levels": resistance_levels, "nearest_support": nearest_support, "support_levels": support_levels, }, "corporate_actions": { "actions": corporate_actions, "last_reverse_split": last_reverse_split, "last_forward_split": last_forward_split, }, "board": extract_board(company_officers), "financials": { "net_income": (financials or {}).get("net_income"), "operating_cash_flow": ( (stockanalysis_data or {}).get("operating_cash_flow") or (financials or {}).get("operating_cash_flow") ), "financing_cash_flow": (financials or {}).get("financing_cash_flow"), }, "stockanalysis": stockanalysis_data or {}, "smallcaplab": smallcaplab_data or {}, "dilution_risk": dilution_risk or {}, "scl_dilution": (smallcaplab_data or {}).get("dilution_risk"), "scl_survival": (smallcaplab_data or {}).get("survival"), "press_releases": press_releases or [], } async def get_stock_catalysts(symbol: str, fast: bool = False) -> dict: """Stock catalysts with layered data fetching. Layer 1 — Cache Layer 2 — Fast: DuckDB + Finviz + StockAnalysis + SmallCapLab Layer 3 — Background: Yahoo board + cashflow + SEC When fast=True: skip StockAnalysis, SmallCapLab, press releases, insider, and catalyst for a faster response. Result is not cached (the full call handles caching). """ if not fast: cached = cache_get(_cache_key_fn(symbol)) if cached is not None: return cached timeout = 5 loop = asyncio.get_event_loop() # Fast tasks (always) duckdb_task = loop.run_in_executor(_PROVIDER_EXECUTOR, fetch_duckdb_data, symbol) finviz_quote_task = loop.run_in_executor(_PROVIDER_EXECUTOR, finviz.fetch_quote_page, symbol) # Slow tasks (skipped in fast mode) finviz_insider_task = None finviz_catalyst_task = None sa_task = None scl_task = None news_task = None if not fast: finviz_insider_task = loop.run_in_executor(_PROVIDER_EXECUTOR, finviz.fetch_insider_officers, symbol) finviz_catalyst_task = loop.run_in_executor(_PROVIDER_EXECUTOR, finviz.fetch_catalyst, symbol) sa_task = loop.run_in_executor(_PROVIDER_EXECUTOR, stockanalysis.fetch_statistics, symbol) scl_task = loop.run_in_executor(_PROVIDER_EXECUTOR, smallcaplab.fetch, symbol) news_task = loop.run_in_executor(_PROVIDER_EXECUTOR, fetch_press_release, symbol) duckdb_data = await duckdb_task finviz_quote = {} catalyst = {} company_officers = [] press_releases: list[dict] = [] try: finviz_quote = await asyncio.wait_for(finviz_quote_task, timeout=timeout) except Exception as e: logger.warning(f"Finviz quote error for {symbol}: {e}") if not fast: try: company_officers = await asyncio.wait_for(finviz_insider_task, timeout=timeout) except Exception as e: logger.warning(f"Finviz insider error for {symbol}: {e}") try: cat = await asyncio.wait_for(finviz_catalyst_task, timeout=timeout) if cat and cat.get("has_catalyst"): catalyst = cat except Exception as e: logger.warning(f"Finviz catalyst error for {symbol}: {e}") try: press_releases = await asyncio.wait_for(news_task, timeout=timeout) except Exception as e: logger.warning(f"Press releases error for {symbol}: {e}") daily_df = duckdb_data.get("daily_df", pd.DataFrame()) # Split-adjust OHLC so 52w-high / SMA / support / resistance are computed # on a consistent share-count scale (DefeatBeta-sourced splits). daily_df = adjust_for_splits(daily_df, splits_provider.get_splits(symbol)) sma_200 = compute_sma(daily_df, 200, symbol, data_manager) sma_50 = compute_sma(daily_df, 50, symbol, data_manager) high_52w = get_52w_high(daily_df) if high_52w is None: high_52w = finviz_quote.get("52w_high") price = finviz_quote.get("price") or (float(daily_df.iloc[-1]["close"]) if not daily_df.empty else 0) prev_close = finviz_quote.get("previous_close") or price finviz_overview = { "symbol": symbol, "company_name": finviz_quote.get("company_name", ""), "market_cap": finviz_quote.get("market_cap", 0), "price": price, "previous_close": prev_close, "volume": finviz_quote.get("volume", 0), "avg_volume": finviz_quote.get("avg_volume", 0), "float_shares": finviz_quote.get("float_shares", 0), "sector": "", "industry": "", "country": "", "city": "", "phone": "", "company_officers": company_officers, "employees": finviz_quote.get("employees", 0), "business_description": finviz_quote.get("company_description", ""), "52w_high": high_52w, "sma_200": sma_200, "sma_50": sma_50, "inst_own": finviz_quote.get("inst_own"), } try: md = data_manager.metadata_df if md is not None and not md.empty: row = md[md["symbol"] == symbol] if not row.empty: r = row.iloc[0] for col in ("sector", "industry", "country"): val = r.get(col) finviz_overview[col] = val if isinstance(val, str) and val else "" except Exception: pass gap_history = duckdb_data.get("gap_history", []) run_history = duckdb_data.get("run_history", []) corporate_actions = get_corporate_actions_dict(symbol) next_earnings_date = parse_earnings_date(finviz_quote.get("earnings_date_raw", "")) nasdaq_non_compliant = None if not fast: try: nc_task = loop.run_in_executor(_PROVIDER_EXECUTOR, non_compliant_info, symbol) nc_info = await asyncio.wait_for(nc_task, timeout=timeout) if nc_info: nasdaq_non_compliant = nc_info except Exception as e: logger.warning(f"NASDAQ non-compliant error for {symbol}: {e}") sec_unlocks = None if not fast: try: su_task = loop.run_in_executor(_PROVIDER_EXECUTOR, fetch_sec_unlocks, symbol) su_info = await asyncio.wait_for(su_task, timeout=timeout) if su_info: sec_unlocks = su_info except Exception as e: logger.warning(f"SEC unlock check error for {symbol}: {e}") net_income = parse_finviz_financial_number(finviz_quote.get("snapshot_Income")) financials = {} dilution_risk = {} if net_income is not None: financials["net_income"] = net_income sa_data = None if not fast: try: sa_result = await asyncio.wait_for(sa_task, timeout=timeout) if sa_result.operating_cash_flow is not None: sa_data = { "operating_cash_flow": sa_result.operating_cash_flow, "free_cash_flow": sa_result.free_cash_flow, "net_income": sa_result.net_income, "revenue": sa_result.revenue, "cash_and_equivalents": sa_result.cash_and_equivalents, "total_debt": sa_result.total_debt, "net_cash": sa_result.net_cash, "book_value": sa_result.book_value, "book_value_per_share": sa_result.book_value_per_share, "pe_ratio": sa_result.pe_ratio, "ps_ratio": sa_result.ps_ratio, "pb_ratio": sa_result.pb_ratio, "current_ratio": sa_result.current_ratio, "debt_equity": sa_result.debt_equity, "roe": sa_result.roe, "roic": sa_result.roic, "shares_outstanding": sa_result.shares_outstanding, "shares_change_yoy": sa_result.shares_change_yoy, "insider_percent": sa_result.insider_percent, "institution_percent": sa_result.institution_percent, "float_shares": sa_result.float_shares, "short_interest": sa_result.short_interest, "short_pct_float": sa_result.short_pct_float, "short_pct_shares": sa_result.short_pct_shares, "short_ratio": sa_result.short_ratio, "altman_z_score": sa_result.altman_z_score, "piotroski_f_score": sa_result.piotroski_f_score, "earnings_date": sa_result.earnings_date, "gross_margin": sa_result.gross_margin, "profit_margin": sa_result.profit_margin, "operating_margin": sa_result.operating_margin, "beta": sa_result.beta, "employee_count": sa_result.employee_count, "market_cap": sa_result.market_cap, "enterprise_value": sa_result.enterprise_value, "working_capital": sa_result.working_capital, } if sa_result.earnings_date and not next_earnings_date: next_earnings_date = sa_result.earnings_date if financials.get("operating_cash_flow") is None: financials["operating_cash_flow"] = sa_result.operating_cash_flow except Exception as e: logger.warning(f"StockAnalysis error for {symbol}: {e}") scl_data = None if not fast: try: scl_result = await asyncio.wait_for(scl_task, timeout=timeout) if scl_result.dilution_risk or scl_result.survival or scl_result.dilution_activity: scl_data = { "dilution_risk": { "rating": scl_result.dilution_risk.rating if scl_result.dilution_risk else None, "structural_score": scl_result.dilution_risk.structural_score if scl_result.dilution_risk else None, "near_term_score": scl_result.dilution_risk.near_term_score if scl_result.dilution_risk else None, "severity_score": scl_result.dilution_risk.severity_score if scl_result.dilution_risk else None, "has_atm_drip": scl_result.dilution_risk.has_atm_drip if scl_result.dilution_risk else None, } if scl_result.dilution_risk else None, "survival": { "cash_on_hand": scl_result.survival.cash_on_hand if scl_result.survival else None, "quarterly_burn": scl_result.survival.quarterly_burn if scl_result.survival else None, "estimated_runway": scl_result.survival.estimated_runway if scl_result.survival else None, "adtv_30d": scl_result.survival.adtv_30d if scl_result.survival else None, } if scl_result.survival else None, "dilution_activity": { "yoy_share_growth_pct": scl_result.dilution_activity.yoy_share_growth_pct if scl_result.dilution_activity else None, "share_cagr_2y_pct": scl_result.dilution_activity.share_cagr_2y_pct if scl_result.dilution_activity else None, "capital_raises_12mo": scl_result.dilution_activity.capital_raises_12mo if scl_result.dilution_activity else None, "capital_raises_24mo": scl_result.dilution_activity.capital_raises_24mo if scl_result.dilution_activity else None, "total_raises": scl_result.dilution_activity.total_raises if scl_result.dilution_activity else None, "last_raise_days_ago": scl_result.dilution_activity.last_raise_days_ago if scl_result.dilution_activity else None, } if scl_result.dilution_activity else None, "capital_structure": { "charter_authorized_shares": scl_result.capital_structure.charter_authorized_shares if scl_result.capital_structure else None, "shares_outstanding": scl_result.capital_structure.shares_outstanding if scl_result.capital_structure else None, "unissued_headroom": scl_result.capital_structure.unissued_headroom if scl_result.capital_structure else None, "has_atm": scl_result.capital_structure.has_atm if scl_result.capital_structure else None, "has_shelf": scl_result.capital_structure.has_shelf if scl_result.capital_structure else None, } if scl_result.capital_structure else None, } except Exception as e: logger.warning(f"SmallCapLab error for {symbol}: {e}") result = build_stock_result( symbol=symbol, finviz_overview=finviz_overview, corporate_actions=corporate_actions, next_earnings_date=next_earnings_date, daily_df=daily_df, gap_history=gap_history, run_history=run_history, company_description=finviz_quote.get("company_description", ""), employees=finviz_quote.get("employees"), float_finviz=finviz_quote.get("float_shares", 0) or 0, nasdaq_non_compliant=nasdaq_non_compliant, sec_unlocks=sec_unlocks, financials=financials, dilution_risk=dilution_risk, stockanalysis_data=sa_data, smallcaplab_data=scl_data, press_releases=press_releases, ) if catalyst: result["catalysts"]["finviz"] = catalyst if not fast: cache_set(_cache_key_fn(symbol), result) asyncio.create_task(background_refresh_slow_data(symbol)) return result async def background_refresh_slow_data(symbol: str) -> None: """Background refresh: Yahoo board + cashflow + SEC dilution risk.""" try: loop = asyncio.get_event_loop() yahoo_task = loop.run_in_executor(_PROVIDER_EXECUTOR, defeatbeta.fetch_consolidated, symbol) sec_task = loop.run_in_executor(_PROVIDER_EXECUTOR, fetch_sec_dilution_risk, symbol) yahoo_data = await asyncio.wait_for(yahoo_task, timeout=15) dr = await asyncio.wait_for(sec_task, timeout=15) except Exception as e: logger.warning(f"Background refresh error for {symbol}: {e}") return from core.chinese_profile import score_chinese_profile from core.shared.cache import cache_get, cache_set from core.shared.cache import cache_key as _ck current = cache_get(_ck(symbol)) if current is None: return # Work on a COPY so concurrent readers of the cached object never # observe a half-mutated dict mid-refresh. current = dict(current) fundamentals = yahoo_data.get("fundamentals", {}) company_officers = fundamentals.get("company_officers", []) city = fundamentals.get("city", "") phone = fundamentals.get("phone", "") cashflow = yahoo_data.get("cashflow", {}) updated = False if company_officers: current["board"] = extract_board(company_officers) updated_company = dict(current.get("company", {})) has_company_update = False if city: updated_company["city"] = city has_company_update = True if phone: updated_company["phone"] = phone has_company_update = True if company_officers and has_company_update: updated_company["chinese_profile"] = score_chinese_profile( { "country": updated_company.get("country", ""), "city": city or updated_company.get("city", ""), "phone": phone or updated_company.get("phone", ""), "business_description": updated_company.get("business_description", ""), "officers": company_officers, } ) if has_company_update: current["company"] = updated_company updated = True current_fin = dict(current.get("financials", {})) current_ocf = current_fin.get("operating_cash_flow") fin_updated = False if cashflow: if "net_income" not in current_fin and cashflow.get("net_income") is not None: current_fin["net_income"] = cashflow["net_income"] fin_updated = True if current_ocf is None and cashflow.get("operating_cash_flow") is not None: current_fin["operating_cash_flow"] = cashflow["operating_cash_flow"] fin_updated = True if cashflow.get("financing_cash_flow") is not None: current_fin["financing_cash_flow"] = cashflow["financing_cash_flow"] fin_updated = True if fin_updated: current["financials"] = current_fin updated = True if dr: current["dilution_risk"] = dr updated = True if updated: cache_set(_ck(symbol), current)