Explorer
/proc/2012/task/2095/root/tmp/aa053-signal_db.py
← Zurück ↓ Download
"""
Signal DB -- Social Media Radar
Canonical Signal Layer + Duplicate Detection
============================================
Neue Logik beim Einfuegen:
  1. URL-Hash: exakter URL-Duplikat -> ablehnen
  2. Jaccard-Similarity: gleiche Story, andere Quelle -> Quelle hinzufuegen
  3. Neu: neues kanonisches Signal anlegen
"""

import sqlite3
import uuid
import re
import hashlib
import json
from datetime import datetime, timezone, timedelta
from pathlib import Path
from sma_database import SIGNALS_DB_PATH
import aa045_schema
import aa045_enrich
from sma_core.signal_intake import SignalIntakeRepository

DB_PATH = SIGNALS_DB_PATH

_STOPWORDS = {
    "de", "die", "das", "der", "und", "oder", "in", "auf", "ist", "von",
    "mit", "zu", "ein", "eine", "sich", "an", "nach", "bei", "hat",
    "the", "a", "an", "and", "or", "in", "on", "is", "of", "for", "to",
    "el", "la", "los", "las", "de", "en", "un", "una", "por", "para",
    "nos", "sus", "del", "con", "que", "se", "nos", "sus",
}

# Jaccard-Schwellenwert fuer Canonical-Matching
JACCARD_THRESHOLD = 0.45
# Zeitfenster fuer Canonical-Suche (Stunden)
CANONICAL_WINDOW_HOURS = 48


def get_conn():
    conn = sqlite3.connect(str(DB_PATH))
    conn.row_factory = sqlite3.Row
    conn.execute("PRAGMA foreign_keys=ON")
    conn.execute("PRAGMA busy_timeout=5000")
    return conn


def _url_hash(url: str) -> str:
    url = url.strip().lower().rstrip("/")
    return hashlib.sha256(url.encode()).hexdigest()[:16]


def _title_hash(topic: str) -> str:
    """16-char SHA256 eines normalisierten Titels fuer Title-Dedup."""
    if not topic:
        return ""
    t = topic.lower().strip()
    import re as _re
    t = _re.sub(r"[^\w\s]", " ", t)
    t = _re.sub(r"\s+", " ", t).strip()
    return hashlib.sha256(t.encode("utf-8")).hexdigest()[:16]


def _find_by_title_hash(conn, title_h: str):
    """Gibt signal_id zurueck wenn title_hash bereits in DB. Sonst None."""
    if not title_h:
        return None
    row = conn.execute(
        "SELECT signal_id FROM signals WHERE title_hash = ?", (title_h,)
    ).fetchone()
    return row["signal_id"] if row else None


def _words_set(text: str) -> set:
    """Normalisierte Wortmenge fuer Jaccard-Vergleich."""
    t = text.lower()
    t = re.sub(r"[^a-z0-9\s]", " ", t)
    return {w for w in t.split() if w not in _STOPWORDS and len(w) > 3}


def _jaccard(a: set, b: set) -> float:
    if not a or not b:
        return 0.0
    inter = len(a & b)
    union = len(a | b)
    return inter / union if union > 0 else 0.0


