""" Crypto News & Narrative Analysis Tool -------------------------------------- Usage: python main.py # single snapshot + dashboard python main.py --live # auto-refresh every N minutes python main.py --export # save JSON report to ./reports/ python main.py --setup # check dependencies and first-run guide """ import argparse import io import json import logging import os import sys import time from datetime import datetime, timezone from pathlib import Path # Force UTF-8 output on Windows (prevents charmap encode errors from Rich) if sys.platform == "win32": sys.stdout = io.TextIOWrapper(sys.stdout.buffer, encoding="utf-8", errors="replace") sys.stderr = io.TextIOWrapper(sys.stderr.buffer, encoding="utf-8", errors="replace") from dotenv import load_dotenv load_dotenv() # --------------------------------------------------------------------------- # Logging setup # --------------------------------------------------------------------------- logging.basicConfig( level=os.getenv("LOG_LEVEL", "WARNING"), format="%(asctime)s [%(levelname)s] %(name)s: %(message)s", ) logger = logging.getLogger(__name__) # --------------------------------------------------------------------------- # Imports (after env is loaded) # --------------------------------------------------------------------------- from collectors import ( RSSCollector, RedditCollector, CoinGeckoCollector, FearGreedCollector, NewsDataCollector, ) from analyzers import ( SentimentAnalyzer, NarrativeDetector, EntityExtractor, TrendScorer, ) from storage import Database from dashboard import Dashboard from config import REFRESH_INTERVAL_MINUTES # --------------------------------------------------------------------------- # Core pipeline # --------------------------------------------------------------------------- def run_pipeline(db: Database) -> dict: """ Full data collection → analysis → scoring pipeline. Returns a dict with all computed artefacts. """ dash = Dashboard() with dash.spinner("Fetching RSS feeds...") as p: p.add_task("") rss_items = RSSCollector().fetch_all() with dash.spinner("Fetching NewsData.io articles...") as p: p.add_task("") newsdata_items = NewsDataCollector().fetch_all() with dash.spinner("Fetching Reddit posts...") as p: p.add_task("") reddit_items = RedditCollector().fetch_all() with dash.spinner("Fetching CoinGecko data...") as p: p.add_task("") cg = CoinGeckoCollector() trending_coins = cg.fetch_trending() global_market = cg.fetch_global() coin_prices = cg.fetch_prices() trending_ids = cg.get_trending_ids() with dash.spinner("Fetching Fear & Greed index...") as p: p.add_task("") fg_collector = FearGreedCollector() fear_greed_history = fg_collector.fetch() fear_greed_current = fg_collector.current() # Combine all text items all_items = rss_items + newsdata_items + reddit_items # Save Fear & Greed and prices to DB db.save_fear_greed(fear_greed_history) db.save_prices(coin_prices) if not all_items: logger.warning("No items collected — check network connectivity") return {} # ---- Analysis ---- with dash.spinner("Running sentiment analysis...") as p: p.add_task("") sentiment_analyzer = SentimentAnalyzer() all_items = sentiment_analyzer.analyze_batch(all_items) with dash.spinner("Detecting narratives...") as p: p.add_task("") narrative_detector = NarrativeDetector() all_items = narrative_detector.detect_batch(all_items) narrative_detector.fit_corpus(all_items) emerging_keywords = narrative_detector.top_emerging_keywords(all_items, top_n=25) with dash.spinner("Extracting entities...") as p: p.add_task("") entity_extractor = EntityExtractor() all_items = entity_extractor.extract_batch(all_items) coin_mentions = entity_extractor.coin_mention_frequency(all_items) with dash.spinner("Scoring narratives...") as p: p.add_task("") scorer = TrendScorer(trending_coin_ids=trending_ids) scored_narratives = scorer.score_narratives(all_items) # Persist to DB db.save_articles(all_items) db.save_snapshot(scored_narratives) db.purge_old() return { "narratives": scored_narratives, "global_market": global_market, "fear_greed": fear_greed_current, "trending_coins": trending_coins, "coin_prices": coin_prices, "coin_mentions": coin_mentions, "emerging_keywords": emerging_keywords, "article_count": len(all_items), "fetch_time": datetime.now(timezone.utc).strftime("%Y-%m-%d %H:%M UTC"), } # --------------------------------------------------------------------------- # Export # --------------------------------------------------------------------------- def export_report(result: dict) -> str: Path("reports").mkdir(exist_ok=True) ts = datetime.now(timezone.utc).strftime("%Y%m%d_%H%M%S") path = f"reports/narrative_report_{ts}.json" with open(path, "w", encoding="utf-8") as f: json.dump(result, f, indent=2, default=str) return path # --------------------------------------------------------------------------- # Setup / dependency check # --------------------------------------------------------------------------- def run_setup() -> None: from rich.console import Console from rich.table import Table from rich import box c = Console() c.print("\n[bold cyan]Crypto Narrative Analysis — Setup Check[/bold cyan]\n") checks = [ ("feedparser", "RSS feed parsing"), ("vaderSentiment", "VADER sentiment"), ("sklearn", "TF-IDF / scikit-learn"), ("rich", "Terminal dashboard"), ("praw", "Reddit API (optional)"), ("spacy", "spaCy NER (optional)"), ("bertopic", "BERTopic clustering (optional)"), ("pytrends", "Google Trends (optional)"), ("plotext", "Terminal charts (optional)"), ] table = Table(box=box.SIMPLE, show_header=True, header_style="bold white") table.add_column("Package", width=25) table.add_column("Purpose", width=30) table.add_column("Status", width=12) for pkg, purpose in checks: try: __import__(pkg) status = "[green]✓ OK[/green]" except ImportError: status = "[red]✗ Missing[/red]" table.add_row(pkg, purpose, status) c.print(table) # spaCy model check try: import spacy spacy.load("en_core_web_sm") c.print("[green]✓ spaCy model en_core_web_sm loaded[/green]") except Exception: c.print("[yellow]spaCy model not installed. Run:[/yellow]") c.print(" python -m spacy download en_core_web_sm") # Reddit credentials c.print() if os.getenv("REDDIT_CLIENT_ID"): c.print("[green]✓ Reddit credentials found in environment[/green]") else: c.print("[yellow]Reddit credentials not set. Copy .env.example → .env and fill in values.[/yellow]") c.print(" Get free credentials at: https://www.reddit.com/prefs/apps") c.print("\n[bold]To run:[/bold] python main.py") c.print("[bold]Live mode:[/bold] python main.py --live") c.print("[bold]Export JSON:[/bold] python main.py --export\n") # --------------------------------------------------------------------------- # CLI entry # --------------------------------------------------------------------------- def parse_args(): p = argparse.ArgumentParser(description="Crypto News & Narrative Analysis") p.add_argument("--live", action="store_true", help="Auto-refresh dashboard") p.add_argument("--export", action="store_true", help="Export JSON report") p.add_argument("--setup", action="store_true", help="Check dependencies") p.add_argument("--no-reddit", action="store_true", help="Skip Reddit collection") p.add_argument("--interval", type=int, default=REFRESH_INTERVAL_MINUTES, help="Refresh interval in minutes (live mode)") p.add_argument("--verbose", action="store_true", help="Enable debug logging") return p.parse_args() def main(): args = parse_args() if args.verbose: logging.getLogger().setLevel(logging.DEBUG) if args.setup: run_setup() return db = Database() dash = Dashboard() def single_run(): try: result = run_pipeline(db) if not result: dash.print_alert("No Data", "Could not collect any items. Check your network.", "error") return dash.render(**result) if args.export: path = export_report(result) dash.print_alert("Report Exported", f"Saved to {path}", "info") except KeyboardInterrupt: raise except Exception as exc: logger.exception("Pipeline error") dash.print_alert("Pipeline Error", str(exc), "error") if args.live: from rich.console import Console Console().print( f"[cyan]Live mode — refreshing every {args.interval} minutes. Ctrl+C to stop.[/cyan]\n" ) try: while True: single_run() time.sleep(args.interval * 60) except KeyboardInterrupt: Console().print("\n[dim]Stopped.[/dim]") else: single_run() db.close() if __name__ == "__main__": main()