Files
rippy/docker/worker/db.py
T
Hitonabi b584cc29ad
Ampel / ampel (push) Successful in 29s
Universal-Sprint: Wizard, UI-Mounts, Encoder-Erkennung, Task-Split, README
Commander-Ziel: All-in-one, universell, weitergebbar.

- First-Run-Wizard: startet automatisch bei neuer Installation (Keys,
  Verarbeitung, erkannte Hardware); /setup + /setup/complete
- Data-Mounts via UI: Einstellungen -> Speicherziele haengt NFS/SMB direkt
  ein (mounts.py, CAP_SYS_ADMIN + rshared-Propagation, Auto-Remount beim
  Start, CIFS-Creds via Datei statt Kommandozeile); nfs-common/cifs-utils
  im api-Image
- Encoder-Erkennung: jeder Worker meldet beim Start ehrlich seine
  Faehigkeiten (caps.py -> workers-Tabelle), GET /capabilities, Anzeige
  in Wizard + Verarbeitung-Tab
- Task-Split: transcode_files als eigener Task auf Queue "transcode"
  (Basis fuer optionale Remote-GPU-Worker, deploy/remote-transcode-worker.yml
  EXPERIMENTELL) + POST /jobs/{id}/retry-transcode + UI-Knopf
  "Neu komprimieren" bei fehlgeschlagenen Jobs
- API-Keys aus der DB: Settings-UI/Wizard ueberstimmen Env — vorher waren
  die Key-Felder im UI reine Dekoration (Clients lasen nur Env)
- README komplett neu: generischer Schnellstart, Laufwerk-Override via
  docker-compose.override.yml, Architektur, Env-Tabelle

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-23 20:49:08 +02:00

129 lines
3.4 KiB
Python

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