def init_db():
    conn = get_conn()
    c = conn.cursor()

    # Haupt-Signale-Tabelle
    c.executescript("""
        CREATE TABLE IF NOT EXISTS signals (
            signal_id            TEXT PRIMARY KEY,
            timestamp            TEXT NOT NULL,
            source_type          TEXT NOT NULL,
            source_name          TEXT NOT NULL,
            source_url           TEXT,
            language             TEXT DEFAULT 'de',
            signal_category      TEXT NOT NULL,
            topic                TEXT,
            short_summary        TEXT,
            extracted_hook       TEXT,
            emotional_direction  TEXT,
            urgency_level        INTEGER DEFAULT 0,
            mallorca_relevance   INTEGER DEFAULT 0,
            seasonal_relevance   INTEGER DEFAULT 0,
            risk_relevance       INTEGER DEFAULT 0,
            estimated_noise_level INTEGER DEFAULT 0,
            duplicate_cluster    TEXT,
            suggested_case_types TEXT,
            suggested_platforms  TEXT,
            suggested_cta        TEXT,
            confidence_score     REAL DEFAULT 0.0,
            radar_score          REAL DEFAULT 0.0,
            decay_rate           REAL DEFAULT 1.0,
            expires_at           TEXT,
            processed            INTEGER DEFAULT 0,
            created_at           TEXT NOT NULL
        );

        CREATE TABLE IF NOT EXISTS signal_sources (
            id          INTEGER PRIMARY KEY AUTOINCREMENT,
            signal_id   TEXT NOT NULL,
            source_name TEXT,
            source_url  TEXT,
            url_hash    TEXT,
            added_at    TEXT NOT NULL,
            FOREIGN KEY (signal_id) REFERENCES signals(signal_id)
        );

        CREATE INDEX IF NOT EXISTS idx_signals_radar_score ON signals(radar_score DESC);
        CREATE INDEX IF NOT EXISTS idx_signals_category    ON signals(signal_category);
        CREATE INDEX IF NOT EXISTS idx_signals_processed   ON signals(processed);
        CREATE INDEX IF NOT EXISTS idx_signals_expires     ON signals(expires_at);
        CREATE INDEX IF NOT EXISTS idx_sources_signal_id   ON signal_sources(signal_id);
        CREATE INDEX IF NOT EXISTS idx_sources_url_hash    ON signal_sources(url_hash);
    """)

    # Migration: fehlende Spalten hinzufuegen
    existing = {row[1] for row in c.execute("PRAGMA table_info(signals)").fetchall()}
    migrations = [
        ("url_hash",           "TEXT"),
        ("title_hash",         "TEXT"),
        ("status",             "TEXT DEFAULT 'neu'"),
        ("score_breakdown",    "TEXT"),
        ("virality_level",     "INTEGER DEFAULT 1"),
        ("emotionality_level", "INTEGER DEFAULT 1"),
        ("comment_potential",  "INTEGER DEFAULT 1"),
        ("content_potential",  "INTEGER DEFAULT 1"),
        ("source_count",       "INTEGER DEFAULT 1"),
        ("topic_de",           "TEXT"),
        ("image_url",          "TEXT"),
        ("image_source",       "TEXT"),
        ("image_found",        "INTEGER DEFAULT 0"),
        ("final_article_url",  "TEXT"),
        ("image_checked_at",   "TEXT"),
        ("old_radar_score",    "REAL"),
        ("cached_image_path",  "TEXT"),
        # Phase 4: Platform Fit + Themenklasse + Cockpit
        ("platform_linkedin",  "INTEGER DEFAULT 1"),
        ("platform_facebook",  "INTEGER DEFAULT 1"),
        ("platform_instagram", "INTEGER DEFAULT 1"),
        ("platform_tiktok",    "INTEGER DEFAULT 1"),
        ("mas_relevant",       "INTEGER DEFAULT 0"),
        ("mas_anchor",         "TEXT DEFAULT ''"),
        ("mas_relevance_reason", "TEXT DEFAULT ''"),
        ("content_blocked",    "INTEGER DEFAULT 0"),
        ("content_block_reason", "TEXT DEFAULT ''"),
        ("topic_class",        "TEXT DEFAULT 'sekundaer'"),
        ("competition_signal",   "INTEGER DEFAULT 0"),
        ("pinned",               "INTEGER DEFAULT 0"),
        # Phase 5: Operative Priorisierung + Aktionsempfehlung
        ("operational_priority", "REAL DEFAULT 0"),
        ("recommended_action",   "TEXT DEFAULT ''"),
        ("action_platform",      "TEXT DEFAULT ''"),
    ]
    for col, col_type in migrations:
        if col not in existing:
            c.execute(f"ALTER TABLE signals ADD COLUMN {col} {col_type}")
            print(f"[DB] Migration: '{col}' hinzugefuegt")

    # AA-045: Multi-Source-Felder + Observability
    aa045_schema.apply_migration(c)

    conn.commit()

    # Indexes fuer neue Spalten
    for stmt in [
        "CREATE INDEX IF NOT EXISTS idx_signals_status     ON signals(status)",
        "CREATE INDEX IF NOT EXISTS idx_signals_created    ON signals(created_at)",
        "CREATE INDEX IF NOT EXISTS idx_signals_url_hash   ON signals(url_hash)",
        "CREATE INDEX IF NOT EXISTS idx_signals_title_hash ON signals(title_hash)",
        "CREATE UNIQUE INDEX IF NOT EXISTS idx_signals_title_hash_uq ON signals(title_hash) WHERE title_hash IS NOT NULL AND title_hash != ''",
    ]:
        try:
            c.execute(stmt)
        except sqlite3.OperationalError:
            pass

    conn.commit()
    conn.close()
    print(f"[DB] Initialisiert: {DB_PATH}")


