You cannot select more than 25 topics
Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
311 lines
11 KiB
Python
311 lines
11 KiB
Python
"""
|
|
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")
|