Explorer
/tmp/multi_now.py
← Zurück ↓ Download
"""
Multi-Source Intake — AA-045
YouTube (RSS-Feeds + Suchseiten-Scraping) und öffentliche Facebook-Signale
(loginfrei via Websuche). Kein Login, keine Cookies, keine Session-Pflicht,
kein Publishing. Fehler einer Quelle stoppen die anderen nicht.

Zentrale Konfiguration: sources_multi.yaml (gleicher Ordner)
Observability: Tabelle source_runs in der Signal-DB.
"""

import sys
from pathlib import Path
sys.path.insert(0, str(Path(__file__).parent.parent))

import re
import json
import unicodedata
import base64
import time
import urllib.parse
import urllib.request
import xml.etree.ElementTree as ET
from datetime import datetime, timezone, timedelta

try:
    import yaml
except ImportError:
    yaml = None

from signal_db import insert_signal, record_source_run
from intake.meta_ads_intake import run_meta_ads
import email.utils
def _pubdate_to_iso(pubdate_str):
    if not pubdate_str:
        return ""
    try:
        dt = email.utils.parsedate_to_datetime(pubdate_str)
        return dt.isoformat()
    except Exception:
        return ""

def normalize_youtube_published(raw: str) -> str:
    """Normalisiert relative YouTube-Zeitangaben; unbekannt bleibt leer."""
    value = (raw or "").strip()
    if not value:
        return ""
    if "T" in value and (value.endswith("Z") or "+" in value or value.count(":") >= 2):
        return value
    folded = "".join(c for c in unicodedata.normalize("NFKD", value.casefold()) if not unicodedata.combining(c))
    m = re.search(r"(?:vor|hace)\s+(\d+)\s+(minuten?|minutos?|hours?|stunden?|horas?|tage?|dias?|weeks?|wochen?|semanas?)|(?:^|\s)(\d+)\s+(minutes?|hours?|days?|weeks?)\s+ago", folded)
    if not m:
        return ""
    number = int(m.group(1) or m.group(3))
    unit = (m.group(2) or m.group(4)).casefold()
    if unit.startswith(("minute", "minut")): delta = timedelta(minutes=number)
    elif unit.startswith(("stund", "hora", "hour")): delta = timedelta(hours=number)
    elif unit.startswith(("tag", "dia", "day")): delta = timedelta(days=number)
    elif unit.startswith(("woch", "semana", "week")): delta = timedelta(weeks=number)
    else: return ""
    return (datetime.now(timezone.utc) - delta).isoformat()


def _parse_relative_time_german(rel_str):
    """Parse a German relative time string like "vor 3 Tagen", "vor 2 Stunden" etc.
    Returns an ISO timestamp string or empty string if not parsable.
    """
    import re
    from datetime import datetime, timezone, timedelta
    rel_str = rel_str.strip().lower()
    if not rel_str.startswith("vor "):
        return ""
    rel_str = rel_str[4:]  # remove "vor "
    # Match patterns: "(\d+) (Tag|Tage|Stunde|Stunden|Monat|Monate|Woche|Wochen)"
    match = re.match(r"(\d+)\s+(Tag|Tage|Stunde|Stunden|Monat|Monate|Woche|Wochen)", rel_str)
    if not match:
        return ""
    value = int(match.group(1))
    unit = match.group(2)
    # Normalize unit to singular
    if unit.endswith("en"):
        unit = unit[:-2]
    now = datetime.now(timezone.utc)
    if unit == "Tag":
        delta = timedelta(days=value)
    elif unit == "Stunde":
        delta = timedelta(hours=value)
    elif unit == "Monat":
        # Approximate a month as 30 days
        delta = timedelta(days=30*value)
    elif unit == "Woche":
        delta = timedelta(weeks=value)
    else:
        return ""
    pub_dt = now - delta
    return pub_dt.isoformat()

from classification.signal_classifier import score_signal

