feat(core): V2-2 — SQLite hinter demselben Port, Auftrags-Queue mit Lease
Ampel / ampel (push) Successful in 39s

WAS: Der Store spricht jetzt beide Datenbanken (PostgreSQL wie bisher,
SQLite fuer den Standalone-Betrieb). Dazu die LocalQueue: Auftraege als
Tabelle, Vergabe per bedingtem UPDATE, Lease statt Zombie-Jagd. Und die
Konfigurationsschicht mit EINER Praezedenz fuer alle drei Betriebsarten.

WARUM: Ohne SQLite und ohne broker-lose Queue gibt es keinen Standalone-
Betrieb — und ohne den keine Windows-App und keine Headless-Variante.
Beides haengt an dieser Etappe.

DER GRUNDSATZ (KONZEPT-V2.md §3.1): Die Datenbank ist die Wahrheit ueber
den Job-Zustand, der Broker ist nur der Wecker. Daraus folgt die ganze
Queue: Ein Auftrag ist eine Zeile mit claimed_by und lease_until, ein
Knoten uebernimmt ihn per bedingtem UPDATE (es gewinnt genau EINER, auch
wenn zehn gleichzeitig fragen), und er haelt ihn per Lease am Leben.
Laeuft die Lease ab, ist der Auftrag frei — egal ob Absturz, Netzausfall
oder gezogener Stecker.

Damit gibt es den Zustand "laeuft, aber niemand arbeitet daran" nicht
mehr, den v1 mit zombies.py (206 Zeilen + 263 Zeilen Tests) nachtraeglich
einsammeln musste. Er kann hoechstens 60 Sekunden bestehen und heilt sich
dann selbst. Beides ist geprueft: zwei Knoten bekommen NICHT denselben
Auftrag, und eine abgelaufene Lease gibt ihn wirklich wieder her.

KEINE QUEUE-BIBLIOTHEK (Taskiq/ARQ/Celery-lite), Begruendung im
Modul-Docstring: Der teure Teil ist kein Task, sondern ein 30-90-Minuten-
Subprozess mit Fortschritts-Parsing. Was die Bibliotheken loesen, ist
nicht das Problem; was das Problem ist, muss man ohnehin selbst bauen.

ABWEICHUNG VOM ENTWURF — KEIN ALEMBIC (in KONZEPT-V2.md §3.2 vermerkt):
Die Datenbank auf der VM gibt es schon, Alembic muesste sie erst
stempeln — ein Handgriff auf einer laufenden Installation, den beim
naechsten Deploy jemand vergisst. Und es gibt hier nichts zu
versionieren: drei nachgetragene Spalten. store.migrieren() fragt
stattdessen per inspect() nach und legt nur Fehlendes an; das laeuft auf
beiden Dialekten. Der alte Weg (ADD COLUMN IF NOT EXISTS) war korrekte
Postgres-Syntax und waere auf SQLite gebrochen.

EIN KONSTRUKTIONSFEHLER, SELBST GEFUNDEN: lokal.py hatte anfangs
`from rippy.store import engine` — ein Import bindet den Wert EINMAL, ein
spaeteres store.verbinden() waere nie angekommen, und das Modul haette
weiter mit der alten Datenbank gesprochen. Jetzt wird store.engine zur
Aufrufzeit gelesen; die Warnung steht an beiden Stellen im Code.

