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
@@ -1,284 +0,0 @@
|
||||
"""Job-, Log- und Settings-Persistenz in PostgreSQL (KONZEPT: Postgres für Job-Logs).
|
||||
|
||||
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.
|
||||
|
||||
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.
|
||||
"""
|
||||
|
||||
import json
|
||||
import os
|
||||
from datetime import datetime, timezone
|
||||
|
||||
from sqlalchemy import (
|
||||
Column,
|
||||
DateTime,
|
||||
Integer,
|
||||
MetaData,
|
||||
String,
|
||||
Table,
|
||||
Text,
|
||||
create_engine,
|
||||
select,
|
||||
)
|
||||
|
||||
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: Jahr/Poster/Plot) für Detail-Popup + 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),
|
||||
)
|
||||
|
||||
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)),
|
||||
)
|
||||
|
||||
storage_mounts = Table(
|
||||
"storage_mounts",
|
||||
metadata,
|
||||
Column("name", String(64), primary_key=True),
|
||||
Column("typ", String(8)), # nfs | cifs
|
||||
Column("quelle", String(255)), # host:/export bzw. //host/share
|
||||
Column("optionen", String(255)),
|
||||
Column("username", String(128)),
|
||||
Column("passwort", String(255)), # Klartext — Heimnetz-Kompromiss, siehe README
|
||||
)
|
||||
|
||||
|
||||
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."""
|
||||
metadata.create_all(engine)
|
||||
# create_all ändert BESTEHENDE Tabellen nicht — neue Spalten hier nachziehen:
|
||||
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"
|
||||
)
|
||||
|
||||
|
||||
def insert_job(
|
||||
job_id: str, device: str, disc_type: str = None, title: str = None,
|
||||
target_dir: str = None, meta: str = None,
|
||||
) -> None:
|
||||
with engine.begin() as conn:
|
||||
conn.execute(
|
||||
jobs.insert().values(
|
||||
id=job_id,
|
||||
device=device,
|
||||
disc_type=disc_type,
|
||||
title=title,
|
||||
target_dir=target_dir,
|
||||
meta=meta,
|
||||
status="pending",
|
||||
progress=0,
|
||||
created_at=utcnow(),
|
||||
)
|
||||
)
|
||||
|
||||
|
||||
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 get_job(job_id: str) -> dict:
|
||||
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 list_jobs(limit: int = 100) -> list:
|
||||
with engine.connect() as conn:
|
||||
zeilen = conn.execute(
|
||||
select(jobs).order_by(jobs.c.created_at.desc()).limit(limit)
|
||||
).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:
|
||||
conn.execute(jobs.delete().where(jobs.c.id == job_id))
|
||||
|
||||
|
||||
def delete_finished_jobs() -> int:
|
||||
"""Räumt alle erledigten Jobs (completed/failed) aus der Liste. Dateien bleiben."""
|
||||
with engine.begin() as conn:
|
||||
ergebnis = conn.execute(
|
||||
jobs.delete().where(jobs.c.status.in_(("completed", "failed")))
|
||||
)
|
||||
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).
|
||||
|
||||
Gebraucht von der Mount-Wache: Eine Freigabe neu zu verbinden bedeutet ein
|
||||
`umount -l` — mitten in einem laufenden Rip oder Encode wäre das ein
|
||||
Datenverlust. `pending` zählt bewusst mit: so ein Job kann jeden Moment
|
||||
anlaufen.
|
||||
"""
|
||||
with engine.connect() as conn:
|
||||
zeile = conn.execute(
|
||||
select(jobs.c.id)
|
||||
.where(jobs.c.status.in_(
|
||||
("pending", "running", "ripping", "transcoding", "canceling")))
|
||||
.limit(1)
|
||||
).first()
|
||||
return zeile is not None
|
||||
|
||||
|
||||
def has_active_job(device: str) -> bool:
|
||||
"""True, wenn auf dem Gerät ein Job läuft oder wartet (Eject-Schutz)."""
|
||||
with engine.connect() as conn:
|
||||
zeile = conn.execute(
|
||||
select(jobs.c.id)
|
||||
.where(jobs.c.device == device)
|
||||
.where(jobs.c.status.in_(("pending", "running")))
|
||||
.limit(1)
|
||||
).first()
|
||||
return zeile is not None
|
||||
|
||||
|
||||
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)
|
||||
)
|
||||
|
||||
|
||||
def list_logs(limit: int = 200) -> list:
|
||||
with engine.connect() as conn:
|
||||
zeilen = conn.execute(
|
||||
select(logs).order_by(logs.c.id.desc()).limit(limit)
|
||||
).mappings().all()
|
||||
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()
|
||||
if not zeile or not zeile[0]:
|
||||
return {}
|
||||
try:
|
||||
return json.loads(zeile[0])
|
||||
except ValueError:
|
||||
return {}
|
||||
|
||||
|
||||
def save_settings(werte: dict, key: str = "ui") -> None:
|
||||
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))
|
||||
@@ -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 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
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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}$")
|
||||
|
||||
Reference in New Issue
Block a user