CONFIG_PATH = Path(__file__).parent.parent / "sources_multi.yaml"
HEADERS = {
    "User-Agent": ("Mozilla/5.0 (Windows NT 10.0; Win64; x64) "
                   "AppleWebKit/537.36 (KHTML, like Gecko) Chrome/124 Safari/537.36"),
    "Accept-Language": "de-DE,de;q=0.9",
}
HTTP_TIMEOUT = 20


def load_config() -> dict:
    if yaml is None:
        raise RuntimeError("PyYAML fehlt – Quellenkonfiguration nicht ladbar")
    with open(CONFIG_PATH, encoding="utf-8") as fh:
        return yaml.safe_load(fh)


# ─────────────────────────── HTTP Helper ────────────────────────────

def _get(url: str, timeout: int = HTTP_TIMEOUT) -> bytes | None:
    try:
        req = urllib.request.Request(url, headers=HEADERS)
        with urllib.request.urlopen(req, timeout=timeout) as resp:
            return resp.read()
    except Exception as e:
        print(f"[MULTI] HTTP-Fehler {url[:80]}: {e}")
        return None


# ─────────────────────────── YouTube ────────────────────────────────

def _yt_rss_channel_videos(channel_id: str) -> list[dict]:
    """Channel-RSS: videoId, Titel, Link, Published, Kanalname."""
    raw = _get(f"https://www.youtube.com/feeds/videos.xml?channel_id={channel_id}")
    if not raw:
        return []
    try:
        root = ET.fromstring(raw)
    except ET.ParseError:
        return []
    ns = {"a": "http://www.w3.org/2005/Atom",
          "yt": "http://www.youtube.com/xml/schemas/2015",
          "media": "http://search.yahoo.com/mrss/"}
    out = []
    for e in root.findall("a:entry", ns):
        vid = e.findtext("yt:videoId", "", ns)
        title = (e.findtext("a:title", "", ns) or "").strip()
        published = e.findtext("a:published", "", ns)
        author = (e.findtext("a:author/a:name", "", ns) or "").strip()
        mg = e.find("media:group", ns)
        desc = (mg.findtext("media:description", "", ns) if mg is not None else "") or ""
        views = None
        if mg is not None:
            cs = mg.find("media:community/media:statistics", ns)
            if cs is not None:
                views = cs.get("views")
        out.append({
            "external_id": vid,
            "title": title,
            "url": f"https://www.youtube.com/watch?v={vid}",
            "published": published,
            "channel": author or channel_id,
            "description": desc,
            "views": views,
        })
    return out


_YT_PATTERNS = {
    "ids":   re.compile(r'"videoRenderer":\{"videoId":"([^"]{11})"'),
    "titles": re.compile(r'"videoRenderer":\{"videoId":"[^"]+".*?"title":\{"runs":\[\{"text":"(.*?)"\}'),
}


def _decode_bing_redirect(u_param: str) -> str | None:
    s = u_param + "=" * (-len(u_param) % 4)
    try:
        d = base64.urlsafe_b64decode(s).decode("utf-8", "ignore")
    except Exception:
        return None
    return d if d.startswith("http") else None


