refactor(core): V2-1 — die vier Ports, ein Store, ein Laufwerks-Treiber
Ampel / ampel (push) Successful in 40s

WAS: rippy/ports.py beschreibt Store/Queue/Bus/Drives als Protocol. Zwei
weitere Doppelungen sind zusammengelegt: db.py (API+Worker) wird
rippy/store, und der Auswurf (api/devices.py + ripping.wirf_disc_aus)
wird rippy/drives/linux. Verhalten unveraendert.

WARUM: Drei Betriebsarten tragen nur, wenn ein Modus die Auswahl der
Treiber hinter vier Nahtstellen ist statt ein eigener Codestand
(KONZEPT-V2.md §1). Diese Etappe zieht die Nahtstellen ein, ohne schon
einen zweiten Treiber zu haben — die kommen in V2-2 (SQLite/LocalQueue)
und V2-4 (Windows).

DIE UNANGENEHMERE DOPPELUNG WAR DER AUSWURF: Er stand zweimal da, mit
UNTERSCHIEDLICHEN Vertraegen — devices.eject wirft OSError, ripping.
wirf_disc_aus gibt False zurueck und wirft nie. Beides ist richtig fuer
seine Seite (Browser-Meldung gegen "ein Rip stirbt nicht an einer
klemmenden Schublade"). Jetzt liegt EINE Mechanik darunter
(auswerfen_mit_grund) und beide Vertraege unveraendert darueber.
Die API-Fassung war ausserdem NIE getestet — jetzt schon, inklusive
"reicht ENOENT/EPERM unveraendert weiter".

get_settings hatte den einzigen echten Verhaltensunterschied der beiden
db.py: die Worker-Fassung schluckte jeden Fehler und gab {} zurueck.
Nicht still entschieden, sondern sichtbar gemacht — der Parameter
bei_fehler_leer steht jetzt in der Signatur, mit der offenen Frage im
Docstring. {} heisst fuer den Aufrufer "nichts gesetzt", nicht "konnte
nicht nachsehen"; das ist dieselbe Klasse wie catch(() => []) im alten
UI. Zu entscheiden in V2-2.

ZWEITER BEINAHE-FEHLER DIESER ETAPPE: linux.py importierte detection
auf Modulebene — und das zieht fcntl. Damit waere ripping.py und ueber
es der NATIVE WINDOWS-WORKER nicht mehr ladbar gewesen. Diesmal haben
die Tests es sofort gefangen (4 Sammelfehler). Behoben an der Wurzel:
Konstanten und die reine classify() leben jetzt in drives/cdrom.py,
ganz ohne fcntl. Nebengewinn — die classify-Tests liefen bisher NUR in
der Ampel ("erst nach dem Push bewiesen") und laufen jetzt ueberall.

Ausserdem: .dockerignore-Testmuster brauchen **, sonst greifen sie nur
in der obersten Ebene. Im laufenden Container nachgezaehlt: 30 test_*.py
lagen in den Images.

GEMESSEN: ruff sauber, 301 Tests gruen + 1 uebersprungen (vorher 290;
+4 neue eject-Tests, +7 classify-Tests die jetzt lokal laufen). Kein
Modul liegt mehr doppelt im Repo.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
Hitonabi
2026-08-28 08:54:28 +02:00
co-authored by Claude Opus 5
parent edfe0d4313
commit dd1d0b7365
19 changed files with 1670 additions and 1436 deletions
+11 -1
View File
@@ -29,11 +29,21 @@ Thumbs.db
.dockerignore
.docker/
# Test
# Test — ** ist Pflicht, sonst greifen die Muster NUR in der obersten Ebene
# des Bau-Kontexts. Ohne ** landeten am 28.08.2026 dreissig test_*.py in beiden
# Laufzeit-Images (im laufenden Container nachgezaehlt: find /app -name
# "test_*.py" | wc -l -> 30). Schadet nichts, gehoert aber nicht ins Image.
tests/
**/tests/
test_*.py
**/test_*.py
*_test.py
**/*_test.py
.pytest_cache/
**/.pytest_cache/
**/__pycache__/
conftest.py
**/conftest.py
# Build
dist/
-110
View File
@@ -1,110 +0,0 @@
"""Laufwerks-Discovery über /sys und ioctls — ehrlich, ohne udev.
Warum kein udevadm mehr: Im Container läuft kein udevd, die udev-Datenbank ist
leer — `udevadm info` lieferte schlicht nichts (der alte Weg zeigte deshalb nie
ein Laufwerk an). /sys/class/block/<name>/device/{vendor,model} kommt dagegen
direkt vom Kernel und funktioniert überall.
"""
import glob
import os
from fcntl import ioctl
from rippy.drives.detection import (
CDS_DISC_OK,
classify,
disc_size_bytes,
disc_status,
drive_status,
)
# include/uapi/linux/cdrom.h
CDROMEJECT = 0x5309
CDROM_LOCKDOOR = 0x5329 # 1 = Tür verriegeln, 0 = entriegeln
CDROM_DRIVE_STATUS = 0x5326
CDS_NO_DISC = 1
CDS_TRAY_OPEN = 2
AUSWURF_WARTEN_SEKUNDEN = 5
def eject(device_path: str) -> None:
"""Wirft die Disc aus und prüft es nach. Wirft OSError, wenn sie drin bleibt.
⚠️ ERST ENTRIEGELN (Befund 26.07.2026, am laufenden System gemessen): Ein
nacktes CDROMEJECT wird von einem verriegelten Laufwerk mit ERFOLG quittiert
und tut nichts. MakeMKV verriegelt die Tür während des Rips
(`CDROM_LOCKDOOR 1`) und entriegelt sie nicht wieder — danach blieb die
Schublade zu, während Rippy „Disc ausgeworfen" ins Log schrieb. Gegenprobe
an derselben Disc: mit `CDROM_LOCKDOOR 0` davor geht sie auf (Status 2).
Deshalb macht das Werkzeug `eject` immer beides.
Gleichlautend in worker/ripping.wirf_disc_aus — es gibt kein geteiltes Paket
zwischen den Containern; dort steht die ausführliche Herleitung.
"""
import time
fd = os.open(device_path, os.O_RDONLY | os.O_NONBLOCK)
try:
try:
ioctl(fd, CDROM_LOCKDOOR, 0)
except OSError:
pass # nicht verriegelt oder ioctl unbekannt — Auswurf trotzdem versuchen
ioctl(fd, CDROMEJECT, 0)
# Nachsehen statt hoffen. „Kein Datenträger" zählt mit: ein
# Slot-Laufwerk hat keine Schublade.
for _ in range(AUSWURF_WARTEN_SEKUNDEN):
if ioctl(fd, CDROM_DRIVE_STATUS, 0) in (CDS_TRAY_OPEN, CDS_NO_DISC):
return
time.sleep(1)
raise OSError(
"Das Laufwerk hat den Auswurf angenommen, die Disc ist aber noch "
"drin. Blockiert etwas die Schublade, oder läuft noch ein Zugriff?"
)
finally:
os.close(fd)
def list_optical_devices() -> list:
"""Alle optischen Laufwerke, die der Container sieht (devices: in compose)."""
return sorted(glob.glob("/dev/sr[0-9]*"))
def read_sys_attr(device_name: str, attr: str) -> str:
pfad = f"/sys/class/block/{device_name}/device/{attr}"
try:
with open(pfad, encoding="ascii", errors="replace") as f:
return f.read().strip()
except OSError:
return ""
def device_info(device_path: str) -> dict:
"""Baut den Geräte-Eintrag fürs UI: Name aus /sys, Disc-Status per ioctl."""
name = os.path.basename(device_path)
vendor = read_sys_attr(name, "vendor")
model = read_sys_attr(name, "model")
try:
status_code = drive_status(device_path)
except OSError:
status_code = -1
disc_type = "unknown"
status = "empty"
if status_code == CDS_DISC_OK:
status = "ready"
try:
disc_type = classify(disc_status(device_path), disc_size_bytes(device_path))
except OSError:
disc_type = "unknown"
return {
"id": name,
"name": " ".join(teil for teil in (vendor, model) if teil) or f"Laufwerk {name}",
"type": disc_type,
"path": device_path,
"status": status,
"model": model,
"serial": read_sys_attr(name, "wwid"),
}
+2 -2
View File
@@ -11,8 +11,8 @@ import shutil
import time
import uuid
import db
import devices as device_discovery
from rippy import store as db
from rippy.drives import linux as device_discovery
import eta
from rippy.rip import makemkv_daten
import makemkv_key
+3 -3
View File
@@ -50,7 +50,7 @@ def fetch_current_key(timeout: int = 15) -> str | None:
def _apply_key(key: str) -> bool:
"""Merge den Key in die Settings. save_settings überschreibt das GANZE JSON,
darum erst lesen, dann setzen. Rueckgabe True, wenn sich der Key geaendert hat."""
import db
from rippy import store as db
einstellungen = db.get_settings()
if (einstellungen.get("makemkvAppKey") or "").strip() == key:
return False
@@ -61,7 +61,7 @@ def _apply_key(key: str) -> bool:
async def refresh_once() -> bool:
"""Einmal holen + anwenden. Best-effort, loggt in die App-Logs. True = geaendert."""
import db
from rippy import store as db
key = await asyncio.to_thread(fetch_current_key)
if not key:
db.add_log("warning", "makemkv-key",
@@ -77,7 +77,7 @@ async def refresh_once() -> bool:
async def refresh_loop(intervall_stunden: int = 24, start_verzoegerung_s: int = 60) -> None:
"""Periodische Erneuerung (Default: taeglich, damit ein Monatswechsel nie verpasst
wird). Fehler werden geloggt, aber nie geworfen -> die API läuft weiter."""
import db
from rippy import store as db
await asyncio.sleep(start_verzoegerung_s)
while True:
try:
+1 -1
View File
@@ -16,7 +16,7 @@ import re
import subprocess
import tempfile
import db
from rippy import store as db
MEDIA_ROOT = "/app/media"
NAME_MUSTER = re.compile(r"^[a-z0-9][a-z0-9-]{1,30}$")
+2 -2
View File
@@ -364,8 +364,8 @@ def werkzeug_versionen() -> dict:
info["handbrake"] = "installiert"
# Woher kommt der MakeMKV-Key? UI-Setting schlägt Env — ehrlich anzeigen.
try:
import db
ui_key = (db.get_settings().get("makemkvAppKey") or "").strip()
from rippy import store as db
ui_key = (db.get_settings(bei_fehler_leer=True).get("makemkvAppKey") or "").strip()
except Exception:
ui_key = ""
if ui_key:
+1 -1
View File
@@ -32,7 +32,7 @@ celery_app.conf.update(
# Import auf Modulebene: im worker_ready-Signal ist /app nicht mehr
# zuverlässig im sys.path (ModuleNotFoundError 'caps', Deploy 23.07.).
import caps # noqa: E402
import db # noqa: E402
from rippy import store as db # noqa: E402
import zombies # noqa: E402
-261
View File
@@ -1,261 +0,0 @@
"""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:<device>'), 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 meta_merken(job_id: str, **felder) -> None:
"""Ergänzt EINZELNE Schlüssel in den Job-Metadaten. Wirft nie.
`update_job(meta=...)` würde die Spalte ersetzen — Poster, Jahr, Titel-Wahl
und Sprachwunsch dieses Rips wären damit fort. Also lesen, mischen,
schreiben.
Fehler werden geschluckt: Diese Funktion vermerkt nur, in welcher Phase ein
Job steht (api/phasen.py). Ein Rip darf daran nicht scheitern — im
schlimmsten Fall fehlt die Marke und der „Neu"-Knopf fragt nach.
"""
import json
try:
zeile = get_job(job_id)
if not zeile:
return
try:
vorher = json.loads(zeile.get("meta") or "{}")
except (ValueError, TypeError):
vorher = {}
if not isinstance(vorher, dict):
vorher = {}
vorher.update(felder)
update_job(job_id, meta=json.dumps(vorher))
except Exception as e:
try:
add_log("warning", "worker", f"Job {job_id}: Metadaten-Vermerk fehlgeschlagen: {e}")
except Exception:
pass
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)
)
+760 -760
View File
File diff suppressed because it is too large Load Diff
+19 -105
View File
@@ -17,6 +17,25 @@ import shutil
import subprocess
import tempfile
# Auswurf und ioctl-Konstanten leben seit V2-1 im gemeinsamen Treiber
# (rippy/drives/linux.py). Vorher stand derselbe Ablauf ZWEIMAL im Repo —
# hier und in docker/api/devices.py, mit unterschiedlichen Vertraegen.
#
# Die Namen werden hier weiter angeboten, weil tasks.py und die Tests sie so
# kennen; `wirf_disc_aus` ist der Worker-Vertrag (gibt False zurueck, wirft nie).
from rippy.drives.linux import ( # noqa: F401
AUSWURF_WARTEN_SEKUNDEN,
CDROM_DRIVE_STATUS,
CDROM_LOCKDOOR,
CDROMCLOSETRAY,
CDROMEJECT,
CDS_DISC_OK,
CDS_DRIVE_NOT_READY,
CDS_NO_DISC,
CDS_TRAY_OPEN,
_auswurf_geglueckt,
)
from rippy.drives.linux import auswerfen_versuchen as wirf_disc_aus # noqa: F401
from winlauf import OHNE_FENSTER
RIP_OUTPUT_DIR = os.getenv("RIP_OUTPUT_DIR", "/app/media")
@@ -27,111 +46,6 @@ class RipAbbruch(Exception):
Nutzer den Job abgebrochen hat (Status 'canceling' in der DB)."""
# include/uapi/linux/cdrom.h — dieselben ioctls wie in api/devices.py
CDROMEJECT = 0x5309
CDROM_LOCKDOOR = 0x5329 # 1 = Tür verriegeln, 0 = entriegeln
CDROM_DRIVE_STATUS = 0x5326
CDROMCLOSETRAY = 0x5319
# Antworten von CDROM_DRIVE_STATUS (cdrom.h)
CDS_NO_DISC = 1
CDS_TRAY_OPEN = 2
CDS_DRIVE_NOT_READY = 3
CDS_DISC_OK = 4
# Wie lange auf die Schublade gewartet wird. Ein Laufwerk braucht dafür ein
# bis zwei Sekunden; fünf sind reichlich und blockieren nichts Wichtiges.
AUSWURF_WARTEN_SEKUNDEN = 5
def _auswurf_geglueckt(status: int) -> bool:
"""Ist die Disc nach dem Auswurf wirklich draußen? (pure Funktion)
Sowohl „Schublade offen" als auch „kein Datenträger" zählen: Ein
Slot-Laufwerk hat keine Schublade und meldet nach dem Auswerfen CDS_NO_DISC.
"""
return status in (CDS_TRAY_OPEN, CDS_NO_DISC)
def wirf_disc_aus(device_path: str, ioctl_fn=None, oeffnen=None,
schliessen=None, warten=None) -> bool:
"""Wirft die Disc aus und PRÜFT, ob sie draußen ist. Wirft NIE.
## Warum `CDROMEJECT` allein nicht genügt (Befund 26.07.2026, gemessen)
Der Commander meldete: „Der Button gibt es in den Settings, aber es passiert
nicht, das Laufwerk geht nicht auf." Am laufenden System nachgestellt:
wirf_disc_aus("/dev/sr0") → True
CDROM_DRIVE_STATUS danach → 4 (Disc drin)
Das ioctl wird also **angenommen und tut nichts**. Ursache: MakeMKV
verriegelt während des Rips die Laufwerkstür (`CDROM_LOCKDOOR 1`) und
entriegelt sie nicht wieder. Ein verriegeltes Laufwerk quittiert den Auswurf
trotzdem mit Erfolg. Deshalb macht das Werkzeug `eject` immer beides:
erst entriegeln, dann auswerfen. Gegenprobe an derselben Disc:
CDROM_LOCKDOOR 0 + CDROMEJECT → Status 2 (SCHUBLADE OFFEN)
## Und deshalb wird das Ergebnis geprüft, nicht geglaubt
Genau diese Sorte Fehler ist zweimal durchgerutscht: In v3.14 stand hier
„Auswurf tat nichts, jetzt entscheidet die Einstellung" — die Einstellung
wurde danach wirklich gelesen, nur ausgeworfen wurde weiterhin nicht, und im
Log stand „Disc ausgeworfen". Ein Rückgabewert eines ioctls beweist nichts;
gefragt wird jetzt das Laufwerk.
Bewusst hier und nicht in detection.py: das Modul ist ein byteweiser
Zwilling der API-Kopie. `fcntl` gibt es nur unter Linux — der native
Windows-Worker lädt ripping.py ebenfalls, rippt dort aber nie.
Die drei Parameter sind nur zum Testen einspritzbar (kein echtes Laufwerk).
"""
if ioctl_fn is None:
try:
from fcntl import ioctl
except ImportError: # Windows — dieser Worker rippt nie
return False
ioctl_fn = ioctl
if warten is None:
import time
warten = time.sleep
oeffnen = oeffnen or (lambda p: os.open(p, os.O_RDONLY | os.O_NONBLOCK))
schliessen = schliessen or os.close
try:
fd = oeffnen(device_path)
except OSError:
return False
try:
# Entriegeln ist der entscheidende Schritt. Scheitert er, wird der
# Auswurf trotzdem versucht — bei einem nicht verriegelten Laufwerk
# (oder einem, das das ioctl nicht kennt) klappt er ohnehin.
try:
ioctl_fn(fd, CDROM_LOCKDOOR, 0)
except OSError:
pass
try:
ioctl_fn(fd, CDROMEJECT, 0)
except OSError:
return False
# Nachsehen statt hoffen: Die Schublade braucht ein bis zwei Sekunden.
for _ in range(AUSWURF_WARTEN_SEKUNDEN):
try:
if _auswurf_geglueckt(ioctl_fn(fd, CDROM_DRIVE_STATUS, 0)):
return True
except OSError:
return False
warten(1)
return False
finally:
try:
schliessen(fd)
except OSError:
pass
def check_makemkv_installed() -> bool:
"""Prüft, ob makemkvcon installiert ist."""
return shutil.which("makemkvcon") is not None
+5 -5
View File
@@ -22,7 +22,7 @@ import time
import requests
import db
from rippy import store as db
from rippy.rip import makemkv_daten
import medien
from rippy.core import notify
@@ -372,7 +372,7 @@ def _serien_episoden_zuordnen(ausgabe: str, serie: str, staffel, meta: dict) ->
def _benachrichtigen(job_id: str, betreff: str, text: str, level: str) -> None:
"""Webhook-Meldung bei Job-Ende — best effort, nie job-entscheidend."""
einstellungen = db.get_settings()
einstellungen = db.get_settings(bei_fehler_leer=True)
url = (einstellungen.get("notificationWebhook") or "").strip()
if not url:
return
@@ -402,7 +402,7 @@ def _job_abschliessen(job_id: str, ergebnis: dict) -> None:
if ergebnis.get("status") == "success":
ausgabe = ergebnis.get("output_dir")
if ausgabe and (job.get("disc_type") in ("dvd", "bluray", "uhd")):
einstellungen = db.get_settings()
einstellungen = db.get_settings(bei_fehler_leer=True)
try:
meta = json.loads(job.get("meta") or "{}")
except ValueError:
@@ -527,7 +527,7 @@ def rip_disc(self, device_path: str, job_id: str, target_dir: str = None):
gesehen.add(text)
db.add_log("info", "makemkv", f"Job {job_id}: {text[:300]}")
einstellungen = db.get_settings()
einstellungen = db.get_settings(bei_fehler_leer=True)
ist_video = disc_type in ("dvd", "bluray", "uhd")
# Je Disc-Typ abwählbar (siehe komprimieren_fuer): 4K verlustfrei behalten,
# DVDs trotzdem schrumpfen — vorher gab es nur alles oder nichts.
@@ -898,7 +898,7 @@ def transcode_files(self, job_id: str, raw_dir: str, final_dir: str):
_job_abschliessen(job_id, ergebnis)
return ergebnis
einstellungen = db.get_settings()
einstellungen = db.get_settings(bei_fehler_leer=True)
# Preset nach Disc-Typ (Befund 25.07.2026): vorher lief JEDE Quelle durch
# dasselbe Preset — eine 4K-UHD wurde damit auf 1080p heruntergerechnet.
job = db.get_job(job_id) or {}
-84
View File
@@ -383,90 +383,6 @@ def test_fehler_ohne_preset_problem_bleibt_der_alte():
assert ergebnis["error"] == "HandBrake endete mit Code 1"
# --- Auswurf: das ioctl meldet Erfolg und tut nichts (Befund 26.07.2026) -----
def _laufwerk(verriegelt=True, kennt_lockdoor=True):
"""Ein nachgebautes Laufwerk, das sich wie das echte verhaelt.
Gemessen am BU40N der Rippy-VM: Nach einem MakeMKV-Rip ist die Tuer
verriegelt. CDROMEJECT wird dann ANGENOMMEN und tut nichts - der Status
bleibt auf 4 (Disc drin). Erst CDROM_LOCKDOOR 0 macht den Auswurf wirksam.
"""
import ripping
zustand = {"verriegelt": verriegelt, "status": ripping.CDS_DISC_OK,
"aufrufe": []}
def ioctl_fn(fd, befehl, arg=0):
zustand["aufrufe"].append(befehl)
if befehl == ripping.CDROM_LOCKDOOR:
if not kennt_lockdoor:
raise OSError("ioctl unbekannt")
zustand["verriegelt"] = bool(arg)
return 0
if befehl == ripping.CDROMEJECT:
if not zustand["verriegelt"]:
zustand["status"] = ripping.CDS_TRAY_OPEN
return 0 # <- auch verriegelt: ERFOLG, aber ohne Wirkung
if befehl == ripping.CDROM_DRIVE_STATUS:
return zustand["status"]
raise OSError("unerwartetes ioctl")
return zustand, ioctl_fn
def test_auswurf_entriegelt_zuerst_und_klappt_dann():
import ripping
zustand, ioctl_fn = _laufwerk(verriegelt=True)
ok = ripping.wirf_disc_aus(
"/dev/sr0", ioctl_fn=ioctl_fn, oeffnen=lambda p: 42,
schliessen=lambda fd: None, warten=lambda s: None)
assert ok is True
assert zustand["status"] == ripping.CDS_TRAY_OPEN
# Reihenfolge: entriegeln VOR auswerfen
assert zustand["aufrufe"][0] == ripping.CDROM_LOCKDOOR
assert zustand["aufrufe"][1] == ripping.CDROMEJECT
def test_auswurf_meldet_fehlschlag_wenn_die_disc_drin_bleibt():
"""Der eigentliche Fehler. Vorher gab wirf_disc_aus True zurueck, weil das
ioctl nicht geworfen hatte - und ins Log kam "Disc ausgeworfen", waehrend
die Schublade zu blieb. Ein ioctl-Rueckgabewert beweist nichts."""
import ripping
# Ein Laufwerk, das LOCKDOOR nicht kennt und verriegelt bleibt
zustand, ioctl_fn = _laufwerk(verriegelt=True, kennt_lockdoor=False)
ok = ripping.wirf_disc_aus(
"/dev/sr0", ioctl_fn=ioctl_fn, oeffnen=lambda p: 42,
schliessen=lambda fd: None, warten=lambda s: None)
assert ok is False
assert zustand["status"] == ripping.CDS_DISC_OK # nie aufgegangen
def test_auswurf_bei_slot_laufwerk_ohne_schublade():
"""Ein Slot-Laufwerk hat keine Schublade und meldet nach dem Auswerfen
CDS_NO_DISC. Das muss als Erfolg zaehlen."""
import ripping
assert ripping._auswurf_geglueckt(ripping.CDS_NO_DISC) is True
assert ripping._auswurf_geglueckt(ripping.CDS_TRAY_OPEN) is True
assert ripping._auswurf_geglueckt(ripping.CDS_DISC_OK) is False
assert ripping._auswurf_geglueckt(ripping.CDS_DRIVE_NOT_READY) is False
def test_auswurf_ohne_laufwerk_wirft_nicht():
import ripping
def oeffnen_kaputt(p):
raise OSError("kein Laufwerk")
assert ripping.wirf_disc_aus(
"/dev/sr9", ioctl_fn=lambda *a: 0, oeffnen=oeffnen_kaputt,
schliessen=lambda fd: None, warten=lambda s: None) is False
# --- Sprachen der Disc: gemessen an der Akira-Blu-ray (26.07.2026) -----------
#
# Woertlich aus `makemkvcon -r --noscan info dev:/dev/sr0` im Worker-Container
+68
View File
@@ -0,0 +1,68 @@
"""Konstanten und reine Logik rund um optische Laufwerke — OHNE `fcntl`.
## Warum es diese Datei gibt (Etappe V2-1, 28.08.2026)
Zwei Gründe, beide praktisch:
**1. `fcntl` gibt es nur unter Linux.** Solange die ioctl-Nummern in
`detection.py` standen, zog jeder Import dieser Nummern auch `fcntl` nach. Als
in V2-1 der Auswurf-Treiber (`linux.py`) sie brauchte, war `ripping.py` auf
einmal unter Windows nicht mehr ladbar — und der native Windows-Worker lädt
genau dieses Modul. Die Tests haben es sofort gefangen; ohne die Trennung hier
wäre es beim nächsten Windows-Start aufgefallen.
**2. Die Zuordnung ist reine Logik und gehört jedem.** `classify()` ist eine
Funktion von zwei Zahlen auf einen String. Sie hing nur an Linux, weil sie in
derselben Datei stand wie die ioctls — ihre Tests liefen deshalb ausschließlich
in der Ampel und nie auf dem Entwicklungsrechner („Tests, die dort landen, sind
erst nach dem Push bewiesen"). Jetzt laufen sie überall. Der Windows-Treiber
aus V2-4 benutzt dieselbe Funktion: Auch dort wird nach Disc-Status und Größe
eingeordnet, nur die Beschaffung der beiden Zahlen unterscheidet sich.
Die Werte stammen aus der Kernel-UAPI (`include/uapi/linux/cdrom.h`, für
BLKGETSIZE64 aus `include/uapi/linux/fs.h`) und sind seit Jahrzehnten stabil.
"""
# ── include/uapi/linux/cdrom.h ──────────────────────────────────────────
CDROMEJECT = 0x5309
CDROMCLOSETRAY = 0x5319
CDROM_DRIVE_STATUS = 0x5326
CDROM_DISC_STATUS = 0x5327
CDROM_LOCKDOOR = 0x5329 # 1 = Tür verriegeln, 0 = entriegeln
# Antworten von CDROM_DRIVE_STATUS
CDS_NO_DISC = 1
CDS_TRAY_OPEN = 2
CDS_DRIVE_NOT_READY = 3
CDS_DISC_OK = 4
# Antworten von CDROM_DISC_STATUS
CDS_AUDIO = 100
CDS_DATA_1 = 101
CDS_DATA_2 = 102
CDS_XA_2_1 = 103
CDS_XA_2_2 = 104
CDS_MIXED = 105
# include/uapi/linux/fs.h: BLKGETSIZE64 = _IOR(0x12, 114, size_t) auf 64-bit
BLKGETSIZE64 = 0x80081272
# ── Schwellen für die Typ-Zuordnung ─────────────────────────────────────
# Eine DVD9 fasst ~8,5 GB; Blu-ray beginnt bei 25 GB (Single Layer).
# Alles ab 10 GB ist also sicher eine Blu-ray.
BLURAY_MIN_BYTES = 10 * 1024**3
# 4K-UHD-Discs sind BD-66 (66 GB) oder BD-100 — eine normale BD-50 bleibt
# unter ~47 GiB. Ab 55 GiB ist es also sicher eine UHD. (Seltene 50-GB-UHDs
# laufen als "bluray" — der Rip-Weg ist ohnehin identisch.)
UHD_MIN_BYTES = 55 * 1024**3
def classify(disc_status_code: int, size_bytes: int) -> str:
"""Pure Zuordnung: Disc-Status + Größe → cd | dvd | bluray | uhd | unknown."""
if disc_status_code in (CDS_AUDIO, CDS_MIXED):
return "cd"
if disc_status_code in (CDS_DATA_1, CDS_DATA_2, CDS_XA_2_1, CDS_XA_2_2):
if size_bytes >= UHD_MIN_BYTES:
return "uhd"
return "bluray" if size_bytes >= BLURAY_MIN_BYTES else "dvd"
return "unknown"
+21 -37
View File
@@ -14,32 +14,27 @@ import os
import struct
from fcntl import ioctl
# include/uapi/linux/cdrom.h
CDROM_DRIVE_STATUS = 0x5326
CDROM_DISC_STATUS = 0x5327
CDS_NO_DISC = 1
CDS_TRAY_OPEN = 2
CDS_DRIVE_NOT_READY = 3
CDS_DISC_OK = 4
CDS_AUDIO = 100
CDS_DATA_1 = 101
CDS_DATA_2 = 102
CDS_XA_2_1 = 103
CDS_XA_2_2 = 104
CDS_MIXED = 105
# include/uapi/linux/fs.h: BLKGETSIZE64 = _IOR(0x12, 114, size_t) auf 64-bit
BLKGETSIZE64 = 0x80081272
# Eine DVD9 fasst ~8,5 GB; Blu-ray beginnt bei 25 GB (Single Layer).
# Alles ab 10 GB ist also sicher eine Blu-ray.
BLURAY_MIN_BYTES = 10 * 1024**3
# 4K-UHD-Discs sind BD-66 (66 GB) oder BD-100 — eine normale BD-50 bleibt
# unter ~47 GiB. Ab 55 GiB ist es also sicher eine UHD. (Seltene 50-GB-UHDs
# laufen als "bluray" — der Rip-Weg ist ohnehin identisch.)
UHD_MIN_BYTES = 55 * 1024**3
# Konstanten und die reine Zuordnung leben in cdrom.py — ohne fcntl, damit
# der Windows-Treiber (V2-4) und der native Windows-Worker sie benutzen
# koennen. Hier werden sie weiter angeboten, weil Aufrufstellen sie so kennen.
from rippy.drives.cdrom import ( # noqa: F401
BLKGETSIZE64,
BLURAY_MIN_BYTES,
CDROM_DISC_STATUS,
CDROM_DRIVE_STATUS,
CDS_AUDIO,
CDS_DATA_1,
CDS_DATA_2,
CDS_DISC_OK,
CDS_DRIVE_NOT_READY,
CDS_MIXED,
CDS_NO_DISC,
CDS_TRAY_OPEN,
CDS_XA_2_1,
CDS_XA_2_2,
UHD_MIN_BYTES,
classify,
)
def _open_nonblock(device_path: str) -> int:
@@ -76,17 +71,6 @@ def disc_size_bytes(device_path: str) -> int:
os.close(fd)
def classify(disc_status_code: int, size_bytes: int) -> str:
"""Pure Zuordnung (testbar): Disc-Status + Größe → cd | dvd | bluray | uhd | unknown."""
if disc_status_code in (CDS_AUDIO, CDS_MIXED):
return "cd"
if disc_status_code in (CDS_DATA_1, CDS_DATA_2, CDS_XA_2_1, CDS_XA_2_2):
if size_bytes >= UHD_MIN_BYTES:
return "uhd"
return "bluray" if size_bytes >= BLURAY_MIN_BYTES else "dvd"
return "unknown"
def detect_disc_type(device_path: str) -> str:
"""Erkennt den Typ der eingelegten Disc; 'no_disc' wenn keine drin ist."""
try:
+245
View File
@@ -0,0 +1,245 @@
"""Laufwerks-Treiber für Linux: finden, Zustand lesen, verriegeln, auswerfen.
Erfüllt `rippy.ports.Drives`. Der Windows-Treiber (Win32 statt ioctl) kommt in
Etappe V2-4 und implementiert denselben Port.
## Diese Datei war zwei — mit ZWEI VERSCHIEDENEN VERTRÄGEN (V2-1, 28.08.2026)
Der Auswurf existierte doppelt, und das war die unangenehmere Sorte Doppelung,
weil die beiden Fassungen sich nicht nur wiederholten, sondern unterschiedlich
ANTWORTETEN:
docker/api/devices.py eject(pfad) wirft OSError
docker/worker/ripping.py wirf_disc_aus(pfad) gibt False zurück, wirft nie
Beide Verhalten sind richtig — für ihre Seite. Die API will einen Fehler, den
sie dem Browser zeigen kann. Der Worker will einen Rip nicht daran scheitern
lassen, dass die Schublade klemmt, und er will die ioctls in Tests einspritzen
können, ohne ein echtes Laufwerk zu haben.
Deshalb liegt hier jetzt EINE Mechanik (`auswerfen_mit_grund`) und darüber
beide Verträge unverändert. Was verschwindet, ist die dritte Kopie der
ioctl-Nummern und der Ablauf-Reihenfolge.
## Warum überhaupt entriegelt wird (Befund 26.07.2026, am System gemessen)
Der Commander meldete: „Der Button gibt es in den Settings, aber es passiert
nicht, das Laufwerk geht nicht auf." Nachgestellt:
CDROMEJECT allein -> Erfolg gemeldet
CDROM_DRIVE_STATUS danach -> 4 (Disc ist noch drin)
CDROM_LOCKDOOR 0 + CDROMEJECT -> Status 2 (Schublade offen)
MakeMKV verriegelt die Tür während des Rips und entriegelt sie nicht wieder.
Ein verriegeltes Laufwerk quittiert den Auswurf **mit Erfolg** und tut nichts.
Deshalb gilt für jeden Treiber dieses Ports: erst entriegeln, dann auswerfen,
dann NACHSEHEN. Ein Rückgabewert ist kein Beweis, wo die Wirkung prüfbar ist.
"""
import glob
import os
# NUR aus cdrom.py — das Modul kommt ohne `fcntl` aus. Wuerde hier
# detection importiert, waere linux.py (und ueber ripping.py auch der native
# Windows-Worker) unter Windows nicht mehr ladbar. Genau das ist beim Bau
# dieser Etappe passiert; die Tests haben es gefangen.
from rippy.drives.cdrom import ( # noqa: F401
CDROM_DRIVE_STATUS,
CDROM_LOCKDOOR,
CDROMCLOSETRAY,
CDROMEJECT,
CDS_DISC_OK,
CDS_DRIVE_NOT_READY,
CDS_NO_DISC,
CDS_TRAY_OPEN,
)
# Wie lange auf die Schublade gewartet wird. Ein Laufwerk braucht dafür ein
# bis zwei Sekunden; fünf sind reichlich und blockieren nichts Wichtiges.
AUSWURF_WARTEN_SEKUNDEN = 5
# Ergebnisse von auswerfen_mit_grund
AUSWURF_OK = "ok"
AUSWURF_KEIN_ZUGRIFF = "kein-zugriff" # Gerät ließ sich nicht öffnen
AUSWURF_ABGELEHNT = "abgelehnt" # ioctl selbst scheiterte
AUSWURF_BLEIBT_DRIN = "bleibt-drin" # angenommen, aber Disc ist noch da
def _auswurf_geglueckt(status: int) -> bool:
"""Ist die Disc nach dem Auswurf wirklich draußen? (pure Funktion)
Sowohl „Schublade offen" als auch „kein Datenträger" zählen: Ein
Slot-Laufwerk hat keine Schublade und meldet nach dem Auswerfen CDS_NO_DISC.
"""
return status in (CDS_TRAY_OPEN, CDS_NO_DISC)
def auswerfen_mit_grund(device_path: str, ioctl_fn=None, oeffnen=None,
schliessen=None, warten=None):
"""Die gemeinsame Mechanik. Wirft NIE.
Rückgabe: `(grund, fehler)` — `grund` ist eine der AUSWURF_*-Konstanten,
`fehler` der ursprüngliche OSError (oder None). Der Fehler wird
weitergereicht statt verschluckt, damit die API ihn unverändert
weiterwerfen kann; vorher stand dort ein eigener `os.open`-Aufruf nur
deshalb, weil diese Unterscheidung fehlte.
Die vier Parameter sind nur zum Testen einspritzbar (kein echtes Laufwerk).
"""
if ioctl_fn is None:
try:
from fcntl import ioctl
except ImportError: # Windows — dieser Worker rippt nie
return AUSWURF_KEIN_ZUGRIFF, None
ioctl_fn = ioctl
if warten is None:
import time
warten = time.sleep
oeffnen = oeffnen or (lambda p: os.open(p, os.O_RDONLY | os.O_NONBLOCK))
schliessen = schliessen or os.close
try:
fd = oeffnen(device_path)
except OSError as e:
return AUSWURF_KEIN_ZUGRIFF, e
try:
# Entriegeln ist der entscheidende Schritt. Scheitert er, wird der
# Auswurf trotzdem versucht — bei einem nicht verriegelten Laufwerk
# (oder einem, das das ioctl nicht kennt) klappt er ohnehin.
try:
ioctl_fn(fd, CDROM_LOCKDOOR, 0)
except OSError:
pass
try:
ioctl_fn(fd, CDROMEJECT, 0)
except OSError as e:
return AUSWURF_ABGELEHNT, e
# Nachsehen statt hoffen: Die Schublade braucht ein bis zwei Sekunden.
for _ in range(AUSWURF_WARTEN_SEKUNDEN):
try:
if _auswurf_geglueckt(ioctl_fn(fd, CDROM_DRIVE_STATUS, 0)):
return AUSWURF_OK, None
except OSError as e:
return AUSWURF_ABGELEHNT, e
warten(1)
return AUSWURF_BLEIBT_DRIN, None
finally:
try:
schliessen(fd)
except OSError:
pass
def auswerfen_versuchen(device_path: str, ioctl_fn=None, oeffnen=None,
schliessen=None, warten=None) -> bool:
"""Wirft die Disc aus und prüft nach. Wirft NIE — Vertrag des Workers.
Ein Rip soll nicht daran sterben, dass die Schublade klemmt: `tasks.py`
protokolliert das Ergebnis und macht weiter.
"""
grund, _ = auswerfen_mit_grund(
device_path, ioctl_fn=ioctl_fn, oeffnen=oeffnen,
schliessen=schliessen, warten=warten,
)
return grund == AUSWURF_OK
def eject(device_path: str) -> None:
"""Wirft die Disc aus. Wirft OSError, wenn sie drin bleibt — Vertrag der API.
Die API zeigt die Meldung im Browser, deshalb muss ein Fehlschlag hier
laut sein und darf nicht als stilles False untergehen.
"""
grund, fehler = auswerfen_mit_grund(device_path)
if grund == AUSWURF_OK:
return
if fehler is not None:
raise fehler # unverändert weiterreichen (ENOENT, EPERM, …)
if grund == AUSWURF_KEIN_ZUGRIFF:
raise OSError(f"Das Laufwerk {device_path} ließ sich nicht ansprechen.")
raise OSError(
"Das Laufwerk hat den Auswurf angenommen, die Disc ist aber noch "
"drin. Blockiert etwas die Schublade, oder läuft noch ein Zugriff?"
)
def verriegeln(device_path: str, an: bool) -> None:
"""Tür verriegeln (True) oder entriegeln (False).
Bisher gab es das nur als Seitenschritt im Auswurf. Als eigene Operation
steht es im Port, weil der Windows-Treiber sie ebenfalls braucht
(IOCTL_STORAGE_MEDIA_REMOVAL) und weil ein Rip die Tür bewusst zuhalten
können soll.
"""
from fcntl import ioctl
fd = os.open(device_path, os.O_RDONLY | os.O_NONBLOCK)
try:
ioctl(fd, CDROM_LOCKDOOR, 1 if an else 0)
finally:
os.close(fd)
def list_optical_devices() -> list:
"""Alle optischen Laufwerke, die dieser Rechner sieht.
Im Container sind das die, die per `devices:` durchgereicht wurden.
"""
return sorted(glob.glob("/dev/sr[0-9]*"))
def read_sys_attr(device_name: str, attr: str) -> str:
pfad = f"/sys/class/block/{device_name}/device/{attr}"
try:
with open(pfad, encoding="ascii", errors="replace") as f:
return f.read().strip()
except OSError:
return ""
def device_info(device_path: str) -> dict:
"""Baut den Geräte-Eintrag fürs UI: Name aus /sys, Disc-Status per ioctl.
Kein udevadm: Im Container läuft kein udevd, die udev-Datenbank ist leer —
`udevadm info` lieferte dort schlicht nichts (deshalb zeigte der alte Weg
nie ein Laufwerk an). `/sys/class/block/<name>/device/*` kommt direkt vom
Kernel und funktioniert überall.
"""
# Erst hier importiert: detection braucht `fcntl`. Auf einer Maschine, die
# ueberhaupt Laufwerke abfragt, ist das vorhanden — und wenn nicht, soll es
# LAUT scheitern statt beim Import des ganzen Moduls.
from rippy.drives.detection import (
classify,
disc_size_bytes,
disc_status,
drive_status,
)
name = os.path.basename(device_path)
vendor = read_sys_attr(name, "vendor")
model = read_sys_attr(name, "model")
try:
status_code = drive_status(device_path)
except OSError:
status_code = -1
disc_type = "unknown"
status = "empty"
if status_code == CDS_DISC_OK:
status = "ready"
try:
disc_type = classify(disc_status(device_path), disc_size_bytes(device_path))
except OSError:
disc_type = "unknown"
return {
"id": name,
"name": " ".join(teil for teil in (vendor, model) if teil) or f"Laufwerk {name}",
"type": disc_type,
"path": device_path,
"status": status,
"model": model,
"serial": read_sys_attr(name, "wwid"),
}
@@ -1,12 +1,16 @@
"""Tests für detection.py — die pure Zuordnung classify().
"""Tests für die reine Zuordnung classify() aus cdrom.py.
Der Vorgänger (`file -L` auf ein Block-Device) konnte strukturell nie etwas
erkennen; sein Test mockte sich die file-Ausgabe passend zurecht. Hier wird
nur echte, deterministische Logik getestet die ioctl-Aufrufe selbst sind
dünne Kernel-Durchreichen und werden im E2E-Test mit echter Disc bewiesen.
Seit V2-1 liegen Konstanten und classify() in cdrom.py ohne `fcntl`.
Damit laufen diese Tests auch auf dem Windows-Entwicklungsrechner und
nicht mehr nur in der Ampel (vorher: erst nach dem Push bewiesen").
"""
from rippy.drives.detection import (
from rippy.drives.cdrom import (
BLURAY_MIN_BYTES,
CDS_AUDIO,
CDS_DATA_1,
+157
View File
@@ -0,0 +1,157 @@
"""Tests fuer den Linux-Laufwerks-Treiber — vor allem den Auswurf.
## Warum diese Tests hier liegen (Etappe V2-1, 28.08.2026)
Sie standen bis hierher in docker/worker/test_ripping_helpers.py, weil der
Auswurf in ripping.py lebte — und ein zweites Mal, anders geschrieben, in
docker/api/devices.py. Der Code liegt jetzt an EINER Stelle
(rippy/drives/linux.py), also liegen die Tests daneben.
Neu gegenueber vorher: Auch der ANDERE Vertrag wird geprueft. Die API-Fassung
(`eject`) muss werfen, wo die Worker-Fassung (`auswerfen_versuchen`) False
zurueckgibt. Genau diese Unterscheidung war vorher ungetestet, weil sie in
einer Datei stand, die niemand testete.
"""
from rippy.drives import linux
def _laufwerk(verriegelt=True, kennt_lockdoor=True):
"""Ein nachgebautes Laufwerk, das sich wie das echte verhaelt.
Gemessen am BU40N der Rippy-VM: Nach einem MakeMKV-Rip ist die Tuer
verriegelt. CDROMEJECT wird dann ANGENOMMEN und tut nichts - der Status
bleibt auf 4 (Disc drin). Erst CDROM_LOCKDOOR 0 macht den Auswurf wirksam.
"""
zustand = {"verriegelt": verriegelt, "status": linux.CDS_DISC_OK,
"aufrufe": []}
def ioctl_fn(fd, befehl, arg=0):
zustand["aufrufe"].append(befehl)
if befehl == linux.CDROM_LOCKDOOR:
if not kennt_lockdoor:
raise OSError("ioctl unbekannt")
zustand["verriegelt"] = bool(arg)
return 0
if befehl == linux.CDROMEJECT:
if not zustand["verriegelt"]:
zustand["status"] = linux.CDS_TRAY_OPEN
return 0 # <- auch verriegelt: ERFOLG, aber ohne Wirkung
if befehl == linux.CDROM_DRIVE_STATUS:
return zustand["status"]
raise OSError("unerwartetes ioctl")
return zustand, ioctl_fn
def test_auswurf_entriegelt_zuerst_und_klappt_dann():
zustand, ioctl_fn = _laufwerk(verriegelt=True)
ok = linux.auswerfen_versuchen(
"/dev/sr0", ioctl_fn=ioctl_fn, oeffnen=lambda p: 42,
schliessen=lambda fd: None, warten=lambda s: None)
assert ok is True
assert zustand["status"] == linux.CDS_TRAY_OPEN
# Reihenfolge: entriegeln VOR auswerfen
assert zustand["aufrufe"][0] == linux.CDROM_LOCKDOOR
assert zustand["aufrufe"][1] == linux.CDROMEJECT
def test_auswurf_meldet_fehlschlag_wenn_die_disc_drin_bleibt():
"""Der eigentliche Fehler. Vorher gab wirf_disc_aus True zurueck, weil das
ioctl nicht geworfen hatte - und ins Log kam "Disc ausgeworfen", waehrend
die Schublade zu blieb. Ein ioctl-Rueckgabewert beweist nichts."""
# Ein Laufwerk, das LOCKDOOR nicht kennt und verriegelt bleibt
zustand, ioctl_fn = _laufwerk(verriegelt=True, kennt_lockdoor=False)
ok = linux.auswerfen_versuchen(
"/dev/sr0", ioctl_fn=ioctl_fn, oeffnen=lambda p: 42,
schliessen=lambda fd: None, warten=lambda s: None)
assert ok is False
assert zustand["status"] == linux.CDS_DISC_OK # nie aufgegangen
def test_auswurf_bei_slot_laufwerk_ohne_schublade():
"""Ein Slot-Laufwerk hat keine Schublade und meldet nach dem Auswerfen
CDS_NO_DISC. Das muss als Erfolg zaehlen."""
assert linux._auswurf_geglueckt(linux.CDS_NO_DISC) is True
assert linux._auswurf_geglueckt(linux.CDS_TRAY_OPEN) is True
assert linux._auswurf_geglueckt(linux.CDS_DISC_OK) is False
assert linux._auswurf_geglueckt(linux.CDS_DRIVE_NOT_READY) is False
def test_auswurf_ohne_laufwerk_wirft_nicht():
def oeffnen_kaputt(p):
raise OSError("kein Laufwerk")
assert linux.auswerfen_versuchen(
"/dev/sr9", ioctl_fn=lambda *a: 0, oeffnen=oeffnen_kaputt,
schliessen=lambda fd: None, warten=lambda s: None) is False
# --- Der ANDERE Vertrag: eject() der API muss WERFEN ------------------------
#
# Diese Fälle waren bis V2-1 ungetestet. Der Auswurf existierte zweimal, und
# die API-Fassung (docker/api/devices.py) hatte gar keine Testdatei — obwohl
# genau sie dem Nutzer im Browser eine Meldung zeigt. Jetzt liegt beides
# nebeneinander, also wird auch beides geprüft.
def test_eject_ist_still_wenn_die_disc_rauskommt(monkeypatch):
monkeypatch.setattr(linux, "auswerfen_mit_grund",
lambda pfad: (linux.AUSWURF_OK, None))
assert linux.eject("/dev/sr0") is None
def test_eject_wirft_wenn_die_disc_drin_bleibt(monkeypatch):
"""Der Unterschied zum Worker-Vertrag: hier ist ein stilles False falsch.
Der Commander drückt im Browser auf „Auswerfen". Passiert nichts und die
API schweigt, sucht er den Fehler am Laufwerk — genau der Verlauf vom
26.07.2026.
"""
monkeypatch.setattr(linux, "auswerfen_mit_grund",
lambda pfad: (linux.AUSWURF_BLEIBT_DRIN, None))
try:
linux.eject("/dev/sr0")
except OSError as e:
assert "noch" in str(e) and "drin" in str(e)
else:
raise AssertionError("eject haette werfen muessen")
def test_eject_reicht_den_urspruenglichen_fehler_unveraendert_weiter(monkeypatch):
"""ENOENT/EPERM sollen NICHT hinter einer eigenen Meldung verschwinden.
Vorher warf devices.eject den rohen OSError von os.open. Wer den ersetzt,
nimmt dem Nutzer die einzige brauchbare Auskunft („Datei nicht gefunden"
gegen „keine Berechtigung" sind zwei ganz verschiedene Probleme).
"""
original = OSError(2, "No such file or directory")
monkeypatch.setattr(linux, "auswerfen_mit_grund",
lambda pfad: (linux.AUSWURF_KEIN_ZUGRIFF, original))
try:
linux.eject("/dev/sr9")
except OSError as e:
assert e is original
else:
raise AssertionError("eject haette werfen muessen")
def test_auswerfen_mit_grund_unterscheidet_die_fehlerarten():
"""Die gemeinsame Mechanik muss sagen KÖNNEN, was schiefging — sonst
kann eject() den Unterschied nicht machen."""
zustand, ioctl_fn = _laufwerk(verriegelt=True, kennt_lockdoor=False)
grund, fehler = linux.auswerfen_mit_grund(
"/dev/sr0", ioctl_fn=ioctl_fn, oeffnen=lambda p: 42,
schliessen=lambda fd: None, warten=lambda s: None)
assert grund == linux.AUSWURF_BLEIBT_DRIN
assert fehler is None
def oeffnen_kaputt(p):
raise OSError(13, "Permission denied")
grund, fehler = linux.auswerfen_mit_grund(
"/dev/sr0", ioctl_fn=lambda *a: 0, oeffnen=oeffnen_kaputt,
schliessen=lambda fd: None, warten=lambda s: None)
assert grund == linux.AUSWURF_KEIN_ZUGRIFF
assert fehler.errno == 13
+153
View File
@@ -0,0 +1,153 @@
"""Die vier Ports — die Nahtstellen, an denen Rippy austauschbar wird.
## Wozu das gut ist (KONZEPT-V2.md § 1 und § 2.3)
Rippy v2 soll in drei Betriebsarten laufen: Docker, native Windows-App,
Headless-Linux-Dienst. Der Trick dabei ist, dass es NICHT drei Programme sind,
sondern eines — und ein Betriebsmodus nur die Auswahl der Treiber hinter diesen
vier Nahtstellen:
Port Standalone Verteilt
──────── ─────────────────────────────── ──────────────────────────
Store SQLite (WAL) PostgreSQL
Queue LocalQueue (Tabelle + Pool) Celery über Redis
Bus In-Process (asyncio) Redis Pub/Sub
Drives LinuxDrives / WindowsDrives LinuxDrives je Knoten
## Warum Protocol und nicht Basisklasse
Ein `typing.Protocol` prüft die FORM, nicht die Abstammung. Ein Treiber muss
nichts erben und nichts importieren — er muss nur die Methoden haben. Damit
lässt sich in Tests ein handgeschriebenes Fake-Objekt einsetzen, ohne
Mocking-Bibliothek. Genau so arbeitet `zombies.raeume_zombies_auf(celery, db)`
heute schon: Die Tests dort geben ein `FakeDb` hinein, und es funktioniert,
weil nur die Form zählt. Diese Datei schreibt die Form auf, die bisher nur
stillschweigend galt.
## Stand (Etappe V2-1, 28.08.2026)
Beschrieben sind alle vier Ports. Treiber gibt es bisher für:
Store -> rippy.store (SQLAlchemy Core; heute Postgres,
ab V2-2 auch SQLite — derselbe Code)
Drives -> rippy.drives.linux (ioctl; Windows folgt in V2-4)
Queue -> Celery, noch ohne Adapter (V2-2)
Bus -> noch keiner (V2-3)
Die Protokolle stehen trotzdem schon hier: Sie sind der Vertrag, gegen den die
nächsten Etappen gebaut werden, und sie machen sichtbar, wo heute noch etwas
fehlt. Ein Port ohne zweiten Treiber ist keine Abstraktion auf Vorrat — er ist
die Stelle, an der der zweite Treiber hinkommt.
"""
from typing import AsyncIterator, Iterable, Optional, Protocol, runtime_checkable
# ───────────────────────────────────────────────────────────── Store ────
@runtime_checkable
class Store(Protocol):
"""Zustand: Jobs, Logs, Einstellungen, Worker, Speicherziele.
WICHTIG (KONZEPT-V2.md § 3.1): Der Store ist die WAHRHEIT über den
Job-Zustand — nicht der Broker. Daraus folgt alles Weitere: Ein Auftrag
gehört dem, der ihn per bedingtem UPDATE übernommen hat, und er ist wieder
frei, wenn dessen Lease abläuft. Ein abgestürzter Knoten hinterlässt damit
keine Leiche, die jemand einsammeln muss.
"""
def init_db(self) -> None: ...
def get_job(self, job_id: str) -> Optional[dict]: ...
def list_jobs(self, limit: int = 100) -> list: ...
def insert_job(self, job_id: str, device: str, **felder) -> None: ...
def update_job(self, job_id: str, **felder) -> None: ...
def delete_job(self, job_id: str) -> None: ...
def get_job_status(self, job_id: str) -> str: ...
def list_jobs_mit_status(self, stati: Iterable[str]) -> list: ...
def add_log(self, level: str, source: str, message: str) -> None: ...
def list_logs(self, limit: int = 200) -> list: ...
def get_settings(self, key: str = "ui") -> dict: ...
def save_settings(self, werte: dict, key: str = "ui") -> None: ...
def save_worker(self, name: str, encoder_liste: list, info: dict = None) -> None: ...
def list_workers(self) -> list: ...
def zaehle_online_worker(self, sekunden: int = 120) -> int: ...
# ───────────────────────────────────────────────────────────── Queue ────
@runtime_checkable
class Queue(Protocol):
"""Aufträge: einreihen, übernehmen, am Leben halten, abschließen.
`uebernehmen` gibt einem Knoten einen Auftrag NUR, wenn er die verlangten
Fähigkeiten hat (`{"drive": "sr0"}` bzw. `{"encoder": "nvenc"}`) — das ist
dieselbe Auswahl, die v1 über getrennte Celery-Queues und `worker_direct`
gelöst hat, nur als Datum statt als Verdrahtung.
`lebenszeichen` ist der Kern: Wer einen Auftrag hält, verlängert alle paar
Sekunden seine Lease. Läuft sie ab, ist der Auftrag frei — egal ob der
Knoten abgestürzt ist, das Netz weg war oder jemand den Stecker gezogen
hat. Das ersetzt die nachträgliche Zombie-Erkennung aus v1 durch einen
Mechanismus, der den Fall gar nicht erst entstehen lässt.
"""
def einreihen(self, auftrag: dict) -> str: ...
def uebernehmen(self, faehigkeiten: set, knoten: str) -> Optional[dict]: ...
def lebenszeichen(self, auftrag_id: str) -> None: ...
def abschliessen(self, auftrag_id: str, ergebnis: dict) -> None: ...
def fehlgeschlagen(self, auftrag_id: str, fehler: str) -> None: ...
def abbrechen(self, auftrag_id: str) -> None: ...
# ─────────────────────────────────────────────────────────────── Bus ────
@runtime_checkable
class Bus(Protocol):
"""Ereignisse — Feuer und vergiss.
⚠️ Der Bus überträgt nur NACHRICHTEN ÜBER ÄNDERUNGEN, nie den Zustand
selbst. Wer den Zustand will, fragt den Store.
Der Grund steht in AGENTS.md: „Ein verpasster Abruf ist keine Nachricht
über die Welt." Im UI stand fünfmal `catch(() => [])` — jeder
fehlgeschlagene Abruf hieß damit „es gibt keine Jobs". Wenn der Bus keinen
Zustand trägt, kann ein verpasstes Ereignis auch keinen Zustand löschen.
`abonnieren(ab_seq=…)` liefert alles ab einer Folgenummer nach. Ist die
Lücke zu groß, schickt der Treiber ausdrücklich ein `snapshot`-Ereignis
statt stillschweigend Deltas.
"""
async def senden(self, ereignis: dict) -> None: ...
def abonnieren(self, ab_seq: Optional[int] = None) -> AsyncIterator[dict]: ...
# ──────────────────────────────────────────────────────────── Drives ────
@runtime_checkable
class Drives(Protocol):
"""Hardware: Laufwerke finden, Zustand lesen, verriegeln, auswerfen.
Der EINZIGE Ort im Projekt mit ioctl- bzw. Win32-Aufrufen.
Zwei Regeln gelten für jeden Treiber (KONZEPT-V2.md § 5):
1. **Auswerfen heißt: entriegeln, auswerfen, NACHSEHEN.** In dieser
Reihenfolge, auf jeder Plattform. Am 26.07.2026 am laufenden System
gemessen: Ein nacktes CDROMEJECT wird von einem verriegelten Laufwerk
mit ERFOLG quittiert und tut nichts — MakeMKV verriegelt die Tür
während des Rips und entriegelt sie nicht wieder.
2. **Jede Operation mit Zeitgrenze, nie im Ereignis-Loop.** Ein hängendes
ioctl auf einem defekten Laufwerk darf die API nicht blockieren —
dieselbe Fehlerklasse wie die CIFS-Blockade aus v1.
"""
def laufwerke(self) -> list: ...
def zustand(self, geraet: str) -> str: ...
def info(self, geraet: str) -> dict: ...
def verriegeln(self, geraet: str, an: bool) -> None: ...
def auswerfen(self, geraet: str) -> None: ...
async def ereignisse(self) -> AsyncIterator[dict]: ...
__all__ = ["Store", "Queue", "Bus", "Drives"]
+216 -62
View File
@@ -1,16 +1,41 @@
"""Job-, Log- und Settings-Persistenz in PostgreSQL (KONZEPT: Postgres für Job-Logs).
"""Store-Treiber: Jobs, Logs, Einstellungen, Worker, Speicherziele.
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.
## Diese Datei war zwei (Etappe V2-1, 28.08.2026)
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.
Bis hierher gab es `docker/api/db.py` UND `docker/worker/db.py` mit
identischen Tabellendefinitionen, überlappenden Funktionen und diesem Satz im
Kopf beider Dateien:
Die Tabellendefinition existiert bewusst identisch in API und Worker
es gibt kein geteiltes Paket zwischen den Containern.
Wer die Struktur ändert, ändert BEIDE Dateien."
Das Paket gibt es seit V2-0. Also gibt es die Datei jetzt einmal, und beide
Seiten importieren sie als `from rippy import store as db` der Aufrufname
`db.` bleibt überall gleich, damit der Umzug keine 60 Aufrufstellen anfasst.
## SQLAlchemy Core, nicht ORM — und warum das ein Glücksfall ist
v1 hat hier von Anfang an mit `Table()`, `select()` und `insert()` gearbeitet
statt mit ORM-Klassen. Genau dieser Code läuft auf SQLite und auf PostgreSQL
gleichermaßen. Der SQLite-Treiber für den Standalone-Betrieb (Etappe V2-2) ist
deshalb kein Neubau, sondern eine andere `DATABASE_URL` plus ein paar PRAGMAs.
## Was hier NOCH nicht stimmt (offen für V2-2)
`init_db()` zieht Mini-Migrationen mit `ALTER TABLE ADD COLUMN IF NOT EXISTS`
nach. Das ist **Postgres-Syntax** auf SQLite bricht es. Es bleibt für diese
Etappe unverändert stehen, weil V2-1 laut Plan KEIN Verhalten ändert; V2-2
ersetzt es durch Alembic.
Erfüllt den Port `rippy.ports.Store` (als Modul, nicht als Klasse die
Aufrufstellen benutzen es seit v1 so, und eine Umstellung auf eine Instanz
wäre eine Verhaltensänderung ohne Gegenwert in dieser Etappe).
"""
import json
import os
from datetime import datetime, timezone
from datetime import datetime, timedelta, timezone
from sqlalchemy import (
Column,
@@ -21,6 +46,7 @@ from sqlalchemy import (
Table,
Text,
create_engine,
func,
select,
)
@@ -86,56 +112,20 @@ storage_mounts = Table(
)
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."""
"""Legt fehlende Tabellen an (idempotent) und zieht Mini-Migrationen nach.
Wer zuerst startet API oder Worker , migriert; der andere findet die
Spalten dann bereits vor. Das war schon vor der Zusammenlegung so und
bleibt es.
"""
metadata.create_all(engine)
# create_all ändert BESTEHENDE Tabellen nicht — neue Spalten hier nachziehen:
# create_all ändert BESTEHENDE Tabellen nicht — neue Spalten hier nachziehen.
# ⚠️ Postgres-Syntax, bricht auf SQLite. Ersatz durch Alembic in V2-2.
with engine.begin() as conn:
conn.exec_driver_sql(
"ALTER TABLE jobs ADD COLUMN IF NOT EXISTS target_dir VARCHAR(255)"
@@ -148,6 +138,7 @@ def init_db() -> None:
)
# ───────────────────────────────────────────────────────────── Jobs ────
def insert_job(
job_id: str, device: str, disc_type: str = None, title: str = None,
target_dir: str = None, meta: str = None,
@@ -174,6 +165,8 @@ def update_job(job_id: str, **fields) -> None:
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)."""
with engine.connect() as conn:
zeile = conn.execute(
select(jobs).where(jobs.c.id == job_id)
@@ -181,6 +174,15 @@ def get_job(job_id: str) -> dict:
return dict(zeile) if zeile else None
def get_job_status(job_id: str) -> str:
"""Nur der Status — der Worker prüft damit kooperative Abbruch-Anfragen."""
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(limit: int = 100) -> list:
with engine.connect() as conn:
zeilen = conn.execute(
@@ -189,6 +191,20 @@ def list_jobs(limit: int = 100) -> list:
return [dict(z) for z in zeilen]
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.
"""
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 delete_job(job_id: str) -> None:
"""Entfernt EINEN Job-Eintrag (nur die DB-Zeile — Dateien bleiben)."""
with engine.begin() as conn:
@@ -204,12 +220,6 @@ def delete_finished_jobs() -> int:
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 hat_arbeit() -> bool:
"""True, wenn IRGENDEIN Job noch nicht durch ist (egal auf welchem Gerät).
@@ -240,6 +250,37 @@ def has_active_job(device: str) -> bool:
return zeile is not None
def meta_merken(job_id: str, **felder) -> None:
"""Ergänzt EINZELNE Schlüssel in den Job-Metadaten. Wirft nie.
`update_job(meta=...)` würde die Spalte ersetzen Poster, Jahr, Titel-Wahl
und Sprachwunsch dieses Rips wären damit fort. Also lesen, mischen,
schreiben.
Fehler werden geschluckt: Diese Funktion vermerkt nur, in welcher Phase ein
Job steht (api/phasen.py). Ein Rip darf daran nicht scheitern im
schlimmsten Fall fehlt die Marke und der Neu"-Knopf fragt nach.
"""
try:
zeile = get_job(job_id)
if not zeile:
return
try:
vorher = json.loads(zeile.get("meta") or "{}")
except (ValueError, TypeError):
vorher = {}
if not isinstance(vorher, dict):
vorher = {}
vorher.update(felder)
update_job(job_id, meta=json.dumps(vorher))
except Exception as e:
try:
add_log("warning", "worker", f"Job {job_id}: Metadaten-Vermerk fehlgeschlagen: {e}")
except Exception:
pass
# ───────────────────────────────────────────────────────────── Logs ────
def add_log(level: str, source: str, message: str) -> None:
with engine.begin() as conn:
conn.execute(
@@ -255,11 +296,37 @@ def list_logs(limit: int = 200) -> list:
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()
# ───────────────────────────────────────────────────── Einstellungen ────
def get_settings(key: str = "ui", bei_fehler_leer: bool = False) -> dict:
"""UI-Einstellungen lesen.
`bei_fehler_leer` ist KEIN Komfort-Schalter, sondern der einzige echte
Unterschied, den die beiden alten db.py-Dateien hatten und er ist eine
offene Frage, keine Entscheidung.
Die Worker-Fassung fing jeden Fehler ab und gab `{}` zurück: Ein Rip sollte
nicht daran sterben, dass die Datenbank kurz zickt. Die API-Fassung ließ
den Fehler durch. Beim Zusammenlegen wäre es leicht gewesen, sich still für
eine Seite zu entscheiden genau das tut diese Etappe NICHT, denn V2-1
ändert kein Verhalten. Deshalb steht der Unterschied jetzt sichtbar in der
Signatur statt unsichtbar in zwei Dateien.
Wichtig zu wissen, wenn das später entschieden wird: `{}` heißt für den
Aufrufer keine Einstellungen gesetzt", nicht „konnte nicht nachsehen". Der
Worker nimmt dann stillschweigend Vorgabewerte womöglich das falsche
Preset oder das falsche Zielverzeichnis. Das ist dieselbe Fehlerklasse wie
`catch(() => [])` im alten UI (AGENTS.md: Ein verpasster Abruf ist keine
Nachricht über die Welt"). Zu klären in V2-2, wenn der Store-Port steht.
"""
try:
with engine.connect() as conn:
zeile = conn.execute(
select(settings_table.c.value).where(settings_table.c.key == key)
).first()
except Exception:
if bei_fehler_leer:
return {}
raise
if not zeile or not zeile[0]:
return {}
try:
@@ -269,6 +336,8 @@ def get_settings(key: str = "ui") -> dict:
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:<device>'), die API liest sie."""
payload = json.dumps(werte)
with engine.begin() as conn:
vorhanden = conn.execute(
@@ -282,3 +351,88 @@ def save_settings(werte: dict, key: str = "ui") -> None:
)
else:
conn.execute(settings_table.insert().values(key=key, value=payload))
# ────────────────────────────────────────────────────────── Worker ────
def save_worker(name: str, encoder_liste: list, info: dict = None) -> None:
"""Worker meldet Name + Encoder-Fähigkeiten + Werkzeug-Versionen (Upsert)."""
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 list_workers() -> list:
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 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 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.
"""
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)
# ──────────────────────────────────────────────────── Speicherziele ────
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))