Spaces:
Running
Running
| import json | |
| import logging | |
| import os | |
| import sys | |
| import threading | |
| from collections import defaultdict | |
| from datetime import datetime, timezone | |
| from flask import Flask, jsonify, request, send_from_directory | |
| from flask_cors import CORS | |
| from config import Config | |
| from main import run_pipeline | |
| from src.models import NewsItem | |
| from src.rss_feed_scraper import RSSFeedScraper | |
| from src.html_scraper import HTMLSiteScraper | |
| logging.basicConfig(level=logging.INFO, format="%(levelname).1s %(message)s", stream=sys.stderr) | |
| logger = logging.getLogger(__name__) | |
| cfg = Config() | |
| app = Flask(__name__) | |
| CORS(app, origins=cfg.cors_origins.split(",") if cfg.cors_origins != "*" else "*") | |
| cache_lock = threading.Lock() | |
| CACHE_FILE = os.getenv("CACHE_FILE", "/tmp/cache_data.json") | |
| CATEGORIES = ["Geopolitical", "World Health", "Tech", "Cybersecurity", "Funny/Weird", "Gaming", "Movies", "Arab World", "Tunisia"] | |
| cached_by_cat: dict[str, list[dict]] = defaultdict(list) | |
| cached_results: list[dict] = [] | |
| cached_status = "starting" | |
| cached_last_refresh = None | |
| def _cluster_to_html_dict(c): | |
| category = "General" | |
| if c.articles and c.articles[0].analysis: | |
| cat = c.articles[0].analysis.category | |
| category = cat if cat in CATEGORIES else "General" | |
| return { | |
| "category": category, | |
| "topic": c.topic, | |
| "score": round(c.final_score, 2), | |
| "trust": f"{c.avg_trustworthiness:.0%}", | |
| "coverage": c.total_coverage, | |
| "image_url": c.image_url, | |
| "top_post_url": c.top_post_url, | |
| "articles": [ | |
| { | |
| "title": a.article.title or a.post.title, | |
| "url": a.post.url, | |
| "domain": a.article.source_domain if a.article else "", | |
| "summary": a.analysis.summary if a.analysis else "", | |
| "topics": a.analysis.topics if a.analysis else [], | |
| "trust": round(a.analysis.trustworthiness_score, 2) if a.analysis else 0, | |
| "sponsor": a.analysis.sponsor.get("display", "") if a.analysis else "", | |
| "sponsor_info": a.analysis.sponsor if a.analysis else {}, | |
| "source_bias": a.analysis.source_bias if a.analysis else "", | |
| "source_factuality": a.analysis.source_factuality if a.analysis else "", | |
| "article_leaning": a.analysis.article_leaning if a.analysis else "", | |
| "sourcing_penalty": a.analysis.sourcing_penalty if a.analysis else 0, | |
| "score": a.post.score, | |
| "comments": a.post.num_comments, | |
| "image": a.article.image_url if a.article else a.post.image_url, | |
| "published": a.article.published or a.post.published, | |
| "published_iso": a.article.published_iso or a.post.published_iso, | |
| } | |
| for a in c.articles[:5] | |
| ], | |
| } | |
| def _cluster_to_api_dict(c): | |
| cat = "General" | |
| if c.articles and c.articles[0].analysis: | |
| ca = c.articles[0].analysis.category | |
| cat = ca if ca in CATEGORIES else "General" | |
| return { | |
| "id": f"cluster-{id(c)}", | |
| "topic": c.topic, | |
| "category": cat, | |
| "final_score": round(c.final_score, 2), | |
| "avg_trustworthiness": round(c.avg_trustworthiness, 2), | |
| "total_coverage": c.total_coverage, | |
| "image_url": c.image_url, | |
| "top_post_url": c.top_post_url, | |
| "articles": [ | |
| { | |
| "title": a.article.title or a.post.title, | |
| "url": a.post.url, | |
| "domain": a.article.source_domain if a.article else "", | |
| "summary": a.analysis.summary if a.analysis else "", | |
| "topics": a.analysis.topics if a.analysis else [], | |
| "trust": round(a.analysis.trustworthiness_score, 2) if a.analysis else 0, | |
| "sponsor": a.analysis.sponsor.get("display", "") if a.analysis else "", | |
| "sponsor_info": a.analysis.sponsor if a.analysis else {}, | |
| "source_bias": a.analysis.source_bias if a.analysis else "", | |
| "source_factuality": a.analysis.source_factuality if a.analysis else "", | |
| "article_leaning": a.analysis.article_leaning if a.analysis else "", | |
| "sourcing_penalty": a.analysis.sourcing_penalty if a.analysis else 0, | |
| "score": a.post.score, | |
| "comments": a.post.num_comments, | |
| "image": a.article.image_url if a.article else a.post.image_url, | |
| "published": a.article.published or a.post.published, | |
| "published_iso": a.article.published_iso or a.post.published_iso, | |
| } | |
| for a in c.articles[:5] | |
| ], | |
| } | |
| def refresh_data(): | |
| global cached_by_cat, cached_results, cached_status, cached_last_refresh | |
| with cache_lock: | |
| cached_status = "running" | |
| try: | |
| rss_posts = RSSFeedScraper(cfg).fetch_posts() | |
| html_posts = HTMLSiteScraper(cfg).fetch_posts() | |
| posts = rss_posts + html_posts | |
| items = [NewsItem(post=p) for p in posts] | |
| new_clusters = run_pipeline(items, cfg) | |
| html_dicts = [_cluster_to_html_dict(c) for c in new_clusters] | |
| api_dicts = [_cluster_to_api_dict(c) for c in new_clusters] | |
| by_cat: dict[str, list[dict]] = defaultdict(list) | |
| for d in html_dicts: | |
| by_cat[d["category"]].append(d) | |
| with cache_lock: | |
| cached_by_cat = by_cat | |
| cached_results = api_dicts | |
| cached_last_refresh = datetime.now(timezone.utc) | |
| cached_status = f"ok — {len(new_clusters)} clusters from {len(items)} articles" | |
| _save_cache(by_cat, api_dicts, cached_status, cached_last_refresh) | |
| logger.info("Refresh complete: %s", cached_status) | |
| except Exception as exc: | |
| with cache_lock: | |
| cached_status = f"error: {exc}" | |
| logger.error("Pipeline failed: %s", exc) | |
| def _save_cache(by_cat, results, status, last_refresh): | |
| try: | |
| data = { | |
| "by_cat": {k: v for k, v in by_cat.items()}, | |
| "results": results, | |
| "status": status, | |
| "last_refresh": last_refresh.isoformat() if last_refresh else None, | |
| } | |
| with open(CACHE_FILE, "w") as f: | |
| json.dump(data, f) | |
| except Exception as exc: | |
| logger.warning("Failed to write cache file: %s", exc) | |
| def _load_cache(): | |
| try: | |
| with open(CACHE_FILE) as f: | |
| data = json.load(f) | |
| return data | |
| except (FileNotFoundError, json.JSONDecodeError): | |
| return None | |
| # ---- Web UI (React SPA) ---- | |
| FRONTEND_DIST = os.path.join(os.path.dirname(__file__), "frontend", "dist") | |
| def index(): | |
| return send_from_directory(FRONTEND_DIST, "index.html") | |
| def serve_assets(filename): | |
| return send_from_directory(os.path.join(FRONTEND_DIST, "assets"), filename) | |
| # ---- REST API ---- | |
| def api_status(): | |
| with cache_lock: | |
| status = cached_status | |
| lr = cached_last_refresh | |
| n_clusters = len(cached_results) | |
| n_articles = sum(c["total_coverage"] for c in cached_results) | |
| return jsonify({ | |
| "status": "ok", | |
| "clusters": n_clusters, | |
| "total_articles": n_articles, | |
| "last_refresh": lr.isoformat() if lr else None, | |
| "pipeline_status": status, | |
| }) | |
| def api_posts(): | |
| category = request.args.get("category") | |
| topic = request.args.get("topic") | |
| limit = request.args.get("limit", type=int) | |
| min_score = request.args.get("min_score", type=float) | |
| with cache_lock: | |
| all_results = list(cached_results) | |
| filtered = all_results | |
| if category: | |
| filtered = [r for r in filtered if r["category"].lower() == category.lower()] | |
| if topic: | |
| filtered = [r for r in filtered if topic.lower() in [t.lower() for t in r.get("articles", [])[:1]]] | |
| if min_score is not None: | |
| filtered = [r for r in filtered if r["final_score"] >= min_score] | |
| if limit: | |
| filtered = filtered[:limit] | |
| return jsonify({"clusters": filtered, "meta": {"total": len(filtered)}}) | |
| def api_categories(): | |
| return jsonify({"categories": CATEGORIES}) | |
| def _build_sponsors(): | |
| with cache_lock: | |
| results = list(cached_results) | |
| groups: dict[str, dict] = {} | |
| for cluster in results: | |
| for article in cluster.get("articles", []): | |
| info = article.get("sponsor_info", {}) | |
| display = info.get("display", "") or article.get("sponsor", "") | |
| if not display: | |
| continue | |
| if display not in groups: | |
| groups[display] = { | |
| "display": display, | |
| "parent": info.get("parent", ""), | |
| "category": info.get("category", ""), | |
| "bias": info.get("bias", ""), | |
| "factuality": info.get("factuality", ""), | |
| "wikipedia": info.get("wikipedia", ""), | |
| "owners": info.get("owners", []), | |
| "owner_wikis": info.get("owner_wikis", {}), | |
| "sources": set(), | |
| } | |
| groups[display]["sources"].add(article.get("domain", "")) | |
| return sorted( | |
| [ | |
| { | |
| "display": g["display"], | |
| "parent": g["parent"], | |
| "category": g["category"], | |
| "bias": g["bias"], | |
| "factuality": g["factuality"], | |
| "wikipedia": g["wikipedia"], | |
| "owners": g["owners"], | |
| "owner_wikis": g["owner_wikis"], | |
| "sources": sorted(g["sources"]), | |
| "source_count": len(g["sources"]), | |
| } | |
| for g in groups.values() | |
| ], | |
| key=lambda x: x["display"].lower(), | |
| ) | |
| def api_sponsors(): | |
| return jsonify({"sponsors": _build_sponsors()}) | |
| def api_refresh(): | |
| threading.Thread(target=refresh_data, daemon=True).start() | |
| return jsonify({"status": "accepted", "message": "refresh started"}), 202 | |
| def health(): | |
| with cache_lock: | |
| ok = "error" not in cached_status | |
| if ok: | |
| return jsonify({"status": "healthy"}), 200 | |
| return jsonify({"status": "unhealthy", "message": cached_status}), 503 | |
| # ---- Startup (runs in both WSGI and dev modes) ---- | |
| _cache_data = _load_cache() | |
| if _cache_data: | |
| cached_by_cat = defaultdict(list, _cache_data.get("by_cat", {})) | |
| cached_results = _cache_data.get("results", []) | |
| cached_status = _cache_data.get("status", "idle") | |
| lr = _cache_data.get("last_refresh") | |
| if lr: | |
| try: | |
| cached_last_refresh = datetime.fromisoformat(lr) | |
| except Exception: | |
| pass | |
| logger.info("Loaded cached data: %s", cached_status) | |
| if not cached_results: | |
| threading.Thread(target=refresh_data, daemon=True).start() | |
| try: | |
| from apscheduler.schedulers.background import BackgroundScheduler | |
| _scheduler = BackgroundScheduler(daemon=True) | |
| _scheduler.add_job( | |
| refresh_data, | |
| trigger="cron", | |
| hour=cfg.refresh_hour, | |
| minute=cfg.refresh_minute, | |
| timezone=cfg.refresh_timezone, | |
| id="daily_refresh", | |
| name="Daily news refresh", | |
| replace_existing=True, | |
| ) | |
| _scheduler.start() | |
| logger.info("Scheduler started — daily refresh at %02d:%02d %s", cfg.refresh_hour, cfg.refresh_minute, cfg.refresh_timezone) | |
| except ImportError: | |
| logger.warning("APScheduler not available — scheduled refresh disabled") | |
| # ---- Dev server (not used in WSGI/gunicorn) ---- | |
| if __name__ == "__main__": | |
| app.run(host=cfg.host, port=cfg.port, debug=(cfg.flask_env != "production")) | |