def _url_exists(conn, url_h: str) -> bool:
    """True wenn URL-Hash bereits als Signal-URL oder Quell-URL existiert."""
    if not url_h:
        return False
    # In signals
    if conn.execute("SELECT 1 FROM signals WHERE url_hash = ?", (url_h,)).fetchone():
        return True
    # In signal_sources
    if conn.execute("SELECT 1 FROM signal_sources WHERE url_hash = ?", (url_h,)).fetchone():
        return True
    return False


def _find_canonical(conn, topic: str) -> str | None:
    """
    Sucht ein semantisch aehnliches kanonisches Signal (letzte CANONICAL_WINDOW_HOURS h).
    Gibt signal_id zurueck wenn Jaccard >= JACCARD_THRESHOLD.
    """
    if not topic:
        return None
    candidate = _words_set(topic)
    if len(candidate) < 3:
        return None  # Zu kurz fuer sinnvollen Vergleich

    cutoff = (datetime.now(timezone.utc) - timedelta(hours=CANONICAL_WINDOW_HOURS)).isoformat()
    rows = conn.execute(
        "SELECT signal_id, topic FROM signals WHERE created_at > ?",
        (cutoff,)
    ).fetchall()

    best_id    = None
    best_score = 0.0
    for row in rows:
        existing = _words_set(row["topic"] or "")
        j = _jaccard(candidate, existing)
        if j > best_score:
            best_score = j
            best_id    = row["signal_id"]

    return best_id if best_score >= JACCARD_THRESHOLD else None


def _add_source(conn, signal_id: str, data: dict):
    """Fuegt eine Quell-URL ueber den autorisierten Intake-Port hinzu."""
    url = data.get("source_url", "") or ""
    url_h = _url_hash(url) if url else ""
    now = datetime.now(timezone.utc).isoformat()
    SignalIntakeRepository(conn).attach_source(
        signal_id,
        data.get("source_name", ""),
        url,
        url_h,
        now,
    )


