"""
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]