Explorer
/root/backups-aa006/20260826-pre/cases_db.py
← Zurück ↓ Download
"""
Cases DB — Operations-Datenbank für Social Media Cockpit
Datei: /opt/struktur/social-media-agent/cases.db

Tabellen:
  cases              — Fallverwaltung
  case_media         — Medien pro Fall
  signal_case_links  — Signal↔Fall-Verknüpfungen
  content_items      — Content-Pipeline
  content_media_links — Content↔Medien-Verknüpfungen
  audit_log          — Vollständiger Verlauf
  publishing_jobs    — Publishing-Queue-Stub
"""

import sqlite3
from datetime import datetime, timezone
from pathlib import Path
from radar_database import RADAR_DB_PATH

CASES_DB = Path("/opt/struktur/social-media-agent/cases.db")


def _conn():
    con = sqlite3.connect(str(CASES_DB))
    con.row_factory = sqlite3.Row
    con.execute("PRAGMA journal_mode=WAL")
    con.execute("PRAGMA foreign_keys=ON")
    return con


def init_cases_db():
    """Legt alle Tabellen an (idempotent)."""
    con = _conn()
    con.executescript("""
    CREATE TABLE IF NOT EXISTS cases (
        case_id        TEXT PRIMARY KEY,
        title          TEXT NOT NULL,
        case_type      TEXT NOT NULL DEFAULT 'sonstiges',
        location_area  TEXT DEFAULT '',
        problem_desc   TEXT DEFAULT '',
        status         TEXT NOT NULL DEFAULT 'offen',
        created_at     TEXT NOT NULL,
        updated_at     TEXT NOT NULL
    );

    CREATE TABLE IF NOT EXISTS case_media (
        media_id       INTEGER PRIMARY KEY AUTOINCREMENT,
        case_id        TEXT NOT NULL REFERENCES cases(case_id) ON DELETE CASCADE,
        filename       TEXT NOT NULL,
        file_type      TEXT NOT NULL DEFAULT 'foto',
        stored_path    TEXT NOT NULL,
        caption        TEXT DEFAULT '',
        uploaded_at    TEXT NOT NULL
    );

    CREATE TABLE IF NOT EXISTS signal_case_links (
        id             INTEGER PRIMARY KEY AUTOINCREMENT,
        signal_id      TEXT NOT NULL,
        case_id        TEXT NOT NULL REFERENCES cases(case_id) ON DELETE CASCADE,
        linked_at      TEXT NOT NULL,
        note           TEXT DEFAULT '',
        UNIQUE(signal_id, case_id)
    );

    CREATE TABLE IF NOT EXISTS content_items (
        content_id         TEXT PRIMARY KEY,
        signal_id          TEXT DEFAULT '',
        case_id            TEXT DEFAULT '',
        title              TEXT NOT NULL,
        platform           TEXT NOT NULL DEFAULT 'facebook',
        content_type       TEXT NOT NULL DEFAULT 'post',
        hook_text          TEXT DEFAULT '',
        body_text          TEXT DEFAULT '',
        cta_text           TEXT DEFAULT '',
        status             TEXT NOT NULL DEFAULT 'idee',
        scheduled_at       TEXT DEFAULT '',
        published_at       TEXT DEFAULT '',
        external_post_id   TEXT DEFAULT '',
        external_post_url  TEXT DEFAULT '',
        platform_account   TEXT DEFAULT '',
        created_at         TEXT NOT NULL,
        updated_at         TEXT NOT NULL,
        notes              TEXT DEFAULT ''
    );

    CREATE TABLE IF NOT EXISTS content_media_links (
        id             INTEGER PRIMARY KEY AUTOINCREMENT,
        content_id     TEXT NOT NULL REFERENCES content_items(content_id) ON DELETE CASCADE,
        media_id       INTEGER NOT NULL REFERENCES case_media(media_id) ON DELETE CASCADE,
        UNIQUE(content_id, media_id)
    );

    CREATE TABLE IF NOT EXISTS audit_log (
        log_id         INTEGER PRIMARY KEY AUTOINCREMENT,
        entity_type    TEXT NOT NULL,
        entity_id      TEXT NOT NULL,
        action         TEXT NOT NULL,
        from_value     TEXT DEFAULT '',
        to_value       TEXT DEFAULT '',
        source         TEXT NOT NULL DEFAULT 'manual',
        note           TEXT DEFAULT '',
        changed_at     TEXT NOT NULL
    );

    CREATE TABLE IF NOT EXISTS publishing_jobs (
        job_id             INTEGER PRIMARY KEY AUTOINCREMENT,
        content_id         TEXT NOT NULL REFERENCES content_items(content_id) ON DELETE CASCADE,
        platform           TEXT NOT NULL,
        platform_account   TEXT DEFAULT '',
        scheduled_at       TEXT DEFAULT '',
        status             TEXT NOT NULL DEFAULT 'pending',
        external_id        TEXT DEFAULT '',
        external_url       TEXT DEFAULT '',
        error_message      TEXT DEFAULT '',
        created_at         TEXT NOT NULL,
        updated_at         TEXT NOT NULL
    );

    CREATE INDEX IF NOT EXISTS idx_cases_status        ON cases(status);
    CREATE INDEX IF NOT EXISTS idx_cases_type          ON cases(case_type);
    CREATE INDEX IF NOT EXISTS idx_media_case          ON case_media(case_id);
    CREATE INDEX IF NOT EXISTS idx_links_case          ON signal_case_links(case_id);
    CREATE INDEX IF NOT EXISTS idx_links_signal        ON signal_case_links(signal_id);
    CREATE INDEX IF NOT EXISTS idx_content_status      ON content_items(status);
    CREATE INDEX IF NOT EXISTS idx_content_platform    ON content_items(platform);
    CREATE INDEX IF NOT EXISTS idx_audit_entity        ON audit_log(entity_type, entity_id);
    CREATE INDEX IF NOT EXISTS idx_audit_changed       ON audit_log(changed_at);
    CREATE INDEX IF NOT EXISTS idx_publishing_content  ON publishing_jobs(content_id);
    CREATE INDEX IF NOT EXISTS idx_publishing_status   ON publishing_jobs(status);
    """)
    con.commit()
    con.close()


