Explorer
/tmp/multi_intake.py
← Zurück ↓ Download
#!/usr/bin/env python3
"""
Multi-Source Intake — AA-045
YouTube (RSS-Feeds + Suchseiten-Scraping + Name-Search-Fallback) 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
import re
import time
import json
import base64
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 pathlib import Path
sys.path.insert(0, str(Path(__file__).parent.parent))

from signal_db import insert_signal, record_source_run
from classification.signal_classifier import score_signal

CONFIG_PATH = Path(__file__).parent.parent / "sources_multi.yaml"
WATCH_ENTITIES_PATH = Path(__file__).parent.parent / "watch_entities.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():
    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)

def load_watch_entities():
    if yaml is None:
        return []
    try:
        with open(WATCH_ENTITIES_PATH, encoding="utf-8") as fh:
            data = yaml.safe_load(fh)
        ents = [e for e in data.get("watch_entities", [])
                if e.get("active") and e.get("verification", {}).get("status") == "verified"]
        return ents
    except Exception:
        return []

# ─────────────────────────── 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)
    html_bytes = _get(f"https://www.youtube.com/results?search_query={q}")
    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

def get_youtube_channel_id_for_entity(entity_name: str) -> str | None:
    """Get YouTube channel ID from watch_entities.yaml by entity name"""
    entities = load_watch_entities()
    for entity in entities:
        if entity.get("canonical_name") == entity_name:
            return entity.get("youtube_channel_id")
    return None

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 = 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

    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", ""),
        "published_at": pub_iso,
        "engagement": 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"],
        "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

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

def run_youtube(cfg: dict) -> int:
    ycfg = cfg["sources"]["youtube"]
    if not ycfg.get("enabled"):
        return 0
    saved = 0
    t0 = time.time()
    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. YouTube name-search fallback for watch entities without channel_id or with failed RSS
    entities = load_watch_entities()
    yt_search_every_runs = int(ycfg.get("yt_search_every_runs", 3))
    day_key = datetime.now(timezone.utc).strftime("%Y%m%d%H")
    offset = int(day_key) % max(len(entities), 1)
    rotated = entities[offset:] + entities[:offset]
    entities_per_run = min(len(rotated), ycfg.get("entities_per_run", 8))
    batch = rotated[:entities_per_run]
    
    for entity in batch:
        entity_id = entity["entity_id"]
        canonical_name = entity["canonical_name"]
        existing_channel_id = entity.get("youtube_channel_id")
        
        # Determine if we should search for this entity
        should_search = False
        if not existing_channel_id:
            # No channel ID, always search
            should_search = True
            search_reason = "no channel ID"
        elif int(day_key) % yt_search_every_runs == offset % yt_search_every_runs:
            # It's time to re-check according to schedule
            should_search = True
            search_reason = "periodic re-check"
        
        if should_search:
            print(f"[YT-SEARCH] Checking {canonical_name} ({entity_id}) - {search_reason}")
            
            # Search queries to try
            search_queries = [
                f'"{canonical_name}" Mallorca',
                f'{canonical_name} Mallorca',
                f'{canonical_name} property management Mallorca',
                f'{canonical_name} house care Mallorca',
                f'{canonical_name} home management Mallorca',
                f'{canonical_name} feuchte Mallorca',
                f'{canonical_name} schimmel Mallorca',
            ]
            
            channel_found = False
            for query in search_queries:
                videos, channels = yt_search(query)
                
                if channels:
                    # Take the first channel and check if it has content
                    channel_id = channels[0]
                    rss_results = _yt_rss_channel_videos(channel_id)
                    
                    if rss_results:
                        print(f"[YT-SEARCH] Found channel for {canonical_name}: {channel_id} ({len(rss_results)} videos)")
                        channel_found = True
                        
                        # Add videos from this channel
                        for v in rss_results:
                            items_seen += 1
                            sid, neu = _store_youtube_video(v, f"YouTube-Suche '{query}'", ycfg.get("min_score", 31))
                            if sid is None:
                                rejected += 1
                            elif neu:
                                items_new += 1; saved += 1
                                print(f"[YT] NEU (Suche): {v['title'][:60]} ({v['channel']})")
                            else:
                                items_dup += 1
                        break
                    else:
                        print(f"[YT-SEARCH] Channel {channel_id} found but no RSS content for {canonical_name}")
                else:
                    print(f"[YT-SEARCH] No channels found for query: {query}")
            
            if not channel_found:
                print(f"[YT-SEARCH] No suitable YouTube channel found for {canonical_name}")
        else:
            # Try existing channel ID if we have one
            if existing_channel_id:
                vids = _yt_rss_channel_videos(existing_channel_id)
                name = f"YouTube-Kanal {canonical_name} (from entity)"
                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

    # C/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()
        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})
    return out

def run_facebook(cfg: dict) -> int:
    fcfg = cfg["sources"]["facebook"]
    if not fcfg.get("enabled"):
        return 0
    t0 = time.time()
    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": datetime.now(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 nacheinander aus; Fehler isoliert je Quelle."""
    cfg = load_config()
    total = 0
    for runner in (run_youtube, run_facebook):
        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()