Universal-Sprint: Wizard, UI-Mounts, Encoder-Erkennung, Task-Split, README
Ampel / ampel (push) Successful in 29s
Ampel / ampel (push) Successful in 29s
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>
This commit is contained in:
@@ -85,4 +85,6 @@ COPY docker/worker/ .
|
||||
RUN chmod +x /app/entrypoint.sh
|
||||
|
||||
ENTRYPOINT ["/app/entrypoint.sh"]
|
||||
CMD ["celery", "-A", "celery_app", "worker", "--loglevel=info"]
|
||||
# Der lokale Worker bedient BEIDE Queues (rip + transcode). Ein optionaler
|
||||
# Remote-GPU-Worker startet dasselbe Image nur mit "-Q transcode".
|
||||
CMD ["celery", "-A", "celery_app", "worker", "--loglevel=info", "-Q", "celery,transcode"]
|
||||
|
||||
@@ -0,0 +1,24 @@
|
||||
"""Encoder-Fähigkeiten dieses Workers — ehrlich erkannt, nicht behauptet.
|
||||
|
||||
Jeder Worker meldet beim Start, was auf SEINER Maschine wirklich verfügbar
|
||||
ist (Commander-Anforderung 23.07.: das UI zeigt an, WAS da ist). Ein
|
||||
Remote-GPU-Worker meldet sich hier genauso wie der eingebaute CPU-Worker.
|
||||
"""
|
||||
|
||||
import os
|
||||
import shutil
|
||||
|
||||
|
||||
def erkenne_encoder() -> list:
|
||||
"""Liste der verfügbaren Encoder-Backends auf dieser Maschine."""
|
||||
gefunden = ["cpu-x264", "cpu-x265"] # HandBrake-Software-Encoder, immer dabei
|
||||
|
||||
# VAAPI: AMD (VCN) und Intel (QuickSync) melden sich über /dev/dri
|
||||
if os.path.exists("/dev/dri/renderD128"):
|
||||
gefunden.append("vaapi")
|
||||
|
||||
# NVENC: NVIDIA-Treiber im Container sichtbar
|
||||
if shutil.which("nvidia-smi") or os.path.exists("/usr/lib/x86_64-linux-gnu/libnvidia-encode.so.1"):
|
||||
gefunden.append("nvenc")
|
||||
|
||||
return gefunden
|
||||
@@ -1,5 +1,8 @@
|
||||
from celery import Celery
|
||||
import os
|
||||
import socket
|
||||
|
||||
from celery import Celery
|
||||
from celery.signals import worker_ready
|
||||
|
||||
REDIS_URL = os.getenv("REDIS_URL", "redis://localhost:6379/0")
|
||||
|
||||
@@ -16,4 +19,18 @@ celery_app.conf.update(
|
||||
result_serializer="json",
|
||||
timezone="UTC",
|
||||
enable_utc=True,
|
||||
broker_connection_retry_on_startup=True,
|
||||
)
|
||||
|
||||
|
||||
@worker_ready.connect
|
||||
def melde_faehigkeiten(**kwargs):
|
||||
"""Beim Start: eigene Encoder-Fähigkeiten in die DB melden (UI-Anzeige)."""
|
||||
import caps
|
||||
import db
|
||||
|
||||
try:
|
||||
db.init_db()
|
||||
db.save_worker(socket.gethostname(), caps.erkenne_encoder())
|
||||
except Exception as e: # DB noch nicht da → nicht den Start verhindern
|
||||
print(f"Fähigkeiten-Meldung fehlgeschlagen: {e}")
|
||||
|
||||
@@ -68,6 +68,35 @@ settings_table = Table(
|
||||
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)."""
|
||||
|
||||
+66
-44
@@ -1,8 +1,10 @@
|
||||
"""Zentraler Rip-Task: erkennt den Disc-Typ, rippt, komprimiert, schreibt Status.
|
||||
"""Rip- und Transcode-Tasks — bewusst GETRENNT (23.07.2026, Basis für Etappe 12/20).
|
||||
|
||||
Die API legt beim POST /jobs die Job-Zeile an und schickt diesen Task los —
|
||||
der Worker hält die Zeile aktuell (running → transcoding → completed/failed)
|
||||
und schreibt Ereignisse ins Log. Das UI liest beides über die API.
|
||||
Warum zwei Tasks: Rippen braucht das Laufwerk (läuft immer lokal), Kompression
|
||||
braucht nur CPU/GPU + Zugriff auf die Rohdatei. Als eigener Celery-Task auf der
|
||||
Queue "transcode" kann die Kompression damit auch ein Remote-Worker mit GPU
|
||||
übernehmen (optionales Add-on) — und fehlgeschlagene Kompressionen lassen sich
|
||||
neu anstoßen, ohne die Disc neu zu rippen (POST /jobs/{id}/retry-transcode).
|
||||
|
||||
Zwei Stufen (Commander-Entscheid 23.07.): MakeMKV rippt verlustfrei (einziger
|
||||
Weg durch AACS — HandBrake kann verschlüsselte Discs nicht lesen), HandBrake
|
||||
@@ -10,6 +12,7 @@ komprimiert danach auf Arbeitsgröße. Die Rohdatei liegt nur temporär in
|
||||
/app/temp und wird nach Erfolg gelöscht (Setting keepOriginal behält sie).
|
||||
"""
|
||||
|
||||
import glob
|
||||
import os
|
||||
import shutil
|
||||
|
||||
@@ -25,8 +28,6 @@ from ripping import (
|
||||
)
|
||||
|
||||
RAW_DIR = os.getenv("RAW_DIR", "/app/temp/raw")
|
||||
|
||||
|
||||
MEDIA_ROOT = "/app/media"
|
||||
|
||||
|
||||
@@ -39,9 +40,30 @@ def _zielbasis(target_dir, disc_type: str) -> str:
|
||||
return os.path.join(RIP_OUTPUT_DIR, disc_type)
|
||||
|
||||
|
||||
def _job_abschliessen(job_id: str, ergebnis: dict) -> None:
|
||||
"""Schreibt den Endzustand eines Jobs (completed/failed) nach Postgres."""
|
||||
if ergebnis.get("status") == "success":
|
||||
db.update_job(
|
||||
job_id,
|
||||
status="completed",
|
||||
progress=100,
|
||||
output_path=ergebnis.get("output_dir"),
|
||||
finished_at=db.utcnow(),
|
||||
)
|
||||
db.add_log("success", "worker", f"Job {job_id}: abgeschlossen → {ergebnis.get('output_dir')}")
|
||||
else:
|
||||
db.update_job(
|
||||
job_id,
|
||||
status="failed",
|
||||
error=ergebnis.get("error", "unbekannter Fehler"),
|
||||
finished_at=db.utcnow(),
|
||||
)
|
||||
db.add_log("error", "worker", f"Job {job_id}: {ergebnis.get('error', 'unbekannter Fehler')}")
|
||||
|
||||
|
||||
@celery_app.task(bind=True, name="worker.tasks.rip_disc")
|
||||
def rip_disc(self, device_path: str, job_id: str, target_dir: str = None):
|
||||
"""Rippt eine Disc basierend auf ihrem Typ; job_id ist die DB-Zeile der API.
|
||||
"""Stufe 1: Rippt die Disc; bei Video folgt die Kompression als eigener Task.
|
||||
|
||||
target_dir (optional): vom Nutzer gewähltes Ablageziel unter /app/media —
|
||||
dort eingehängte Shares (NFS/SMB) sind damit direkt wählbar.
|
||||
@@ -54,9 +76,7 @@ def rip_disc(self, device_path: str, job_id: str, target_dir: str = None):
|
||||
"Keine Disc im Laufwerk" if disc_type == "no_disc"
|
||||
else "Disc-Typ nicht erkennbar"
|
||||
)
|
||||
db.update_job(
|
||||
job_id, status="failed", error=fehler, finished_at=db.utcnow()
|
||||
)
|
||||
db.update_job(job_id, status="failed", error=fehler, finished_at=db.utcnow())
|
||||
db.add_log("error", "worker", f"Job {job_id}: {fehler} ({device_path})")
|
||||
return {"status": "error", "error": fehler, "disc_type": disc_type}
|
||||
|
||||
@@ -87,7 +107,7 @@ def rip_disc(self, device_path: str, job_id: str, target_dir: str = None):
|
||||
if disc_type == "cd":
|
||||
ergebnis = rip_cd(device_path, job_id, progress_cb=fortschritt, output_dir=final_dir)
|
||||
elif transcode_an:
|
||||
# Stufe 1: Roh-Rip nach /app/temp (wird nach der Kompression gelöscht)
|
||||
# Roh-Rip nach /app/temp (wird nach erfolgreicher Kompression gelöscht)
|
||||
ergebnis = rip_video(
|
||||
device_path, job_id, disc_type,
|
||||
progress_cb=fortschritt,
|
||||
@@ -100,48 +120,48 @@ def rip_disc(self, device_path: str, job_id: str, target_dir: str = None):
|
||||
)
|
||||
|
||||
if ergebnis.get("status") == "success" and transcode_an:
|
||||
ergebnis = _komprimiere(job_id, final_dir, ergebnis, einstellungen)
|
||||
|
||||
if ergebnis.get("status") == "success":
|
||||
db.update_job(
|
||||
job_id,
|
||||
status="completed",
|
||||
progress=100,
|
||||
output_path=ergebnis.get("output_dir"),
|
||||
finished_at=db.utcnow(),
|
||||
# Kompression als eigener Task auf der transcode-Queue — kann vom
|
||||
# lokalen Worker ODER einem Remote-GPU-Worker übernommen werden.
|
||||
db.update_job(job_id, status="transcoding", progress=0)
|
||||
db.add_log("info", "worker", f"Job {job_id}: Rip fertig, Kompression eingereiht")
|
||||
transcode_files.apply_async(
|
||||
args=[job_id, os.path.join(RAW_DIR, job_id), final_dir],
|
||||
queue="transcode",
|
||||
)
|
||||
db.add_log("success", "worker", f"Job {job_id}: Rip abgeschlossen → {ergebnis.get('output_dir')}")
|
||||
else:
|
||||
db.update_job(
|
||||
job_id,
|
||||
status="failed",
|
||||
error=ergebnis.get("error", "unbekannter Fehler"),
|
||||
finished_at=db.utcnow(),
|
||||
)
|
||||
db.add_log("error", "worker", f"Job {job_id}: {ergebnis.get('error', 'unbekannter Fehler')}")
|
||||
return {"status": "ripped", "raw_dir": os.path.join(RAW_DIR, job_id)}
|
||||
|
||||
_job_abschliessen(job_id, ergebnis)
|
||||
return ergebnis
|
||||
|
||||
|
||||
def _komprimiere(job_id: str, final_dir: str, rip_ergebnis: dict, einstellungen: dict) -> dict:
|
||||
@celery_app.task(bind=True, name="worker.tasks.transcode_files")
|
||||
def transcode_files(self, job_id: str, raw_dir: str, final_dir: str):
|
||||
"""Stufe 2: HandBrake komprimiert die Roh-MKVs auf Arbeitsgröße.
|
||||
|
||||
Erst wenn ALLE Dateien sauber komprimiert sind, wird das Roh-Verzeichnis
|
||||
gelöscht — bricht die Kompression ab, bleibt das Original in /app/temp
|
||||
liegen (kein Datenverlust wie bei ARMs berüchtigtem Move-Bug #1530).
|
||||
Über POST /jobs/{id}/retry-transcode jederzeit neu anstoßbar.
|
||||
"""
|
||||
quellen = rip_ergebnis.get("files", [])
|
||||
os.makedirs(final_dir, exist_ok=True)
|
||||
db.init_db()
|
||||
quellen = sorted(glob.glob(os.path.join(raw_dir, "*.mkv")))
|
||||
if not quellen:
|
||||
ergebnis = {"status": "error", "error": f"Keine Roh-MKVs in {raw_dir} gefunden"}
|
||||
_job_abschliessen(job_id, ergebnis)
|
||||
return ergebnis
|
||||
|
||||
einstellungen = db.get_settings()
|
||||
preset = einstellungen.get("transcodePreset") or DEFAULT_HB_PRESET
|
||||
original_behalten = einstellungen.get("keepOriginal", False)
|
||||
|
||||
db.update_job(job_id, status="transcoding", progress=0)
|
||||
os.makedirs(final_dir, exist_ok=True)
|
||||
db.update_job(job_id, status="transcoding", progress=0, error=None)
|
||||
db.add_log(
|
||||
"info", "worker",
|
||||
f"Job {job_id}: Kompression gestartet ({len(quellen)} Datei(en), Preset '{preset}')",
|
||||
)
|
||||
|
||||
anzahl = max(1, len(quellen))
|
||||
anzahl = len(quellen)
|
||||
for index, quelle in enumerate(quellen):
|
||||
ziel = os.path.join(final_dir, os.path.basename(quelle))
|
||||
|
||||
@@ -151,21 +171,23 @@ def _komprimiere(job_id: str, final_dir: str, rip_ergebnis: dict, einstellungen:
|
||||
|
||||
hb = run_handbrake(quelle, ziel, preset=preset, progress_cb=datei_fortschritt)
|
||||
if hb.get("status") != "success":
|
||||
return {
|
||||
ergebnis = {
|
||||
"status": "error",
|
||||
"error": (
|
||||
f"Kompression fehlgeschlagen bei {os.path.basename(quelle)}: "
|
||||
f"{hb.get('error')} — Roh-Datei bleibt in /app/temp erhalten"
|
||||
),
|
||||
}
|
||||
_job_abschliessen(job_id, ergebnis)
|
||||
return ergebnis
|
||||
|
||||
raw_dir = os.path.dirname(quellen[0]) if quellen else None
|
||||
if raw_dir:
|
||||
if original_behalten:
|
||||
ziel_original = os.path.join(final_dir, "original")
|
||||
shutil.move(raw_dir, ziel_original)
|
||||
db.add_log("info", "worker", f"Job {job_id}: Original behalten unter {ziel_original}")
|
||||
else:
|
||||
shutil.rmtree(raw_dir, ignore_errors=True)
|
||||
if original_behalten:
|
||||
ziel_original = os.path.join(final_dir, "original")
|
||||
shutil.move(raw_dir, ziel_original)
|
||||
db.add_log("info", "worker", f"Job {job_id}: Original behalten unter {ziel_original}")
|
||||
else:
|
||||
shutil.rmtree(raw_dir, ignore_errors=True)
|
||||
|
||||
return {"status": "success", "output_dir": final_dir}
|
||||
ergebnis = {"status": "success", "output_dir": final_dir}
|
||||
_job_abschliessen(job_id, ergebnis)
|
||||
return ergebnis
|
||||
|
||||
Reference in New Issue
Block a user