def _now():
    return datetime.now(timezone.utc).isoformat()


def _next_case_id(con):
    year = datetime.now().year
    row = con.execute(
        "SELECT case_id FROM cases WHERE case_id LIKE ? ORDER BY case_id DESC LIMIT 1",
        (f"CASE-{year}-%",)
    ).fetchone()
    if row:
        try:
            n = int(row["case_id"].split("-")[-1]) + 1
        except (ValueError, IndexError):
            n = 1
    else:
        n = 1
    return f"CASE-{year}-{n:04d}"


def _next_content_id(con):
    year = datetime.now().year
    row = con.execute(
        "SELECT content_id FROM content_items WHERE content_id LIKE ? ORDER BY content_id DESC LIMIT 1",
        (f"POST-{year}-%",)
    ).fetchone()
    if row:
        try:
            n = int(row["content_id"].split("-")[-1]) + 1
        except (ValueError, IndexError):
            n = 1
    else:
        n = 1
    return f"POST-{year}-{n:04d}"


# ---------------------------------------------------------------------------
# Audit
# ---------------------------------------------------------------------------

def add_audit(entity_type, entity_id, action, from_value="", to_value="",
              source="manual", note=""):
    con = _conn()
    con.execute(
        "INSERT INTO audit_log(entity_type,entity_id,action,from_value,to_value,source,note,changed_at)"
        " VALUES(?,?,?,?,?,?,?,?)",
        (entity_type, entity_id, action, str(from_value), str(to_value), source, note, _now())
    )
    con.commit()
    con.close()


def get_audit(entity_id, limit=50):
    con = _conn()
    rows = con.execute(
        "SELECT * FROM audit_log WHERE entity_id = ? ORDER BY changed_at DESC LIMIT ?",
        (entity_id, limit)
    ).fetchall()
    con.close()
    return [dict(r) for r in rows]


# ---------------------------------------------------------------------------
# Cases
# ---------------------------------------------------------------------------