def insert_signal(data: dict):
    """
    Fuegt ein Signal ein oder fuegt es einem kanonischen Signal hinzu.
    Gibt signal_id zurueck (entweder neu oder canonical).
    Gibt None zurueck wenn die URL ein exaktes Duplikat ist.
    """
    url     = data.get("source_url", "") or ""
    topic   = data.get("topic", "") or ""
    url_h   = _url_hash(url) if url else ""
    title_h = _title_hash(topic) if topic else ""

    conn = get_conn()
    try:
        # ── Check 1: Exakte URL ────────────────────────────────────────────
        if _url_exists(conn, url_h):
            return None  # Exakter URL-Duplikat

        # ── Check 1c (AA-045): Externe ID (z.B. YouTube-Video-ID) ─────────
        ext_platform = data.get("source_platform", "") or data.get("source_type", "")
        ext_id = data.get("external_id", "") or ""
        dup_id = aa045_enrich.find_by_external_id(conn, ext_platform, ext_id)
        if dup_id:
            _add_source(conn, dup_id, data)
            conn.commit()
            return dup_id  # gleiches externes Objekt, andere Quelle/Feed

        # ── Check 1b: Title-Hash (gleicher normalisierter Titel) ──────────
        if title_h:
            existing_id = _find_by_title_hash(conn, title_h)
            if existing_id:
                _add_source(conn, existing_id, data)
                conn.commit()
                return existing_id  # Gleicher Titel, andere Quelle

        # ── Check 2: Jaccard-Kanonisches Signal (semantisch gleich) ───────
        canonical_id = _find_canonical(conn, topic)
        if canonical_id:
            _add_source(conn, canonical_id, data)
            conn.commit()
            return canonical_id  # Quelle zu bestehendem Signal hinzugefuegt

        # ── Check 3: Neues kanonisches Signal ─────────────────────────────
        signal_id = str(uuid.uuid4())
        now       = datetime.now(timezone.utc).isoformat()

        conn.execute("""
            INSERT INTO signals (
                signal_id, timestamp, source_type, source_name, source_url,
                language, signal_category, topic, short_summary, extracted_hook,
                emotional_direction, urgency_level, mallorca_relevance,
                seasonal_relevance, risk_relevance, estimated_noise_level,
                duplicate_cluster, suggested_case_types, suggested_platforms,
                suggested_cta, confidence_score, radar_score, decay_rate,
                expires_at, created_at, url_hash, title_hash, status,
                score_breakdown, virality_level, emotionality_level,
                comment_potential, content_potential, source_count, topic_de, image_url,
                platform_linkedin, platform_facebook, platform_instagram, platform_tiktok,
                topic_class, competition_signal,
                operational_priority, recommended_action, action_platform,
                mas_relevant, mas_anchor, mas_relevance_reason, content_blocked, content_block_reason
            ) VALUES (
                :signal_id, :timestamp, :source_type, :source_name, :source_url,
                :language, :signal_category, :topic, :short_summary, :extracted_hook,
                :emotional_direction, :urgency_level, :mallorca_relevance,
                :seasonal_relevance, :risk_relevance, :estimated_noise_level,
                :duplicate_cluster, :suggested_case_types, :suggested_platforms,
                :suggested_cta, :confidence_score, :radar_score, :decay_rate,
                :expires_at, :created_at, :url_hash, :title_hash, :status,
                :score_breakdown, :virality_level, :emotionality_level,
                :comment_potential, :content_potential, 1, :topic_de, :image_url,
                :platform_linkedin, :platform_facebook, :platform_instagram, :platform_tiktok,
                :topic_class, :competition_signal,
                :operational_priority, :recommended_action, :action_platform,
                :mas_relevant, :mas_anchor, :mas_relevance_reason, :content_blocked, :content_block_reason
            )
        """, {
            "signal_id":            signal_id,
            "timestamp":            data.get("timestamp", now),
            "source_type":          data.get("source_type", "unknown"),
            "source_name":          data.get("source_name", ""),
            "source_url":           url,
            "language":             data.get("language", "de"),
            "signal_category":      data.get("signal_category", "uncategorized"),
            "topic":                topic,
            "short_summary":        data.get("short_summary", ""),
            "extracted_hook":       data.get("extracted_hook", ""),
            "emotional_direction":  data.get("emotional_direction", "neutral"),
            "urgency_level":        data.get("urgency_level", 0),
            "mallorca_relevance":   data.get("mallorca_relevance", 0),
            "seasonal_relevance":   data.get("seasonal_relevance", 0),
            "risk_relevance":       data.get("risk_relevance", 0),
            "estimated_noise_level": data.get("estimated_noise_level", 0),
            "duplicate_cluster":    data.get("duplicate_cluster", ""),
            "suggested_case_types": data.get("suggested_case_types", ""),
            "suggested_platforms":  data.get("suggested_platforms", ""),
            "suggested_cta":        data.get("suggested_cta", ""),
            "confidence_score":     data.get("confidence_score", 0.0),
            "radar_score":          data.get("radar_score", 0.0),
            "decay_rate":           data.get("decay_rate", 1.0),
            "expires_at":           data.get("expires_at", ""),
            "created_at":           now,
            "url_hash":             url_h,
            "status":               data.get("status", "neu"),
            "score_breakdown":      json.dumps(data.get("score_breakdown", {})),
            "virality_level":       data.get("virality_level", 1),
            "emotionality_level":   data.get("emotionality_level", 1),
            "comment_potential":    data.get("comment_potential", 1),
            "content_potential":    data.get("content_potential", 1),
            "topic_de":             data.get("topic_de", ""),
            "image_url":            data.get("image_url", ""),
            "platform_linkedin":    data.get("platform_linkedin", 1),
            "platform_facebook":    data.get("platform_facebook", 1),
            "platform_instagram":   data.get("platform_instagram", 1),
            "platform_tiktok":      data.get("platform_tiktok", 1),
            "topic_class":          data.get("topic_class", "sekundaer"),
            "competition_signal":   data.get("competition_signal", 0),
            "title_hash":           title_h,
            "operational_priority": data.get("operational_priority", 0),
            "recommended_action":   data.get("recommended_action", ""),
            "action_platform":      data.get("action_platform", ""),
            "mas_relevant":          data.get("mas_relevant", 0),
            "mas_anchor":            data.get("mas_anchor", ""),
            "mas_relevance_reason":  data.get("mas_relevance_reason", ""),
            "content_blocked":       data.get("content_blocked", 0),
            "content_block_reason":  data.get("content_block_reason", ""),
        })

        # AA-045: Felder des einheitlichen Signalmodells befuellen
        aa045_enrich.enrich_signal_row(conn, signal_id, data)

        # Bildfelder des einheitlichen Signalmodells persistieren.
        # image_found darf nur vom Intake/Refresh gesetzt werden, nachdem eine
        # belastbare Bildquelle ermittelt bzw. validiert wurde.
        if data.get("image_url"):
            conn.execute("""
                UPDATE signals SET image_url=?, image_source=?, image_found=?,
                    final_article_url=?, cached_image_path=?, image_checked_at=?
                WHERE signal_id=?
            """, (
                data.get("image_url"), data.get("image_source") or "",
                int(bool(data.get("image_found"))),
                data.get("final_article_url"), data.get("cached_image_path"),
                data.get("image_checked_at"), signal_id,
            ))

        # AA-045-F5: Watch-Entity-Felder
        if data.get("watch_entity_id"):
            conn.execute(
                "UPDATE signals SET watch_entity_id=?, watch_sector=?, watch_priority=? WHERE signal_id=?",
                (data["watch_entity_id"], data.get("watch_sector", ""),
                 data.get("watch_priority", ""), signal_id))

        # Erste Quelle eintragen
        _add_source(conn, signal_id, data)
        conn.commit()
        return signal_id

    except sqlite3.IntegrityError:
        return None
    finally:
        conn.close()


