"""AA-045: Multi-Source-Anreicherung von insert_signal."""
import json
import sqlite3
from datetime import datetime, timezone
from published_at_normalizer import normalize_published_at
def enrich_signal_row(conn: sqlite3.Connection, signal_id: str, data: dict):
now = datetime.now(timezone.utc).isoformat()
observed_at = data.get("observed_at", "") or now
published = normalize_published_at(data.get("published_at"), observed_at)
ad_start = data.get("ad_start_at") or (published if data.get("source_type") == "meta_ad" else None)
last_seen = data.get("last_seen_at") if data.get("source_type") == "meta_ad" else None
active_status = data.get("active_status", "") if data.get("source_type") == "meta_ad" else ""
conn.execute("""UPDATE signals SET
source_platform=COALESCE(NULLIF(?, ''), source_type), external_id=?,
observed_at=?, published_at=?, content_hash=?, processing_state='neu',
ad_start_at=?, last_seen_at=?,
active_status=COALESCE(NULLIF(?, ''), active_status, 'unknown')
WHERE signal_id=?""", (
data.get("source_platform", "") or data.get("source_type", ""),
data.get("external_id", "") or "", observed_at, published,
data.get("content_hash", "") or "", ad_start, last_seen,
active_status, signal_id))
def find_by_external_id(conn: sqlite3.Connection, platform: str, external_id: str):
if not external_id:
return None
row = conn.execute("SELECT signal_id FROM signals WHERE external_id=? AND (source_platform=? OR ?='')",
(external_id, platform, platform)).fetchone()
return row["signal_id"] if row else None
def update_cluster_meta(conn: sqlite3.Connection, cluster_id: str, signal_ids: list, bonus: float):
now = datetime.now(timezone.utc).isoformat()
for sid in signal_ids:
conn.execute("UPDATE signals SET cluster_id=?, radar_score=MIN(radar_score+?,100.0), score_breakdown=COALESCE(score_breakdown,'{}') WHERE signal_id=?",
(cluster_id, bonus, sid))
return now