VALID_CASE_TYPES = [
    "feuchte", "schimmel", "lueftung", "kondensation",
    "waermebruecke", "luftqualitaet", "sonstiges"
]

VALID_CASE_STATUS = ["offen", "in_arbeit", "abgeschlossen", "archiviert"]


def get_cases(status_filter=None, limit=200):
    con = _conn()
    if status_filter == "no_media":
        rows = con.execute(
            "SELECT c.*, "
            " (SELECT COUNT(*) FROM case_media m WHERE m.case_id=c.case_id) as media_count,"
            " (SELECT COUNT(*) FROM signal_case_links l WHERE l.case_id=c.case_id) as signal_count"
            " FROM cases c WHERE c.status != 'archiviert'"
            " AND (SELECT COUNT(*) FROM case_media m WHERE m.case_id=c.case_id) = 0"
            " ORDER BY c.updated_at DESC LIMIT ?",
            (limit,)
        ).fetchall()
    elif status_filter and status_filter != "alle":
        rows = con.execute(
            "SELECT c.*, "
            " (SELECT COUNT(*) FROM case_media m WHERE m.case_id=c.case_id) as media_count,"
            " (SELECT COUNT(*) FROM signal_case_links l WHERE l.case_id=c.case_id) as signal_count"
            " FROM cases c WHERE c.status=? ORDER BY c.updated_at DESC LIMIT ?",
            (status_filter, limit)
        ).fetchall()
    else:
        rows = con.execute(
            "SELECT c.*, "
            " (SELECT COUNT(*) FROM case_media m WHERE m.case_id=c.case_id) as media_count,"
            " (SELECT COUNT(*) FROM signal_case_links l WHERE l.case_id=c.case_id) as signal_count"
            " FROM cases c ORDER BY c.updated_at DESC LIMIT ?",
            (limit,)
        ).fetchall()
    con.close()
    return [dict(r) for r in rows]


def get_case(case_id):
    con = _conn()
    row = con.execute("SELECT * FROM cases WHERE case_id=?", (case_id,)).fetchone()
    if not row:
        con.close()
        return None
    case = dict(row)
    case["media"] = [dict(r) for r in con.execute(
        "SELECT * FROM case_media WHERE case_id=? ORDER BY uploaded_at DESC", (case_id,)
    ).fetchall()]
    case["signal_links"] = [dict(r) for r in con.execute(
        "SELECT * FROM signal_case_links WHERE case_id=? ORDER BY linked_at DESC", (case_id,)
    ).fetchall()]
    case["audit"] = get_audit(case_id)
    case["content_items"] = [dict(r) for r in con.execute(
        "SELECT content_id,title,platform,status,created_at FROM content_items WHERE case_id=? ORDER BY created_at DESC",
        (case_id,)
    ).fetchall()]
    con.close()
    return case


def create_case(title, case_type, location_area="", problem_desc=""):
    con = _conn()
    case_id = _next_case_id(con)
    now = _now()
    con.execute(
        "INSERT INTO cases(case_id,title,case_type,location_area,problem_desc,status,created_at,updated_at)"
        " VALUES(?,?,?,?,?,?,?,?)",
        (case_id, title, case_type, location_area, problem_desc, "offen", now, now)
    )
    con.commit()
    con.close()
    add_audit("case", case_id, "erstellt", "", "offen", "manual", f"Titel: {title}")
    return case_id


def update_case_field(case_id, field, value):
    """Aktualisiert ein einzelnes Feld und schreibt Audit-Eintrag."""
    allowed = {"title", "case_type", "location_area", "problem_desc"}
    if field not in allowed:
        raise ValueError(f"Feld nicht erlaubt: {field}")
    con = _conn()
    old = con.execute(f"SELECT {field} FROM cases WHERE case_id=?", (case_id,)).fetchone()
    old_val = old[0] if old else ""
    con.execute(f"UPDATE cases SET {field}=?, updated_at=? WHERE case_id=?",
                (value, _now(), case_id))
    con.commit()
    con.close()
    add_audit("case", case_id, f"feld_{field}", old_val, value)


