"""
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()