"""RSS ingestion Prefect flow""" import asyncio import logging from typing import List import feedparser from prefect import flow, task from prefect.task_runners import ConcurrentTaskRunner from sentiment_engine.ingestion.rss import RSSConnector from sentiment_engine.schemas.config import RSSConnectorConfig from sentiment_engine.utils.config import get_settings logger = logging.getLogger(__name__) @task(retries=3, retry_delay_seconds=30) async def fetch_rss_feed(feed_url: str, config: RSSConnectorConfig) -> List[dict]: """Fetch and parse a single RSS feed""" try: feed = feedparser.parse(feed_url) items = [] for entry in feed.entries[:config.max_items_per_feed]: title = getattr(entry, "title", "").strip() summary = getattr(entry, "summary", getattr(entry, "description", "")).strip() raw_text = f"{title}\n\n{summary}" if not raw_text.strip(): continue items.append({ "source_id": f"rss:{feed_url}", "source_type": "news", "raw_text": raw_text, "title": title, "url": getattr(entry, "link", ""), "author": getattr(entry, "author", ""), "publish_ts": getattr(entry, "published_parsed", None), "metadata": {"feed_url": feed_url} }) return items except Exception as e: logger.error(f"Error fetching {feed_url}: {e}") raise @flow( name="rss_ingest", task_runner=ConcurrentTaskRunner(max_workers=10), log_prints=True ) async def rss_ingest_flow(feed_urls: List[str] = None): """Main RSS ingestion flow""" settings = get_settings() if feed_urls is None: # Default crypto news feeds feed_urls = [ "https://www.coindesk.com/arc/outboundfeeds/rss/", "https://cointelegraph.com/rss", "https://www.theblock.co/rss", "https://decrypt.co/feed", "https://messari.io/feed", "https://cryptoslate.com/feed/", "https://bitcoinmagazine.com/feed/", ] config = RSSConnectorConfig( name="prefect_rss", source_type="news", feed_urls=feed_urls, max_items_per_feed=50 ) # Fetch all feeds concurrently results = await asyncio.gather( *[fetch_rss_feed(url, config) for url in feed_urls], return_exceptions=True ) all_items = [] for i, result in enumerate(results): if isinstance(result, Exception): logger.error(f"Feed {feed_urls[i]} failed: {result}") else: all_items.extend(result) logger.info(f"Fetched {len(all_items)} items from {len(feed_urls)} feeds") # In production, publish to NATS # For now, return items return all_items if __name__ == "__main__": asyncio.run(rss_ingest_flow())