refactor(core): V2-1 — die vier Ports, ein Store, ein Laufwerks-Treiber
Ampel / ampel (push) Successful in 40s
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:
co-authored by
Claude Opus 5
parent
edfe0d4313
commit
dd1d0b7365
+11
-1
@@ -29,11 +29,21 @@ Thumbs.db
|
|||||||
.dockerignore
|
.dockerignore
|
||||||
.docker/
|
.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/
|
||||||
|
**/tests/
|
||||||
test_*.py
|
test_*.py
|
||||||
|
**/test_*.py
|
||||||
*_test.py
|
*_test.py
|
||||||
|
**/*_test.py
|
||||||
.pytest_cache/
|
.pytest_cache/
|
||||||
|
**/.pytest_cache/
|
||||||
|
**/__pycache__/
|
||||||
|
conftest.py
|
||||||
|
**/conftest.py
|
||||||
|
|
||||||
# Build
|
# Build
|
||||||
dist/
|
dist/
|
||||||
|
|||||||
@@ -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
@@ -11,8 +11,8 @@ import shutil
|
|||||||
import time
|
import time
|
||||||
import uuid
|
import uuid
|
||||||
|
|
||||||
import db
|
from rippy import store as db
|
||||||
import devices as device_discovery
|
from rippy.drives import linux as device_discovery
|
||||||
import eta
|
import eta
|
||||||
from rippy.rip import makemkv_daten
|
from rippy.rip import makemkv_daten
|
||||||
import makemkv_key
|
import makemkv_key
|
||||||
|
|||||||
@@ -50,7 +50,7 @@ def fetch_current_key(timeout: int = 15) -> str | None:
|
|||||||
def _apply_key(key: str) -> bool:
|
def _apply_key(key: str) -> bool:
|
||||||
"""Merge den Key in die Settings. save_settings überschreibt das GANZE JSON,
|
"""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."""
|
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()
|
einstellungen = db.get_settings()
|
||||||
if (einstellungen.get("makemkvAppKey") or "").strip() == key:
|
if (einstellungen.get("makemkvAppKey") or "").strip() == key:
|
||||||
return False
|
return False
|
||||||
@@ -61,7 +61,7 @@ def _apply_key(key: str) -> bool:
|
|||||||
|
|
||||||
async def refresh_once() -> bool:
|
async def refresh_once() -> bool:
|
||||||
"""Einmal holen + anwenden. Best-effort, loggt in die App-Logs. True = geaendert."""
|
"""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)
|
key = await asyncio.to_thread(fetch_current_key)
|
||||||
if not key:
|
if not key:
|
||||||
db.add_log("warning", "makemkv-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:
|
async def refresh_loop(intervall_stunden: int = 24, start_verzoegerung_s: int = 60) -> None:
|
||||||
"""Periodische Erneuerung (Default: taeglich, damit ein Monatswechsel nie verpasst
|
"""Periodische Erneuerung (Default: taeglich, damit ein Monatswechsel nie verpasst
|
||||||
wird). Fehler werden geloggt, aber nie geworfen -> die API läuft weiter."""
|
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)
|
await asyncio.sleep(start_verzoegerung_s)
|
||||||
while True:
|
while True:
|
||||||
try:
|
try:
|
||||||
|
|||||||
@@ -16,7 +16,7 @@ import re
|
|||||||
import subprocess
|
import subprocess
|
||||||
import tempfile
|
import tempfile
|
||||||
|
|
||||||
import db
|
from rippy import store as db
|
||||||
|
|
||||||
MEDIA_ROOT = "/app/media"
|
MEDIA_ROOT = "/app/media"
|
||||||
NAME_MUSTER = re.compile(r"^[a-z0-9][a-z0-9-]{1,30}$")
|
NAME_MUSTER = re.compile(r"^[a-z0-9][a-z0-9-]{1,30}$")
|
||||||
|
|||||||
@@ -364,8 +364,8 @@ def werkzeug_versionen() -> dict:
|
|||||||
info["handbrake"] = "installiert"
|
info["handbrake"] = "installiert"
|
||||||
# Woher kommt der MakeMKV-Key? UI-Setting schlägt Env — ehrlich anzeigen.
|
# Woher kommt der MakeMKV-Key? UI-Setting schlägt Env — ehrlich anzeigen.
|
||||||
try:
|
try:
|
||||||
import db
|
from rippy import store as db
|
||||||
ui_key = (db.get_settings().get("makemkvAppKey") or "").strip()
|
ui_key = (db.get_settings(bei_fehler_leer=True).get("makemkvAppKey") or "").strip()
|
||||||
except Exception:
|
except Exception:
|
||||||
ui_key = ""
|
ui_key = ""
|
||||||
if ui_key:
|
if ui_key:
|
||||||
|
|||||||
@@ -32,7 +32,7 @@ celery_app.conf.update(
|
|||||||
# Import auf Modulebene: im worker_ready-Signal ist /app nicht mehr
|
# Import auf Modulebene: im worker_ready-Signal ist /app nicht mehr
|
||||||
# zuverlässig im sys.path (ModuleNotFoundError 'caps', Deploy 23.07.).
|
# zuverlässig im sys.path (ModuleNotFoundError 'caps', Deploy 23.07.).
|
||||||
import caps # noqa: E402
|
import caps # noqa: E402
|
||||||
import db # noqa: E402
|
from rippy import store as db # noqa: E402
|
||||||
import zombies # noqa: E402
|
import zombies # noqa: E402
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -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
File diff suppressed because it is too large
Load Diff
+19
-105
@@ -17,6 +17,25 @@ import shutil
|
|||||||
import subprocess
|
import subprocess
|
||||||
import tempfile
|
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
|
from winlauf import OHNE_FENSTER
|
||||||
|
|
||||||
RIP_OUTPUT_DIR = os.getenv("RIP_OUTPUT_DIR", "/app/media")
|
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)."""
|
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:
|
def check_makemkv_installed() -> bool:
|
||||||
"""Prüft, ob makemkvcon installiert ist."""
|
"""Prüft, ob makemkvcon installiert ist."""
|
||||||
return shutil.which("makemkvcon") is not None
|
return shutil.which("makemkvcon") is not None
|
||||||
|
|||||||
@@ -22,7 +22,7 @@ import time
|
|||||||
|
|
||||||
import requests
|
import requests
|
||||||
|
|
||||||
import db
|
from rippy import store as db
|
||||||
from rippy.rip import makemkv_daten
|
from rippy.rip import makemkv_daten
|
||||||
import medien
|
import medien
|
||||||
from rippy.core import notify
|
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:
|
def _benachrichtigen(job_id: str, betreff: str, text: str, level: str) -> None:
|
||||||
"""Webhook-Meldung bei Job-Ende — best effort, nie job-entscheidend."""
|
"""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()
|
url = (einstellungen.get("notificationWebhook") or "").strip()
|
||||||
if not url:
|
if not url:
|
||||||
return
|
return
|
||||||
@@ -402,7 +402,7 @@ def _job_abschliessen(job_id: str, ergebnis: dict) -> None:
|
|||||||
if ergebnis.get("status") == "success":
|
if ergebnis.get("status") == "success":
|
||||||
ausgabe = ergebnis.get("output_dir")
|
ausgabe = ergebnis.get("output_dir")
|
||||||
if ausgabe and (job.get("disc_type") in ("dvd", "bluray", "uhd")):
|
if ausgabe and (job.get("disc_type") in ("dvd", "bluray", "uhd")):
|
||||||
einstellungen = db.get_settings()
|
einstellungen = db.get_settings(bei_fehler_leer=True)
|
||||||
try:
|
try:
|
||||||
meta = json.loads(job.get("meta") or "{}")
|
meta = json.loads(job.get("meta") or "{}")
|
||||||
except ValueError:
|
except ValueError:
|
||||||
@@ -527,7 +527,7 @@ def rip_disc(self, device_path: str, job_id: str, target_dir: str = None):
|
|||||||
gesehen.add(text)
|
gesehen.add(text)
|
||||||
db.add_log("info", "makemkv", f"Job {job_id}: {text[:300]}")
|
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")
|
ist_video = disc_type in ("dvd", "bluray", "uhd")
|
||||||
# Je Disc-Typ abwählbar (siehe komprimieren_fuer): 4K verlustfrei behalten,
|
# Je Disc-Typ abwählbar (siehe komprimieren_fuer): 4K verlustfrei behalten,
|
||||||
# DVDs trotzdem schrumpfen — vorher gab es nur alles oder nichts.
|
# 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)
|
_job_abschliessen(job_id, ergebnis)
|
||||||
return 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
|
# Preset nach Disc-Typ (Befund 25.07.2026): vorher lief JEDE Quelle durch
|
||||||
# dasselbe Preset — eine 4K-UHD wurde damit auf 1080p heruntergerechnet.
|
# dasselbe Preset — eine 4K-UHD wurde damit auf 1080p heruntergerechnet.
|
||||||
job = db.get_job(job_id) or {}
|
job = db.get_job(job_id) or {}
|
||||||
|
|||||||
@@ -383,90 +383,6 @@ def test_fehler_ohne_preset_problem_bleibt_der_alte():
|
|||||||
assert ergebnis["error"] == "HandBrake endete mit Code 1"
|
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) -----------
|
# --- Sprachen der Disc: gemessen an der Akira-Blu-ray (26.07.2026) -----------
|
||||||
#
|
#
|
||||||
# Woertlich aus `makemkvcon -r --noscan info dev:/dev/sr0` im Worker-Container
|
# Woertlich aus `makemkvcon -r --noscan info dev:/dev/sr0` im Worker-Container
|
||||||
|
|||||||
@@ -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"
|
||||||
@@ -14,32 +14,27 @@ import os
|
|||||||
import struct
|
import struct
|
||||||
from fcntl import ioctl
|
from fcntl import ioctl
|
||||||
|
|
||||||
# include/uapi/linux/cdrom.h
|
# Konstanten und die reine Zuordnung leben in cdrom.py — ohne fcntl, damit
|
||||||
CDROM_DRIVE_STATUS = 0x5326
|
# der Windows-Treiber (V2-4) und der native Windows-Worker sie benutzen
|
||||||
CDROM_DISC_STATUS = 0x5327
|
# koennen. Hier werden sie weiter angeboten, weil Aufrufstellen sie so kennen.
|
||||||
|
from rippy.drives.cdrom import ( # noqa: F401
|
||||||
CDS_NO_DISC = 1
|
BLKGETSIZE64,
|
||||||
CDS_TRAY_OPEN = 2
|
BLURAY_MIN_BYTES,
|
||||||
CDS_DRIVE_NOT_READY = 3
|
CDROM_DISC_STATUS,
|
||||||
CDS_DISC_OK = 4
|
CDROM_DRIVE_STATUS,
|
||||||
|
CDS_AUDIO,
|
||||||
CDS_AUDIO = 100
|
CDS_DATA_1,
|
||||||
CDS_DATA_1 = 101
|
CDS_DATA_2,
|
||||||
CDS_DATA_2 = 102
|
CDS_DISC_OK,
|
||||||
CDS_XA_2_1 = 103
|
CDS_DRIVE_NOT_READY,
|
||||||
CDS_XA_2_2 = 104
|
CDS_MIXED,
|
||||||
CDS_MIXED = 105
|
CDS_NO_DISC,
|
||||||
|
CDS_TRAY_OPEN,
|
||||||
# include/uapi/linux/fs.h: BLKGETSIZE64 = _IOR(0x12, 114, size_t) auf 64-bit
|
CDS_XA_2_1,
|
||||||
BLKGETSIZE64 = 0x80081272
|
CDS_XA_2_2,
|
||||||
|
UHD_MIN_BYTES,
|
||||||
# Eine DVD9 fasst ~8,5 GB; Blu-ray beginnt bei 25 GB (Single Layer).
|
classify,
|
||||||
# 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 _open_nonblock(device_path: str) -> int:
|
def _open_nonblock(device_path: str) -> int:
|
||||||
@@ -76,17 +71,6 @@ def disc_size_bytes(device_path: str) -> int:
|
|||||||
os.close(fd)
|
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:
|
def detect_disc_type(device_path: str) -> str:
|
||||||
"""Erkennt den Typ der eingelegten Disc; 'no_disc' wenn keine drin ist."""
|
"""Erkennt den Typ der eingelegten Disc; 'no_disc' wenn keine drin ist."""
|
||||||
try:
|
try:
|
||||||
|
|||||||
@@ -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
|
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
|
erkennen; sein Test mockte sich die file-Ausgabe passend zurecht. Hier wird
|
||||||
nur echte, deterministische Logik getestet — die ioctl-Aufrufe selbst sind
|
nur echte, deterministische Logik getestet — die ioctl-Aufrufe selbst sind
|
||||||
dünne Kernel-Durchreichen und werden im E2E-Test mit echter Disc bewiesen.
|
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,
|
BLURAY_MIN_BYTES,
|
||||||
CDS_AUDIO,
|
CDS_AUDIO,
|
||||||
CDS_DATA_1,
|
CDS_DATA_1,
|
||||||
@@ -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
|
||||||
@@ -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"]
|
||||||
@@ -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
|
## Diese Datei war zwei (Etappe V2-1, 28.08.2026)
|
||||||
(docker/worker/db.py) — es gibt kein geteiltes Paket zwischen den Containern.
|
|
||||||
Wer die Struktur ändert, ändert BEIDE Dateien. create_all ist idempotent.
|
|
||||||
|
|
||||||
Bis 23.07. lag Postgres komplett brach: GET /jobs gab hart [] zurück, nichts
|
Bis hierher gab es `docker/api/db.py` UND `docker/worker/db.py` — mit
|
||||||
schrieb je eine Zeile — der Job-Verlauf im UI war ein Placebo.
|
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 json
|
||||||
import os
|
import os
|
||||||
from datetime import datetime, timezone
|
from datetime import datetime, timedelta, timezone
|
||||||
|
|
||||||
from sqlalchemy import (
|
from sqlalchemy import (
|
||||||
Column,
|
Column,
|
||||||
@@ -21,6 +46,7 @@ from sqlalchemy import (
|
|||||||
Table,
|
Table,
|
||||||
Text,
|
Text,
|
||||||
create_engine,
|
create_engine,
|
||||||
|
func,
|
||||||
select,
|
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:
|
def utcnow() -> datetime:
|
||||||
return datetime.now(timezone.utc)
|
return datetime.now(timezone.utc)
|
||||||
|
|
||||||
|
|
||||||
def init_db() -> None:
|
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)
|
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:
|
with engine.begin() as conn:
|
||||||
conn.exec_driver_sql(
|
conn.exec_driver_sql(
|
||||||
"ALTER TABLE jobs ADD COLUMN IF NOT EXISTS target_dir VARCHAR(255)"
|
"ALTER TABLE jobs ADD COLUMN IF NOT EXISTS target_dir VARCHAR(255)"
|
||||||
@@ -148,6 +138,7 @@ def init_db() -> None:
|
|||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
# ───────────────────────────────────────────────────────────── Jobs ────
|
||||||
def insert_job(
|
def insert_job(
|
||||||
job_id: str, device: str, disc_type: str = None, title: str = None,
|
job_id: str, device: str, disc_type: str = None, title: str = None,
|
||||||
target_dir: str = None, meta: 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:
|
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:
|
with engine.connect() as conn:
|
||||||
zeile = conn.execute(
|
zeile = conn.execute(
|
||||||
select(jobs).where(jobs.c.id == job_id)
|
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
|
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:
|
def list_jobs(limit: int = 100) -> list:
|
||||||
with engine.connect() as conn:
|
with engine.connect() as conn:
|
||||||
zeilen = conn.execute(
|
zeilen = conn.execute(
|
||||||
@@ -189,6 +191,20 @@ def list_jobs(limit: int = 100) -> list:
|
|||||||
return [dict(z) for z in zeilen]
|
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:
|
def delete_job(job_id: str) -> None:
|
||||||
"""Entfernt EINEN Job-Eintrag (nur die DB-Zeile — Dateien bleiben)."""
|
"""Entfernt EINEN Job-Eintrag (nur die DB-Zeile — Dateien bleiben)."""
|
||||||
with engine.begin() as conn:
|
with engine.begin() as conn:
|
||||||
@@ -204,12 +220,6 @@ def delete_finished_jobs() -> int:
|
|||||||
return ergebnis.rowcount or 0
|
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:
|
def hat_arbeit() -> bool:
|
||||||
"""True, wenn IRGENDEIN Job noch nicht durch ist (egal auf welchem Gerät).
|
"""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
|
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:
|
def add_log(level: str, source: str, message: str) -> None:
|
||||||
with engine.begin() as conn:
|
with engine.begin() as conn:
|
||||||
conn.execute(
|
conn.execute(
|
||||||
@@ -255,11 +296,37 @@ def list_logs(limit: int = 200) -> list:
|
|||||||
return [dict(z) for z in zeilen]
|
return [dict(z) for z in zeilen]
|
||||||
|
|
||||||
|
|
||||||
def get_settings(key: str = "ui") -> dict:
|
# ───────────────────────────────────────────────────── Einstellungen ────
|
||||||
with engine.connect() as conn:
|
def get_settings(key: str = "ui", bei_fehler_leer: bool = False) -> dict:
|
||||||
zeile = conn.execute(
|
"""UI-Einstellungen lesen.
|
||||||
select(settings_table.c.value).where(settings_table.c.key == key)
|
|
||||||
).first()
|
⚠️ `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]:
|
if not zeile or not zeile[0]:
|
||||||
return {}
|
return {}
|
||||||
try:
|
try:
|
||||||
@@ -269,6 +336,8 @@ def get_settings(key: str = "ui") -> dict:
|
|||||||
|
|
||||||
|
|
||||||
def save_settings(werte: dict, key: str = "ui") -> 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."""
|
||||||
payload = json.dumps(werte)
|
payload = json.dumps(werte)
|
||||||
with engine.begin() as conn:
|
with engine.begin() as conn:
|
||||||
vorhanden = conn.execute(
|
vorhanden = conn.execute(
|
||||||
@@ -282,3 +351,88 @@ def save_settings(werte: dict, key: str = "ui") -> None:
|
|||||||
)
|
)
|
||||||
else:
|
else:
|
||||||
conn.execute(settings_table.insert().values(key=key, value=payload))
|
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))
|
||||||
Reference in New Issue
Block a user