""" APScheduler entry point for News v2 ingest. Jobs: alpaca_news_realtime_poll every 5 min: latest Alpaca News stream alpaca_news_daily_backfill 04:30 ET: prior 26h Alpaca News refill stocktwits_universe_refresh 09:00 ET: rebuild dynamic poll list stocktwits_poll every 5 min: per-ticker stream pull finnhub_daily_backfill 05:00 ET: prior day Finnhub /company-news Activation gate: NEWS_INGEST_ENABLED must be true. Per-source key presence determines which jobs are registered — Alpaca needs ALPACA_API_KEY/SECRET, Finnhub needs FINNHUB_API_KEY, StockTwits needs nothing. If a source's key is missing, its jobs are skipped with a WARNING. If NO sources are configured the scheduler aborts (fail-fast) so the misconfiguration is loud. """ from __future__ import annotations import logging from datetime import datetime, timedelta, timezone from app.core.config import settings logger = logging.getLogger(__name__) _scheduler = None def _get_scheduler(): global _scheduler if _scheduler is None: try: from apscheduler.schedulers.asyncio import AsyncIOScheduler _scheduler = AsyncIOScheduler(timezone=settings.NEWS_INGEST_TIMEZONE) except ImportError: logger.warning("apscheduler not installed; news ingest scheduling disabled") return None return _scheduler # --------------------------------------------------------------------------- # Job bodies # --------------------------------------------------------------------------- async def _run_alpaca_realtime_poll() -> None: """Poll the latest 5 minutes of Alpaca News (no symbol filter).""" from app.services.news.alpaca_news_client import AlpacaNewsClient from app.services.news.headline_ingest_service import ( alpaca_articles_to_rows, insert_headline_rows, ) client = AlpacaNewsClient() if not client.is_configured(): return end = datetime.now(timezone.utc) start = end - timedelta(minutes=10) # 5-min poll + 5-min overlap for safety try: articles = [] async for art in client.fetch_news(symbols=None, start=start, end=end, page_limit=50): articles.append(art) rows = alpaca_articles_to_rows(articles) if rows: inserted = await insert_headline_rows(rows) if inserted: logger.info(f"[News] Alpaca realtime: +{inserted} new ({len(articles)} articles)") except Exception as e: logger.error(f"[News] Alpaca realtime poll failed: {e}") finally: await client.close() async def _run_alpaca_daily_backfill() -> None: """Refill prior 26h Alpaca News stream to catch any gaps.""" from app.services.news.alpaca_news_client import AlpacaNewsClient from app.services.news.headline_ingest_service import ( alpaca_articles_to_rows, insert_headline_rows, ) client = AlpacaNewsClient() if not client.is_configured(): return end = datetime.now(timezone.utc) start = end - timedelta(hours=26) try: articles = [] async for art in client.fetch_news(symbols=None, start=start, end=end, page_limit=50): articles.append(art) rows = alpaca_articles_to_rows(articles) inserted = await insert_headline_rows(rows) if rows else 0 logger.info(f"[News] Alpaca daily backfill: +{inserted} new ({len(articles)} articles, 26h)") except Exception as e: logger.error(f"[News] Alpaca daily backfill failed: {e}") finally: await client.close() async def _run_stocktwits_universe_refresh() -> None: from app.services.news.stocktwits_universe import compute_stocktwits_universe try: await compute_stocktwits_universe() except Exception as e: logger.error(f"[News] StockTwits universe refresh failed: {e}") async def _run_stocktwits_poll() -> None: """Pull recent messages for every ticker in the cached universe.""" from app.services.news.stocktwits_client import StocktwitsClient from app.services.news.stocktwits_universe import get_cached_universe from app.services.news.headline_ingest_service import ( stocktwits_messages_to_rows, insert_headline_rows, ) universe = await get_cached_universe() if not universe: # On first start there's no cached universe — compute on the fly so # the very first poll has something to do. from app.services.news.stocktwits_universe import compute_stocktwits_universe try: universe = await compute_stocktwits_universe() except Exception as e: logger.warning(f"[News] StockTwits initial universe build failed: {e}") return if not universe: return client = StocktwitsClient() total_inserted = 0 try: for ticker in universe: try: payload = await client.fetch_symbol_stream(ticker, max_results=30) except Exception as e: logger.warning(f"[News] StockTwits {ticker} failed: {e}") continue messages = payload.get("messages") or [] if not messages: continue rows = stocktwits_messages_to_rows(ticker, messages) if rows: total_inserted += await insert_headline_rows(rows) finally: await client.close() if total_inserted: logger.info(f"[News] StockTwits poll: +{total_inserted} new across {len(universe)} tickers") async def _run_finnhub_daily_backfill() -> None: """Pull yesterday's Finnhub /company-news for the active V49 universe.""" from app.services.news.finnhub_client import FinnhubClient from app.services.news.headline_ingest_service import ( finnhub_articles_to_rows, insert_headline_rows, ) from sqlalchemy import select from app.core.database import AsyncSessionLocal from app.models.universe_snapshot import UniverseTickerRegistry client = FinnhubClient() if not client.is_configured(): return today = datetime.now(timezone.utc).date() yesterday = today - timedelta(days=1) # Pull active universe tickers async with AsyncSessionLocal() as db: rows = (await db.execute( select(UniverseTickerRegistry.ticker).where(UniverseTickerRegistry.is_active.is_(True)) )).scalars().all() tickers = [t.upper() for t in rows if t] if not tickers: logger.info("[News] Finnhub daily: no active universe tickers") return total_inserted = 0 try: for ticker in tickers: try: payload = await client.fetch_company_news(ticker, yesterday, today) except Exception as e: logger.warning(f"[News] Finnhub {ticker} failed: {e}") continue if not payload: continue news_rows = finnhub_articles_to_rows(ticker, payload) if news_rows: total_inserted += await insert_headline_rows(news_rows) finally: await client.close() logger.info(f"[News] Finnhub daily: +{total_inserted} new across {len(tickers)} tickers") # --------------------------------------------------------------------------- # Lifecycle # --------------------------------------------------------------------------- def start_news_ingest_scheduler() -> None: """Wire jobs into the AsyncIOScheduler. Idempotent.""" if not settings.NEWS_INGEST_ENABLED: logger.info("[News] NEWS_INGEST_ENABLED=false — scheduler not started") return sched = _get_scheduler() if sched is None: return # Source key presence checks (fail-fast on no-keys-at-all) alpaca_ok = bool(settings.ALPACA_API_KEY and settings.ALPACA_SECRET_KEY) finnhub_ok = bool(settings.FINNHUB_API_KEY) stocktwits_ok = True # public API enabled_sources = [] if alpaca_ok: enabled_sources.append("alpaca_benzinga") else: logger.warning("[News] ALPACA_API_KEY/SECRET missing — Alpaca News jobs skipped") if finnhub_ok: enabled_sources.append("finnhub") else: logger.warning("[News] FINNHUB_API_KEY missing — Finnhub jobs skipped") if stocktwits_ok: enabled_sources.append("stocktwits") # Fail-fast: if neither Alpaca nor Finnhub is configured, the only source # is StockTwits which is low-signal on its own. Abort instead of silently # running a degenerate ingest. if not (alpaca_ok or finnhub_ok): logger.error( "[News] NEWS_INGEST_ENABLED=true but no premium-source keys configured " "(need at least ALPACA_API_KEY/SECRET or FINNHUB_API_KEY) — aborting scheduler" ) return try: from apscheduler.triggers.cron import CronTrigger from apscheduler.triggers.interval import IntervalTrigger if alpaca_ok: sched.add_job( _run_alpaca_realtime_poll, trigger=IntervalTrigger(minutes=5), id="alpaca_news_realtime_poll", replace_existing=True, max_instances=1, misfire_grace_time=120, coalesce=True, ) sched.add_job( _run_alpaca_daily_backfill, trigger=CronTrigger(hour=4, minute=30), id="alpaca_news_daily_backfill", replace_existing=True, max_instances=1, misfire_grace_time=600, coalesce=True, ) if stocktwits_ok: sched.add_job( _run_stocktwits_universe_refresh, trigger=CronTrigger(day_of_week="mon-fri", hour=9, minute=0), id="stocktwits_universe_refresh", replace_existing=True, max_instances=1, misfire_grace_time=600, coalesce=True, ) sched.add_job( _run_stocktwits_poll, trigger=IntervalTrigger(minutes=5), id="stocktwits_poll", replace_existing=True, max_instances=1, misfire_grace_time=120, coalesce=True, ) if finnhub_ok: sched.add_job( _run_finnhub_daily_backfill, trigger=CronTrigger(hour=5, minute=0), id="finnhub_daily_backfill", replace_existing=True, max_instances=1, misfire_grace_time=600, coalesce=True, ) if not sched.running: sched.start() logger.info(f"[News] ingest scheduler started — sources: {enabled_sources}") except Exception as e: logger.error(f"[News] ingest scheduler start failed: {e}") def stop_news_ingest_scheduler() -> None: sched = _get_scheduler() if sched and sched.running: sched.shutdown(wait=False) logger.info("[News] ingest scheduler stopped")