"""
AA-045 Patch-Modul: Schema-Migration + Observability für den Multi-Source-Radar.
Wird von signal_db.init_db() importiert (idempotent).
"""
import json
import hashlib
from datetime import datetime, timezone
AA045_COLUMNS = [
# Einheitliches Signalmodell §7
("source_platform", "TEXT DEFAULT ''"), # youtube | facebook | web | rss | weather
("external_id", "TEXT DEFAULT ''"), # z.B. YouTube-Video-ID
("observed_at", "TEXT"), # Zeitpunkt der Erfassung
("published_at", "TEXT"), # Original-Veröffentlichung (soweit bekannt)
("content_hash", "TEXT DEFAULT ''"), # SHA-256 über normalisierten Inhalt
("cluster_id", "TEXT DEFAULT ''"), # Multi-Source-Themencluster
("engagement", "TEXT DEFAULT ''"), # JSON: views/comments etc.
("processing_state", "TEXT DEFAULT 'neu'"),
]
AA045_INDEXES = [
"CREATE INDEX IF NOT EXISTS idx_signals_external_id ON signals(external_id)",
"CREATE INDEX IF NOT EXISTS idx_signals_cluster_id ON signals(cluster_id)",
"CREATE INDEX IF NOT EXISTS idx_signals_source_platform ON signals(source_platform)",
]
def apply_migration(c) -> list[str]:
applied = []
existing = {row[1] for row in c.execute("PRAGMA table_info(signals)").fetchall()}
for col, col_type in AA045_COLUMNS:
if col not in existing:
c.execute(f"ALTER TABLE signals ADD COLUMN {col} {col_type}")
applied.append(col)
for stmt in AA045_INDEXES:
try:
c.execute(stmt)
except Exception:
pass
# Observability-Tabelle je Quelle (§18)
c.execute("""
CREATE TABLE IF NOT EXISTS source_runs (
id INTEGER PRIMARY KEY AUTOINCREMENT,
source_name TEXT NOT NULL,
run_at TEXT NOT NULL,
last_success TEXT,
last_error TEXT,
items_seen INTEGER DEFAULT 0,
items_new INTEGER DEFAULT 0,
items_duplicate INTEGER DEFAULT 0,
items_clustered INTEGER DEFAULT 0,
items_rejected INTEGER DEFAULT 0,
runtime_seconds REAL DEFAULT 0,
quota_usage TEXT DEFAULT ''
)
""")
return applied
def content_hash_for(title: str, summary: str = "") -> str:
norm = " ".join(f"{title} {summary}".lower().split())
return hashlib.sha256(norm.encode("utf-8")).hexdigest()[:16]
def record_run(conn, source_name: str, t0_monotonic: float, items_seen: int,
items_new: int, items_duplicate: int, items_clustered: int,
items_rejected: int, error: str | None, quota_usage: str = ""):
now = datetime.now(timezone.utc).isoformat()
conn.execute("""
INSERT INTO source_runs (
source_name, run_at, last_success, last_error,
items_seen, 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_new, items_duplicate, items_clustered,
items_rejected,
0.0, # runtime wird vom Aufrufer gesetzt, wenn messbar
quota_usage,
))