GEMESSEN: ruff sauber, 346 Tests gruen + 3 uebersprungen (vorher 301).
Neu: 30 Tests gegen ECHTES SQLite (kein Fake) — inklusive der Gegenprobe,
dass WAL und busy_timeout wirklich gesetzt sind, und dass migrieren()
eine weggenommene Spalte zurueckholt.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
Hitonabi
2026-08-28 09:03:46 +02:00
co-authored by Claude Opus 5
parent 7558be4854
commit d2645f87e8
9 changed files with 1346 additions and 83 deletions
+15 -1
View File
@@ -305,12 +305,26 @@ Was wirklich angefasst werden muss:
| Thema | v1 | v2 |
|-------|-----|-----|
| Migrationen | `ALTER TABLE … ADD COLUMN IF NOT EXISTS` von Hand in `init_db()`**Postgres-only, bricht auf SQLite** | Alembic, ein Verzeichnis, beide Dialekte |
| Migrationen | `ALTER TABLE … ADD COLUMN IF NOT EXISTS` von Hand in `init_db()`**Postgres-only, bricht auf SQLite** | dialektneutral: erst `inspect()` fragen, dann nur fehlende Spalten anlegen — **kein Alembic**, siehe Kasten unten |
| Nebenläufigkeit | Postgres regelt es | SQLite: `PRAGMA journal_mode=WAL`, `busy_timeout=5000`, **ein** Schreiber-Kontext |
| Zeitstempel | `DateTime(timezone=True)` | unverändert; SQLite speichert ISO-8601 UTC |
| JSON-Felder | `Text` + `json.dumps` von Hand | unverändert (portabel), Serialisierung in den Store gezogen |
| Duplikat | `api/db.py` **und** `worker/db.py` | eine Datei |
> **Abweichung vom Entwurf (beim Bauen entschieden, 28.08.2026):** Hier stand
> ursprünglich Alembic. Zwei Dinge sprachen beim Umsetzen dagegen. Erstens
> **gibt es die Datenbank auf der VM schon** — Alembic müsste sie erst
> „stempeln" (`alembic stamp head`), sonst hält es sie für leer und versucht
> vorhandene Tabellen anzulegen; das ist ein Handgriff auf einer laufenden
> Installation und damit genau die Sorte Schritt, die beim nächsten Deploy
> jemand vergisst. Zweitens **gibt es hier nichts zu versionieren**: Die
> gesamte Migrationslast des Projekts sind drei nachgetragene Spalten.
> Stattdessen fragt `store.migrieren()` per `inspect()` nach, welche Spalten
> existieren, und legt nur die fehlenden an — das läuft auf beiden Dialekten.
> Sollte das Schema je wirklich wandern (Spalten umbenennen, Daten
> umschichten), ist Alembic die richtige Antwort, dann aber mit einer
> bewussten Stempel-Runde.
**SQLite-Grenzen, ehrlich benannt:** WAL erlaubt viele Leser und *einen*
Schreiber. Für Rippy reicht das mit großem Abstand — ein Rip schreibt etwa alle
2 s einen Fortschrittswert. Wer mehr als ~4 gleichzeitige Rip-Knoten fährt,
+235
View File
@@ -0,0 +1,235 @@
"""Konfiguration — eine Präzedenz, in allen drei Betriebsarten dieselbe.
CLI-Flag > Umgebungsvariable (RIPPY_*) > Konfigurationsdatei > Vorgabe
## Warum das ein eigenes Modul ist
v1 liest Einstellungen an drei verschiedenen Stellen und auf drei Arten: aus
`os.getenv` (verstreut über tasks.py, ripping.py, db.py), aus der
`settings`-Tabelle (`db.get_settings()`) und aus der `.env` über Compose. Wer
wissen will, woher ein Wert wirklich kommt, muss alle drei durchsuchen — und
`.env.example` warnt an zwei Stellen davor, dass ein vorbelegter Wert „gesund
aussieht" und trotzdem falsch ist.
v2 hat drei Betriebsarten, und jede hat eine andere natürliche Quelle:
Docker Umgebungsvariablen (kein Dateisystem-Zugriff nötig)
Headless /etc/rippy/rippy.toml (dort erwartet man Konfiguration)
Windows %ProgramData%\\Rippy\\rippy.toml (der Installer schreibt sie)
Ohne eine festgelegte Reihenfolge wäre nicht entscheidbar, wer gewinnt. Mit ihr
ist es eine Zeile Doku und ein Test.
## Verschachtelung in Umgebungsvariablen
`[store] pfad = ""` heißt `RIPPY_STORE__PFAD` — doppelter Unterstrich als
Trenner. Der einfache scheidet aus, weil Schlüssel selbst welche enthalten
(`poll_sekunden`).
## Was hier NICHT hingehört
Laufzeit-Einstellungen, die der Nutzer im Browser ändert (Presets,
Sprachwünsche, Medienserver-Adresse), bleiben in der `settings`-Tabelle. Hier
steht nur, was VOR dem Start feststehen muss: wo die Datenbank liegt, welcher
Port, welche Treiber. Faustregel: Wenn es einen Neustart braucht, gehört es
hierher; wenn nicht, in die Datenbank.
"""
import os
from typing import Any
try: # Python 3.11+
import tomllib
except ModuleNotFoundError: # 3.10 — nur der Entwicklungsrechner
tomllib = None
PRAEFIX = "RIPPY_"
TRENNER = "__"
# Die Vorgaben. Sie sind zugleich die Liste dessen, was es überhaupt gibt —
# ein Schlüssel, der hier fehlt, wird aus Datei und Umgebung NICHT übernommen
# (siehe `_zusammenfuehren`). Das fängt Tippfehler ab, statt sie stillschweigend
# zu ignorieren: Ein `RIPPY_STORE__PAFD` soll auffallen, nicht wirkungslos sein.
VORGABEN: dict = {
"profil": "standalone", # standalone | api | node
"server": {
"host": "0.0.0.0",
"port": 7788,
"ui": True, # statische UI-Assets mitausliefern
},
"store": {
"treiber": "sqlite", # sqlite | postgres
"pfad": "", # bei sqlite; leer = Vorgabe je Plattform
"url": "", # bei postgres
},
"queue": {
"treiber": "lokal", # lokal | celery
"broker": "", # bei celery
"rip_slots": 1,
"encode_slots": 2,
"knoten": "", # Anzeigename; leer = Rechnername
},
"drives": {
"erkennung": "poll", # udev | poll
"poll_sekunden": 3,
"auto_auswurf": True,
},
"storage": {
"medien": "",
"temp": "",
"pfad_timeout_sekunden": 5,
},
"metadata": {
"tmdb_key": "",
"tvdb_key": "",
"omdb_key": "",
"sprachen": ["de", "en"],
},
"transcode": {
"engine": "handbrake", # handbrake | ffmpeg
"encoder": "auto",
},
"keys": {
# Bezugsadresse für Stufe 2 der Schlüsselkette (KONZEPT.md § 10,
# Commander-Entscheid 28.08.2026). BEWUSST LEER: Eine vorbelegte
# Adresse, die irgendwann tot ist, lässt die Konfiguration gesund
# aussehen und scheitert erst im Betrieb — genau die Falle, die in
# .env.example beim MakeMKV-Download dokumentiert ist. Ohne Eintrag
# ist Stufe 2 schlicht übersprungen.
"quelle_url": "",
},
}
def _als_typ(wert: str, vorbild: Any) -> Any:
"""String aus der Umgebung auf den Typ der Vorgabe bringen.
Ohne das wäre `RIPPY_SERVER__PORT=8000` der String "8000", und ein
Vergleich oder eine Bindung an den Port schlüge fehl — mit einer Meldung,
die nach allem Möglichen aussieht, nur nicht nach einem Typfehler.
"""
if isinstance(vorbild, bool):
return wert.strip().lower() in ("1", "true", "ja", "yes", "on")
if isinstance(vorbild, int):
return int(wert)
if isinstance(vorbild, list):
return [teil.strip() for teil in wert.split(",") if teil.strip()]
return wert
def aus_umgebung(umgebung: dict = None) -> dict:
"""Baut aus RIPPY_*-Variablen einen verschachtelten dict.
Unbekannte Schlüssel werden ÜBERGANGEN, nicht übernommen — sie tauchen in
`unbekannte_schluessel()` auf, damit ein Tippfehler sichtbar wird statt
wirkungslos zu bleiben.
"""
umgebung = os.environ if umgebung is None else umgebung
ergebnis: dict = {}
for name, wert in umgebung.items():
if not name.startswith(PRAEFIX):
continue
pfad = name[len(PRAEFIX):].lower().split(TRENNER)
vorbild = VORGABEN
for teil in pfad:
if not isinstance(vorbild, dict) or teil not in vorbild:
vorbild = None
break
vorbild = vorbild[teil]
if vorbild is None or isinstance(vorbild, dict):
continue # unbekannt oder kein Blattwert
ziel = ergebnis
for teil in pfad[:-1]:
ziel = ziel.setdefault(teil, {})
ziel[pfad[-1]] = _als_typ(wert, vorbild)
return ergebnis
def aus_datei(pfad: str) -> dict:
"""Liest eine TOML-Datei. Fehlt sie, ist das kein Fehler — leer zurück.
Eine KAPUTTE Datei ist dagegen sehr wohl ein Fehler und fliegt hoch: Wer
seine Konfiguration verschreibt, soll das beim Start erfahren und nicht
stundenlang rätseln, warum eine Einstellung nicht greift.
"""
if not pfad or not os.path.isfile(pfad):
return {}
if tomllib is None:
raise RuntimeError(
"TOML-Konfiguration braucht Python 3.11 oder neuer "
"(auf 3.10 fehlt tomllib). Bis dahin die RIPPY_*-Variablen nutzen."
)
with open(pfad, "rb") as f:
return tomllib.load(f)
def _zusammenfuehren(basis: dict, oben: dict, vorbild: dict = None) -> dict:
"""Legt `oben` über `basis` — rekursiv, und nur bekannte Schlüssel."""
vorbild = VORGABEN if vorbild is None else vorbild
ergebnis = dict(basis)
for schluessel, wert in (oben or {}).items():
if schluessel not in vorbild:
continue
if isinstance(vorbild[schluessel], dict) and isinstance(wert, dict):
ergebnis[schluessel] = _zusammenfuehren(
ergebnis.get(schluessel, {}), wert, vorbild[schluessel])
else:
ergebnis[schluessel] = wert
return ergebnis
def unbekannte_schluessel(quelle: dict, vorbild: dict = None, pfad: str = "") -> list:
"""Alles, was `_zusammenfuehren` verworfen hätte — für `rippy doctor`.
Ein stillschweigend ignorierter Schlüssel ist dieselbe Fehlerklasse wie ein
verschluckter Fehler: Die Konfiguration sieht gesund aus, wirkt aber nicht.
"""
vorbild = VORGABEN if vorbild is None else vorbild
gefunden = []
for schluessel, wert in (quelle or {}).items():
voll = f"{pfad}.{schluessel}" if pfad else schluessel
if schluessel not in vorbild:
gefunden.append(voll)
elif isinstance(vorbild[schluessel], dict) and isinstance(wert, dict):
gefunden.extend(unbekannte_schluessel(wert, vorbild[schluessel], voll))
return gefunden
def laden(datei: str = None, flags: dict = None, umgebung: dict = None) -> dict:
"""Die vollständige Konfiguration nach der Präzedenz oben."""
werte = _zusammenfuehren(VORGABEN, aus_datei(datei) if datei else {})
werte = _zusammenfuehren(werte, aus_umgebung(umgebung))
werte = _zusammenfuehren(werte, flags or {})
return werte
def datenbank_url(werte: dict, standard_pfad: str = None) -> str:
"""Übersetzt den store-Abschnitt in eine SQLAlchemy-URL.
Für `store.verbinden()` — damit die Entscheidung „welche Datenbank" an
genau EINER Stelle getroffen wird und nicht an jeder Aufrufstelle neu.
"""
store = werte.get("store", {})
if store.get("treiber") == "postgres":
url = (store.get("url") or "").strip()
if not url:
raise ValueError(
"store.treiber ist 'postgres', aber store.url ist leer. "
"Ohne Adresse kann Rippy die Datenbank nicht finden."
)
return url
pfad = (store.get("pfad") or "").strip() or standard_pfad or standard_datenbankpfad()
return "sqlite:///" + pfad.replace("\\", "/")
def standard_datenbankpfad() -> str:
"""Wo die SQLite-Datei liegt, wenn niemand etwas anderes sagt.
Je Plattform dort, wo Dienste ihre Daten ablegen — nicht neben dem
Programm. Ein Programmordner ist unter Windows für einen Dienst nicht
zuverlässig beschreibbar, und unter Linux widerspräche es dem FHS.
"""
if os.name == "nt":
basis = os.environ.get("PROGRAMDATA") or os.path.expanduser("~")
return os.path.join(basis, "Rippy", "rippy.db")
return "/var/lib/rippy/rippy.db"
+10
View File
@@ -0,0 +1,10 @@
"""Queue-Schicht: Aufträge verteilen (KONZEPT-V2.md § 3.4).
Zwei Treiber, ein Port (`rippy.ports.Queue`):
lokal.py Auftragstabelle + Lease, ohne Broker — Standalone-Betrieb
celery.py Celery über Redis — verteilter Betrieb (V2-5)
Beide stehen auf demselben Grundsatz: Die Datenbank ist die Wahrheit über den
Job-Zustand, der Broker ist nur der Wecker.
"""
+290
View File
@@ -0,0 +1,290 @@
"""LocalQueue — Aufträge ohne Broker. Erfüllt `rippy.ports.Queue`.
## Der Grundsatz, aus dem alles folgt (KONZEPT-V2.md § 3.1)
Die Datenbank ist die Wahrheit über den Job-Zustand.
Der Broker ist nur der Wecker.
In v1 war es umgekehrt gedacht: Celery hielt den Auftrag, die Datenbank
spiegelte ihn nach. Deshalb gibt es dort `zombies.py` (206 Zeilen) plus
`test_zombies.py` (263 Zeilen) einen nachträglichen Reparaturmechanismus für
Jobs, die in der Datenbank laufen, während in Celery niemand mehr an ihnen
arbeitet. Nach einer Gnadenfrist werden sie ehrlich auf failed gesetzt".
Wenn die Datenbank die Wahrheit ist, verschwindet dieses Problem, statt
repariert zu werden:
* Ein Auftrag ist eine Zeile mit `status`, `claimed_by`, `lease_until`.
* Ein Knoten übernimmt ihn per **bedingtem UPDATE** atomar, es gewinnt genau
einer, auch wenn zehn gleichzeitig fragen.
* Er hält ihn per Lease am Leben: alle `LEASE_ERNEUERN_SEKUNDEN` wird
`lease_until` vorgerückt.
* Läuft die Lease ab, ist der Auftrag frei. Egal ob der Knoten abgestürzt ist,
das Netz weg war oder jemand den Stecker gezogen hat.
Es gibt also keinen Zustand läuft, aber niemand arbeitet daran" mehr, den
jemand erkennen müsste. Er kann höchstens `LEASE_DAUER_SEKUNDEN` lang bestehen
und heilt sich dann selbst.
## Warum keine Queue-Bibliothek (Taskiq, ARQ, Celery-lite)
Drei Gründe, alle am konkreten Fall gemessen:
1. **Der teure Teil ist kein Task, sondern ein Subprozess.** Ein Rip ist ein
3090-Minuten-`makemkvcon`-Aufruf, dessen stdout zeilenweise geparst wird
(`ripping.get_progress_from_prgv`). Was die Bibliotheken liefern
Serialisierung, Retry, Broker-Anbindung löst davon nichts. Was gebraucht
wird Fortschritts-Streaming, Abbruch mitten im Lauf, Lease muss man
ohnehin selbst bauen.
2. **Zwei Bibliotheken heißen zwei Programmiermodelle.** Der verteilte Modus
behält Celery (funktioniert, der Remote-Windows-Worker ist bewiesen). Eine
zweite Queue-Bibliothek daneben bedeutet zwei Fehlerbilder, zwei
Retry-Semantiken, zwei Sorten Doku.
3. **Mit dem Grundsatz oben ist der Treiber trivial** die eigentliche
Auswahl unten ist ein SELECT und ein bedingtes UPDATE.
## Wozu die Zeitgrenzen
`LEASE_DAUER_SEKUNDEN` muss deutlich größer sein als
`LEASE_ERNEUERN_SEKUNDEN`, sonst gilt ein arbeitender Knoten bei jedem
Aussetzer als tot und ein zweiter greift sich denselben Auftrag. Faktor 4 ist
reichlich: Drei Erneuerungen dürfen ausfallen, bevor jemand übernimmt.
"""
import json
import uuid
from datetime import timedelta
from sqlalchemy import or_
from rippy import store
from rippy.store import auftraege
# ⚠️ NICHT `from rippy.store import engine`. Ein Import bindet den Wert EINMAL,
# beim Laden des Moduls — ein späteres `store.verbinden(...)` (Standalone-Start,
# Tests) käme hier nie an, und dieses Modul spräche weiter mit der alten
# Datenbank. Über `store.engine` wird bei jedem Aufruf nachgesehen.
# Wie lange ein übernommener Auftrag als „in Arbeit" gilt, ohne Lebenszeichen.
LEASE_DAUER_SEKUNDEN = 60
# In welchem Abstand ein arbeitender Knoten die Lease verlängert.
LEASE_ERNEUERN_SEKUNDEN = 15
# Zustände
WARTEND = "wartend"
LAUFEND = "laufend"
FERTIG = "fertig"
FEHLER = "fehler"
ABGEBROCHEN = "abgebrochen"
# Wie oft ein Auftrag höchstens neu vergeben wird, bevor er als Fehler gilt.
# Ohne Obergrenze würde ein Auftrag, der den bearbeitenden Knoten jedes Mal
# umbringt, endlos im Kreis laufen und dabei jeden freien Knoten mitnehmen.
MAX_VERSUCHE = 3
def _jetzt():
return store.utcnow()
def einreihen(job_id: str, art: str, faehigkeiten: dict = None,
payload: dict = None, prioritaet: int = 0) -> str:
"""Legt einen Auftrag an und gibt seine Kennung zurück."""
auftrag_id = str(uuid.uuid4())
with store.engine.begin() as conn:
conn.execute(auftraege.insert().values(
id=auftrag_id,
job_id=job_id,
art=art,
status=WARTEND,
faehigkeiten=json.dumps(faehigkeiten or {}),
payload=json.dumps(payload or {}),
prioritaet=prioritaet,
versuche=0,
created_at=_jetzt(),
))
return auftrag_id
def _passt(auftrag_faehigkeiten: str, knoten_kann: set) -> bool:
"""Kann dieser Knoten den Auftrag bedienen? (pure Funktion)
Ein Auftrag verlangt Fähigkeiten als Schlüssel-Wert-Paare, z. B.
`{"art": "rip", "drive": "sr0"}`. Der Knoten meldet, was er kann, als
Menge von `"schluessel=wert"`-Strings bzw. blanken Schlüsseln.
Leere Anforderung heißt: jeder darf. Das ist Absicht ein Auftrag ohne
besondere Ansprüche soll nicht liegen bleiben, nur weil niemand explizit
kann nichts Besonderes" gemeldet hat.
"""
try:
verlangt = json.loads(auftrag_faehigkeiten or "{}")
except (ValueError, TypeError):
verlangt = {}
if not verlangt:
return True
for schluessel, wert in verlangt.items():
if f"{schluessel}={wert}" in knoten_kann or schluessel in knoten_kann:
continue
return False
return True
def uebernehmen(knoten: str, kann: set = None):
"""Holt EINEN passenden Auftrag und markiert ihn als übernommen.
Gibt den Auftrag als dict zurück oder None, wenn nichts passt.
## Warum das UPDATE eine Bedingung trägt
Zwischen dem SELECT und dem UPDATE kann ein anderer Knoten denselben
Auftrag genommen haben. Das `WHERE AND claimed_by IS NULL` (bzw. eine
abgelaufene Lease) sorgt dafür, dass genau EINER gewinnt: Der Verlierer
ändert null Zeilen und sucht weiter. Das ist derselbe Mechanismus auf
SQLite wie auf PostgreSQL beide führen ein UPDATE atomar aus.
Ohne die Bedingung hätte man das klassische Doppel-Rip: zwei Knoten am
selben Laufwerk, und der zweite überschreibt die Ausgabe des ersten.
"""
kann = kann or set()
jetzt = _jetzt()
with store.engine.begin() as conn:
kandidaten = conn.execute(
auftraege.select()
.where(auftraege.c.status.in_((WARTEND, LAUFEND)))
.where(or_(
auftraege.c.claimed_by.is_(None),
auftraege.c.lease_until.is_(None),
auftraege.c.lease_until < jetzt, # Lease abgelaufen = frei
))
.where(auftraege.c.versuche < MAX_VERSUCHE)
.order_by(auftraege.c.prioritaet.desc(), auftraege.c.created_at)
.limit(20)
).mappings().all()
for kandidat in kandidaten:
if not _passt(kandidat["faehigkeiten"], kann):
continue
betroffen = conn.execute(
auftraege.update()
.where(auftraege.c.id == kandidat["id"])
.where(or_(
auftraege.c.claimed_by.is_(None),
auftraege.c.lease_until.is_(None),
auftraege.c.lease_until < jetzt,
))
.values(
status=LAUFEND,
claimed_by=knoten,
lease_until=jetzt + timedelta(seconds=LEASE_DAUER_SEKUNDEN),
versuche=(kandidat["versuche"] or 0) + 1,
)
).rowcount
if betroffen == 1:
eintrag = dict(kandidat)
eintrag.update({
"status": LAUFEND,
"claimed_by": knoten,
"versuche": (kandidat["versuche"] or 0) + 1,
})
return eintrag
return None
def lebenszeichen(auftrag_id: str, knoten: str) -> bool:
"""Verlängert die Lease. False, wenn der Auftrag inzwischen weg ist.
Der `knoten`-Vergleich ist kein Beiwerk: Hat ein anderer den Auftrag
übernommen (weil diese Lease einmal abgelaufen war), darf der alte
Bearbeiter ihn NICHT wieder an sich ziehen. Ein False heißt für den
Aufrufer: aufhören, jemand anderes macht das jetzt.
"""
with store.engine.begin() as conn:
betroffen = conn.execute(
auftraege.update()
.where(auftraege.c.id == auftrag_id)
.where(auftraege.c.claimed_by == knoten)
.where(auftraege.c.status == LAUFEND)
.values(lease_until=_jetzt() + timedelta(seconds=LEASE_DAUER_SEKUNDEN))
).rowcount
return betroffen == 1
def abschliessen(auftrag_id: str, ergebnis: dict = None) -> None:
with store.engine.begin() as conn:
conn.execute(
auftraege.update()
.where(auftraege.c.id == auftrag_id)
.values(status=FERTIG, finished_at=_jetzt(),
payload=json.dumps(ergebnis or {}), lease_until=None)
)
def fehlgeschlagen(auftrag_id: str, fehler: str, erneut: bool = True) -> None:
"""Auftrag ist gescheitert.
`erneut=True` gibt ihn wieder frei, solange die Versuche reichen ein
abgestürzter Encoder soll den Job nicht endgültig verlieren. Ist
MAX_VERSUCHE erreicht, bleibt er auf `fehler` stehen: Ein Auftrag, der
dreimal denselben Knoten umgebracht hat, wird beim vierten Mal nicht
plötzlich funktionieren.
"""
with store.engine.begin() as conn:
zeile = conn.execute(
auftraege.select().where(auftraege.c.id == auftrag_id)
).mappings().first()
if not zeile:
return
wieder_frei = erneut and (zeile["versuche"] or 0) < MAX_VERSUCHE
conn.execute(
auftraege.update()
.where(auftraege.c.id == auftrag_id)
.values(
status=WARTEND if wieder_frei else FEHLER,
claimed_by=None,
lease_until=None,
fehler=(fehler or "")[:2000],
finished_at=None if wieder_frei else _jetzt(),
)
)
def abbrechen(auftrag_id: str) -> None:
with store.engine.begin() as conn:
conn.execute(
auftraege.update()
.where(auftraege.c.id == auftrag_id)
.values(status=ABGEBROCHEN, claimed_by=None, lease_until=None,
finished_at=_jetzt())
)
def offene_auftraege(job_id: str = None) -> list:
"""Alles, was noch nicht durch ist — für `rippy queue` und das UI."""
abfrage = auftraege.select().where(
auftraege.c.status.in_((WARTEND, LAUFEND))
).order_by(auftraege.c.prioritaet.desc(), auftraege.c.created_at)
if job_id:
abfrage = abfrage.where(auftraege.c.job_id == job_id)
with store.engine.connect() as conn:
return [dict(z) for z in conn.execute(abfrage).mappings().all()]
def verwaiste_freigeben() -> int:
"""Gibt Aufträge frei, deren Lease abgelaufen ist. Gibt die Anzahl zurück.
Braucht man streng genommen nicht `uebernehmen` behandelt abgelaufene
Leases ohnehin als frei. Diese Funktion macht es nur SICHTBAR: Nach einem
Absturz soll im UI und im Log stehen, dass etwas zurückgefallen ist, statt
dass es stillschweigend wieder auftaucht.
"""
jetzt = _jetzt()
with store.engine.begin() as conn:
betroffen = conn.execute(
auftraege.update()
.where(auftraege.c.status == LAUFEND)
.where(auftraege.c.lease_until.isnot(None))
.where(auftraege.c.lease_until < jetzt)
.values(status=WARTEND, claimed_by=None, lease_until=None)
).rowcount
return betroffen or 0
+215
View File
@@ -0,0 +1,215 @@
"""LocalQueue: übernehmen, Lease, Rückfall — gegen echtes SQLite (V2-2).
Der Grundsatz, den diese Tests prüfen: **Die Datenbank ist die Wahrheit, der
Broker ist nur der Wecker.** Konkret heißt das, dass zwei Dinge stimmen müssen,
und beide sind hier nachgestellt:
1. Zwei Knoten dürfen sich denselben Auftrag NICHT teilen. Ohne diese Zusage
liefen zwei Rips auf demselben Laufwerk, und der zweite überschriebe die
Ausgabe des ersten.
2. Ein abgestürzter Knoten muss seinen Auftrag von SELBST wieder hergeben.
Genau das konnte v1 nicht dort brauchte es `zombies.py` (206 Zeilen plus
263 Zeilen Tests), um hinterher aufzuräumen.
"""
from datetime import timedelta
import pytest
from rippy import store
from rippy.queue import lokal
@pytest.fixture
def q(tmp_path):
"""Frische SQLite-Datenbank je Test, danach zurückdrehen."""
vorher_url, vorher_engine = store.DATABASE_URL, store.engine
store.verbinden(f"sqlite:///{(tmp_path / 'q.db').as_posix()}")
store.init_db()
yield lokal
store.engine.dispose()
store.DATABASE_URL, store.engine = vorher_url, vorher_engine
def _lease_zuruecksetzen(auftrag_id, sekunden):
"""Schiebt die Lease künstlich in die Vergangenheit — so sieht ein
abgestürzter Knoten aus, ohne dass ein Test eine Minute warten muss."""
with store.engine.begin() as conn:
conn.execute(
store.auftraege.update()
.where(store.auftraege.c.id == auftrag_id)
.values(lease_until=store.utcnow() - timedelta(seconds=sekunden))
)
# ── Grundfall ───────────────────────────────────────────────────────────
def test_eingereihter_auftrag_wird_uebernommen(q):
q.einreihen("job1", "rip")
auftrag = q.uebernehmen("knoten-a")
assert auftrag is not None
assert auftrag["art"] == "rip"
assert auftrag["claimed_by"] == "knoten-a"
assert auftrag["status"] == lokal.LAUFEND
def test_ohne_auftrag_kommt_none(q):
assert q.uebernehmen("knoten-a") is None
def test_hoehere_prioritaet_kommt_zuerst(q):
q.einreihen("job1", "transcode", prioritaet=0)
q.einreihen("job2", "rip", prioritaet=5)
assert q.uebernehmen("knoten-a")["job_id"] == "job2"
# ── Die zentrale Zusage: genau EINER gewinnt ────────────────────────────
def test_zwei_knoten_bekommen_nicht_denselben_auftrag(q):
"""DIE Zusage der ganzen Queue. Ohne sie: zwei Rips auf einem Laufwerk."""
q.einreihen("job1", "rip")
erster = q.uebernehmen("knoten-a")
zweiter = q.uebernehmen("knoten-b")
assert erster is not None
assert zweiter is None, "Der zweite Knoten hat denselben Auftrag bekommen"
def test_zwei_auftraege_gehen_an_zwei_knoten(q):
"""Die Gegenprobe — sonst wäre die Sperre oben nur eine kaputte Queue."""
q.einreihen("job1", "rip")
q.einreihen("job2", "rip")
a = q.uebernehmen("knoten-a")
b = q.uebernehmen("knoten-b")
assert a and b and a["id"] != b["id"]
# ── Lease: ein abgestürzter Knoten gibt von selbst her ──────────────────
def test_abgelaufene_lease_macht_den_auftrag_wieder_frei(q):
"""Das ist der Ersatz für zombies.py.
In v1 blieb ein Job nach einem Worker-Absturz auf `transcoding` stehen,
und ein eigener Mechanismus musste ihn nach einer Gnadenfrist ehrlich auf
failed setzen". Hier passiert nichts dergleichen: Die Lease läuft ab, und
der Auftrag ist wieder zu haben.
"""
auftrag_id = q.einreihen("job1", "rip")
erster = q.uebernehmen("knoten-a")
assert erster is not None
assert q.uebernehmen("knoten-b") is None # solange die Lease lebt
_lease_zuruecksetzen(auftrag_id, 10) # Knoten A ist abgestürzt
zweiter = q.uebernehmen("knoten-b")
assert zweiter is not None
assert zweiter["claimed_by"] == "knoten-b"
assert zweiter["versuche"] == 2
def test_lebenszeichen_haelt_den_auftrag(q):
auftrag_id = q.einreihen("job1", "rip")
q.uebernehmen("knoten-a")
_lease_zuruecksetzen(auftrag_id, 10)
assert q.lebenszeichen(auftrag_id, "knoten-a") is True
assert q.uebernehmen("knoten-b") is None, "Lease wurde nicht verlaengert"
def test_fremder_knoten_kann_die_lease_nicht_verlaengern(q):
"""Sonst könnte ein Knoten, dessen Lease einmal abgelaufen ist, den Auftrag
dem neuen Bearbeiter wieder wegnehmen und beide arbeiteten daran."""
auftrag_id = q.einreihen("job1", "rip")
q.uebernehmen("knoten-a")
assert q.lebenszeichen(auftrag_id, "knoten-b") is False
def test_lebenszeichen_auf_unbekannten_auftrag_ist_false(q):
assert q.lebenszeichen("gibt-es-nicht", "knoten-a") is False
def test_verwaiste_freigeben_macht_den_rueckfall_sichtbar(q):
"""Nötig ist die Funktion nicht — uebernehmen() behandelt abgelaufene
Leases ohnehin als frei. Sie macht den Rückfall aber SICHTBAR, statt ihn
stillschweigend passieren zu lassen."""
auftrag_id = q.einreihen("job1", "rip")
q.uebernehmen("knoten-a")
assert q.verwaiste_freigeben() == 0
_lease_zuruecksetzen(auftrag_id, 10)
assert q.verwaiste_freigeben() == 1
assert q.offene_auftraege()[0]["status"] == lokal.WARTEND
# ── Fähigkeiten ─────────────────────────────────────────────────────────
def test_passt_ohne_anforderung_darf_jeder():
assert lokal._passt("{}", set()) is True
assert lokal._passt(None, set()) is True
def test_passt_verlangt_alle_faehigkeiten():
verlangt = '{"art": "transcode", "encoder": "nvenc"}'
assert lokal._passt(verlangt, {"art=transcode", "encoder=nvenc"}) is True
assert lokal._passt(verlangt, {"art=transcode"}) is False
def test_passt_bei_kaputtem_json_nicht_abstuerzen():
"""Ein beschädigtes Feld darf die Auftragsvergabe nicht anhalten —
sonst blockiert eine einzige kaputte Zeile die ganze Queue."""
assert lokal._passt("{kein json", {"art=rip"}) is True
def test_knoten_ohne_die_faehigkeit_bekommt_nichts(q):
q.einreihen("job1", "transcode", faehigkeiten={"encoder": "nvenc"})
assert q.uebernehmen("nur-cpu", kann={"encoder=x265"}) is None
assert q.uebernehmen("gpu-knoten", kann={"encoder=nvenc"}) is not None
def test_passender_auftrag_wird_gefunden_obwohl_ein_fremder_vorne_liegt(q):
"""Der GPU-Knoten darf nicht leer ausgehen, nur weil ein CPU-Auftrag mit
höherer Prioritaet davor steht."""
q.einreihen("job1", "transcode", faehigkeiten={"encoder": "x265"}, prioritaet=9)
q.einreihen("job2", "transcode", faehigkeiten={"encoder": "nvenc"})
auftrag = q.uebernehmen("gpu", kann={"encoder=nvenc"})
assert auftrag is not None and auftrag["job_id"] == "job2"
# ── Abschluss und Fehlschlag ────────────────────────────────────────────
def test_abgeschlossener_auftrag_kommt_nicht_zurueck(q):
auftrag_id = q.einreihen("job1", "rip")
q.uebernehmen("knoten-a")
q.abschliessen(auftrag_id, {"pfad": "/app/media/Akira"})
assert q.uebernehmen("knoten-b") is None
assert q.offene_auftraege() == []
def test_fehlschlag_gibt_den_auftrag_wieder_frei(q):
auftrag_id = q.einreihen("job1", "transcode")
q.uebernehmen("knoten-a")
q.fehlgeschlagen(auftrag_id, "Encoder abgestuerzt")
zweiter = q.uebernehmen("knoten-b")
assert zweiter is not None and zweiter["id"] == auftrag_id
def test_nach_max_versuchen_ist_schluss(q):
"""Ohne Obergrenze liefe ein Auftrag, der jeden Knoten umbringt, endlos im
Kreis und nähme dabei jeden freien Knoten mit."""
auftrag_id = q.einreihen("job1", "transcode")
for _ in range(lokal.MAX_VERSUCHE):
q.uebernehmen("knoten-a")
q.fehlgeschlagen(auftrag_id, "immer derselbe Fehler")
assert q.uebernehmen("knoten-a") is None
with store.engine.connect() as conn:
zeile = conn.execute(
store.auftraege.select().where(store.auftraege.c.id == auftrag_id)
).mappings().first()
assert zeile["status"] == lokal.FEHLER
assert "immer derselbe Fehler" in zeile["fehler"]
def test_abbrechen_nimmt_den_auftrag_aus_dem_verkehr(q):
auftrag_id = q.einreihen("job1", "rip")
q.abbrechen(auftrag_id)
assert q.uebernehmen("knoten-a") is None
def test_offene_auftraege_lassen_sich_nach_job_filtern(q):
q.einreihen("job1", "rip")
q.einreihen("job2", "rip")
assert len(q.offene_auftraege()) == 2
assert len(q.offene_auftraege(job_id="job1")) == 1
+137 -82
View File
@@ -21,12 +21,12 @@ 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)
## Beide Datenbanken seit 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.
`DATABASE_URL` entscheidet, welche: `postgresql://` im verteilten Betrieb,
`sqlite:///` im Standalone-Betrieb. Der Unterschied steckt allein in
`_engine_bauen()` (vier PRAGMAs) und in `migrieren()` der Rest dieser Datei
weiß nicht, mit wem er spricht.
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
@@ -37,105 +37,160 @@ import json
import os
from datetime import datetime, timedelta, timezone
from sqlalchemy import (
Column,
DateTime,
Integer,
MetaData,
String,
Table,
Text,
create_engine,
func,
select,
from sqlalchemy import create_engine, event, func, inspect, select
from rippy.store import schema
from rippy.store.schema import ( # noqa: F401 — Aufrufstellen kennen sie so
auftraege,
ereignisse,
jobs,
laufwerke,
logs,
metadata,
settings_table,
storage_mounts,
workers,
)
# Standard bleibt PostgreSQL: In beiden Containern ist DATABASE_URL gesetzt,
# und ein falscher Standard wäre hier gefährlicher als gar keiner. Der
# Standalone-Betrieb setzt stattdessen z. B.
# sqlite:////var/lib/rippy/rippy.db (Linux)
# sqlite:///C:/ProgramData/Rippy/rippy.db (Windows)
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)),
)
def ist_sqlite(url: str = None) -> bool:
return (url or DATABASE_URL).startswith("sqlite")
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),
)
def _engine_bauen(url: str):
"""Eine Engine, zwei Datenbanken — der Unterschied steckt nur hier.
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)),
)
## Die vier SQLite-PRAGMAs, und warum jedes einzelne nötig ist
`journal_mode=WAL` ohne WAL sperrt SCHON EIN LESER die ganze Datei.
Rippy hat aber immer mindestens zwei Zugriffe gleichzeitig: die API
beantwortet Anfragen, während der Rip alle paar Sekunden seinen
Fortschritt schreibt. Mit WAL lesen beliebig viele parallel, und ein
Schreiber stört sie nicht.
`busy_timeout=5000` SQLite gibt bei einer belegten Sperre SOFORT auf
(database is locked"), statt zu warten. Fünf Sekunden Geduld machen
aus einem harten Fehler eine kurze Verzögerung. Ohne diesen Wert wäre
ein Fortschritts-Schreiben, das mit einem Auftragswechsel
zusammenfällt, ein abgebrochener Rip.
`synchronous=NORMAL` mit WAL sicher gegen Programmabstürze (nur ein
Stromausfall im falschen Moment kann die letzten Transaktionen
kosten). `FULL` würde jeden Fortschrittswert einzeln auf die Platte
zwingen; das ist für eine Fortschrittsanzeige verschwendet.
`foreign_keys=ON` SQLite prüft Fremdschlüssel sonst NICHT, obwohl es
sie versteht. Standardmäßig aus, aus Rücksicht auf alte Programme.
PostgreSQL braucht nichts davon, nur `pool_pre_ping` gegen Verbindungen,
die der Server in der Zwischenzeit geschlossen hat.
"""
if not url.startswith("sqlite"):
return create_engine(url, pool_pre_ping=True)
# Ordner anlegen, sonst scheitert SQLite mit "unable to open database file"
# — eine Meldung, die nach einem kaputten Pfad klingt und nur "das
# Verzeichnis gibt es nicht" heißt.
pfad = url.split("///", 1)[-1].lstrip("/")
if pfad and pfad != ":memory:":
ordner = os.path.dirname(os.path.abspath(pfad))
if ordner:
os.makedirs(ordner, exist_ok=True)
eng = create_engine(url, connect_args={"timeout": 30}, future=True)
@event.listens_for(eng, "connect")
def _pragmas(dbapi_conn, _):
cur = dbapi_conn.cursor()
cur.execute("PRAGMA journal_mode=WAL")
cur.execute("PRAGMA busy_timeout=5000")
cur.execute("PRAGMA synchronous=NORMAL")
cur.execute("PRAGMA foreign_keys=ON")
cur.close()
return eng
engine = _engine_bauen(DATABASE_URL)
def verbinden(url: str) -> None:
"""Stellt den Store auf eine andere Datenbank um.
Gebraucht an zwei Stellen: beim Start im Standalone-Betrieb (dort steht die
URL in der Konfiguration, nicht in der Umgebung) und in Tests, die gegen
eine SQLite-Datei laufen.
Wer dieses Modul anderswo benutzt, muss `store.engine` zur AUFRUFZEIT
lesen ein `from rippy.store import engine` bindet den alten Wert fest und
merkt von dieser Umstellung nichts. Genau darauf achtet rippy/queue/lokal.py.
"""
global engine, DATABASE_URL
DATABASE_URL = url
engine = _engine_bauen(url)
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 utcnow() -> datetime:
return datetime.now(timezone.utc)
def migrieren() -> None:
"""Trägt fehlende Spalten nach — auf PostgreSQL wie auf SQLite.
## Warum hier kein Alembic steht (Abweichung von KONZEPT-V2.md § 3.2)
Der Entwurf sah Alembic vor. Beim Bauen sprachen zwei Dinge dagegen:
1. **Die Datenbank auf der VM gibt es schon.** Alembic müsste sie erst
stempeln" (`alembic stamp head`), sonst hält es sie für leer und
versucht, vorhandene Tabellen anzulegen. Das ist ein Handgriff auf
einer laufenden Installation genau die Sorte Schritt, die beim
nächsten Deploy jemand vergisst.
2. **Es gibt hier nichts zu versionieren.** Die gesamte Migrationslast
dieses Projekts sind drei nachgetragene Spalten. Ein Werkzeug mit
Versionsgraph, Downgrade-Pfaden und eigenem Verzeichnis wäre mehr
Maschinerie als Inhalt.
Was v1 hier stehen hatte, war `ALTER TABLE ADD COLUMN IF NOT EXISTS`
korrekt, aber **Postgres-Syntax**; auf SQLite bricht das. Statt der
Syntax wird jetzt VORHER nachgesehen, welche Spalten es gibt. Das
funktioniert auf beiden Dialekten und ist obendrein ehrlicher: Der Code
sagt, was er prüft, statt es der Datenbank zu überlassen.
Sollte das Schema je wirklich wandern (Spalten umbenennen, Daten
umschichten), ist Alembic die richtige Antwort dann aber mit einer
bewussten Stempel-Runde.
"""
inspektor = inspect(engine)
vorhandene_tabellen = set(inspektor.get_table_names())
for tabelle, spalte, typ in schema.NACHZUTRAGENDE_SPALTEN:
if tabelle not in vorhandene_tabellen:
continue # create_all hat sie gerade frisch angelegt
spalten = {s["name"] for s in inspektor.get_columns(tabelle)}
if spalte in spalten:
continue
with engine.begin() as conn:
conn.exec_driver_sql(f"ALTER TABLE {tabelle} ADD COLUMN {spalte} {typ}")
def init_db() -> None:
"""Legt fehlende Tabellen an (idempotent) und zieht Mini-Migrationen nach.
"""Legt fehlende Tabellen an (idempotent) und zieht 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.
Spalten dann bereits vor. Das war schon vor der Zusammenlegung so.
"""
metadata.create_all(engine)
# create_all ändert BESTEHENDE Tabellen nicht — neue Spalten hier nachziehen.
# ⚠️ Postgres-Syntax, bricht auf SQLite. Ersatz durch Alembic in V2-2.
with engine.begin() as conn:
conn.exec_driver_sql(
"ALTER TABLE jobs ADD COLUMN IF NOT EXISTS target_dir VARCHAR(255)"
)
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"
)
migrieren()
# ───────────────────────────────────────────────────────────── Jobs ────
+154
View File
@@ -0,0 +1,154 @@
"""Tabellen — dieselben auf PostgreSQL wie auf SQLite.
## Warum das ohne Klimmzüge geht
v1 hat hier von Anfang an SQLAlchemy **Core** benutzt (`Table()`, `select()`,
`insert()`) statt ORM-Klassen. Genau dieser Code läuft auf beiden Datenbanken.
Der SQLite-Treiber für den Standalone-Betrieb ist deshalb kein Neubau, sondern
eine andere `DATABASE_URL` plus ein paar PRAGMAs (siehe `__init__.py`).
## Neu in Etappe V2-2 (28.08.2026)
Drei Tabellen kommen dazu. Sie sind noch nicht überall in Gebrauch sie sind
die Grundlage, auf der V2-3 (Ereignisse) und V2-7 (mehrere Laufwerke) stehen:
`auftraege`
Aufträge getrennt von Jobs. Ein Job (Akira rippen") besteht aus mehreren
Aufträgen: scan rip transcode ablegen. v1 hatte diese Kette implizit
in tasks.py verdrahtet und konnte deshalb nicht gezielt bei einem Schritt
neu ansetzen Neu komprimieren" war ein Sonderweg, kein Normalfall.
Die drei Spalten `claimed_by`, `lease_until` und `versuche` tragen den
Grundsatz aus KONZEPT-V2.md § 3.1: **Die Datenbank ist die Wahrheit über
den Job-Zustand, der Broker ist nur der Wecker.** Wer einen Auftrag hält,
verlängert seine Lease; läuft sie ab, ist der Auftrag frei. Das ersetzt die
nachträgliche Zombie-Erkennung durch einen Mechanismus, der den Fall gar
nicht erst entstehen lässt.
`laufwerke`
Laufwerke als erstklassige Objekte, mit STABILER Kennung (WWID/Serial)
statt `sr0`. Grund steht in .env.example: nach USB-Reconnect zur Laufzeit
kann die sg-Nummer wandern". Ein Job, der an `sr0` hängt, hängt danach am
falschen Gerät.
`ereignisse`
Ringpuffer mit monoton steigender `seq` für die Wiederaufnahme des
SSE-Stroms (V2-3). Damit holt ein UI, das an Knoten A hängt, auch
Ereignisse von Knoten B lückenlos nach.
"""
from sqlalchemy import (
Column,
DateTime,
Integer,
MetaData,
String,
Table,
Text,
)
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
)
# ── neu ab V2-2 ─────────────────────────────────────────────────────────
auftraege = Table(
"auftraege",
metadata,
Column("id", String(36), primary_key=True),
Column("job_id", String(36), nullable=False),
Column("art", String(16), nullable=False), # scan | rip | transcode | ablegen
Column("status", String(16), nullable=False), # wartend | laufend | fertig | fehler | abgebrochen
Column("faehigkeiten", Text), # JSON: {"drive":"sr0"} / {"encoder":"nvenc"}
Column("prioritaet", Integer, nullable=False, server_default="0"),
Column("claimed_by", String(128)), # Knotenname
Column("lease_until", DateTime(timezone=True)),
Column("versuche", Integer, nullable=False, server_default="0"),
Column("payload", Text), # JSON
Column("fehler", Text),
Column("created_at", DateTime(timezone=True)),
Column("finished_at", DateTime(timezone=True)),
)
laufwerke = Table(
"laufwerke",
metadata,
Column("id", String(128), primary_key=True), # stabil: WWID/Serial, NICHT sr0
Column("knoten", String(128), nullable=False),
Column("pfad", String(64), nullable=False), # /dev/sr0 bzw. D:
Column("anzeigename", String(255)),
Column("profil", Text), # JSON: Zero-Click je Disc-Typ
Column("last_seen", DateTime(timezone=True)),
)
ereignisse = Table(
"ereignisse",
metadata,
Column("seq", Integer, primary_key=True, autoincrement=True),
Column("ts", DateTime(timezone=True)),
Column("typ", String(32), nullable=False),
Column("entitaet", String(32)),
Column("entitaet_id", String(64)),
Column("daten", Text), # JSON
)
# Spalten, die `create_all` an BESTEHENDEN Tabellen nicht mehr nachträgt.
# Siehe `migrieren()` in __init__.py — bewusst dialektneutral gehalten.
NACHZUTRAGENDE_SPALTEN = [
("jobs", "target_dir", "VARCHAR(255)"),
("jobs", "meta", "TEXT"),
("workers", "info", "TEXT"),
]
+136
View File
@@ -0,0 +1,136 @@
"""Der Store auf SQLite — die Grundlage des Standalone-Betriebs (V2-2).
Diese Tests laufen gegen eine ECHTE SQLite-Datei im temporären Verzeichnis,
nicht gegen ein Fake. Das ist der Punkt: Die Behauptung von KONZEPT-V2.md § 3.2
lautet, derselbe SQLAlchemy-Core-Code laufe auf beiden Datenbanken. Eine solche
Behauptung prüft man, indem man sie ausführt.
Kein Postgres nötig deshalb laufen sie auch auf dem Windows-Entwicklungs-
rechner und nicht nur in der Ampel.
"""
import pytest
from sqlalchemy import inspect, select
from rippy import store
@pytest.fixture
def sqlite_store(tmp_path):
"""Store auf eine frische SQLite-Datei umstellen und danach zurückdrehen."""
vorher_url, vorher_engine = store.DATABASE_URL, store.engine
store.verbinden(f"sqlite:///{(tmp_path / 'rippy.db').as_posix()}")
store.init_db()
yield store
store.engine.dispose()
store.DATABASE_URL, store.engine = vorher_url, vorher_engine
def test_init_db_legt_alle_tabellen_an(sqlite_store):
"""Wenn eine Tabelle fehlt, merkt es sonst erst der Nutzer im Betrieb."""
vorhanden = set(inspect(sqlite_store.engine).get_table_names())
erwartet = {"jobs", "logs", "settings", "workers", "storage_mounts",
"auftraege", "laufwerke", "ereignisse"}
assert erwartet <= vorhanden, f"fehlen: {erwartet - vorhanden}"
def test_wal_ist_wirklich_eingeschaltet(sqlite_store):
"""WAL ist keine Kosmetik: ohne ihn sperrt schon EIN Leser die Datei.
Rippy liest und schreibt immer gleichzeitig (API beantwortet Anfragen,
während der Rip alle paar Sekunden Fortschritt schreibt). Deshalb wird das
PRAGMA hier nicht geglaubt, sondern zurückgefragt.
"""
with sqlite_store.engine.connect() as conn:
modus = conn.exec_driver_sql("PRAGMA journal_mode").scalar()
wartezeit = conn.exec_driver_sql("PRAGMA busy_timeout").scalar()
assert str(modus).lower() == "wal"
assert int(wartezeit) >= 5000
def test_job_kommt_so_zurueck_wie_er_reinging(sqlite_store):
sqlite_store.insert_job("j1", "/dev/sr0", disc_type="bluray", title="Akira")
zeile = sqlite_store.get_job("j1")
assert zeile["title"] == "Akira"
assert zeile["status"] == "pending"
assert zeile["progress"] == 0
sqlite_store.update_job("j1", status="ripping", progress=42)
assert sqlite_store.get_job_status("j1") == "ripping"
assert sqlite_store.get_job("j1")["progress"] == 42
def test_einstellungen_ueberleben_das_speichern(sqlite_store):
sqlite_store.save_settings({"transcodePreset": "HQ 1080p30 Surround"})
assert sqlite_store.get_settings()["transcodePreset"] == "HQ 1080p30 Surround"
def test_meta_merken_ergaenzt_statt_zu_ersetzen(sqlite_store):
"""Der Fehler, gegen den meta_merken gebaut wurde: update_job(meta=…)
ERSETZT die Spalte Poster, Jahr und Sprachwunsch des Rips wären fort."""
import json
sqlite_store.insert_job("j2", "/dev/sr0", meta=json.dumps({"year": 1988}))
sqlite_store.meta_merken("j2", rip_fertig=True)
meta = json.loads(sqlite_store.get_job("j2")["meta"])
assert meta == {"year": 1988, "rip_fertig": True}
def test_migrieren_traegt_eine_fehlende_spalte_nach(sqlite_store):
"""Der eigentliche Grund, warum hier kein `IF NOT EXISTS` mehr steht.
v1 benutzte `ALTER TABLE ADD COLUMN IF NOT EXISTS` korrekte
Postgres-Syntax, die auf SQLite mit einem Syntaxfehler bricht. Hier wird
eine Spalte weggenommen und nachgesehen, ob migrieren() sie zurückholt.
"""
# DROP COLUMN kann SQLite erst ab 3.35 (März 2021). Ein älterer Runner soll
# hier eine klare Auskunft geben statt eines rätselhaften OperationalError.
import sqlite3
if tuple(int(x) for x in sqlite3.sqlite_version.split(".")) < (3, 35, 0):
pytest.skip(f"SQLite {sqlite3.sqlite_version} kann kein DROP COLUMN")
with sqlite_store.engine.begin() as conn:
conn.exec_driver_sql("ALTER TABLE jobs DROP COLUMN meta")
assert "meta" not in {s["name"] for s in inspect(sqlite_store.engine).get_columns("jobs")}
sqlite_store.migrieren()
assert "meta" in {s["name"] for s in inspect(sqlite_store.engine).get_columns("jobs")}
def test_migrieren_darf_zweimal_laufen(sqlite_store):
"""init_db() läuft bei JEDEM Start von API und Worker. Wäre migrieren()
nicht wiederholbar, würde der zweite Container beim Hochfahren sterben."""
sqlite_store.migrieren()
sqlite_store.migrieren()
sqlite_store.insert_job("j3", "/dev/sr0")
assert sqlite_store.get_job("j3") is not None
def test_hat_arbeit_zaehlt_wartende_jobs_mit(sqlite_store):
"""Die Mount-Wache fragt das, bevor sie ein `umount -l` macht. Ein
`pending`-Job kann jeden Moment anlaufen er MUSS mitzählen, sonst wird
die Freigabe mitten in einem gerade startenden Rip weggezogen."""
assert sqlite_store.hat_arbeit() is False
sqlite_store.insert_job("j4", "/dev/sr0")
assert sqlite_store.hat_arbeit() is True
sqlite_store.update_job("j4", status="completed")
assert sqlite_store.hat_arbeit() is False
def test_logs_kommen_in_umgekehrter_reihenfolge(sqlite_store):
for i in range(3):
sqlite_store.add_log("info", "test", f"Zeile {i}")
zeilen = sqlite_store.list_logs(limit=10)
assert [z["message"] for z in zeilen] == ["Zeile 2", "Zeile 1", "Zeile 0"]
def test_worker_upsert_legt_nicht_zweimal_an(sqlite_store):
"""Der Herzschlag läuft jede Minute. Ohne Upsert hätte die Worker-Liste
nach einem Tag 1440 Einträge desselben Rechners."""
sqlite_store.save_worker("pc", ["x265"], {"makemkv": "1.18.4"})
sqlite_store.save_worker("pc", ["x265", "nvenc"], {"makemkv": "1.18.5"})
with sqlite_store.engine.connect() as conn:
anzahl = len(conn.execute(select(sqlite_store.workers)).all())
assert anzahl == 1
assert sqlite_store.list_workers()[0]["encoders"] == ["x265", "nvenc"]
+154
View File
@@ -0,0 +1,154 @@
"""Konfiguration: die Präzedenz muss stimmen, sonst gewinnt der falsche Wert.
CLI-Flag > Umgebungsvariable > Datei > Vorgabe
Das ist keine Geschmacksfrage: Ein Docker-Nutzer setzt Umgebungsvariablen, ein
Server-Nutzer eine TOML-Datei, und beim Debuggen will man beides mit einem Flag
übersteuern können. Wenn die Reihenfolge nicht festliegt, sucht man Fehler an
der falschen Stelle.
"""
import pytest
from rippy import config
def test_ohne_alles_gelten_die_vorgaben():
werte = config.laden(umgebung={})
assert werte["server"]["port"] == 7788
assert werte["store"]["treiber"] == "sqlite"
assert werte["queue"]["treiber"] == "lokal"
def test_umgebung_schlaegt_vorgabe():
werte = config.laden(umgebung={"RIPPY_SERVER__PORT": "9000"})
assert werte["server"]["port"] == 9000
def test_flag_schlaegt_umgebung():
werte = config.laden(
umgebung={"RIPPY_SERVER__PORT": "9000"},
flags={"server": {"port": 1234}},
)
assert werte["server"]["port"] == 1234
# TOML lesen kann erst Python 3.11 (tomllib). Beide Container laufen 3.12, die
# Ampel ebenfalls — nur der Entwicklungsrechner hat 3.10 als Standard-`python`.
# Die PRÄZEDENZ selbst (Flag > Umgebung > Datei > Vorgabe) ist oben ohne Datei
# geprüft und läuft überall; hier hängt nur das Einlesen der Datei dran.
braucht_toml = pytest.mark.skipif(
config.tomllib is None,
reason="tomllib gibt es erst ab Python 3.11 (Container und Ampel: 3.12)",
)
@braucht_toml
def test_datei_wird_von_umgebung_geschlagen(tmp_path):
datei = tmp_path / "rippy.toml"
datei.write_text('[server]\nport = 5000\nhost = "127.0.0.1"\n', encoding="utf-8")
werte = config.laden(str(datei), umgebung={"RIPPY_SERVER__PORT": "9000"})
assert werte["server"]["port"] == 9000 # Umgebung gewinnt
assert werte["server"]["host"] == "127.0.0.1" # aus der Datei, unangetastet
def test_fehlende_datei_ist_kein_fehler():
assert config.aus_datei("/gibt/es/nicht/rippy.toml") == {}
@braucht_toml
def test_kaputte_datei_fliegt_hoch(tmp_path):
"""Wer seine Konfiguration verschreibt, soll das beim START erfahren.
Ein stilles Zurückfallen auf die Vorgaben wäre die schlimmere Variante:
Rippy liefe dann mit falschen Werten und säße im falschen Verzeichnis,
ohne dass irgendwo etwas stünde.
"""
datei = tmp_path / "kaputt.toml"
datei.write_text("[server\nport = ", encoding="utf-8")
with pytest.raises(Exception):
config.aus_datei(str(datei))
# ── Typen ───────────────────────────────────────────────────────────────
def test_zahlen_kommen_als_zahlen_an():
"""Sonst wäre RIPPY_SERVER__PORT=8000 der String "8000" — und das
Binden an den Port schlüge mit einer Meldung fehl, die nach allem
aussieht, nur nicht nach einem Typfehler."""
werte = config.laden(umgebung={"RIPPY_DRIVES__POLL_SEKUNDEN": "10"})
assert werte["drives"]["poll_sekunden"] == 10
assert isinstance(werte["drives"]["poll_sekunden"], int)
def test_wahrheitswerte_verstehen_deutsch_und_englisch():
for wahr in ("1", "true", "ja", "yes", "on", "TRUE"):
assert config.laden(umgebung={"RIPPY_SERVER__UI": wahr})["server"]["ui"] is True
for falsch in ("0", "false", "nein", "no", "off"):
assert config.laden(umgebung={"RIPPY_SERVER__UI": falsch})["server"]["ui"] is False
def test_listen_kommen_kommagetrennt():
werte = config.laden(umgebung={"RIPPY_METADATA__SPRACHEN": "de, en, ja"})
assert werte["metadata"]["sprachen"] == ["de", "en", "ja"]
# ── Tippfehler dürfen nicht stillschweigend wirkungslos sein ────────────
def test_unbekannte_variable_wird_nicht_uebernommen():
werte = config.laden(umgebung={"RIPPY_STORE__PAFD": "/tmp/x.db"})
assert "pafd" not in werte["store"]
assert werte["store"]["pfad"] == ""
def test_unbekannte_schluessel_werden_gemeldet():
"""Für `rippy doctor`. Ein stillschweigend ignorierter Schlüssel ist
dieselbe Fehlerklasse wie ein verschluckter Fehler: Die Konfiguration
sieht gesund aus und wirkt nicht."""
gefunden = config.unbekannte_schluessel({
"server": {"port": 80, "prot": 81},
"quatsch": {"a": 1},
})
assert set(gefunden) == {"server.prot", "quatsch"}
def test_bekannte_schluessel_werden_nicht_gemeldet():
assert config.unbekannte_schluessel({"server": {"port": 80}}) == []
# ── Datenbank-URL ───────────────────────────────────────────────────────
def test_sqlite_url_aus_pfad():
werte = config.laden(umgebung={"RIPPY_STORE__PFAD": "/var/lib/rippy/r.db"})
assert config.datenbank_url(werte) == "sqlite:////var/lib/rippy/r.db"
def test_windows_pfad_wird_zu_schraegstrichen():
"""SQLAlchemy-URLs kennen keine Backslashes — ohne die Umschrift landete
die Datenbank unter einem Namen mit Backslashes drin."""
werte = config.laden(umgebung={"RIPPY_STORE__PFAD": r"C:\ProgramData\Rippy\r.db"})
assert config.datenbank_url(werte) == "sqlite:///C:/ProgramData/Rippy/r.db"
def test_postgres_url_wird_durchgereicht():
werte = config.laden(umgebung={
"RIPPY_STORE__TREIBER": "postgres",
"RIPPY_STORE__URL": "postgresql://rippy:geheim@db:5432/rippy",
})
assert config.datenbank_url(werte) == "postgresql://rippy:geheim@db:5432/rippy"
def test_postgres_ohne_adresse_meldet_klartext():
"""„could not translate host name" ist keine brauchbare Auskunft.
Wer treiber=postgres setzt und die URL vergisst, soll das lesen können."""
werte = config.laden(umgebung={"RIPPY_STORE__TREIBER": "postgres"})
with pytest.raises(ValueError, match="store.url ist leer"):
config.datenbank_url(werte)
def test_schluesselquelle_ist_leer_vorbelegt():
"""Commander-Entscheid 28.08.2026, Regel 1: Die Bezugsadresse steht in der
Konfiguration und ist LEER vorbelegt. Eine vorbelegte Adresse, die
irgendwann tot ist, lässt die Konfiguration gesund aussehen und scheitert
erst im Betrieb dieselbe Falle wie beim MakeMKV-Download in .env.example.
Wer hier je einen Standard einträgt, muss an diesem Test vorbei.
"""
assert config.VORGABEN["keys"]["quelle_url"] == ""
assert config.laden(umgebung={})["keys"]["quelle_url"] == ""