def yt_search(query: str) -> list[dict]:
    """YouTube-Suchseite ohne API-Key; liefert videoId/Titel/Kanal/Views."""
    q = urllib.parse.quote(query)
    # CAI= is YouTube's public upload-date sort; keep the search source fresh.
    html_bytes = _get(f"https://www.youtube.com/results?search_query={q}&sp=CAI%253D")
    if not html_bytes:
        return []
    h = html_bytes.decode("utf-8", errors="ignore")

    videos = []
    seen = set()
    # Blockweise parsen für konsistente Zuordnung
    for block in re.findall(r'"videoRenderer":\{"videoId":"[^"]+".*?(?="videoRenderer"|\Z)', h):
        mid = re.match(r'"videoRenderer":\{"videoId":"([^"]{11})"', block)
        mt = re.search(r'"title":\{"runs":\[\{"text":"(.*?)"\}', block)
        mc = re.search(r'"ownerText":\{"runs":\[\{"text":"(.*?)"\}', block)
        mp = re.search(r'"publishedTimeText":\{"simpleText":"(.*?)"\}', block)
        mv = re.search(r'"viewCountText":\{"simpleText":"([\d\.]+)', block)
        md = re.search(r'"detailedMetadataSnippets".*?"text":"(.*?)"', block)
        if not (mid and mt):
            continue
        vid = mid.group(1)
        if vid in seen:
            continue
        seen.add(vid)

        def _unesc(s: str) -> str:
            try:
                s = s.encode("utf-8").decode("unicode_escape").encode("latin-1", "ignore").decode("utf-8", "ignore")
            except Exception:
                pass  # roh lassen bei ungueltigen Escape-Sequenzen
            return s.replace("\\/", "/")

        channel = _unesc(mc.group(1)).split('","')[0] if mc else ""
        views = mv.group(1).replace(".", "") if mv else ""
        videos.append({
            "external_id": vid,
            "title": _unesc(mt.group(1)),
            "url": f"https://www.youtube.com/watch?v={vid}",
            "published": mp.group(1) if mp else "",   # relativ ('vor 3 Tagen')
            "channel": channel,
            "description": _unesc(md.group(1))[:400] if md else "",
            "views": views,
        })
    return videos


# ─────────────── YouTube → Signale (Kanäle + Suche) ─────────────────

_DECAY_H = 168  # YouTube-Themen halten länger als News


def _store_youtube_video(v: dict, source_name: str, min_score: float) -> tuple[str | None, bool]:
    """Speichert ein YouTube-Video als Signal. Rückgabe (signal_id, ist_neu)."""
    text = f"{v['title']} {v.get('description','')}"
    pub_iso = normalize_youtube_published(v.get("published") or "")
    sd = score_signal(text, language="de", base_category="feuchte",
                      source_url=v["url"], timestamp=pub_iso if "T" in pub_iso else "")
    if sd["radar_score"] < min_score:
        return None, False

    observed_at = datetime.now(timezone.utc).isoformat()
    expires = (datetime.now(timezone.utc) + timedelta(hours=_DECAY_H)).isoformat()
    # AA-045-F1: oeffentliches YouTube-HQ-Thumbnail (kein Download, nur URL)
    yt_thumb = ""
    if v.get("external_id"):
        cand = f"https://i.ytimg.com/vi/{v['external_id']}/hqdefault.jpg"
        if _get(cand, timeout=8):
            yt_thumb = cand
    signal = {
        "source_type":       "youtube",
        "source_platform":   "youtube",
        "external_id":       v.get("external_id", ""),
        "observed_at":       observed_at,
        "published_at":      (pub_iso or None),
        "engagement":        __import__("json").dumps({"views": v.get("views", "")}),
        "image_url":         yt_thumb,
        "image_source":      ("youtube" if yt_thumb else ""),
        "image_found":       (1 if yt_thumb else 0),
        "source_name":       source_name,
        "source_url":        v["url"],
        "language":          "de",
        "signal_category":   sd["signal_category"] or "youtube_thema",
        "topic":             v["title"][:200],
        "short_summary":     (v.get("description") or v["title"])[:500],
        "extracted_hook":    sd["extracted_hook"],
        "emotional_direction": sd["emotional_direction"],
        "urgency_level":     sd["urgency_level"],
        "mallorca_relevance": sd["mallorca_relevance"],
        "seasonal_relevance": 0,
        "risk_relevance":    sd["risk_relevance"],
        "estimated_noise_level": sd["estimated_noise_level"],
        "suggested_case_types": sd["suggested_case_types"],
        "suggested_platforms": sd["suggested_platforms"],
        "suggested_cta":     sd["suggested_cta"],
        "confidence_score":  sd["confidence_score"],
        "radar_score":       sd["radar_score"],
        "decay_rate":        1.2,
        "expires_at":        expires,
        "virality_level":    sd["virality_level"],
        "emotionality_level": sd["emotionality_level"],
        "comment_potential": sd["comment_potential"],
        "content_potential": sd["content_potential"],
        "score_breakdown":   sd["score_breakdown"],
        "platform_facebook": 1,
        "platform_instagram": 1,
        "platform_linkedin": 1,
        "topic_class":       "primär",
    }
    sid = insert_signal(signal)
    return sid, sid is not None