def update_case_status(case_id, new_status):
    if new_status not in VALID_CASE_STATUS:
        raise ValueError(f"Ungültiger Status: {new_status}")
    con = _conn()
    row = con.execute("SELECT status FROM cases WHERE case_id=?", (case_id,)).fetchone()
    if not row:
        con.close()
        raise ValueError("Fall nicht gefunden")
    old_status = row["status"]
    con.execute("UPDATE cases SET status=?, updated_at=? WHERE case_id=?",
                (new_status, _now(), case_id))
    con.commit()
    con.close()
    add_audit("case", case_id, "status_change", old_status, new_status)


# ---------------------------------------------------------------------------
# Media
# ---------------------------------------------------------------------------

def add_media(case_id, filename, file_type, stored_path, caption=""):
    con = _conn()
    con.execute(
        "INSERT INTO case_media(case_id,filename,file_type,stored_path,caption,uploaded_at)"
        " VALUES(?,?,?,?,?,?)",
        (case_id, filename, file_type, stored_path, caption, _now())
    )
    con.execute("UPDATE cases SET updated_at=? WHERE case_id=?", (_now(), case_id))
    con.commit()
    con.close()
    add_audit("case", case_id, "media_upload", "", filename, "manual", f"Typ: {file_type}")


def delete_media(media_id, case_id):
    con = _conn()
    row = con.execute("SELECT filename,stored_path FROM case_media WHERE media_id=? AND case_id=?",
                      (media_id, case_id)).fetchone()
    if not row:
        con.close()
        return False
    filename = row["filename"]
    stored_path = row["stored_path"]
    con.execute("DELETE FROM case_media WHERE media_id=?", (media_id,))
    con.execute("UPDATE cases SET updated_at=? WHERE case_id=?", (_now(), case_id))
    con.commit()
    con.close()
    # Datei löschen
    try:
        Path(stored_path).unlink(missing_ok=True)
    except Exception:
        pass
    add_audit("case", case_id, "media_delete", filename, "", "manual")
    return True


# ---------------------------------------------------------------------------
# Signal↔Fall-Verknüpfung
# ---------------------------------------------------------------------------

def link_signal_to_case(signal_id, case_id, note=""):
    con = _conn()
    try:
        con.execute(
            "INSERT OR IGNORE INTO signal_case_links(signal_id,case_id,linked_at,note)"
            " VALUES(?,?,?,?)",
            (signal_id, case_id, _now(), note)
        )
        con.execute("UPDATE cases SET updated_at=? WHERE case_id=?", (_now(), case_id))
        con.commit()
        linked = con.execute(
            "SELECT changes() as n"
        ).fetchone()
    finally:
        con.close()
    add_audit("case", case_id, "signal_linked", "", signal_id)
    add_audit("signal", signal_id, "fall_zugeordnet", "", case_id)
    return True


def unlink_signal(signal_id, case_id):
    con = _conn()
    con.execute("DELETE FROM signal_case_links WHERE signal_id=? AND case_id=?",
                (signal_id, case_id))
    con.commit()
    con.close()
    add_audit("case", case_id, "signal_unlinked", signal_id, "")


def get_signals_for_case(case_id):
    """Holt Signal-Details für einen Fall (join mit signals.db)."""
    import sqlite3 as _sq
    RADAR_DB = RADAR_DB_PATH
    con = _conn()
    links = con.execute(
        "SELECT signal_id, linked_at, note FROM signal_case_links WHERE case_id=?",
        (case_id,)
    ).fetchall()
    con.close()
    if not links:
        return []
    result = []
    if RADAR_DB.exists():
        rcon = _sq.connect(str(RADAR_DB))
        rcon.row_factory = _sq.Row
        for lnk in links:
            row = rcon.execute(
                "SELECT signal_id, topic, topic_de, signal_category, radar_score,"
                " operational_priority, status, source_name, created_at"
                " FROM signals WHERE signal_id=?",
                (lnk["signal_id"],)
            ).fetchone()
            if row:
                d = dict(row)
                d["linked_at"] = lnk["linked_at"]
                d["link_note"] = lnk["note"]
            else:
                d = {"signal_id": lnk["signal_id"], "topic": "Signal nicht mehr vorhanden",
                     "linked_at": lnk["linked_at"], "link_note": lnk["note"]}
            result.append(d)
        rcon.close()
    else:
        result = [{"signal_id": l["signal_id"], "linked_at": l["linked_at"]} for l in links]
    return result


