"""Publish-ready package builder for the existing SMA cockpit.
The module reads the canonical signal database and persists publication packages
in the existing social-media-agent ``cases.db``. It never calls a platform API.
The daily radar shell script is the only scheduler; this module is a child step.
"""
from __future__ import annotations
import argparse
import hashlib
import json
import re
import shutil
import sqlite3
import sys
import uuid
from dataclasses import asdict, dataclass
from datetime import datetime, timezone
from pathlib import Path
from urllib.parse import urlparse
BASE_DIR = Path(__file__).resolve().parent
RADAR_DIR = Path(__import__("os").environ.get("SMA_RADAR_DIR", "/opt/struktur/social-media-radar")).resolve()
CASES_DB = Path(__import__("os").environ.get("SMA_CASES_DB", str(BASE_DIR / "cases.db"))).resolve()
RADAR_DB = Path(__import__("os").environ.get("SMA_DATABASE_PATH", "/var/lib/sma-data/signals.db")).resolve()
PACKAGE_DIR = Path(__import__("os").environ.get("SMA_PACKAGES_DIR", str(BASE_DIR / "generated_packages"))).resolve()
EVERGREEN_SOURCE = Path("/opt/struktur/mirofish/lueftungsprofi/produktblatt.md")
EVERGREEN_ASSET_SOURCE = Path("/opt/struktur/kurse/lueftung-pro/content-assets/asset-library.json")
EVERGREEN_LIBRARY_PATH = Path(__import__("os").environ.get("SMA_EVERGREEN_LIBRARY", str(BASE_DIR / "evergreen_library.json"))).resolve()
RELEVANT_CATEGORIES = {"feuchte", "schimmel_risiko", "wetter_klima", "ferienimmobilie"}
RELEVANT_TERMS = ("feucht", "schimmel", "lüft", "luftqualität", "co2", "voc", "kondens", "raumklima", "hitze", "leerstand", "immobil")
BAD_FREIGABE = {"verwerfen", "zurueckstellen", "zurückstellen"}
PILLAR_BY_CATEGORY = {
"feuchte": "feuchte_schimmel", "schimmel_risiko": "feuchte_schimmel",
"wetter_klima": "mallorca_gebaeudesituationen", "ferienimmobilie": "mallorca_gebaeudesituationen",
}
CTA_TEXTS = {
"beratung": "Wenn Sie die Situation Ihres Objekts fachlich einordnen möchten, sprechen Sie uns an.",
"messung": "Bei sichtbaren Feuchtezeichen kann eine professionelle Messung der sinnvolle nächste Schritt sein.",
"objektpruefung": "Sammeln Sie die wichtigsten Objektdaten und lassen Sie die Situation fachlich prüfen.",
"kontakt": "Fragen zu Ihrer Mallorca-Immobilie? Starten Sie mit einer persönlichen Erstberatung.",
"erfahrung": "Welche Erfahrung haben Sie mit Feuchte oder Raumklima in Ihrer Immobilie gemacht?",
"weiterfuehrend": "Speichern Sie den Beitrag und prüfen Sie die nächsten Schritte in Ruhe.",
}
TRANSFORMABLE_SOURCE_TERMS = ("maintenance", "property care", "villa management", "ferienimmobil", "leerstand", "hausverwaltung", "feuchte", "lueft", "luft")
REJECT_TITLE_TERMS = ("privacy", "privacidad", "legal notice", "aviso legal", "avís legal", "cookie", "terms and conditions", "impressum", "for sale", "luxury villa", "property tour", "sea view villa")
try:
import yaml
_editorial_config = yaml.safe_load((RADAR_DIR / "sources_multi.yaml").read_text(encoding="utf-8")).get("editorial_filters", {})
REJECT_TITLE_TERMS = tuple(_editorial_config.get("reject_title_patterns", REJECT_TITLE_TERMS))
TRANSFORMABLE_SOURCE_TERMS = tuple(_editorial_config.get("transformable_source_terms", TRANSFORMABLE_SOURCE_TERMS))
except Exception:
pass
@dataclass(frozen=True)
class Assessment:
publish_ready: bool
block_reason: str = ""
def _text(value) -> str:
return str(value or "").strip()
def _now() -> str:
return datetime.now(timezone.utc).isoformat()
def _direct_url(value: str) -> bool:
try:
parsed = urlparse(_text(value))
return parsed.scheme in {"http", "https"} and bool(parsed.netloc) and "news.google.com" not in parsed.netloc.lower()
except ValueError:
return False
def _relevant(signal: dict) -> bool:
category = _text(signal.get("signal_category")).lower()
haystack = " ".join(_text(signal.get(k)) for k in ("topic", "topic_de", "short_summary", "briefing_winkel")).lower()
return category in RELEVANT_CATEGORIES or any(term in haystack for term in RELEVANT_TERMS)
def load_evergreen_library() -> list[dict]:
return json.loads(EVERGREEN_LIBRARY_PATH.read_text(encoding="utf-8"))
def content_pillar(signal: dict) -> str:
return _text(signal.get("content_pillar")) or PILLAR_BY_CATEGORY.get(_text(signal.get("signal_category")).lower(), "praxiswissen")
def _transformable(signal: dict) -> bool:
"""A source may provide an occasion; it must not be mistaken for the post itself."""
title = _text(signal.get("topic")).lower()
source = " ".join((_text(signal.get("source_name")), _text(signal.get("source_url")), _text(signal.get("short_summary")))).lower()
if any(term in title for term in REJECT_TITLE_TERMS):
return False
if not _direct_url(signal.get("final_article_url") or signal.get("source_url")):
return False
if float(signal.get("mallorca_relevance") or 0) <= 0:
return False
if not any(term in source for term in TRANSFORMABLE_SOURCE_TERMS):
return False
category = _text(signal.get("signal_category")).lower()
return category in RELEVANT_CATEGORIES and (category != "allgemein")
def transform_signal(signal: dict) -> dict | None:
if not _transformable(signal):
return None
category = _text(signal.get("signal_category")).lower()
pillar = PILLAR_BY_CATEGORY.get(category, "mallorca_gebaeudesituationen")
if category == "wetter_klima":
angle = "Das externe Wetterereignis ist nur der Anlass: Für Innenräume zählt, wie Außenfeuchte, Regen und anschließende Hitze im Objekt bewertet werden."
advice = "Nicht automatisch lange lüften, sondern Innen- und Außenbedingungen vergleichen und Feuchtezeichen nach dem Ereignis prüfen."
hook = "Wetter ist der Anlass – das Raumklima entscheidet."
cta_type = "objektpruefung"
elif category in {"feuchte", "schimmel_risiko"}:
angle = "Das Signal zeigt einen konkreten Anlass, aber die fachliche Frage lautet: Welche Feuchtequelle liegt vor und was passiert im Raum?"
advice = "Muffigen Geruch, Kondensat oder sichtbare Veränderungen dokumentieren und die Ursache messen lassen, bevor nur überstrichen wird."
hook = "Nicht das Symptom überstreichen: Erst die Feuchteursache verstehen."
cta_type = "messung"
else:
angle = "Bei einer Mallorca-Ferienimmobilie gehört das Raumklima zur Objektbetreuung – besonders dann, wenn Eigentümer nicht dauerhaft vor Ort sind."
advice = "Leerstandszeiten, Raumdaten und sichtbare Veränderungen in eine regelmäßige Objektprüfung einbeziehen."
hook = "Ferienimmobilien brauchen auch während der Abwesenheit einen Plan."
cta_type = "objektpruefung"
result = dict(signal)
result.update({
"content_pillar": pillar, "transformed": True, "transformation_angle": angle,
"transformation_advice": advice, "transformation_hook": hook, "cta_type": cta_type,
"briefing_freigabe": "geeignet", "briefing_winkel": hook,
"briefing_nutzen": "Aus einem externen Anlass wird konkretes, fachliches Handlungswissen für Mallorca-Immobilien.",
"briefing_risiko": "", "image_prompt": _text(signal.get("image_url")) or "Fachlich passende, fotorealistische Mallorca-Immobilie mit sichtbarem Mess- oder Lüftungskontext; keine Texte, keine Logos.",
})
return result
def build_evergreen_candidate(entry: dict | None = None) -> dict:
"""Build one traceable fallback from the persisted internal knowledge library."""
entry = entry or {
"evergreen_id": "ev-012", "title": "Warum Feuchte- und Luftqualitätsmonitoring bei Mallorca-Immobilien sinnvoll ist",
"content_pillar": "monitoring_homeprotect", "core_claim": "Das Produktblatt beschreibt KI HomeProtect zur Überwachung von Feuchte, Temperatur und Luftqualität für abwesende Eigentümer.",
"target_group": "Abwesende Eigentümer von Ferienimmobilien auf Mallorca", "benefit": "Eigentümer erhalten einen nachvollziehbaren Ansatz für Feuchte- und Luftqualitätskontrolle.",
"source_refs": ["produktblatt.md"], "hooks": ["Was passiert in Ihrem Haus, wenn Sie nicht vor Ort sind?"], "cta_type": "beratung",
"media_idea": "Fotorealistisches, helles mediterranes Ferienhaus-Interieur auf Mallorca mit unaufdringlichem Feuchte- und Luftqualitätssensor, keine Texte, keine Logos, sachliche Premium-Fotografie.", "platforms": ["facebook", "linkedin"]
}
source_path = str(EVERGREEN_SOURCE if "produktblatt" in " ".join(entry.get("source_refs", [])) else EVERGREEN_ASSET_SOURCE)
return {
"signal_id": "evergreen-" + entry["evergreen_id"], "evergreen_id": entry["evergreen_id"], "evergreen": True,
"topic": entry["title"], "topic_de": "", "signal_category": entry["content_pillar"], "content_pillar": entry["content_pillar"],
"short_summary": entry["core_claim"], "source_name": "evergreen:knowledge-library", "source_url": "", "final_article_url": "",
"provenance_path": source_path, "source_refs": entry.get("source_refs", []), "radar_score": 60, "mallorca_relevance": 30,
"briefing_winkel": entry["hooks"][0] if entry.get("hooks") else entry["title"], "briefing_nutzen": entry["benefit"],
"briefing_zielgruppe": entry["target_group"], "briefing_post_typ": "evergreen", "briefing_freigabe": "geeignet", "briefing_risiko": "",
"image_url": "", "image_found": 0, "image_prompt": entry["media_idea"], "cta_type": entry.get("cta_type", "beratung"), "publish_status": "", "source_count": 1,
"knowledge_entry": entry,
}
def platform_variants(signal: dict, *, cta_text: str | None = None) -> dict:
"""Create different presentation, never different facts."""
topic = _text(signal.get("topic"))
summary = _text(signal.get("short_summary"))
hook = _text(signal.get("transformation_hook") or signal.get("briefing_winkel") or topic)
angle = _text(signal.get("transformation_angle"))
advice = _text(signal.get("transformation_advice"))
cta = cta_text or CTA_TEXTS.get(_text(signal.get("cta_type")), CTA_TEXTS["beratung"])
fb = f"{hook}\n\n{summary}\n\n{angle or 'Für Mallorca-Immobilien ist die Einordnung des konkreten Objekts entscheidend.'}\n\n{advice or _text(signal.get('briefing_nutzen'))}\n\n{cta}"
li = f"{topic}\n\nFachliche Einordnung: {angle or summary}\n\nPraktischer Nutzen: {advice or _text(signal.get('briefing_nutzen'))}\n\n{cta}"
return {"facebook": fb.strip(), "linkedin": li.strip(), "cta": cta, "hook": hook}
def _package_text(signal: dict, *, cta_text: str | None = None) -> tuple[str, str, str, str]:
if signal.get("transformed") or signal.get("evergreen"):
variants = platform_variants(signal, cta_text=cta_text)
return variants["facebook"], variants["linkedin"], variants["hook"], variants["cta"]
result = _render_text(signal)
return result[0], result[1], _text(signal.get("briefing_winkel") or signal.get("topic")), cta_text or CTA_TEXTS["beratung"]
def assess_candidate(signal: dict, *, has_final_text: bool = True, has_media: bool | None = None) -> Assessment:
"""Apply the business publish-ready gate without changing signal state."""
if not _text(signal.get("signal_id")):
return Assessment(False, "signal_id_missing")
if _text(signal.get("publish_status")).lower() in {"gepostet", "veroeffentlicht", "manuell_veroeffentlicht"}:
return Assessment(False, "already_published")
if not _relevant(signal) or float(signal.get("mallorca_relevance") or 0) <= 0:
return Assessment(False, "business_relevance_missing")
if float(signal.get("radar_score") or 0) < 45:
return Assessment(False, "score_below_publish_threshold")
freigabe = _text(signal.get("briefing_freigabe")).lower()
if freigabe in BAD_FREIGABE:
return Assessment(False, "editorial_freigabe_not_suitable")
if freigabe not in {"geeignet", "freigegeben", "approved"} and not _text(signal.get("provenance_path")):
return Assessment(False, "editorial_freigabe_missing")
if not _text(signal.get("briefing_winkel")) or not _text(signal.get("briefing_nutzen")):
return Assessment(False, "briefing_incomplete")
risk = _text(signal.get("briefing_risiko"))
if risk and "kein erhöhtes risiko" not in risk.lower() and "kein erhohtes risiko" not in risk.lower():
return Assessment(False, "editorial_warning_present")
if not _text(signal.get("provenance_path")) and not _direct_url(signal.get("final_article_url") or signal.get("source_url")):
return Assessment(False, "source_url_not_direct")
if not has_final_text:
return Assessment(False, "final_text_missing")
if has_media is False:
return Assessment(False, "media_not_ready")
if has_media is None and not (_text(signal.get("image_url")) or int(signal.get("image_found") or 0) or _text(signal.get("image_prompt"))):
return Assessment(False, "media_not_ready")
return Assessment(True)
def _load_signal_rows(limit: int) -> list[dict]:
con = sqlite3.connect(f"file:{RADAR_DB.as_posix()}?mode=ro", uri=True)
con.row_factory = sqlite3.Row
rows = [dict(row) for row in con.execute(
"SELECT * FROM signals WHERE COALESCE(status,'') NOT IN ('ignoriert','veroeffentlicht') "
"AND COALESCE(publish_status,'') NOT IN ('gepostet','veroeffentlicht','manuell_veroeffentlicht') "
"ORDER BY radar_score DESC, created_at DESC LIMIT ?", (limit,)
)]
con.close()
return rows
def _render_text(signal: dict) -> tuple[str, str, str, str]:
sys.path.insert(0, str(RADAR_DIR))
sys.path.insert(0, str(BASE_DIR))
from caption_render import render_caption
if signal.get("evergreen"):
topic = _text(signal["topic"])
summary = _text(signal["short_summary"])
fb = (f"{topic}\n\n{summary}\n\n"
"Bei einer längeren Abwesenheit lohnt es sich, Feuchte, Temperatur und Luftqualität nicht nur nach Gefühl zu beurteilen. "
"Ein passendes Monitoring kann Veränderungen früher sichtbar machen.\n\n"
"👉 Mehr über präventive Immobilienbetreuung auf Mallorca erfahren → mallorca-airservices.com")
li = (f"{topic}\n\n{summary}\n\n"
"Für Eigentümer, die nicht dauerhaft vor Ort sind, ist die kontinuierliche Beobachtung des Raumklimas ein nachvollziehbarer Präventionsansatz. "
"Die konkrete Lösung muss zum Objekt und zur Nutzung passen.\n\n"
"Mallorca AirServices bietet dazu Feuchte-, Temperatur- und Luftqualitätsmonitoring im HomeProtect-Kontext an.\n\n"
"→ mallorca-airservices.com")
return fb, li, topic, "Mehr über präventive Immobilienbetreuung auf Mallorca erfahren → mallorca-airservices.com"
result = render_caption(signal)
if not result.get("success"):
return "", "", "", ""
topic = _text(signal.get("topic_de") or signal.get("topic"))
cta = _text(signal.get("suggested_cta")) or "mallorca-airservices.com"
return result["facebook"], result["linkedin"], topic, cta
def _ensure_column(con: sqlite3.Connection, table: str, column: str, definition: str) -> None:
columns = {row[1] for row in con.execute(f"PRAGMA table_info({table})")}
if column not in columns:
con.execute(f"ALTER TABLE {table} ADD COLUMN {column} {definition}")
def init_schema(con: sqlite3.Connection) -> None:
con.executescript("""
CREATE TABLE IF NOT EXISTS publication_packages (
package_id TEXT PRIMARY KEY,
content_id TEXT NOT NULL UNIQUE REFERENCES content_items(content_id),
signal_id TEXT NOT NULL,
canonical_identity TEXT NOT NULL,
topic TEXT NOT NULL,
topic_category TEXT NOT NULL DEFAULT '',
source_names_json TEXT NOT NULL DEFAULT '[]',
source_urls_json TEXT NOT NULL DEFAULT '[]',
provenance_json TEXT NOT NULL DEFAULT '{}',
evidence_summary TEXT NOT NULL DEFAULT '',
relevance TEXT NOT NULL DEFAULT '',
freshness TEXT NOT NULL DEFAULT '',
score REAL NOT NULL DEFAULT 0,
target_platforms_json TEXT NOT NULL DEFAULT '[]',
hook TEXT NOT NULL DEFAULT '',
final_text TEXT NOT NULL DEFAULT '',
cta TEXT NOT NULL DEFAULT '',
media_refs_json TEXT NOT NULL DEFAULT '[]',
image_prompt TEXT NOT NULL DEFAULT '',
alt_text TEXT NOT NULL DEFAULT '',
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
approval_status TEXT NOT NULL DEFAULT 'blocked',
block_reason TEXT NOT NULL DEFAULT '',
package_status TEXT NOT NULL DEFAULT 'blocked',
evergreen INTEGER NOT NULL DEFAULT 0
);
CREATE UNIQUE INDEX IF NOT EXISTS idx_publication_packages_signal ON publication_packages(signal_id);
CREATE INDEX IF NOT EXISTS idx_publication_packages_status ON publication_packages(package_status);
CREATE INDEX IF NOT EXISTS idx_publication_packages_updated ON publication_packages(updated_at);
""")
for column, definition in {
"content_pillar": "TEXT NOT NULL DEFAULT ''", "evergreen_id": "TEXT NOT NULL DEFAULT ''",
"topic_key": "TEXT NOT NULL DEFAULT ''", "cta_type": "TEXT NOT NULL DEFAULT ''",
"media_strategy": "TEXT NOT NULL DEFAULT ''", "facebook_text": "TEXT NOT NULL DEFAULT ''",
"linkedin_text": "TEXT NOT NULL DEFAULT ''", "platform_variants_json": "TEXT NOT NULL DEFAULT '{}'",
"publication_state": "TEXT NOT NULL DEFAULT 'publish_ready'", "last_published_at": "TEXT DEFAULT ''",
}.items():
_ensure_column(con, "publication_packages", column, definition)
con.executescript("""
CREATE TABLE IF NOT EXISTS evergreen_library (
evergreen_id TEXT PRIMARY KEY, title TEXT NOT NULL, content_pillar TEXT NOT NULL,
core_claim TEXT NOT NULL, target_group TEXT NOT NULL, benefit TEXT NOT NULL,
source_refs_json TEXT NOT NULL DEFAULT '[]', hooks_json TEXT NOT NULL DEFAULT '[]',
cta_type TEXT NOT NULL DEFAULT 'beratung', media_idea TEXT NOT NULL DEFAULT '',
platforms_json TEXT NOT NULL DEFAULT '[]', last_published_at TEXT DEFAULT '',
reuse_status TEXT NOT NULL DEFAULT 'available', created_at TEXT NOT NULL, updated_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS publication_events (
event_id TEXT PRIMARY KEY, package_id TEXT NOT NULL REFERENCES publication_packages(package_id),
platform TEXT NOT NULL, event_status TEXT NOT NULL,
published_at TEXT DEFAULT '', external_url TEXT DEFAULT '', external_post_id TEXT DEFAULT '',
final_text TEXT NOT NULL DEFAULT '', content_version TEXT NOT NULL DEFAULT '', note TEXT DEFAULT '', created_at TEXT NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_publication_events_package ON publication_events(package_id);
CREATE INDEX IF NOT EXISTS idx_publication_events_platform_status ON publication_events(platform,event_status);
""")
library = load_evergreen_library() if EVERGREEN_LIBRARY_PATH.exists() else []
now = _now()
for entry in library:
con.execute("""INSERT INTO evergreen_library(evergreen_id,title,content_pillar,core_claim,target_group,benefit,source_refs_json,hooks_json,cta_type,media_idea,platforms_json,created_at,updated_at)
VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?) ON CONFLICT(evergreen_id) DO UPDATE SET title=excluded.title,content_pillar=excluded.content_pillar,core_claim=excluded.core_claim,target_group=excluded.target_group,benefit=excluded.benefit,source_refs_json=excluded.source_refs_json,hooks_json=excluded.hooks_json,cta_type=excluded.cta_type,media_idea=excluded.media_idea,platforms_json=excluded.platforms_json,updated_at=excluded.updated_at""",
(entry["evergreen_id"], entry["title"], entry["content_pillar"], entry["core_claim"], entry["target_group"], entry["benefit"], json.dumps(entry.get("source_refs", []), ensure_ascii=False), json.dumps(entry.get("hooks", []), ensure_ascii=False), entry.get("cta_type", "beratung"), entry.get("media_idea", ""), json.dumps(entry.get("platforms", [])), now, now))
con.execute("UPDATE publication_packages SET evergreen_id='ev-012', content_pillar='monitoring_homeprotect', publication_state=COALESCE(publication_state,'publish_ready') WHERE signal_id='evergreen-produktblatt-homeprotect' AND evergreen_id=''")
def _ensure_content(con: sqlite3.Connection, signal: dict, fb: str, cta: str, status: str) -> str:
existing = con.execute("SELECT content_id FROM content_items WHERE signal_id=? ORDER BY updated_at DESC LIMIT 1", (_text(signal["signal_id"]),)).fetchone()
now = _now()
if existing:
content_id = existing[0]
con.execute("UPDATE content_items SET title=?,hook_text=?,body_text=?,cta_text=?,status=?,updated_at=?,notes=? WHERE content_id=?",
(_text(signal["topic"]), _text(signal["topic"]), fb, cta, status, now, "AA-006 publication package", content_id))
return content_id
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()
number = int(row[0].rsplit("-", 1)[1]) + 1 if row else 1
content_id = f"POST-{year}-{number:04d}"
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, _text(signal["signal_id"]), "", _text(signal["topic"]), "facebook", "post", _text(signal["topic"]), fb, cta, status, now, now, "AA-006 publication package"))
return content_id
def _package_record(signal: dict, content_id: str, fb: str, li: str, topic: str, cta: str, assessment: Assessment, media_refs: list[str]) -> dict:
now = _now()
source_url = _text(signal.get("final_article_url") or signal.get("source_url"))
source_names = [_text(signal.get("source_name"))] if _text(signal.get("source_name")) else []
return {
"package_id": f"PKG-{uuid.uuid4()}", "content_id": content_id, "signal_id": _text(signal["signal_id"]),
"canonical_identity": _text(signal.get("content_hash") or signal.get("url_hash") or signal["signal_id"]),
"topic": topic or _text(signal.get("topic")), "topic_category": _text(signal.get("signal_category")),
"source_names_json": json.dumps(source_names, ensure_ascii=False),
"source_urls_json": json.dumps([source_url] if source_url else [], ensure_ascii=False),
"provenance_json": json.dumps({"source_name": _text(signal.get("source_name")), "source_platform": _text(signal.get("source_platform")), "provenance_path": _text(signal.get("provenance_path")), "evergreen": bool(signal.get("evergreen"))}, ensure_ascii=False),
"evidence_summary": _text(signal.get("short_summary")),
"relevance": _text(signal.get("briefing_nutzen")), "freshness": _text(signal.get("observed_at") or signal.get("created_at")),
"score": float(signal.get("radar_score") or 0), "target_platforms_json": json.dumps(["facebook", "linkedin"]),
"hook": _text(signal.get("briefing_winkel") or topic), "final_text": fb, "cta": cta,
"media_refs_json": json.dumps(media_refs, ensure_ascii=False), "image_prompt": _text(signal.get("image_prompt")),
"alt_text": f"{topic} – Mallorca AirServices", "created_at": now, "updated_at": now,
"approval_status": "unreviewed", "block_reason": assessment.block_reason,
"package_status": "publish_ready" if assessment.publish_ready else "blocked", "evergreen": int(bool(signal.get("evergreen"))),
"content_pillar": content_pillar(signal), "evergreen_id": _text(signal.get("evergreen_id")),
"topic_key": hashlib.sha256(_text(signal.get("topic_key") or signal.get("topic")).casefold().encode("utf-8")).hexdigest()[:16],
"cta_type": _text(signal.get("cta_type") or "beratung"), "media_strategy": "existing_source_image" if _text(signal.get("image_url")) else "prepared_image_prompt",
"facebook_text": fb, "linkedin_text": li, "platform_variants_json": json.dumps({"facebook": fb, "linkedin": li}, ensure_ascii=False),
"publication_state": "publish_ready", "last_published_at": "",
}
def _persist(signal: dict, fb: str, li: str, topic: str, cta: str, assessment: Assessment, media_refs: list[str]) -> dict:
con = sqlite3.connect(CASES_DB)
con.execute("PRAGMA foreign_keys=ON")
init_schema(con)
content_status = "bereit_zum_posten" if assessment.publish_ready else "entwurf"
content_id = _ensure_content(con, signal, fb, cta, content_status)
record = _package_record(signal, content_id, fb, li, topic, cta, assessment, media_refs)
existing = con.execute("SELECT package_id FROM publication_packages WHERE signal_id=?", (record["signal_id"],)).fetchone()
if existing:
record["package_id"] = existing[0]
record["created_at"] = con.execute("SELECT created_at FROM publication_packages WHERE package_id=?", (existing[0],)).fetchone()[0]
cols = [k for k in record if k not in {"package_id", "created_at"}]
con.execute(f"UPDATE publication_packages SET {','.join(k+'=?' for k in cols)} WHERE package_id=?", [record[k] for k in cols] + [record["package_id"]])
else:
cols = list(record)
con.execute(f"INSERT INTO publication_packages({','.join(cols)}) VALUES({','.join('?' for _ in cols)})", [record[k] for k in cols])
con.commit()
con.close()
return {"package_id": record["package_id"], "content_id": content_id, "package_status": record["package_status"], "block_reason": record["block_reason"]}
def get_packages(*, package_status: str | None = None, limit: int = 100) -> list[dict]:
con = sqlite3.connect(f"file:{CASES_DB.as_posix()}?mode=ro", uri=True)
con.row_factory = sqlite3.Row
if package_status:
rows = con.execute("SELECT * FROM publication_packages WHERE package_status=? ORDER BY updated_at DESC LIMIT ?", (package_status, limit)).fetchall()
else:
rows = con.execute("SELECT * FROM publication_packages ORDER BY updated_at DESC LIMIT ?", (limit,)).fetchall()
result = []
for row in rows:
item = dict(row)
for key in ("source_names_json", "source_urls_json", "target_platforms_json", "media_refs_json"):
try:
item[key.removesuffix("_json")] = json.loads(item.pop(key))
except (TypeError, json.JSONDecodeError):
item[key.removesuffix("_json")] = []
try:
item["provenance"] = json.loads(item.pop("provenance_json"))
except (TypeError, json.JSONDecodeError):
item["provenance"] = {}
result.append(item)
con.close()
return result
def decide_package(package_id: str, decision: str, note: str = "") -> dict:
"""Record a human decision; this never creates a publication call."""
mapping = {"freigeben": ("approved", "publish_ready", ""), "ablehnen": ("rejected", "blocked", "human_rejected"), "ueberarbeiten": ("revision_requested", "blocked", "human_revision_requested")}
if decision not in mapping:
raise ValueError("decision must be freigeben, ablehnen or ueberarbeiten")
approval_status, package_status, reason = mapping[decision]
con = sqlite3.connect(CASES_DB)
con.row_factory = sqlite3.Row
row = con.execute("SELECT content_id FROM publication_packages WHERE package_id=?", (package_id,)).fetchone()
if not row:
con.close()
raise ValueError("publication package not found")
now = _now()
publication_state = "approved" if decision == "freigeben" else ("rejected" if decision == "ablehnen" else "revision_requested")
con.execute("UPDATE publication_packages SET approval_status=?,package_status=?,publication_state=?,block_reason=?,updated_at=? WHERE package_id=?", (approval_status, package_status, publication_state, reason, now, package_id))
content_status = "bereit_zum_posten" if decision == "freigeben" else "entwurf"
con.execute("UPDATE content_items SET status=?,updated_at=?,notes=? WHERE content_id=?", (content_status, now, note or decision, row["content_id"]))
con.commit()
con.close()
return {"package_id": package_id, "approval_status": approval_status, "package_status": package_status, "content_status": content_status}
def _existing_evergreen_ids() -> set[str]:
try:
con = sqlite3.connect(f"file:{CASES_DB.as_posix()}?mode=ro", uri=True)
ids = {row[0] for row in con.execute("SELECT evergreen_id FROM publication_packages WHERE evergreen_id<>'' AND package_status IN ('publish_ready','blocked')")}
con.close()
return ids
except sqlite3.OperationalError:
return set()
def select_evergreen_entries(limit: int = 2) -> list[dict]:
used = _existing_evergreen_ids()
selected, pillars = [], set()
for entry in load_evergreen_library():
if entry["evergreen_id"] in used or entry["content_pillar"] in pillars:
continue
selected.append(entry); pillars.add(entry["content_pillar"])
if len(selected) >= limit:
break
return selected
def select_current_transformations(rows: list[dict], limit: int = 1) -> list[dict]:
selected, seen_topics = [], set()
for row in rows:
transformed = transform_signal(row)
if not transformed:
continue
key = hashlib.sha256(_text(transformed.get("topic")).casefold().encode("utf-8")).hexdigest()[:16]
if key in seen_topics:
continue
transformed["topic_key"] = key
selected.append(transformed); seen_topics.add(key)
if len(selected) >= limit:
break
return selected
def record_publication_event(package_id: str, platform: str, event_status: str, *, published_at: str = "", external_url: str = "", external_post_id: str = "", note: str = "") -> dict:
if platform not in {"facebook", "linkedin"}:
raise ValueError("platform must be facebook or linkedin")
if event_status not in {"published", "not_published"}:
raise ValueError("event_status must be published or not_published")
con = sqlite3.connect(CASES_DB); con.row_factory = sqlite3.Row
row = con.execute("SELECT * FROM publication_packages WHERE package_id=?", (package_id,)).fetchone()
if not row: con.close(); raise ValueError("publication package not found")
if row["approval_status"] != "approved":
con.close(); raise ValueError("package must be approved before publication event")
text = row["facebook_text"] if platform == "facebook" else row["linkedin_text"]
now = _now(); event_time = published_at or now
con.execute("INSERT INTO publication_events(event_id,package_id,platform,event_status,published_at,external_url,external_post_id,final_text,content_version,note,created_at) VALUES(?,?,?,?,?,?,?,?,?,?,?)",
(str(uuid.uuid4()), package_id, platform, event_status, event_time if event_status == "published" else "", external_url, external_post_id, text, _text(row["topic_key"]), note, now))
if event_status == "published":
con.execute("UPDATE publication_packages SET last_published_at=?,publication_state='published',updated_at=? WHERE package_id=?", (event_time, now, package_id))
if row["evergreen_id"]:
con.execute("UPDATE evergreen_library SET last_published_at=?,reuse_status='used',updated_at=? WHERE evergreen_id=?", (event_time, now, row["evergreen_id"]))
con.commit(); con.close()
return {"package_id": package_id, "platform": platform, "event_status": event_status, "published_at": event_time if event_status == "published" else ""}
def run(*, limit: int = 50, dry_run: bool = True) -> dict:
rows = _load_signal_rows(limit)
current = select_current_transformations(rows, limit=1)
evergreen_entries = select_evergreen_entries(limit=2)
selected = current + [build_evergreen_candidate(entry) for entry in evergreen_entries]
packages, blocked = [], []
for signal in selected:
cta_type = _text(signal.get("cta_type") or "beratung")
cta_text = CTA_TEXTS.get(cta_type, CTA_TEXTS["beratung"])
fb, li, hook, cta = _package_text(signal, cta_text=cta_text)
final_assess = assess_candidate(signal, has_final_text=bool(fb and li), has_media=True)
item = {"signal_id": signal["signal_id"], "topic": signal.get("topic"), "content_pillar": content_pillar(signal), "publish_ready": final_assess.publish_ready, "block_reason": final_assess.block_reason, "evergreen": bool(signal.get("evergreen")), "transformed": bool(signal.get("transformed")), "cta_type": cta_type, "facebook": fb, "linkedin": li, "cta": cta, "media": _text(signal.get("image_url")) or _text(signal.get("image_prompt"))}
if final_assess.publish_ready:
media_refs = ["source:" + _text(signal.get("image_url"))] if signal.get("image_url") else ["image_prompt:" + _text(signal.get("image_prompt"))]
if not dry_run: item.update(_persist(signal, fb, li, signal.get("topic", ""), cta, final_assess, media_refs))
packages.append(item)
else:
blocked.append(item)
from collections import Counter
return {"dry_run": dry_run, "considered": len(rows), "direct_current": len(current), "evergreen_candidates": len(evergreen_entries), "publish_ready": len(packages), "evergreen_fallbacks": len(evergreen_entries), "blocked": len(blocked), "persisted_blocked": 0, "block_reasons": dict(Counter(x["block_reason"] for x in blocked)), "pillars": dict(Counter(x["content_pillar"] for x in packages)), "packages": packages, "blocked_examples": blocked[:10]}
def main() -> int:
parser = argparse.ArgumentParser()
parser.add_argument("--limit", type=int, default=50)
parser.add_argument("--production", action="store_true", help="Persist packages in existing cases.db; no external publishing")
args = parser.parse_args()
print(json.dumps(run(limit=args.limit, dry_run=not args.production), ensure_ascii=False, indent=2))
return 0
if __name__ == "__main__":
raise SystemExit(main())