def run_youtube(cfg: dict) -> int:
    ycfg = cfg["sources"]["youtube"]
    if not ycfg.get("enabled"):
        return 0
    saved = 0
    t0 = time.monotonic()
    items_seen = items_new = items_dup = rejected = 0

    # A. Beobachtete Kanäle via RSS
    for ch in ycfg.get("channels", []):
        cid = ch["channel_id"]
        vids = _yt_rss_channel_videos(cid)
        name = f"YouTube-Kanal {ch.get('name', cid)}"
        for v in vids:
            items_seen += 1
            sid, neu = _store_youtube_video(v, name, ycfg.get("min_score", 31))
            if sid is None:
                rejected += 1
            elif neu:
                items_new += 1; saved += 1
                print(f"[YT] NEU (Feed): {v['title'][:60]}")
            else:
                items_dup += 1

    # B/C. Suchthemen rotiert
    topics = list(dict.fromkeys(
        ycfg.get("search_topics", []) + ycfg.get("product_topics", [])))
    per_run = int(ycfg.get("searches_per_run", 4))
    day_key = datetime.now(timezone.utc).strftime("%Y%m%d%H")
    offset = int(day_key) % max(len(topics), 1)
    rotated = [topics[(offset + i) % len(topics)] for i in range(min(per_run, len(topics)))]

    for topic in rotated:
        vids = yt_search(topic)
        print(f"[YT] Suche '{topic}': {len(vids)} Videos")
        for v in vids:
            items_seen += 1
            sid, neu = _store_youtube_video(v, f"YouTube-Suche '{topic}'",
                                            ycfg.get("min_score", 31))
            if sid is None:
                rejected += 1
            elif neu:
                items_new += 1; saved += 1
                print(f"[YT] NEU: {v['title'][:60]} ({v['channel']})")
            else:
                items_dup += 1

    record_source_run("youtube", t0, items_seen, items_new, items_dup,
                      clustered=0, rejected=rejected, error=None)
    print(f"[YT] Gesamt: gesehen={items_seen} neu={items_new} dup={items_dup} reject={rejected}")
    return saved


# ─────────────── Facebook loginfrei über Bing-RSS-Websuche ──────────

_FB_SKIP = re.compile(r"/login|register|m\.me|l\.facebook|/facebook/?$|sharer|plugins")


def fb_websuche(query: str) -> list[dict]:
    """Loginfreie Facebook-Signalfindung über Bing-RSS (site:-Queries).
    Keine Cookies, keine Session, kein Playwright."""
    url = ("https://www.bing.com/search?q=" + urllib.parse.quote(query)
           + "&format=rss&count=25")
    raw = _get(url)
    if not raw:
        return []
    try:
        root = ET.fromstring(raw)
    except ET.ParseError:
        return []
    out = []
    for item in root.findall(".//item"):
        link = (item.findtext("link", "") or "").strip()
        title = (item.findtext("title", "") or "").strip()
        desc = (item.findtext("description", "") or "").strip()
        pubDate = (item.findtext("pubDate", "") or "").strip()
        if not link or "facebook.com" not in link:
            continue
        if _FB_SKIP.search(link) or link.rstrip("/") == "https://www.facebook.com":
            continue
        out.append({"title": title, "url": link, "description": desc, "pubDate": pubDate})
    return out