# ---------------------------------------------------------------------------
# Content Items
# ---------------------------------------------------------------------------

VALID_PLATFORMS     = ["linkedin", "facebook", "instagram", "reels"]
VALID_CONTENT_TYPES = ["post", "carousel", "reel", "story", "artikel"]
VALID_CONTENT_STATUS = ["idee", "entwurf", "pruefung", "freigegeben",
                         "bereit_zum_posten", "geplant",
                         "manuell_veroeffentlicht", "veroeffentlicht", "archiviert"]


def create_content_item(title, platform, content_type="post",
                        signal_id="", case_id="",
                        hook_text="", body_text="", cta_text="", notes=""):
    con = _conn()
    content_id = _next_content_id(con)
    now = _now()
    con.execute(
        "INSERT INTO content_items"
        "(content_id,signal_id,case_id,title,platform,content_type,"
        " hook_text,body_text,cta_text,status,created_at,updated_at,notes)"
        " VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?)",
        (content_id, signal_id, case_id, title, platform, content_type,
         hook_text, body_text, cta_text, "idee", now, now, notes)
    )
    con.commit()
    con.close()
    add_audit("content", content_id, "erstellt", "", "idee", "manual", f"Plattform: {platform}")
    return content_id


def get_content_items(status_filter=None, platform=None, limit=100):
    con = _conn()
    clauses = []
    params = []
    if status_filter and status_filter != "alle":
        clauses.append("status=?")
        params.append(status_filter)
    if platform:
        clauses.append("platform=?")
        params.append(platform)
    where = ("WHERE " + " AND ".join(clauses)) if clauses else ""
    rows = con.execute(
        f"SELECT * FROM content_items {where} ORDER BY updated_at DESC LIMIT ?",
        params + [limit]
    ).fetchall()
    con.close()
    return [dict(r) for r in rows]


def update_content_status(content_id, new_status):
    if new_status not in VALID_CONTENT_STATUS:
        raise ValueError(f"Ungültiger Status: {new_status}")
    con = _conn()
    row = con.execute("SELECT status FROM content_items WHERE content_id=?",
                      (content_id,)).fetchone()
    if not row:
        con.close()
        raise ValueError("Content-Item nicht gefunden")
    old = row["status"]
    now = _now()
    extra = {}
    if new_status in ("veroeffentlicht", "manuell_veroeffentlicht"):
        extra = {"published_at": now}
    updates = ", ".join([f"{k}=?" for k in ["status", "updated_at"] + list(extra.keys())])
    con.execute(
        f"UPDATE content_items SET {updates} WHERE content_id=?",
        [new_status, now] + list(extra.values()) + [content_id]
    )
    con.commit()
    con.close()
    add_audit("content", content_id, "status_change", old, new_status)


# ---------------------------------------------------------------------------
# Dashboard-Stats
# ---------------------------------------------------------------------------

def get_ops_stats():
    """Übersichts-Zahlen für das Dashboard."""
    con = _conn()
    stats = {}
    for status in VALID_CASE_STATUS:
        n = con.execute("SELECT COUNT(*) FROM cases WHERE status=?", (status,)).fetchone()[0]
        stats[f"cases_{status}"] = n
    stats["cases_total"] = con.execute("SELECT COUNT(*) FROM cases").fetchone()[0]
    stats["media_total"] = con.execute("SELECT COUNT(*) FROM case_media").fetchone()[0]
    stats["signal_links_total"] = con.execute("SELECT COUNT(*) FROM signal_case_links").fetchone()[0]
    for cstatus in VALID_CONTENT_STATUS:
        n = con.execute("SELECT COUNT(*) FROM content_items WHERE status=?", (cstatus,)).fetchone()[0]
        stats[f"content_{cstatus}"] = n
    stats["content_total"] = con.execute("SELECT COUNT(*) FROM content_items").fetchone()[0]
    stats["publishing_pending"] = con.execute(
        "SELECT COUNT(*) FROM publishing_jobs WHERE status='pending'"
    ).fetchone()[0]
    con.close()
    return stats


