feat(windows): Rippy arbeitet eigenstaendig — Rippen, Werkzeuge, Fähigkeiten
Ampel / ampel (push) Failing after 48s
Ampel / ampel (push) Failing after 48s
WAS: Der Ablauf ist aus tasks.py heraus (ablauf.py, ohne Celery), die
LocalQueue wird bedient, ein Laeufer arbeitet Auftraege im selben Prozess
ab, und Rippy meldet sich mit gemessenen Faehigkeiten selbst als Arbeiter.
WARUM: "Rippy fuer Windows soll standalone funktionieren" (Commander). Bis
hierher konnte die Windows-App alles ANZEIGEN und nichts TUN — ein Rip waere
eingereiht worden und fuer immer liegengeblieben, weil niemand ihn holt.
DIE TRENNUNG: rip_disc hing an GENAU DREI Celery-Stellen in 234 Zeilen —
self.update_state, _transcode_queue, transcode_files.apply_async. Alle drei
sind Fragen der ZUSTELLUNG, nicht des Ablaufs. Sie sind jetzt Rueckrufe:
tasks.py reicht die Celery-Fassung herein, standalone.py die lokale. OHNE
Rueckruf komprimiert derselbe Prozess weiter — genau das, was ein
Ein-Prozess-Rippy braucht. Der Ablauf selbst ist Zeile fuer Zeile derselbe;
der Docker-Betrieb merkt vom Umbau nichts (Task-Namen, Argumente, Queues
unveraendert).
EINE ZUSTELL-STELLE statt drei: celery_client.abschicken() bedient alle
Auftragsarten. Vorher rief jede Stelle send_task selbst auf — der
Standalone-Betrieb haette an drei Stellen umgebogen werden muessen, beim
naechsten Auftragstyp an einer vierten.
WEITERER BLOCKER GEFUNDEN: ablauf.py holte detect_disc_type fest aus dem
LINUX-Treiber, in einem try/except. Unter Windows waere es damit IMMER None
gewesen und Rippen "hart verriegelt" — Rippy haette alles angezeigt und
nichts gerippt, ohne dass irgendwo ein Fehler stuende. Jetzt fragt es den
Treiber-Port.
WERKZEUGE: ripping.py und caps.py suchten nur im PATH. Auf dem Commander-PC
gemessen, vorher/nachher:
vorher check_makemkv_installed() -> False (obwohl installiert)
erkenne_encoder() -> nur CPU
nachher MakeMKV 1.18.4 C:\Program Files (x86)\MakeMKV\makemkvcon64.exe
HandBrake 1.11.2 ueber die API geholt, in 2,7 s
Encoder cpu-x264, cpu-x265, cpu-av1, VCE, VCE-AV1
107 Presets, Ryzen 7 9700X, 16 Kerne, avx512f
VCE ist die Hardwarebeschleunigung der Radeon — die hat Rippy auf diesem
Rechner vorher nie gesehen, weil es HandBrake gar nicht fand.
NEUE ROUTEN: GET /system/werkzeuge (was liegt wo, in welcher Fassung, gibt
es Neueres) und POST /system/werkzeuge/{name}/holen. HandBrake kommt
vollautomatisch von GitHub. MakeMKV wird NICHT mitgeliefert — Rippy laedt
die offizielle Datei und startet sie (Black-Box-Trennung, KONZEPT.md § 6).
makemkv.com antwortete beim Bauen mit HTTP 525; das wird im Klartext
gemeldet, und eine selbst geholte Datei bleibt moeglich.
HERZSCHLAG: /capabilities las die workers-Tabelle, die bisher nur der
Celery-Herzschlag fuellte. Im Standalone-Betrieb stand dort "0 Worker" und
die Encoder-Auswahl im UI blieb LEER — auf einem Rechner, der alles kann.
Jetzt meldet sich der Prozess selbst, mit dem, was caps.py MISST.
GEMESSEN, aus der fertigen EXE (29,0 MB):
bereit nach 1 s, keine Fehler im Log
Werkzeuge: beide gefunden, mit Version und Pfad
Worker: 1 (TobisNicerPC), 5 Encoder, 107 Presets
Die ganze Kette ist als Test festgehalten (test_kette.py): zustellen ->
einreihen -> Laeufer -> ablauf -> Job endet in einem EHRLICHEN Zustand.
Ohne Laufwerk geprueft, und das ist der wichtigere Fall: Ein Rip auf ein
totes Geraet muss zuegig scheitern, nicht auf "pending" haengenbleiben.
GEMESSEN: ruff sauber, 512 Tests gruen + 15 uebersprungen (vorher 489).
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5
parent
288f9eeb8d
commit
f4a8d77598
@@ -31,6 +31,11 @@ REDIS_URL = os.getenv("REDIS_URL", "redis://localhost:6379/0")
|
||||
|
||||
_client = None
|
||||
|
||||
# Wer Auftraege zustellt. None = Celery (verteilter Betrieb).
|
||||
# Der Standalone-Betrieb setzt hier seine eigene Zustellung ein — siehe
|
||||
# `zusteller_setzen()`.
|
||||
_zusteller = None
|
||||
|
||||
|
||||
class KeinBroker(RuntimeError):
|
||||
"""Es gibt hier kein Celery — mit Ansage statt mit Importfehler."""
|
||||
@@ -64,6 +69,40 @@ def __getattr__(name):
|
||||
raise AttributeError(f"module {__name__!r} has no attribute {name!r}")
|
||||
|
||||
|
||||
def zusteller_setzen(funktion) -> None:
|
||||
"""Legt fest, WER Auftraege bekommt.
|
||||
|
||||
Ohne Aufruf geht alles an Celery — das ist der Docker-Betrieb, unveraendert.
|
||||
Der Standalone-Betrieb (Windows-App, Headless-Linux) setzt hier seine
|
||||
lokale Auftrags-Queue ein; dann laeuft die Arbeit im selben Prozess.
|
||||
|
||||
Die Unterschrift ist die von `abschicken`:
|
||||
|
||||
funktion(task_name, args, queue=None) -> irgendetwas
|
||||
|
||||
Das ist bewusst dieselbe Form, die Celery hat. So muss keine Aufrufstelle
|
||||
wissen, in welchem Betrieb sie gerade laeuft.
|
||||
"""
|
||||
global _zusteller
|
||||
_zusteller = funktion
|
||||
|
||||
|
||||
def abschicken(task_name: str, args: list, queue: str = None):
|
||||
"""Einen Auftrag zustellen — an Celery oder an die lokale Queue.
|
||||
|
||||
EINE Stelle fuer alle drei Auftragsarten (rip_disc, transcode_files,
|
||||
scan_tracks). Vorher rief jede Aufrufstelle `send_task` selbst auf; damit
|
||||
haette der Standalone-Betrieb an drei Stellen umgebogen werden muessen —
|
||||
und beim naechsten Auftragstyp an einer vierten.
|
||||
"""
|
||||
if _zusteller is not None:
|
||||
return _zusteller(task_name, args, queue)
|
||||
kwargs = {"args": args}
|
||||
if queue:
|
||||
kwargs["queue"] = queue
|
||||
return hole_client().send_task(task_name, **kwargs)
|
||||
|
||||
|
||||
def transcode_queue(node: str = None):
|
||||
"""Ziel-Queue für die Kompression: der GEWÄHLTE Worker (worker_direct)
|
||||
wenn er gerade online ist, sonst die geteilte transcode-Queue.
|
||||
@@ -87,7 +126,5 @@ def transcode_queue(node: str = None):
|
||||
|
||||
|
||||
def start_rip(device_path: str, job_id: str, target_dir: str = None):
|
||||
"""Schickt den Rip-Task an den Worker (Task-Name aus worker/tasks.py)."""
|
||||
return hole_client().send_task(
|
||||
"worker.tasks.rip_disc", args=[device_path, job_id, target_dir]
|
||||
)
|
||||
"""Schickt den Rip-Auftrag los (Task-Name aus worker/tasks.py)."""
|
||||
return abschicken("worker.tasks.rip_disc", [device_path, job_id, target_dir])
|
||||
|
||||
+88
-7
@@ -1214,8 +1214,7 @@ async def scan_tracks_starten(name: str):
|
||||
raise HTTPException(status_code=409, detail="Auf diesem Laufwerk läuft gerade ein Job")
|
||||
|
||||
await asyncio.to_thread(db.save_settings, {"status": "running"}, f"tracks:{device_path}")
|
||||
celery_anbindung.hole_client().send_task(
|
||||
"worker.tasks.scan_tracks", args=[device_path])
|
||||
celery_anbindung.abschicken("worker.tasks.scan_tracks", [device_path])
|
||||
return {"status": "scanning"}
|
||||
|
||||
|
||||
@@ -1273,11 +1272,9 @@ async def retry_transcode(job_id: str):
|
||||
except ValueError:
|
||||
meta = {}
|
||||
from celery_client import transcode_queue
|
||||
celery_anbindung.hole_client().send_task(
|
||||
"worker.tasks.transcode_files",
|
||||
args=[job_id, raw_dir, final_dir],
|
||||
queue=transcode_queue(meta.get("transcode_node")),
|
||||
)
|
||||
celery_anbindung.abschicken(
|
||||
"worker.tasks.transcode_files", [job_id, raw_dir, final_dir],
|
||||
queue=transcode_queue(meta.get("transcode_node")))
|
||||
await asyncio.to_thread(db.update_job, job_id, status="transcoding", progress=0, error=None)
|
||||
await asyncio.to_thread(db.add_log, "info", "api", f"Job {job_id}: Kompression neu eingereiht")
|
||||
return {"id": job_id, "status": "transcoding"}
|
||||
@@ -1833,6 +1830,90 @@ async def browse_mkdir(request: MkdirRequest):
|
||||
return {"path": ziel}
|
||||
|
||||
|
||||
# ═══════════════════════════════════════════════════════════════════════
|
||||
# Werkzeuge: Bestand, Update-Stand, Beschaffung — Etappe V2-4
|
||||
# ═══════════════════════════════════════════════════════════════════════
|
||||
#
|
||||
# WOFUER: Damit Rippy unter Windows eigenstaendig arbeiten kann, muss es
|
||||
# seine Werkzeuge nicht nur FINDEN, sondern auch HOLEN und aktuell halten
|
||||
# koennen. Im Container ist das anders — dort stecken sie im Image, und ein
|
||||
# Update ist ein Rebuild (siehe /system/updates, das bleibt unveraendert).
|
||||
#
|
||||
# WARUM EIGENE ROUTEN statt /system/updates zu erweitern: /system/updates
|
||||
# beantwortet "welche Version gibt es draussen". Diese hier beantworten
|
||||
# "was liegt auf DIESER Maschine, wo, und was kann ich dagegen tun". Zwei
|
||||
# verschiedene Fragen; eine Route, die beides taete, koennte keine davon
|
||||
# ehrlich beantworten.
|
||||
|
||||
|
||||
@app.get("/system/werkzeuge")
|
||||
async def system_werkzeuge():
|
||||
"""Was ist installiert, wo, in welcher Fassung — und gibt es Neueres?
|
||||
|
||||
`neueste` ist leer, wenn die Quelle nicht antwortet. Dann steht dort
|
||||
ausdruecklich NICHT "Update verfuegbar": Eine Nichtauskunft ist keine
|
||||
Aussage ueber die Welt.
|
||||
"""
|
||||
from rippy.tools import beschaffen as werkzeug_beschaffung
|
||||
|
||||
einstellungen = await asyncio.to_thread(db.get_settings)
|
||||
eingestellt = (einstellungen or {}).get("werkzeugPfade") or {}
|
||||
lage = await asyncio.to_thread(werkzeug_beschaffung.lage, eingestellt)
|
||||
return {
|
||||
"werkzeuge": lage,
|
||||
# Wohin Rippy selbst installiert — im UI sichtbar, damit niemand
|
||||
# raten muss, wo die geholten Dateien landen.
|
||||
"ordner": werkzeug_beschaffung.katalog.werkzeug_ordner(),
|
||||
"windows": os.name == "nt",
|
||||
}
|
||||
|
||||
|
||||
@app.post("/system/werkzeuge/{name}/holen")
|
||||
async def system_werkzeug_holen(name: str):
|
||||
"""Holt ein Werkzeug und installiert es.
|
||||
|
||||
HandBrake wird vollstaendig automatisch geholt (GitHub-Release, ein ZIP
|
||||
mit einer .exe — kein Installer, keine Administratorrechte).
|
||||
|
||||
MakeMKV wird NICHT mitgeliefert: Rippy laedt die offizielle Datei vom
|
||||
Hersteller und startet sie. Das ist dieselbe Black-Box-Trennung, die
|
||||
KONZEPT.md § 6 fuer den Container festhaelt — MakeMKV bleibt ein fremdes
|
||||
Programm, Rippy nimmt nur die Handgriffe ab.
|
||||
"""
|
||||
from rippy.tools import beschaffen as werkzeug_beschaffung
|
||||
|
||||
if name not in werkzeug_beschaffung.katalog.WERKZEUGE:
|
||||
raise HTTPException(
|
||||
status_code=404,
|
||||
detail=f"Unbekanntes Werkzeug '{name}'. Bekannt: "
|
||||
+ ", ".join(sorted(werkzeug_beschaffung.katalog.WERKZEUGE)))
|
||||
|
||||
meldungen = []
|
||||
|
||||
def fortschritt(text, anteil=None):
|
||||
if text:
|
||||
meldungen.append(text)
|
||||
|
||||
try:
|
||||
if name == "handbrake":
|
||||
pfad = await asyncio.to_thread(
|
||||
werkzeug_beschaffung.handbrake_holen, None, fortschritt)
|
||||
else:
|
||||
pfad = await asyncio.to_thread(
|
||||
werkzeug_beschaffung.makemkv_holen, "", None, fortschritt, False)
|
||||
except werkzeug_beschaffung.BeschaffungsFehler as e:
|
||||
# Klartext statt Traceback: Der haeufigste Grund ist eine nicht
|
||||
# erreichbare Quelle, und das ist nichts, was der Nutzer im Code
|
||||
# suchen sollte.
|
||||
await asyncio.to_thread(db.add_log, "error", "werkzeuge", str(e))
|
||||
raise HTTPException(status_code=502, detail=str(e)) from e
|
||||
|
||||
await asyncio.to_thread(
|
||||
db.add_log, "info", "werkzeuge",
|
||||
f"{werkzeug_beschaffung.katalog.WERKZEUGE[name]['titel']} beschafft: {pfad}")
|
||||
return {"name": name, "pfad": pfad, "meldungen": meldungen}
|
||||
|
||||
|
||||
@app.get("/system/updates")
|
||||
async def system_updates():
|
||||
"""Update-Check für die Kern-Werkzeuge (Einstellungen → System).
|
||||
|
||||
@@ -145,7 +145,11 @@ def test_erzeugtes_mapping_uebersetzt_den_echten_fehlerfall(monkeypatch):
|
||||
# Der Worker liegt neben der API im Repo; kein geteiltes Paket zwischen
|
||||
# den Containern, deshalb per Pfad laden statt importieren.
|
||||
worker_tasks = os.path.join(
|
||||
os.path.dirname(os.path.dirname(os.path.abspath(__file__))), "worker", "tasks.py"
|
||||
# Seit V2-4 steht der Ablauf in ablauf.py; tasks.py ist nur noch
|
||||
# die Celery-Huelle. Dieser Waechter muss dorthin schauen, wo der
|
||||
# Code WIRKLICH steht — sonst prueft er eine leere Datei und ist
|
||||
# gruen, ohne etwas zu beweisen.
|
||||
os.path.dirname(os.path.dirname(os.path.abspath(__file__))), "worker", "ablauf.py"
|
||||
)
|
||||
spec = importlib.util.spec_from_file_location("_worker_tasks_pfad", worker_tasks)
|
||||
quelltext = open(worker_tasks, encoding="utf-8").read()
|
||||
|
||||
@@ -55,7 +55,11 @@ def test_der_schluessel_heisst_im_worker_genauso():
|
||||
import re
|
||||
|
||||
quelle = (
|
||||
pathlib.Path(__file__).resolve().parents[1] / "worker" / "tasks.py"
|
||||
# Seit V2-4 steht der Ablauf in ablauf.py; tasks.py ist nur noch
|
||||
# die Celery-Huelle. Dieser Waechter muss dorthin schauen, wo der
|
||||
# Code WIRKLICH steht — sonst prueft er eine leere Datei und ist
|
||||
# gruen, ohne etwas zu beweisen.
|
||||
pathlib.Path(__file__).resolve().parents[1] / "worker" / "ablauf.py"
|
||||
).read_text(encoding="utf-8")
|
||||
# Der Worker schreibt die Marke über eine Konstante — deren Wert muss hier
|
||||
# ankommen.
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
+60
-959
File diff suppressed because it is too large
Load Diff
@@ -7,7 +7,7 @@ erreicht das Ziel nicht", obwohl die Freigabe erreichbar war. Es fehlten nur
|
||||
zwei noch nie angelegte Ordner.
|
||||
"""
|
||||
|
||||
import tasks
|
||||
import ablauf as tasks
|
||||
|
||||
|
||||
# --- Der eigentliche Fehler: nur EINE fehlende Ebene war erlaubt -------------
|
||||
|
||||
@@ -105,7 +105,7 @@ def test_matche_episoden_ohne_treffer_gibt_none():
|
||||
|
||||
|
||||
def test_pfad_lokal_uebersetzt_fuer_windows_worker():
|
||||
from tasks import pfad_lokal
|
||||
from ablauf import pfad_lokal
|
||||
|
||||
mapping = "/app/media=Z:\\media;/app/temp=Y:\\temp"
|
||||
assert pfad_lokal("/app/temp/raw/abc", mapping) == "Y:\\temp\\raw\\abc"
|
||||
@@ -128,7 +128,7 @@ def test_arbeitsverzeichnis_wahl_des_rips_schlaegt_die_einstellung():
|
||||
Setting-Wert ist genau der, der bei Vollautomatik-Rips greift, weil dort
|
||||
niemand gefragt wird.
|
||||
"""
|
||||
import tasks
|
||||
import ablauf as tasks
|
||||
|
||||
einst = {"workDir": "/app/media/movies"}
|
||||
assert tasks._arbeitsverzeichnis(einst, "/app/media/rippy") == "/app/media/rippy"
|
||||
@@ -144,7 +144,7 @@ def test_arbeitsverzeichnis_wahl_des_rips_schlaegt_die_einstellung():
|
||||
def test_unter_wurzel_faellt_nicht_auf_praefix_namen_herein():
|
||||
"""Befund 25.07.2026: Elf Stellen prüften mit nacktem startswith().
|
||||
„/app/media-boese/x" beginnt mit „/app/media", liegt aber außerhalb."""
|
||||
import tasks
|
||||
import ablauf as tasks
|
||||
|
||||
assert tasks.unter_wurzel("/app/media", "/app/media") is True
|
||||
assert tasks.unter_wurzel("/app/media/movies", "/app/media") is True
|
||||
@@ -160,7 +160,7 @@ def test_unter_wurzel_faellt_nicht_auf_praefix_namen_herein():
|
||||
|
||||
|
||||
def test_zielbasis_lehnt_praefix_ausbruch_ab():
|
||||
import tasks
|
||||
import ablauf as tasks
|
||||
|
||||
assert tasks._zielbasis("/app/media/movies", "bluray") == "/app/media/movies"
|
||||
# Ausbruch per Praefix-Namen fällt auf den Standard zurück
|
||||
@@ -169,7 +169,7 @@ def test_zielbasis_lehnt_praefix_ausbruch_ab():
|
||||
|
||||
|
||||
def test_arbeitsverzeichnis_lehnt_praefix_ausbruch_ab():
|
||||
import tasks
|
||||
import ablauf as tasks
|
||||
|
||||
assert tasks._arbeitsverzeichnis({}, "/app/media-boese") == tasks.RAW_DIR
|
||||
assert tasks._arbeitsverzeichnis({"workDir": "/app/mediaX"}) == tasks.RAW_DIR
|
||||
|
||||
@@ -13,7 +13,7 @@ import os
|
||||
|
||||
import pytest
|
||||
|
||||
import tasks
|
||||
import ablauf as tasks
|
||||
|
||||
|
||||
class FakeDb:
|
||||
|
||||
@@ -40,7 +40,11 @@ def test_handbrake_cmd_arbeitet_auf_datei_nicht_geraet():
|
||||
MakeMKV-Rip, nie das Laufwerk (die alte Direkt-am-Gerät-Pipeline war
|
||||
für Blu-rays prinzipiell funktionsunfähig)."""
|
||||
cmd = build_handbrake_cmd("/app/temp/raw/x/t00.mkv", "/app/media/bluray/x/t00.mkv")
|
||||
assert cmd[0] == "HandBrakeCLI"
|
||||
# Wie beim MakeMKV-Befehl: seit V2-4 steht hier der GEFUNDENE Pfad. Unter
|
||||
# Windows liegt HandBrakeCLI in Program Files oder in Rippys eigenem
|
||||
# Werkzeug-Ordner und nicht im PATH — mit dem nackten Namen faende
|
||||
# `subprocess` es dort nie.
|
||||
assert "handbrakecli" in cmd[0].lower()
|
||||
assert cmd[cmd.index("--input") + 1] == "/app/temp/raw/x/t00.mkv"
|
||||
assert cmd[cmd.index("--output") + 1] == "/app/media/bluray/x/t00.mkv"
|
||||
assert "--preset" in cmd
|
||||
|
||||
@@ -221,7 +221,11 @@ def test_arbeitsstati_deckt_ab_was_der_worker_wirklich_schreibt():
|
||||
import os
|
||||
import re
|
||||
|
||||
pfad = os.path.join(os.path.dirname(os.path.abspath(__file__)), "tasks.py")
|
||||
# Seit V2-4 steht der Ablauf in ablauf.py; tasks.py ist nur noch
|
||||
# die Celery-Huelle. Dieser Waechter muss dorthin schauen, wo der
|
||||
# Code WIRKLICH steht — sonst prueft er eine leere Datei und ist
|
||||
# gruen, ohne etwas zu beweisen.
|
||||
pfad = os.path.join(os.path.dirname(os.path.abspath(__file__)), "ablauf.py")
|
||||
quelle = open(pfad, encoding="utf-8").read()
|
||||
|
||||
# Nur die Aufrufe, die wirklich die Job-Zeile aendern.
|
||||
|
||||
@@ -87,6 +87,9 @@ def pruefen() -> list:
|
||||
maengel.append("Eine Windows-.exe laesst sich nur unter Windows bauen.")
|
||||
if not os.path.isfile(os.path.join(REPO, "docker", "api", "main.py")):
|
||||
maengel.append("docker/api/main.py fehlt.")
|
||||
if not os.path.isfile(os.path.join(REPO, "docker", "worker", "ablauf.py")):
|
||||
maengel.append("docker/worker/ablauf.py fehlt — ohne den Ablauf kann "
|
||||
"die App nicht rippen.")
|
||||
if not os.path.isdir(os.path.join(REPO, "docker", "ui", "dist")):
|
||||
maengel.append(
|
||||
"docker/ui/dist fehlt — die Oberflaeche ist nicht gebaut. "
|
||||
@@ -127,7 +130,11 @@ def bauen(ausgabe: str, version: str) -> str:
|
||||
"--specpath", arbeit,
|
||||
"--paths", os.path.join(REPO, "src"),
|
||||
"--paths", os.path.join(REPO, "docker", "api"),
|
||||
"--paths", os.path.join(REPO, "docker", "worker"),
|
||||
"--add-data", os.path.join(REPO, "docker", "api") + trenner + "api",
|
||||
# Die Worker-Module — ohne sie kann Rippy anzeigen, aber nicht rippen.
|
||||
# ablauf.py, ripping.py, medien.py, caps.py und Nachbarn.
|
||||
"--add-data", os.path.join(REPO, "docker", "worker") + trenner + "worker",
|
||||
"--add-data", os.path.join(REPO, "docker", "ui", "dist") + trenner + "ui",
|
||||
"--add-data", os.path.join(REPO, "deploy", "worker-windows", "rippy.ico") + trenner + ".",
|
||||
# Was PyInstaller nicht von allein findet: dynamisch importierte Module.
|
||||
|
||||
@@ -262,6 +262,23 @@ def starten(argv=None) -> int:
|
||||
app, ui = anwendung_bauen(werte)
|
||||
print(f" Oberflaeche {ui or '(nicht mitgeliefert)'}")
|
||||
|
||||
# Im Standalone-Betrieb arbeitet DIESER Prozess die Auftraege ab —
|
||||
# es gibt keinen Worker daneben. Ohne diesen Schritt koennte Rippy alles
|
||||
# anzeigen und nichts tun: Ein Rip wuerde eingereiht und laege dann fuer
|
||||
# immer da, weil niemand ihn holt.
|
||||
if werte.get("profil") == "standalone" and werte["queue"]["treiber"] == "lokal":
|
||||
from rippy import standalone
|
||||
from rippy.bus.memory import bus as ereignis_bus
|
||||
|
||||
laeufer = standalone.starten(
|
||||
knoten=werte["queue"].get("knoten") or None,
|
||||
bus=ereignis_bus, store=store,
|
||||
slots=werte["queue"].get("rip_slots", 1),
|
||||
)
|
||||
print(f" Laeufer {len(laeufer)} (Auftraege werden hier abgearbeitet)")
|
||||
else:
|
||||
print(" Laeufer keiner — Auftraege gehen an Celery")
|
||||
|
||||
host = werte["server"]["host"]
|
||||
port = werte["server"]["port"]
|
||||
print(f" Adresse http://{'localhost' if host in ('0.0.0.0', '') else host}:{port}")
|
||||
|
||||
@@ -0,0 +1,152 @@
|
||||
"""Der Läufer: nimmt Aufträge aus der LocalQueue und arbeitet sie ab.
|
||||
|
||||
## Wofür er da ist
|
||||
|
||||
Im verteilten Betrieb macht Celery das: Ein Worker-Prozess wartet auf
|
||||
Nachrichten und ruft die Task-Funktion auf. Im Standalone-Betrieb gibt es
|
||||
keinen Broker und keinen zweiten Prozess — also übernimmt dieser Läufer die
|
||||
Rolle. Er ist das letzte Stück, das gefehlt hat, damit Rippy unter Windows
|
||||
nicht nur anzeigen, sondern auch **arbeiten** kann.
|
||||
|
||||
## Die Lease ist der Kern, nicht ein Detail
|
||||
|
||||
Ein Rip läuft 30 bis 90 Minuten. In dieser Zeit muss zweierlei stimmen:
|
||||
|
||||
* **Niemand sonst darf denselben Auftrag anfassen.** Zwei Rips auf einem
|
||||
Laufwerk wären zwei kaputte Dateien.
|
||||
* **Stürzt der Läufer ab, muss der Auftrag zurückfallen.** Sonst bliebe er
|
||||
für immer als „läuft" stehen — genau der Zustand, den v1 mit
|
||||
`zombies.py` nachträglich einsammeln musste.
|
||||
|
||||
Beides erledigt die Lease aus `rippy.queue.lokal`: Der Läufer verlängert sie
|
||||
alle 15 Sekunden aus einem eigenen Thread. Hört er auf zu leben, läuft sie
|
||||
nach 60 Sekunden ab und der Auftrag ist wieder frei.
|
||||
|
||||
## Warum der Herzschlag NICHT im Arbeits-Thread liegt
|
||||
|
||||
Weil er dann mit der Arbeit stillstünde. Der Rip verbringt seine Zeit in
|
||||
`subprocess.communicate()` — dort läuft kein Python-Code, der nebenbei eine
|
||||
Lease verlängern könnte. Der Herzschlag braucht deshalb einen eigenen Thread,
|
||||
der nichts anderes tut.
|
||||
|
||||
## Was der Läufer NICHT tut
|
||||
|
||||
Er kennt den Rip-Vorgang nicht. Er bekommt eine Funktion `ausfuehren(auftrag)`
|
||||
und ruft sie auf. Das hält ihn testbar (die Tests hier rippen nichts) und
|
||||
erlaubt dieselbe Mechanik später für andere Auftragsarten.
|
||||
"""
|
||||
|
||||
import threading
|
||||
import time
|
||||
import traceback
|
||||
|
||||
from rippy.queue import lokal
|
||||
|
||||
|
||||
class Laeufer:
|
||||
"""Ein Arbeiter: holt einen Auftrag, hält ihn, arbeitet ihn ab."""
|
||||
|
||||
def __init__(self, ausfuehren, knoten: str, kann: set = None,
|
||||
bus=None, store=None,
|
||||
herzschlag_sekunden: float = lokal.LEASE_ERNEUERN_SEKUNDEN):
|
||||
self._ausfuehren = ausfuehren
|
||||
self.knoten = knoten
|
||||
self.kann = kann or set()
|
||||
self._bus = bus
|
||||
self._store = store
|
||||
self._herzschlag_sekunden = herzschlag_sekunden
|
||||
self._laeuft = False
|
||||
self.aktueller_auftrag = None
|
||||
|
||||
# ── Ein Auftrag ─────────────────────────────────────────────────────
|
||||
def einmal(self) -> bool:
|
||||
"""Holt EINEN Auftrag und arbeitet ihn ab. False, wenn nichts da war.
|
||||
|
||||
Blockiert für die Dauer der Arbeit — bei einem Rip also Stunden. Das
|
||||
ist Absicht: Der Aufrufer entscheidet, ob das in einem eigenen Thread
|
||||
passiert.
|
||||
"""
|
||||
auftrag = lokal.uebernehmen(self.knoten, self.kann)
|
||||
if auftrag is None:
|
||||
return False
|
||||
|
||||
self.aktueller_auftrag = auftrag
|
||||
schlagen = self._herzschlag_starten(auftrag["id"])
|
||||
try:
|
||||
ergebnis = self._ausfuehren(auftrag)
|
||||
lokal.abschliessen(auftrag["id"], ergebnis if isinstance(ergebnis, dict) else {})
|
||||
self._melden("job.finished", auftrag, {"status": "fertig"})
|
||||
except Exception as e:
|
||||
# Ein Fehlschlag darf den Laeufer NICHT mitnehmen — sonst bleibt
|
||||
# nach einem kaputten Auftrag jeder weitere liegen. Der Auftrag
|
||||
# geht zurueck in die Queue (bis MAX_VERSUCHE), der Grund wird
|
||||
# protokolliert, und zwar mit Rueckverfolgung: Ein "Fehler" ohne
|
||||
# Stelle ist beim naechsten Mal wertlos.
|
||||
self._protokollieren(auftrag, e)
|
||||
lokal.fehlgeschlagen(auftrag["id"], f"{type(e).__name__}: {e}")
|
||||
self._melden("job.finished", auftrag, {"status": "fehler", "fehler": str(e)})
|
||||
finally:
|
||||
schlagen.set()
|
||||
self.aktueller_auftrag = None
|
||||
return True
|
||||
|
||||
def _herzschlag_starten(self, auftrag_id: str) -> threading.Event:
|
||||
"""Verlängert die Lease, solange gearbeitet wird. Eigener Thread."""
|
||||
fertig = threading.Event()
|
||||
|
||||
def schlagen():
|
||||
while not fertig.wait(self._herzschlag_sekunden):
|
||||
if not lokal.lebenszeichen(auftrag_id, self.knoten):
|
||||
# Der Auftrag gehoert uns nicht mehr — jemand anderes hat
|
||||
# ihn uebernommen, weil unsere Lease abgelaufen war. Dann
|
||||
# ist Weiterarbeiten falsch, aber Abbrechen mitten im Rip
|
||||
# waere schlimmer. Also: laut sagen und weiterlaufen.
|
||||
self._protokollieren_text(
|
||||
f"Auftrag {auftrag_id} ist an einen anderen Knoten "
|
||||
"gefallen — die Lease war abgelaufen.")
|
||||
return
|
||||
|
||||
threading.Thread(target=schlagen, daemon=True,
|
||||
name=f"lease-{auftrag_id[:8]}").start()
|
||||
return fertig
|
||||
|
||||
# ── Die Schleife ────────────────────────────────────────────────────
|
||||
def schleife(self, takt: float = 1.0) -> None:
|
||||
"""Läuft, bis `stoppen()` gerufen wird. Für einen eigenen Thread."""
|
||||
self._laeuft = True
|
||||
while self._laeuft:
|
||||
try:
|
||||
if not self.einmal():
|
||||
time.sleep(takt)
|
||||
except Exception as e: # die Schleife selbst darf nie sterben
|
||||
self._protokollieren_text(f"Läufer-Fehler: {e}")
|
||||
time.sleep(takt * 5)
|
||||
|
||||
def stoppen(self) -> None:
|
||||
self._laeuft = False
|
||||
|
||||
# ── Melden ──────────────────────────────────────────────────────────
|
||||
def _melden(self, typ: str, auftrag: dict, daten: dict) -> None:
|
||||
if not self._bus:
|
||||
return
|
||||
try:
|
||||
self._bus.senden(typ, {**daten, "art": auftrag.get("art")},
|
||||
entitaet="job", entitaet_id=auftrag.get("job_id"))
|
||||
except Exception:
|
||||
pass # eine Meldung darf die Arbeit nicht aufhalten
|
||||
|
||||
def _protokollieren(self, auftrag: dict, fehler: Exception) -> None:
|
||||
self._protokollieren_text(
|
||||
f"Auftrag {auftrag.get('art')} für Job {auftrag.get('job_id')} "
|
||||
f"gescheitert: {type(fehler).__name__}: {fehler}\n"
|
||||
+ "".join(traceback.format_exception_only(type(fehler), fehler)).strip()
|
||||
)
|
||||
|
||||
def _protokollieren_text(self, text: str) -> None:
|
||||
if self._store:
|
||||
try:
|
||||
self._store.add_log("error", "laeufer", text)
|
||||
return
|
||||
except Exception:
|
||||
pass
|
||||
print("[Läufer] " + text)
|
||||
@@ -0,0 +1,185 @@
|
||||
"""Der Läufer: hält er den Auftrag, gibt er ihn zurück, überlebt er Fehler?
|
||||
|
||||
Gerippt wird hier nichts — die Arbeit ist eine eingespritzte Funktion. Geprüft
|
||||
wird das Drumherum, und genau dort sitzen die Fehler, die man im Betrieb
|
||||
teuer bezahlt: ein Läufer, der nach dem ersten kaputten Auftrag stehenbleibt,
|
||||
oder eine Lease, die während eines dreistündigen Encodes abläuft.
|
||||
"""
|
||||
|
||||
import threading
|
||||
import time
|
||||
|
||||
import pytest
|
||||
|
||||
from rippy import store
|
||||
from rippy.queue import laeufer as laeufer_modul
|
||||
from rippy.queue import lokal
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def q(tmp_path):
|
||||
vorher = store.zustand_sichern()
|
||||
store.verbinden(f"sqlite:///{(tmp_path / 'l.db').as_posix()}")
|
||||
store.init_db()
|
||||
yield lokal
|
||||
store.engine_holen().dispose()
|
||||
store.zustand_wiederherstellen(vorher)
|
||||
|
||||
|
||||
class FakeBus:
|
||||
def __init__(self):
|
||||
self.gesendet = []
|
||||
|
||||
def senden(self, typ, daten=None, entitaet=None, entitaet_id=None):
|
||||
self.gesendet.append((typ, entitaet_id, daten or {}))
|
||||
|
||||
|
||||
class FakeStore:
|
||||
def __init__(self):
|
||||
self.protokoll = []
|
||||
|
||||
def add_log(self, level, quelle, text):
|
||||
self.protokoll.append((level, quelle, text))
|
||||
|
||||
|
||||
# ── Grundfall ───────────────────────────────────────────────────────────
|
||||
def test_ohne_auftrag_passiert_nichts(q):
|
||||
lauf = laeufer_modul.Laeufer(lambda a: {}, "knoten-a")
|
||||
assert lauf.einmal() is False
|
||||
|
||||
|
||||
def test_auftrag_wird_geholt_und_abgearbeitet(q):
|
||||
gesehen = []
|
||||
q.einreihen("job1", "rip", payload={"device": "/dev/sr0"})
|
||||
lauf = laeufer_modul.Laeufer(lambda a: gesehen.append(a) or {"ok": True}, "knoten-a")
|
||||
|
||||
assert lauf.einmal() is True
|
||||
assert len(gesehen) == 1
|
||||
assert gesehen[0]["art"] == "rip"
|
||||
assert q.offene_auftraege() == [] # abgeschlossen, nicht mehr offen
|
||||
|
||||
|
||||
def test_zweiter_laeufer_bekommt_ihn_nicht(q):
|
||||
"""Die zentrale Zusage — sonst zwei Rips auf einem Laufwerk."""
|
||||
q.einreihen("job1", "rip")
|
||||
lauf_a = laeufer_modul.Laeufer(lambda a: {}, "knoten-a")
|
||||
lauf_b = laeufer_modul.Laeufer(lambda a: {}, "knoten-b")
|
||||
|
||||
langsam = threading.Event()
|
||||
lauf_a._ausfuehren = lambda a: langsam.wait(2) or {}
|
||||
t = threading.Thread(target=lauf_a.einmal, daemon=True)
|
||||
t.start()
|
||||
time.sleep(0.3)
|
||||
try:
|
||||
assert lauf_b.einmal() is False, "Knoten B hat denselben Auftrag bekommen"
|
||||
finally:
|
||||
langsam.set()
|
||||
t.join(timeout=5)
|
||||
|
||||
|
||||
# ── Fehler ──────────────────────────────────────────────────────────────
|
||||
def test_ein_kaputter_auftrag_nimmt_den_laeufer_nicht_mit(q):
|
||||
"""Ohne das bliebe nach dem ersten Fehlschlag JEDER weitere Auftrag
|
||||
liegen — und im UI saehe es aus, als tue Rippy nichts mehr."""
|
||||
protokoll = FakeStore()
|
||||
q.einreihen("job1", "rip")
|
||||
|
||||
def kaputt(_):
|
||||
raise RuntimeError("MakeMKV ist abgestuerzt")
|
||||
|
||||
lauf = laeufer_modul.Laeufer(kaputt, "knoten-a", store=protokoll)
|
||||
assert lauf.einmal() is True # er hat gearbeitet, nicht geworfen
|
||||
assert any("MakeMKV ist abgestuerzt" in text for _, _, text in protokoll.protokoll)
|
||||
|
||||
|
||||
def test_gescheiterter_auftrag_geht_zurueck_in_die_queue(q):
|
||||
"""Ein abgestuerzter Encoder soll den Job nicht endgueltig verlieren."""
|
||||
auftrag_id = q.einreihen("job1", "transcode")
|
||||
|
||||
def kaputt(_):
|
||||
raise RuntimeError("weg")
|
||||
|
||||
laeufer_modul.Laeufer(kaputt, "knoten-a", store=FakeStore()).einmal()
|
||||
offen = q.offene_auftraege()
|
||||
assert len(offen) == 1 and offen[0]["id"] == auftrag_id
|
||||
|
||||
|
||||
def test_fehler_wird_mit_typ_gemeldet(q):
|
||||
"""„Fehler" allein ist beim naechsten Mal wertlos. Der Typ gehoert dazu."""
|
||||
q.einreihen("job1", "rip")
|
||||
protokoll = FakeStore()
|
||||
|
||||
def kaputt(_):
|
||||
raise FileNotFoundError("makemkvcon fehlt")
|
||||
|
||||
laeufer_modul.Laeufer(kaputt, "knoten-a", store=protokoll).einmal()
|
||||
text = " ".join(t for _, _, t in protokoll.protokoll)
|
||||
assert "FileNotFoundError" in text
|
||||
assert "makemkvcon fehlt" in text
|
||||
|
||||
|
||||
# ── Lease ───────────────────────────────────────────────────────────────
|
||||
def test_lease_wird_waehrend_der_arbeit_verlaengert(q):
|
||||
"""DER Punkt bei einem Rip: Er dauert Stunden, die Lease 60 Sekunden.
|
||||
Ohne Herzschlag fiele der Auftrag mitten im Lauf zurueck — und ein
|
||||
zweiter Laeufer finge an, dieselbe Disc zu lesen.
|
||||
"""
|
||||
from datetime import timedelta
|
||||
|
||||
auftrag_id = q.einreihen("job1", "rip")
|
||||
weiter = threading.Event()
|
||||
|
||||
def langsam(_):
|
||||
# Waehrend der Arbeit die Lease kuenstlich ablaufen lassen —
|
||||
# der Herzschlag muss sie zurueckholen.
|
||||
with store.engine_holen().begin() as conn:
|
||||
conn.execute(
|
||||
store.auftraege.update()
|
||||
.where(store.auftraege.c.id == auftrag_id)
|
||||
.values(lease_until=store.utcnow() - timedelta(seconds=10)))
|
||||
weiter.wait(3)
|
||||
return {}
|
||||
|
||||
lauf = laeufer_modul.Laeufer(langsam, "knoten-a", herzschlag_sekunden=0.2)
|
||||
t = threading.Thread(target=lauf.einmal, daemon=True)
|
||||
t.start()
|
||||
time.sleep(1.0) # dem Herzschlag Zeit geben
|
||||
try:
|
||||
assert laeufer_modul.lokal.uebernehmen("knoten-b") is None, \
|
||||
"Die Lease wurde nicht verlaengert — ein zweiter Knoten kam ran"
|
||||
finally:
|
||||
weiter.set()
|
||||
t.join(timeout=5)
|
||||
|
||||
|
||||
def test_aktueller_auftrag_ist_sichtbar_und_danach_wieder_leer(q):
|
||||
"""Fuer das UI und fuer die Diagnose: Woran arbeitet dieser Knoten gerade?"""
|
||||
q.einreihen("job1", "rip")
|
||||
gesehen = {}
|
||||
lauf = laeufer_modul.Laeufer(
|
||||
lambda a: gesehen.update(id=lauf.aktueller_auftrag["id"]) or {}, "knoten-a")
|
||||
lauf.einmal()
|
||||
assert gesehen["id"]
|
||||
assert lauf.aktueller_auftrag is None
|
||||
|
||||
|
||||
# ── Meldungen ───────────────────────────────────────────────────────────
|
||||
def test_abschluss_wird_gemeldet(q):
|
||||
bus = FakeBus()
|
||||
q.einreihen("job1", "rip")
|
||||
laeufer_modul.Laeufer(lambda a: {}, "knoten-a", bus=bus).einmal()
|
||||
assert [t for t, _, _ in bus.gesendet] == ["job.finished"]
|
||||
assert bus.gesendet[0][1] == "job1"
|
||||
|
||||
|
||||
def test_ein_kaputter_bus_haelt_die_arbeit_nicht_auf(q):
|
||||
"""Eine Meldung ist Beiwerk. Wenn sie scheitert, ist der Rip trotzdem
|
||||
fertig — und darf nicht als Fehlschlag gelten."""
|
||||
class KaputterBus:
|
||||
def senden(self, *a, **k):
|
||||
raise RuntimeError("Bus weg")
|
||||
|
||||
q.einreihen("job1", "rip")
|
||||
lauf = laeufer_modul.Laeufer(lambda a: {"ok": True}, "knoten-a", bus=KaputterBus())
|
||||
assert lauf.einmal() is True
|
||||
assert q.offene_auftraege() == []
|
||||
@@ -0,0 +1,243 @@
|
||||
"""Standalone-Betrieb: Auftraege lokal zustellen und im selben Prozess abarbeiten.
|
||||
|
||||
## Was hier zusammenkommt
|
||||
|
||||
Das ist das letzte Stueck, damit Rippy unter Windows nicht nur ANZEIGEN,
|
||||
sondern auch ARBEITEN kann:
|
||||
|
||||
celery_client.zusteller_setzen(…) die API reiht lokal ein statt an Celery
|
||||
rippy.queue.lokal die Auftragstabelle mit Lease
|
||||
rippy.queue.laeufer holt sie und arbeitet sie ab
|
||||
ablauf.rippen / .komprimieren der Ablauf selbst, ohne Celery
|
||||
|
||||
Der Ablauf ist derselbe wie im Docker-Betrieb — Zeile fuer Zeile dieselbe
|
||||
Datei. Was sich unterscheidet, ist allein die Zustellung.
|
||||
|
||||
## Warum die Task-NAMEN erhalten bleiben
|
||||
|
||||
Die API schickt `"worker.tasks.rip_disc"` los. Dieser Name ist der Vertrag
|
||||
zwischen API und Worker; `test_api_smoke.py` nagelt ihn ausdruecklich fest.
|
||||
Hier wird er deshalb NICHT durch etwas Eigenes ersetzt, sondern uebersetzt —
|
||||
so bleibt beides gueltig, und ein Betrieb laesst sich umstellen, ohne dass
|
||||
die API etwas davon merkt.
|
||||
|
||||
## Warum die Kompression NICHT weitergereicht wird
|
||||
|
||||
`ablauf.rippen()` nimmt einen Rueckruf `weiterreichen`, mit dem die
|
||||
Kompression an einen anderen Knoten geht. Im Standalone-Betrieb gibt es
|
||||
keinen anderen Knoten — hier komprimiert derselbe Prozess weiter. Deshalb
|
||||
wird der Rueckruf bewusst NICHT gestellt: `ablauf` ruft dann selbst
|
||||
`komprimieren()` auf, im selben Auftrag, ohne Umweg ueber die Queue.
|
||||
|
||||
Das ist auch das Ehrlichere: Eine Queue mit genau einem Bearbeiter, die
|
||||
Arbeit an sich selbst weiterreicht, ist nur Zeremonie.
|
||||
"""
|
||||
|
||||
import os
|
||||
import sys
|
||||
import threading
|
||||
|
||||
# Die Auftragsarten, die die API kennt — auf die Formen, die `ablauf` versteht.
|
||||
AUFTRAGSARTEN = {
|
||||
"worker.tasks.rip_disc": "rip",
|
||||
"worker.tasks.transcode_files": "transcode",
|
||||
"worker.tasks.scan_tracks": "scan",
|
||||
}
|
||||
|
||||
|
||||
def _worker_pfad() -> str:
|
||||
"""Wo liegen die Worker-Module (ablauf.py, ripping.py, medien.py)?
|
||||
|
||||
Dieselben drei Faelle wie bei der API (siehe `daemon._api_pfad`): aus dem
|
||||
Repo gestartet, im PyInstaller-Paket, oder installiert.
|
||||
"""
|
||||
gebuendelt = getattr(sys, "_MEIPASS", None)
|
||||
if gebuendelt:
|
||||
return os.path.join(gebuendelt, "worker")
|
||||
|
||||
hier = os.path.dirname(os.path.abspath(__file__))
|
||||
repo = os.path.dirname(os.path.dirname(hier))
|
||||
kandidaten = [
|
||||
os.path.join(repo, "docker", "worker"),
|
||||
os.path.join(os.path.dirname(sys.executable), "worker"),
|
||||
]
|
||||
for pfad in kandidaten:
|
||||
if os.path.isfile(os.path.join(pfad, "ablauf.py")):
|
||||
return pfad
|
||||
raise RuntimeError(
|
||||
"Die Worker-Module wurden nicht gefunden. Gesucht wurde in:\n "
|
||||
+ "\n ".join(kandidaten)
|
||||
+ "\nOhne sie kann Rippy nicht rippen."
|
||||
)
|
||||
|
||||
|
||||
def _ablauf():
|
||||
"""`ablauf` importieren — erst beim ersten Bedarf.
|
||||
|
||||
Nicht auf Modulebene: Das Modul zieht `ripping`, `medien` und `db` nach,
|
||||
und die brauchen einen gesetzten sys.path. Ausserdem soll der Import
|
||||
nicht schon beim Start des Servers passieren, sondern dann, wenn wirklich
|
||||
gearbeitet wird.
|
||||
"""
|
||||
pfad = _worker_pfad()
|
||||
if pfad not in sys.path:
|
||||
sys.path.insert(0, pfad)
|
||||
import ablauf
|
||||
|
||||
return ablauf
|
||||
|
||||
|
||||
# ── Zustellung ──────────────────────────────────────────────────────────
|
||||
def zustellung_bauen(bus=None):
|
||||
"""Gibt die Funktion zurueck, die `celery_client.zusteller_setzen` erwartet."""
|
||||
from rippy.queue import lokal
|
||||
|
||||
def zustellen(task_name: str, args: list, queue: str = None):
|
||||
art = AUFTRAGSARTEN.get(task_name)
|
||||
if art is None:
|
||||
raise ValueError(
|
||||
f"Unbekannte Auftragsart {task_name!r}. Bekannt sind: "
|
||||
+ ", ".join(sorted(AUFTRAGSARTEN))
|
||||
+ ". Wer eine neue einfuehrt, traegt sie in "
|
||||
"rippy/standalone.py AUFTRAGSARTEN ein — sonst landet sie "
|
||||
"im Standalone-Betrieb nirgends, und zwar lautlos."
|
||||
)
|
||||
if art == "rip":
|
||||
device_path, job_id, target_dir = (list(args) + [None, None, None])[:3]
|
||||
auftrag_id = lokal.einreihen(
|
||||
job_id, "rip", faehigkeiten={"art": "rip"},
|
||||
payload={"device_path": device_path, "target_dir": target_dir})
|
||||
elif art == "transcode":
|
||||
job_id, raw_dir, final_dir = (list(args) + [None, None, None])[:3]
|
||||
auftrag_id = lokal.einreihen(
|
||||
job_id, "transcode", faehigkeiten={"art": "transcode"},
|
||||
payload={"raw_dir": raw_dir, "final_dir": final_dir})
|
||||
else: # scan
|
||||
device_path = args[0] if args else ""
|
||||
auftrag_id = lokal.einreihen(
|
||||
"", "scan", faehigkeiten={"art": "scan"},
|
||||
payload={"device_path": device_path})
|
||||
|
||||
if bus:
|
||||
try:
|
||||
bus.senden("system.notice",
|
||||
{"level": "info", "text": f"{art} eingereiht"})
|
||||
except Exception:
|
||||
pass
|
||||
return auftrag_id
|
||||
|
||||
return zustellen
|
||||
|
||||
|
||||
# ── Ausfuehrung ─────────────────────────────────────────────────────────
|
||||
def ausfuehren(auftrag: dict):
|
||||
"""Arbeitet EINEN Auftrag ab — der Laeufer ruft das auf."""
|
||||
import json
|
||||
|
||||
ablauf = _ablauf()
|
||||
try:
|
||||
payload = json.loads(auftrag.get("payload") or "{}")
|
||||
except (ValueError, TypeError):
|
||||
payload = {}
|
||||
|
||||
art = auftrag.get("art")
|
||||
if art == "rip":
|
||||
# KEIN `weiterreichen`: Im Standalone-Betrieb komprimiert derselbe
|
||||
# Prozess weiter (siehe Modul-Docstring).
|
||||
return ablauf.rippen(payload.get("device_path"), auftrag["job_id"],
|
||||
payload.get("target_dir"))
|
||||
if art == "transcode":
|
||||
return ablauf.komprimieren(auftrag["job_id"], payload.get("raw_dir"),
|
||||
payload.get("final_dir"))
|
||||
if art == "scan":
|
||||
return ablauf.scan_tracks(payload.get("device_path"))
|
||||
raise ValueError(f"Unbekannte Auftragsart im Auftrag: {art!r}")
|
||||
|
||||
|
||||
# ── Herzschlag ──────────────────────────────────────────────────────────
|
||||
HERZSCHLAG_SEKUNDEN = 60
|
||||
|
||||
|
||||
def herzschlag_starten(knoten: str, store, takt: float = HERZSCHLAG_SEKUNDEN):
|
||||
"""Meldet DIESEN Prozess als Arbeiter — mit seinen Faehigkeiten.
|
||||
|
||||
## Warum das noetig ist
|
||||
|
||||
`/capabilities` liest die `workers`-Tabelle. Gefuellt hat sie bisher der
|
||||
Celery-Herzschlag in `celery_app.py`. Im Standalone-Betrieb gibt es den
|
||||
nicht — also stand dort „0 Worker", und im UI blieb die Encoder-Auswahl
|
||||
LEER. Rippy haette auf einem Rechner mit AMD-Hardwarebeschleunigung
|
||||
behauptet, es koenne nichts.
|
||||
|
||||
Gemeldet wird, was `caps.py` MISST: die Encoder, die HandBrake auf dieser
|
||||
Maschine wirklich anbietet, die Preset-Namen dieser HandBrake-Fassung,
|
||||
CPU-Modell, Kerne und Vektorbefehle. Nichts davon wird behauptet — das
|
||||
ist die Lehre aus Etappe 19 („Encoder wurden behauptet statt gemessen").
|
||||
"""
|
||||
import threading
|
||||
import time
|
||||
|
||||
def schlagen():
|
||||
while True:
|
||||
try:
|
||||
pfad = _worker_pfad()
|
||||
if pfad not in sys.path:
|
||||
sys.path.insert(0, pfad)
|
||||
import caps
|
||||
|
||||
store.save_worker(knoten, caps.erkenne_encoder(),
|
||||
caps.werkzeug_versionen())
|
||||
except Exception as e:
|
||||
# Nicht sterben — aber auch nicht schweigen.
|
||||
try:
|
||||
store.add_log("warning", "standalone",
|
||||
f"Faehigkeiten-Meldung fehlgeschlagen: {e}")
|
||||
except Exception:
|
||||
pass
|
||||
time.sleep(takt)
|
||||
|
||||
threading.Thread(target=schlagen, daemon=True, name="rippy-herzschlag").start()
|
||||
|
||||
|
||||
# ── Alles zusammen ──────────────────────────────────────────────────────
|
||||
def starten(knoten: str = None, bus=None, store=None, slots: int = 1) -> list:
|
||||
"""Zustellung umbiegen und die Laeufer starten. Gibt die Laeufer zurueck.
|
||||
|
||||
`slots` ist die Zahl gleichzeitiger Auftraege. Vorgabe 1: Ein Rip belegt
|
||||
das Laufwerk, und zwei gleichzeitige Encodes auf einem Rechner sind
|
||||
langsamer als zwei nacheinander.
|
||||
"""
|
||||
import socket
|
||||
|
||||
from rippy.queue.laeufer import Laeufer
|
||||
|
||||
knoten = knoten or socket.gethostname()
|
||||
|
||||
# Die API reiht ab jetzt lokal ein statt an Celery.
|
||||
sys.path.insert(0, os.path.join(
|
||||
os.path.dirname(os.path.dirname(os.path.dirname(
|
||||
os.path.abspath(__file__)))), "docker", "api"))
|
||||
try:
|
||||
import celery_client
|
||||
except ImportError: # im Paket liegt sie woanders
|
||||
from rippy.daemon import _api_pfad
|
||||
|
||||
sys.path.insert(0, _api_pfad())
|
||||
import celery_client
|
||||
|
||||
celery_client.zusteller_setzen(zustellung_bauen(bus))
|
||||
|
||||
# Sich selbst als Arbeiter melden — sonst zeigt das UI "0 Worker" und
|
||||
# eine leere Encoder-Auswahl, obwohl dieser Rechner alles kann.
|
||||
if store is not None:
|
||||
herzschlag_starten(knoten, store)
|
||||
|
||||
laeufer = []
|
||||
for nummer in range(max(1, slots)):
|
||||
lauf = Laeufer(ausfuehren, f"{knoten}#{nummer}",
|
||||
kann={"art=rip", "art=transcode", "art=scan"},
|
||||
bus=bus, store=store)
|
||||
threading.Thread(target=lauf.schleife, daemon=True,
|
||||
name=f"rippy-laeufer-{nummer}").start()
|
||||
laeufer.append(lauf)
|
||||
return laeufer
|
||||
@@ -0,0 +1,137 @@
|
||||
"""Die ganze Kette im Standalone-Betrieb — von der Zustellung bis zum Job-Ende.
|
||||
|
||||
## Was hier bewiesen wird
|
||||
|
||||
celery_client.abschicken(...) die API stellt zu
|
||||
-> standalone.zustellung uebersetzt in einen Auftrag
|
||||
-> lokal.einreihen Auftragstabelle mit Lease
|
||||
-> Laeufer.einmal() holt ihn und arbeitet ihn ab
|
||||
-> ablauf.rippen() der echte Ablauf, ohne Celery
|
||||
-> db.update_job() der Job endet in einem ehrlichen Zustand
|
||||
|
||||
Das ist der Weg, der unter Windows bis V2-4 gar nicht existierte: Dort haette
|
||||
die API einen Rip eingereiht, und er waere fuer immer liegengeblieben, weil
|
||||
niemand ihn holt.
|
||||
|
||||
## Warum ohne Laufwerk getestet wird
|
||||
|
||||
Es gibt hier keins (es haengt an der VM). Das ist kein Mangel, sondern der
|
||||
wichtigere Fall: Ein Rip auf ein Geraet, das nicht antwortet, muss **ehrlich
|
||||
scheitern** — mit einem Zustand und einer Begruendung, die im UI ankommen.
|
||||
Was NICHT passieren darf: dass der Job auf „pending" stehenbleibt und
|
||||
niemand erfaehrt, warum nichts geschieht.
|
||||
"""
|
||||
|
||||
import time
|
||||
|
||||
import pytest
|
||||
|
||||
from rippy import standalone, store
|
||||
from rippy.queue import laeufer as laeufer_modul
|
||||
from rippy.queue import lokal
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def umgebung(tmp_path):
|
||||
vorher = store.zustand_sichern()
|
||||
store.verbinden(f"sqlite:///{(tmp_path / 'kette.db').as_posix()}")
|
||||
store.init_db()
|
||||
yield store
|
||||
store.engine_holen().dispose()
|
||||
store.zustand_wiederherstellen(vorher)
|
||||
|
||||
|
||||
class FakeBus:
|
||||
def __init__(self):
|
||||
self.gesendet = []
|
||||
|
||||
def senden(self, typ, daten=None, entitaet=None, entitaet_id=None):
|
||||
self.gesendet.append((typ, entitaet_id, daten or {}))
|
||||
|
||||
|
||||
def test_die_ganze_kette_bis_zum_ehrlichen_fehlschlag(umgebung):
|
||||
"""Ein Rip auf ein Geraet, das es nicht gibt.
|
||||
|
||||
Erwartet wird KEIN Erfolg — erwartet wird, dass der Job in einem
|
||||
ENDZUSTAND landet und die Begruendung im Log steht. Ein Job, der auf
|
||||
„pending" haengenbleibt, waere der schlimmere Ausgang: Der Nutzer saehe
|
||||
einen Balken, der sich nie bewegt, und nirgends stuende warum.
|
||||
"""
|
||||
umgebung.insert_job("job-kette", "/dev/gibt-es-nicht", disc_type=None)
|
||||
|
||||
bus = FakeBus()
|
||||
zustellen = standalone.zustellung_bauen(bus)
|
||||
zustellen("worker.tasks.rip_disc", ["/dev/gibt-es-nicht", "job-kette", None])
|
||||
|
||||
assert len(lokal.offene_auftraege()) == 1, "Der Auftrag wurde nicht eingereiht"
|
||||
|
||||
lauf = laeufer_modul.Laeufer(standalone.ausfuehren, "test-knoten",
|
||||
kann={"art=rip"}, bus=bus, store=umgebung)
|
||||
assert lauf.einmal() is True, "Der Laeufer hat den Auftrag nicht geholt"
|
||||
|
||||
job = umgebung.get_job("job-kette")
|
||||
assert job["status"] in ("failed", "completed"), (
|
||||
f"Der Job steht auf {job['status']!r} — er muss in einem Endzustand "
|
||||
"landen, sonst wartet der Nutzer auf etwas, das nie passiert.")
|
||||
assert job["status"] == "failed"
|
||||
assert job["error"], "Ein Fehlschlag ohne Begruendung ist im UI wertlos"
|
||||
|
||||
meldungen = " ".join(z["message"] for z in umgebung.list_logs(limit=50))
|
||||
assert "job-kette" in meldungen, "Im Log steht nichts ueber diesen Job"
|
||||
|
||||
assert lokal.offene_auftraege() == [], "Der Auftrag blieb in der Queue liegen"
|
||||
|
||||
|
||||
def test_der_laeufer_bleibt_danach_arbeitsfaehig(umgebung):
|
||||
"""Nach einem Fehlschlag muss der naechste Auftrag trotzdem laufen —
|
||||
sonst legt ein einziger kaputter Rip die ganze Installation lahm."""
|
||||
umgebung.insert_job("job-a", "/dev/gibt-es-nicht")
|
||||
umgebung.insert_job("job-b", "/dev/gibt-es-auch-nicht")
|
||||
zustellen = standalone.zustellung_bauen()
|
||||
zustellen("worker.tasks.rip_disc", ["/dev/gibt-es-nicht", "job-a", None])
|
||||
zustellen("worker.tasks.rip_disc", ["/dev/gibt-es-auch-nicht", "job-b", None])
|
||||
|
||||
lauf = laeufer_modul.Laeufer(standalone.ausfuehren, "test-knoten",
|
||||
kann={"art=rip"}, store=umgebung)
|
||||
assert lauf.einmal() is True
|
||||
assert lauf.einmal() is True
|
||||
assert umgebung.get_job("job-a")["status"] == "failed"
|
||||
assert umgebung.get_job("job-b")["status"] == "failed"
|
||||
|
||||
|
||||
def test_die_api_stellt_ueber_dieselbe_stelle_zu(umgebung, monkeypatch):
|
||||
"""`celery_client.start_rip` ist der Weg, den POST /jobs geht. Er MUSS
|
||||
im Standalone-Betrieb in der lokalen Queue landen — sonst reiht die API
|
||||
an einen Broker ein, den es nicht gibt, und der Job verschwindet."""
|
||||
import sys
|
||||
|
||||
from rippy.daemon import _api_pfad
|
||||
|
||||
if _api_pfad() not in sys.path:
|
||||
sys.path.insert(0, _api_pfad())
|
||||
import celery_client
|
||||
|
||||
umgebung.insert_job("job-api", "/dev/sr0")
|
||||
monkeypatch.setattr(celery_client, "_zusteller", standalone.zustellung_bauen())
|
||||
|
||||
celery_client.start_rip("/dev/sr0", "job-api", "/app/media/x")
|
||||
|
||||
offen = lokal.offene_auftraege()
|
||||
assert len(offen) == 1
|
||||
assert offen[0]["job_id"] == "job-api"
|
||||
assert offen[0]["art"] == "rip"
|
||||
|
||||
|
||||
def test_ein_rip_haengt_nicht_ewig(umgebung):
|
||||
"""Zeitgrenze als Zusage: Ein Auftrag auf ein totes Geraet muss ZUEGIG
|
||||
scheitern. Haengt er, merkt es im Betrieb niemand — der Job steht auf
|
||||
„laeuft", und der Laeufer nimmt keinen weiteren an."""
|
||||
umgebung.insert_job("job-zeit", "/dev/nichts")
|
||||
standalone.zustellung_bauen()("worker.tasks.rip_disc", ["/dev/nichts", "job-zeit", None])
|
||||
|
||||
lauf = laeufer_modul.Laeufer(standalone.ausfuehren, "k", kann={"art=rip"},
|
||||
store=umgebung)
|
||||
start = time.monotonic()
|
||||
lauf.einmal()
|
||||
dauer = time.monotonic() - start
|
||||
assert dauer < 30, f"Der Fehlschlag brauchte {dauer:.1f}s — das ist zu lang"
|
||||
@@ -0,0 +1,100 @@
|
||||
"""Die lokale Zustellung: kommt der Auftrag an, und in der richtigen Form?
|
||||
|
||||
Gerippt wird nichts — geprueft wird die Uebersetzung von Celery-Task-Namen in
|
||||
Auftraege und die Frage, was bei einem unbekannten Namen passiert.
|
||||
"""
|
||||
|
||||
import json
|
||||
|
||||
import pytest
|
||||
|
||||
from rippy import standalone, store
|
||||
from rippy.queue import lokal
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def q(tmp_path):
|
||||
vorher = store.zustand_sichern()
|
||||
store.verbinden(f"sqlite:///{(tmp_path / 's.db').as_posix()}")
|
||||
store.init_db()
|
||||
yield lokal
|
||||
store.engine_holen().dispose()
|
||||
store.zustand_wiederherstellen(vorher)
|
||||
|
||||
|
||||
def test_rip_wird_zum_rip_auftrag(q):
|
||||
zustellen = standalone.zustellung_bauen()
|
||||
zustellen("worker.tasks.rip_disc", ["/dev/sr0", "job1", "/app/media/x"])
|
||||
|
||||
offen = q.offene_auftraege()
|
||||
assert len(offen) == 1
|
||||
auftrag = offen[0]
|
||||
assert auftrag["art"] == "rip"
|
||||
assert auftrag["job_id"] == "job1"
|
||||
nutzlast = json.loads(auftrag["payload"])
|
||||
assert nutzlast["device_path"] == "/dev/sr0"
|
||||
assert nutzlast["target_dir"] == "/app/media/x"
|
||||
|
||||
|
||||
def test_transcode_wird_zum_transcode_auftrag(q):
|
||||
standalone.zustellung_bauen()(
|
||||
"worker.tasks.transcode_files", ["job2", "/raw", "/final"])
|
||||
auftrag = q.offene_auftraege()[0]
|
||||
assert auftrag["art"] == "transcode"
|
||||
nutzlast = json.loads(auftrag["payload"])
|
||||
assert nutzlast["raw_dir"] == "/raw" and nutzlast["final_dir"] == "/final"
|
||||
|
||||
|
||||
def test_scan_braucht_keinen_job(q):
|
||||
"""Ein Track-Scan gehoert zu keinem Job — er passiert VOR dem Anlegen."""
|
||||
standalone.zustellung_bauen()("worker.tasks.scan_tracks", ["/dev/sr0"])
|
||||
auftrag = q.offene_auftraege()[0]
|
||||
assert auftrag["art"] == "scan"
|
||||
assert auftrag["job_id"] == ""
|
||||
|
||||
|
||||
def test_fehlendes_zielverzeichnis_ist_kein_fehler(q):
|
||||
"""POST /jobs darf target_dir weglassen — dann gilt der Standard."""
|
||||
standalone.zustellung_bauen()("worker.tasks.rip_disc", ["/dev/sr0", "job3"])
|
||||
assert json.loads(q.offene_auftraege()[0]["payload"])["target_dir"] is None
|
||||
|
||||
|
||||
def test_unbekannter_taskname_fliegt_auf(q):
|
||||
"""Wuerde er stillschweigend verworfen, waere der Auftrag weg und niemand
|
||||
wuesste warum — der Nutzer sieht nur einen Job, der nie anfaengt."""
|
||||
with pytest.raises(ValueError, match="Unbekannte Auftragsart"):
|
||||
standalone.zustellung_bauen()("worker.tasks.gibts_nicht", ["x"])
|
||||
|
||||
|
||||
def test_faehigkeiten_werden_gesetzt(q):
|
||||
"""Damit ein reiner Encode-Knoten spaeter keine Rip-Auftraege bekommt."""
|
||||
zustellen = standalone.zustellung_bauen()
|
||||
zustellen("worker.tasks.rip_disc", ["/dev/sr0", "job1"])
|
||||
zustellen("worker.tasks.transcode_files", ["job1", "/raw", "/final"])
|
||||
arten = {a["art"]: json.loads(a["faehigkeiten"]) for a in q.offene_auftraege()}
|
||||
assert arten["rip"] == {"art": "rip"}
|
||||
assert arten["transcode"] == {"art": "transcode"}
|
||||
|
||||
|
||||
def test_ein_kaputter_bus_haelt_die_zustellung_nicht_auf(q):
|
||||
class KaputterBus:
|
||||
def senden(self, *a, **k):
|
||||
raise RuntimeError("weg")
|
||||
|
||||
standalone.zustellung_bauen(KaputterBus())(
|
||||
"worker.tasks.rip_disc", ["/dev/sr0", "job1"])
|
||||
assert len(q.offene_auftraege()) == 1
|
||||
|
||||
|
||||
def test_ausfuehren_lehnt_unbekannte_art_ab():
|
||||
with pytest.raises(ValueError, match="Unbekannte Auftragsart"):
|
||||
standalone.ausfuehren({"art": "quatsch", "job_id": "j", "payload": "{}"})
|
||||
|
||||
|
||||
def test_worker_pfad_wird_gefunden():
|
||||
"""Ohne die Worker-Module kann Rippy nicht rippen. Wenn sich die
|
||||
Verzeichnisstruktur aendert, soll das HIER auffallen."""
|
||||
import os
|
||||
|
||||
pfad = standalone._worker_pfad()
|
||||
assert os.path.isfile(os.path.join(pfad, "ablauf.py"))
|
||||
@@ -262,6 +262,13 @@ class Dienst:
|
||||
store.verbinden(config.datenbank_url(werte))
|
||||
app, _ = daemon.anwendung_bauen(werte)
|
||||
|
||||
# Ohne diesen Schritt koennte Rippy alles anzeigen und nichts tun.
|
||||
from rippy import standalone
|
||||
from rippy.bus.memory import bus as ereignis_bus
|
||||
|
||||
standalone.starten(bus=ereignis_bus, store=store,
|
||||
slots=werte["queue"].get("rip_slots", 1))
|
||||
|
||||
import uvicorn
|
||||
|
||||
self._server = uvicorn.Server(uvicorn.Config(
|
||||
|
||||
Reference in New Issue
Block a user