def update_signal_status(signal_id: str, status: str):
    valid = {"neu", "analysiert", "fall_gefunden", "fall_fehlt",
             "content_vorbereitet", "veroeffentlicht", "ignoriert", "beobachten", "spaeter"}
    if status not in valid:
        raise ValueError(f"Ungueltiger Status: {status}")
    conn = get_conn()
    conn.execute("UPDATE signals SET status = ? WHERE signal_id = ?", (status, signal_id))
    conn.commit()
    conn.close()


def get_active_signals(min_score: float = 35.0, limit: int = 20, days: int | None = None):
    """Get active signals, optionally filtered by published_at days.
    
    Args:
        min_score: Minimum radar score threshold
        limit: Maximum number of results
        days: Optional. If given, only return signals with published_at within 
              this many days. None = show all signals (ALLE view).
    """
    conn = get_conn()
    now = datetime.now(timezone.utc)
    
    # Base query
    query = """
        SELECT * FROM signals
        WHERE radar_score >= ?
          AND (expires_at = '' OR expires_at > ?)"""
    params = [min_score, now.isoformat()]
    
    # Add published_at filter only when days parameter is explicitly given
    if days is not None and days > 0:
        cutoff = (now - timedelta(days=days)).isoformat()
        query += """
            AND published_at IS NOT NULL 
            AND published_at != '' 
            AND published_at >= ?
        """
        params.append(cutoff)
    
    # No status filter - ALL signals visible in ALLE view
    query += """
        ORDER BY timestamp DESC
        LIMIT ?"""
    params.append(limit)
    
    rows = conn.execute(query, params).fetchall()
    conn.close()
    return [dict(r) for r in rows]