def run_facebook(cfg: dict) -> int:
    fcfg = cfg["sources"]["facebook"]
    if not fcfg.get("enabled"):
        return 0
    t0 = time.monotonic()
    saved = items_seen = items_new = items_dup = rejected = 0
    queries = fcfg.get("queries", [])
    per_run = int(fcfg.get("queries_per_run", 3))
    day_key = datetime.now(timezone.utc).strftime("%Y%m%d%H")
    offset = int(day_key) % max(len(queries), 1)
    rotated = [queries[(offset + i) % len(queries)] for i in range(min(per_run, len(queries)))]

    expires = (datetime.now(timezone.utc) + timedelta(hours=240)).isoformat()
    err = None
    for q in rotated:
        results = fb_websuche(q)
        print(f"[FB] '{q}': {len(results)} FB-Treffer")
        for r in results:
            items_seen += 1
            # Nur Titel/Beschreibung als Signal – KEIN fremder Posttext wird kopiert
            text = f"{r['title']} {r['description'][:300]}"
            sd = score_signal(text, language="de", base_category="feuchte",
                              source_url=r["url"])
            if sd["radar_score"] < 31:
                rejected += 1
                continue
            signal = {
                "source_type":       "facebook_public",
                "source_platform":   "facebook",
                "observed_at":       __import__("datetime").datetime.now(__import__("datetime").timezone.utc).isoformat(),
                "source_name":       f"FB-Signal '{q}'",
                "source_url":        r["url"],
                "language":          "de",
                "signal_category":   sd["signal_category"] or "fb_oeffentlich",
                "topic":             r["title"][:200] or "(Facebook-Diskussion)",
                "short_summary":     (r["description"] or r["title"])[:500],
                "extracted_hook":    sd["extracted_hook"],
                "emotional_direction": sd["emotional_direction"],
                "urgency_level":     sd["urgency_level"],
                "mallorca_relevance": sd["mallorca_relevance"],
                "risk_relevance":    sd["risk_relevance"],
                "estimated_noise_level": sd["estimated_noise_level"],
                "suggested_case_types": sd["suggested_case_types"],
                "suggested_platforms": sd["suggested_platforms"],
                "suggested_cta":     sd["suggested_cta"],
                "confidence_score":  sd["confidence_score"],
                "radar_score":       sd["radar_score"],
                "decay_rate":        1.0,
                "expires_at":        expires,
                "virality_level":    sd["virality_level"],
                "emotionality_level": sd["emotionality_level"],
                "comment_potential": sd["comment_potential"],
                "content_potential": sd["content_potential"],
                "score_breakdown":   sd["score_breakdown"],
                "platform_facebook": 1,
                "topic_class":       "sozial",
            }
            sid = insert_signal(signal)
            if sid is None:
                items_dup += 1
            else:
                items_new += 1; saved += 1
                print(f"[FB] NEU: {r['title'][:60]} -> {r['url'][:80]}")

    record_source_run("facebook_public", t0, items_seen, items_new, items_dup,
                      clustered=0, rejected=rejected, error=err)
    print(f"[FB] Gesamt: gesehen={items_seen} neu={items_new} dup={items_dup} reject={rejected}")
    return saved


# ─────────────────────── Orchestrator-Einstieg ──────────────────────

def run_all() -> int:
    """Führt YouTube, Facebook Organic und Meta Ads isoliert aus."""
    cfg = load_config()
    total = 0
    for runner in (run_youtube, run_facebook, run_meta_ads):
        try:
            total += runner(cfg)
        except Exception as e:
            src = getattr(runner, "__name__", "?")
            print(f"[MULTI] QUELLENFEHLER {src}: {e}")
            record_source_run(src, time.time(), 0, 0, 0, 0, 0, error=str(e))
    return total


if __name__ == "__main__":
    run_all()