# ---------------------------------------------------------------------------
# Content Items — Phase 2 Ergänzungen
# ---------------------------------------------------------------------------

def get_content_item(content_id):
    """Gibt ein Content-Item mit verknüpften Daten zurück."""
    con = _conn()
    row = con.execute("SELECT * FROM content_items WHERE content_id=?",
                      (content_id,)).fetchone()
    if not row:
        con.close()
        return None
    item = dict(row)
    # Verknüpfte Medien (über case_media)
    item["media"] = [dict(r) for r in con.execute(
        "SELECT cm.*, cml.content_id FROM case_media cm"
        " JOIN content_media_links cml ON cm.media_id = cml.media_id"
        " WHERE cml.content_id=? ORDER BY cm.uploaded_at DESC",
        (content_id,)
    ).fetchall()]
    # Audit
    item["audit"] = [dict(r) for r in con.execute(
        "SELECT * FROM audit_log WHERE entity_id=? AND entity_type='content'"
        " ORDER BY changed_at DESC LIMIT 30",
        (content_id,)
    ).fetchall()]
    # Fall-Info wenn verknüpft
    item["case_info"] = None
    if item.get("case_id"):
        crow = con.execute(
            "SELECT case_id, title, case_type, location_area, status FROM cases WHERE case_id=?",
            (item["case_id"],)
        ).fetchone()
        if crow:
            item["case_info"] = dict(crow)
            # Medien des Falls für Auswahl
            item["case_media_all"] = [dict(r) for r in con.execute(
                "SELECT * FROM case_media WHERE case_id=? ORDER BY uploaded_at DESC",
                (item["case_id"],)
            ).fetchall()]
        else:
            item["case_media_all"] = []
    else:
        item["case_media_all"] = []
    con.close()
    return item


def update_content_item(content_id, fields):
    """Aktualisiert Content-Item-Felder. fields = dict mit erlaubten Schlüsseln."""
    allowed = {"title", "platform", "content_type", "hook_text", "body_text",
               "cta_text", "scheduled_at", "platform_account", "notes",
               "external_post_url", "external_post_id"}
    filtered = {k: v for k, v in fields.items() if k in allowed}
    if not filtered:
        return
    con = _conn()
    sets = ", ".join(f"{k}=?" for k in filtered)
    vals = list(filtered.values()) + [_now(), content_id]
    con.execute(f"UPDATE content_items SET {sets}, updated_at=? WHERE content_id=?", vals)
    con.commit()
    con.close()
    add_audit("content", content_id, "bearbeitet", "",
              ", ".join(filtered.keys()), "manual")


def get_pipeline_board():
    """Gibt alle Content-Items gruppiert nach Status zurück."""
    con = _conn()
    result = {}
    for status in VALID_CONTENT_STATUS:
        rows = con.execute(
            "SELECT content_id, title, platform, content_type, case_id, signal_id,"
            " status, updated_at, scheduled_at, notes,"
            " (SELECT COUNT(*) FROM content_media_links cml WHERE cml.content_id=content_items.content_id) as media_count"
            " FROM content_items WHERE status=? ORDER BY updated_at DESC LIMIT 50",
            (status,)
        ).fetchall()
        result[status] = [dict(r) for r in rows]
    con.close()
    return result


def link_media_to_content(content_id, media_id):
    con = _conn()
    con.execute(
        "INSERT OR IGNORE INTO content_media_links(content_id, media_id) VALUES(?,?)",
        (content_id, media_id)
    )
    con.execute("UPDATE content_items SET updated_at=? WHERE content_id=?",
                (_now(), content_id))
    con.commit()
    con.close()
    add_audit("content", content_id, "media_verknuepft", "", str(media_id))


def unlink_media_from_content(content_id, media_id):
    con = _conn()
    con.execute(
        "DELETE FROM content_media_links WHERE content_id=? AND media_id=?",
        (content_id, media_id)
    )
    con.execute("UPDATE content_items SET updated_at=? WHERE content_id=?",
                (_now(), content_id))
    con.commit()
    con.close()