def get_score_stats() -> dict:
    conn = get_conn()
    rows = conn.execute(
        "SELECT radar_score FROM signals WHERE status NOT IN ('ignoriert', 'veroeffentlicht')"
    ).fetchall()
    conn.close()
    if not rows:
        return {}
    scores = sorted(r[0] for r in rows)
    n = len(scores)
    total = sum(scores)
    median = scores[n // 2] if n % 2 else (scores[n // 2 - 1] + scores[n // 2]) / 2
    buckets = {"irrelevant": 0, "beobachten": 0, "interessant": 0, "operativ": 0, "prioritaet": 0}
    for s in scores:
        if s <= 30:   buckets["irrelevant"]  += 1
        elif s <= 50: buckets["beobachten"]  += 1
        elif s <= 69: buckets["interessant"] += 1
        elif s <= 84: buckets["operativ"]    += 1
        else:         buckets["prioritaet"]  += 1
    clustered = max(buckets.values()) / n if n > 0 else 0
    return {
        "count":   n,
        "min":     round(scores[0], 1),
        "max":     round(scores[-1], 1),
        "avg":     round(total / n, 1),
        "median":  round(median, 1),
        "buckets": buckets,
        "warning": clustered > 0.5,
    }


def mark_processed(signal_id: str):
    conn = get_conn()
    conn.execute("UPDATE signals SET processed = 1 WHERE signal_id = ?", (signal_id,))
    conn.commit()
    conn.close()


if __name__ == "__main__":
    init_db()
    print("[DB] Setup abgeschlossen.")
    stats = get_score_stats()
    if stats:
        print(f"[DB] Stats: min={stats['min']} max={stats['max']} avg={stats['avg']} median={stats['median']}")
        print(f"[DB] Verteilung: {stats['buckets']}")
        if stats.get("warning"):
            print("[DB] WARNUNG: Score-Clustering erkannt!")


def record_source_run(source_name: str, t0_monotonic: float, items_seen: int,
                      items_new: int, items_duplicate: int, clustered: int = 0,
                      rejected: int = 0, error=None, quota_usage: str = "",
                      items_parsed: int = 0, items_canonicalized: int = 0):
    """AA-045 Observability: eine Zeile je Quellenlauf in source_runs."""
    import time as _t
    conn = get_conn()
    try:
        now = datetime.now(timezone.utc).isoformat()
        conn.execute("""
            INSERT INTO source_runs (
                source_name, run_at, last_success, last_error,
                items_seen, items_parsed, items_canonicalized, items_new,
                items_duplicate, items_clustered, items_rejected, runtime_seconds, quota_usage
            ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
        """, (
            source_name, now, None if error else now, error,
            items_seen, items_parsed, items_canonicalized, items_new,
            items_duplicate, clustered, rejected,
            round(_t.monotonic() - t0_monotonic, 2), quota_usage,
        ))
        conn.commit()
    finally:
        conn.close()