"""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("target_dir", String(255)), Column("error", Text), Column("meta", Text), # Disc-Metadaten (JSON) — Quelle für Ordnernamen + 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), ) def utcnow() -> datetime: return datetime.now(timezone.utc) def init_db() -> None: """Legt fehlende Tabellen an (idempotent) und zieht Mini-Migrationen nach. Die ALTERs stehen identisch in docker/api/db.py — wer zuerst startet, migriert; der andere findet die Spalten dann bereits vor. """ metadata.create_all(engine) 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" ) 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)), ) def save_worker(name: str, encoder_liste: list, info: dict = None) -> None: """Worker meldet Name + Encoder-Fähigkeiten + Werkzeug-Versionen (Upsert).""" import json payload = json.dumps(encoder_liste) info_payload = json.dumps(info or {}) 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, info=info_payload, last_seen=utcnow() ) ) else: conn.execute( workers.insert().values( name=name, encoders=payload, info=info_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 get_job(job_id: str) -> dict: """Ganze Job-Zeile — der Worker braucht Titel + Metadaten für die Ordner-Benennung und die Media-Server-Aufbereitung (NFO/Poster).""" from sqlalchemy import select 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 save_settings(werte: dict, key: str = "ui") -> None: """Upsert in die settings-Tabelle — der Worker legt hier z. B. die Track-Scan-Ergebnisse ab (key 'tracks:'), die API liest sie.""" import json from sqlalchemy import select 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)) def get_job_status(job_id: str) -> str: """Nur der Status — der Worker prüft damit kooperative Abbruch-Anfragen.""" from sqlalchemy import select with engine.connect() as conn: zeile = conn.execute( select(jobs.c.status).where(jobs.c.id == job_id) ).first() return zeile[0] if zeile else "" def list_jobs_mit_status(stati) -> list: """Alle Jobs in einem der genannten Zustände (id/status/title/created_at). Basis der Zombie-Erkennung: Jobs, die behaupten, es arbeite gerade jemand an ihnen. Bewusst NUR diese schmale Auswahl statt der ganzen Zeile — die Erkennung braucht nichts weiter. """ from sqlalchemy import select with engine.connect() as conn: zeilen = conn.execute( select(jobs.c.id, jobs.c.status, jobs.c.title, jobs.c.created_at) .where(jobs.c.status.in_(list(stati))) ).mappings().all() return [dict(z) for z in zeilen] def zaehle_online_worker(sekunden: int = 120) -> int: """Wie viele Worker gelten laut Herzschlag gerade als online? Die Zombie-Erkennung vergleicht das mit der Zahl der Celery-Antworten: melden sich weniger Worker als bekannt sind, ist die Auskunft unvollständig — dann wird NICHTS als Leiche gewertet. """ from datetime import timedelta from sqlalchemy import func, select grenze = utcnow() - timedelta(seconds=sekunden) with engine.connect() as conn: anzahl = conn.execute( select(func.count()).select_from(workers).where(workers.c.last_seen >= grenze) ).scalar() return int(anzahl or 0) 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) )