# ---------------------------------------------------------------------------
# Publishing Jobs
# ---------------------------------------------------------------------------


def get_publishing_job(job_id):
    """Einzelner Publishing-Job mit vollem Content-Item (inkl. Medien)."""
    con = _conn()
    row = con.execute(
        "SELECT pj.*, ci.title as content_title, ci.platform as ci_platform,"
        "  ci.hook_text, ci.body_text, ci.cta_text, ci.notes,"
        "  ci.status as content_status, ci.signal_id, ci.case_id,"
        "  ci.published_at as content_published_at, ci.external_post_url,"
        "  ci.content_type, ci.scheduled_at as ci_scheduled_at"
        " FROM publishing_jobs pj"
        " LEFT JOIN content_items ci ON pj.content_id = ci.content_id"
        " WHERE pj.job_id=?",
        (job_id,)
    ).fetchone()
    if not row:
        con.close()
        return None
    job = dict(row)
    # Medien des Content-Items
    if job.get("content_id"):
        job["media"] = [dict(r) for r in con.execute(
            "SELECT cm.* FROM case_media cm"
            " JOIN content_media_links cml ON cm.media_id = cml.media_id"
            " WHERE cml.content_id=? ORDER BY cm.uploaded_at",
            (job["content_id"],)
        ).fetchall()]
    else:
        job["media"] = []
    # Fall-Info
    if job.get("case_id"):
        crow = con.execute(
            "SELECT case_id, title, case_type, location_area, status FROM cases WHERE case_id=?",
            (job["case_id"],)
        ).fetchone()
        job["case_info"] = dict(crow) if crow else None
    else:
        job["case_info"] = None
    con.close()
    return job


def mark_publishing_job_published(job_id, external_url=""):
    """Markiert Job als manuell_veroeffentlicht + setzt Content-Status."""
    con = _conn()
    now = _now()
    row = con.execute(
        "SELECT content_id FROM publishing_jobs WHERE job_id=?", (job_id,)
    ).fetchone()
    if not row:
        con.close()
        raise ValueError("Job nicht gefunden")
    content_id = row["content_id"]
    # Job updaten
    con.execute(
        "UPDATE publishing_jobs SET status=?, external_url=?, updated_at=? WHERE job_id=?",
        ("manuell_veroeffentlicht", external_url, now, job_id)
    )
    # Content updaten
    con.execute(
        "UPDATE content_items SET status=?, published_at=?, external_post_url=?, updated_at=?"
        " WHERE content_id=?",
        ("manuell_veroeffentlicht", now, external_url, now, content_id)
    )
    con.commit()
    con.close()
    add_audit("content", content_id, "status_change", "bereit_zum_posten",
              "manuell_veroeffentlicht", "manual", f"Job {job_id}, URL: {external_url[:80]}")
    add_audit("publishing_job", str(job_id), "veroeffentlicht", "pending",
              "manuell_veroeffentlicht", "manual", external_url[:80])

def get_publishing_jobs(status_filter=None, limit=100):
    con = _conn()
    if status_filter and status_filter != "alle":
        rows = con.execute(
            "SELECT pj.*, ci.title as content_title, ci.platform as content_platform"
            " FROM publishing_jobs pj"
            " LEFT JOIN content_items ci ON pj.content_id = ci.content_id"
            " WHERE pj.status=? ORDER BY pj.created_at DESC LIMIT ?",
            (status_filter, limit)
        ).fetchall()
    else:
        rows = con.execute(
            "SELECT pj.*, ci.title as content_title, ci.platform as content_platform"
            " FROM publishing_jobs pj"
            " LEFT JOIN content_items ci ON pj.content_id = ci.content_id"
            " ORDER BY pj.created_at DESC LIMIT ?",
            (limit,)
        ).fetchall()
    con.close()
    return [dict(r) for r in rows]


