"""Job-, Log- und Settings-Persistenz in PostgreSQL (KONZEPT: Postgres für Job-Logs). Die jobs/logs-Tabellendefinition existiert bewusst identisch im Worker (docker/worker/db.py) — es gibt kein geteiltes Paket zwischen den Containern. Wer die Struktur ändert, ändert BEIDE Dateien. create_all ist idempotent. Bis 23.07. lag Postgres komplett brach: GET /jobs gab hart [] zurück, nichts schrieb je eine Zeile — der Job-Verlauf im UI war ein Placebo. """ import json import os from datetime import datetime, timezone from sqlalchemy import ( Column, DateTime, Integer, MetaData, String, Table, Text, create_engine, select, ) 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("target_dir", String(255)), Column("error", Text), Column("meta", Text), # Disc-Metadaten (JSON: Jahr/Poster/Plot) für Detail-Popup + NFO 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), ) 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("info", Text), # Werkzeug-Versionen (JSON: makemkv/handbrake/key-Quelle) Column("last_seen", DateTime(timezone=True)), ) storage_mounts = Table( "storage_mounts", metadata, Column("name", String(64), primary_key=True), Column("typ", String(8)), # nfs | cifs Column("quelle", String(255)), # host:/export bzw. //host/share Column("optionen", String(255)), Column("username", String(128)), Column("passwort", String(255)), # Klartext — Heimnetz-Kompromiss, siehe README ) def list_workers() -> list: import json with engine.connect() as conn: zeilen = conn.execute(select(workers)).mappings().all() ergebnis = [] for z in zeilen: eintrag = dict(z) try: eintrag["encoders"] = json.loads(eintrag.get("encoders") or "[]") except ValueError: eintrag["encoders"] = [] try: eintrag["info"] = json.loads(eintrag.get("info") or "{}") except ValueError: eintrag["info"] = {} if eintrag.get("last_seen"): eintrag["last_seen"] = eintrag["last_seen"].isoformat() ergebnis.append(eintrag) return ergebnis def list_mounts() -> list: with engine.connect() as conn: return [dict(z) for z in conn.execute(select(storage_mounts)).mappings().all()] def save_mount(name: str, typ: str, quelle: str, optionen: str, username: str, passwort: str) -> None: with engine.begin() as conn: conn.execute( storage_mounts.insert().values( name=name, typ=typ, quelle=quelle, optionen=optionen, username=username, passwort=passwort, ) ) def delete_mount(name: str) -> None: with engine.begin() as conn: conn.execute(storage_mounts.delete().where(storage_mounts.c.name == name)) def utcnow() -> datetime: return datetime.now(timezone.utc) def init_db() -> None: """Legt fehlende Tabellen an (idempotent) und zieht Mini-Migrationen nach.""" metadata.create_all(engine) # create_all ändert BESTEHENDE Tabellen nicht — neue Spalten hier nachziehen: with engine.begin() as conn: conn.exec_driver_sql( "ALTER TABLE jobs ADD COLUMN IF NOT EXISTS target_dir VARCHAR(255)" ) conn.exec_driver_sql( "ALTER TABLE jobs ADD COLUMN IF NOT EXISTS meta TEXT" ) conn.exec_driver_sql( "ALTER TABLE workers ADD COLUMN IF NOT EXISTS info TEXT" ) def insert_job( job_id: str, device: str, disc_type: str = None, title: str = None, target_dir: str = None, meta: str = None, ) -> None: with engine.begin() as conn: conn.execute( jobs.insert().values( id=job_id, device=device, disc_type=disc_type, title=title, target_dir=target_dir, meta=meta, status="pending", progress=0, created_at=utcnow(), ) ) 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 get_job(job_id: str) -> dict: with engine.connect() as conn: zeile = conn.execute( select(jobs).where(jobs.c.id == job_id) ).mappings().first() return dict(zeile) if zeile else None def list_jobs(limit: int = 100) -> list: with engine.connect() as conn: zeilen = conn.execute( select(jobs).order_by(jobs.c.created_at.desc()).limit(limit) ).mappings().all() return [dict(z) for z in zeilen] def delete_job(job_id: str) -> None: """Entfernt EINEN Job-Eintrag (nur die DB-Zeile — Dateien bleiben).""" with engine.begin() as conn: conn.execute(jobs.delete().where(jobs.c.id == job_id)) def delete_finished_jobs() -> int: """Räumt alle erledigten Jobs (completed/failed) aus der Liste. Dateien bleiben.""" with engine.begin() as conn: ergebnis = conn.execute( jobs.delete().where(jobs.c.status.in_(("completed", "failed"))) ) return ergebnis.rowcount or 0 def delete_worker(name: str) -> None: """Entfernt einen (verwaisten) Worker-Eintrag aus der Liste.""" with engine.begin() as conn: conn.execute(workers.delete().where(workers.c.name == name)) def has_active_job(device: str) -> bool: """True, wenn auf dem Gerät ein Job läuft oder wartet (Eject-Schutz).""" with engine.connect() as conn: zeile = conn.execute( select(jobs.c.id) .where(jobs.c.device == device) .where(jobs.c.status.in_(("pending", "running"))) .limit(1) ).first() return zeile is not None 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) ) def list_logs(limit: int = 200) -> list: with engine.connect() as conn: zeilen = conn.execute( select(logs).order_by(logs.c.id.desc()).limit(limit) ).mappings().all() return [dict(z) for z in zeilen] def get_settings(key: str = "ui") -> dict: with engine.connect() as conn: zeile = conn.execute( select(settings_table.c.value).where(settings_table.c.key == key) ).first() if not zeile or not zeile[0]: return {} try: return json.loads(zeile[0]) except ValueError: return {} def save_settings(werte: dict, key: str = "ui") -> None: payload = json.dumps(werte) with engine.begin() as conn: vorhanden = conn.execute( select(settings_table.c.key).where(settings_table.c.key == key) ).first() if vorhanden: conn.execute( settings_table.update() .where(settings_table.c.key == key) .values(value=payload) ) else: conn.execute(settings_table.insert().values(key=key, value=payload))