"""Job- und Log-Persistenz in PostgreSQL (KONZEPT: Postgres für Job-Logs). Die Tabellendefinition existiert bewusst identisch in API und Worker (docker/api/db.py) — es gibt kein geteiltes Paket zwischen den Containern. Wer die Struktur ändert, ändert BEIDE Dateien. create_all ist idempotent. """ import os from datetime import datetime, timezone from sqlalchemy import ( Column, DateTime, Integer, MetaData, String, Table, Text, create_engine, ) DATABASE_URL = os.getenv( "DATABASE_URL", "postgresql://rippy:rippy@localhost:5432/rippy" ) engine = create_engine(DATABASE_URL, pool_pre_ping=True) metadata = MetaData() jobs = Table( "jobs", metadata, Column("id", String(36), primary_key=True), Column("disc_type", String(16)), Column("device", String(64)), Column("title", String(255)), Column("status", String(16), nullable=False, server_default="pending"), Column("progress", Integer, nullable=False, server_default="0"), Column("output_path", Text), Column("error", Text), Column("created_at", DateTime(timezone=True)), Column("finished_at", DateTime(timezone=True)), ) logs = Table( "logs", metadata, Column("id", Integer, primary_key=True, autoincrement=True), Column("ts", DateTime(timezone=True)), Column("level", String(16)), Column("source", String(32)), Column("message", Text), ) def utcnow() -> datetime: return datetime.now(timezone.utc) def init_db() -> None: """Legt fehlende Tabellen an (idempotent).""" metadata.create_all(engine) settings_table = Table( "settings", metadata, Column("key", String(64), primary_key=True), Column("value", Text), ) workers = Table( "workers", metadata, Column("name", String(128), primary_key=True), Column("encoders", Text), Column("last_seen", DateTime(timezone=True)), ) def save_worker(name: str, encoder_liste: list) -> None: """Worker meldet Name + Encoder-Fähigkeiten (Upsert).""" import json payload = json.dumps(encoder_liste) with engine.begin() as conn: vorhanden = conn.execute( workers.select().where(workers.c.name == name) ).first() if vorhanden: conn.execute( workers.update().where(workers.c.name == name).values( encoders=payload, last_seen=utcnow() ) ) else: conn.execute( workers.insert().values(name=name, encoders=payload, last_seen=utcnow()) ) def get_settings(key: str = "ui") -> dict: """UI-Einstellungen lesen (der Worker respektiert Transcode-Optionen).""" import json from sqlalchemy import select try: with engine.connect() as conn: zeile = conn.execute( select(settings_table.c.value).where(settings_table.c.key == key) ).first() if zeile and zeile[0]: return json.loads(zeile[0]) except Exception: pass return {} def update_job(job_id: str, **fields) -> None: with engine.begin() as conn: conn.execute(jobs.update().where(jobs.c.id == job_id).values(**fields)) def add_log(level: str, source: str, message: str) -> None: with engine.begin() as conn: conn.execute( logs.insert().values(ts=utcnow(), level=level, source=source, message=message) )