def create_publishing_job(content_id, platform, scheduled_at="", platform_account=""):
    con = _conn()
    now = _now()
    con.execute(
        "INSERT INTO publishing_jobs(content_id,platform,platform_account,"
        " scheduled_at,status,created_at,updated_at)"
        " VALUES(?,?,?,?,?,?,?)",
        (content_id, platform, platform_account, scheduled_at, "pending", now, now)
    )
    con.commit()
    con.close()
    add_audit("content", content_id, "publishing_job_erstellt", "", platform)


def update_publishing_job_status(job_id, new_status, external_url="", error_message=""):
    con = _conn()
    now = _now()
    con.execute(
        "UPDATE publishing_jobs SET status=?, external_url=?, error_message=?, updated_at=?"
        " WHERE job_id=?",
        (new_status, external_url, error_message, now, job_id)
    )
    con.commit()
    con.close()


# ---------------------------------------------------------------------------
# Dashboard — Stale Signals
# ---------------------------------------------------------------------------

def get_stale_signals(hours=48, limit=10):
    """Signale ohne Aktion seit X Stunden aus signals.db."""
    import sqlite3 as _sq
    from datetime import timedelta
    RADAR_DB = RADAR_DB_PATH
    if not RADAR_DB.exists():
        return []
    try:
        cutoff = (datetime.now(timezone.utc) - timedelta(hours=hours)).isoformat()
        rcon = _sq.connect(str(RADAR_DB))
        rcon.row_factory = _sq.Row
        rows = rcon.execute(
            "SELECT signal_id, topic, topic_de, signal_category, radar_score,"
            " operational_priority, created_at, status"
            " FROM signals"
            " WHERE status = 'neu' AND radar_score >= 50 AND created_at < ?"
            " ORDER BY operational_priority DESC, radar_score DESC LIMIT ?",
            (cutoff, limit)
        ).fetchall()
        rcon.close()
        return [dict(r) for r in rows]
    except Exception:
        return []


def get_top_priority_signals(limit=5):
    """Top-Signale nach operational_priority aus signals.db."""
    import sqlite3 as _sq
    RADAR_DB = RADAR_DB_PATH
    if not RADAR_DB.exists():
        return []
    try:
        rcon = _sq.connect(str(RADAR_DB))
        rcon.row_factory = _sq.Row
        rows = rcon.execute(
            "SELECT signal_id, topic, topic_de, signal_category, radar_score,"
            " operational_priority, recommended_action, action_platform,"
            " status, image_url, image_found, topic_class, created_at"
            " FROM signals"
            " WHERE status NOT IN ('ignoriert','veroeffentlicht')"
            " AND radar_score >= 50"
            " ORDER BY pinned DESC, operational_priority DESC, radar_score DESC LIMIT ?",
            (limit,)
        ).fetchall()
        rcon.close()
        return [dict(r) for r in rows]
    except Exception:
        return []


def get_recent_cases(limit=5):
    """Zuletzt aktualisierte Fälle."""
    con = _conn()
    rows = con.execute(
        "SELECT c.case_id, c.title, c.case_type, c.location_area, c.status, c.updated_at,"
        " (SELECT COUNT(*) FROM case_media m WHERE m.case_id=c.case_id) as media_count,"
        " (SELECT COUNT(*) FROM signal_case_links l WHERE l.case_id=c.case_id) as signal_count"
        " FROM cases c"
        " WHERE c.status != 'archiviert'"
        " ORDER BY c.updated_at DESC LIMIT ?",
        (limit,)
    ).fetchall()
    con.close()
    return [dict(r) for r in rows]


def get_content_in_progress(limit=8):
    """Content-Items im aktiven Bearbeitungsstatus."""
    con = _conn()
    rows = con.execute(
        "SELECT ci.content_id, ci.title, ci.platform, ci.status, ci.case_id,"
        " ci.signal_id, ci.updated_at,"
        " c.title as case_title"
        " FROM content_items ci"
        " LEFT JOIN cases c ON ci.case_id = c.case_id"
        " WHERE ci.status IN ('entwurf','pruefung','freigegeben','geplant')"
        " ORDER BY ci.updated_at DESC LIMIT ?",
        (limit,)
    ).fetchall()
    con.close()
    return [dict(r) for r in rows]