Files
rippy/docker/worker/celery_app.py
T
HitonabiandClaude Opus 5 dd1d0b7365
Ampel / ampel (push) Successful in 40s
refactor(core): V2-1 — die vier Ports, ein Store, ein Laufwerks-Treiber
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>
2026-08-28 08:54:28 +02:00

91 lines
2.9 KiB
Python

import os
import socket
from celery import Celery
from celery.signals import worker_ready
REDIS_URL = os.getenv("REDIS_URL", "redis://localhost:6379/0")
celery_app = Celery(
"rippy_worker",
broker=REDIS_URL,
backend=REDIS_URL,
include=["tasks"]
)
celery_app.conf.update(
task_serializer="json",
accept_content=["json"],
result_serializer="json",
timezone="UTC",
enable_utc=True,
broker_connection_retry_on_startup=True,
# worker_direct: jeder Worker konsumiert zusätzlich eine eigene
# Direkt-Queue (<nodename>.dq). Damit kann die API die Kompression
# gezielt an EINEN gewählten Worker schicken (Encoder-Auswahl im
# Rip-Dialog) — ohne Wahl läuft sie weiter über die geteilte
# transcode-Queue (irgendein freier Worker).
worker_direct=True,
)
# Import auf Modulebene: im worker_ready-Signal ist /app nicht mehr
# zuverlässig im sys.path (ModuleNotFoundError 'caps', Deploy 23.07.).
import caps # noqa: E402
from rippy import store as db # noqa: E402
import zombies # noqa: E402
@worker_ready.connect
def raeume_job_leichen_auf(**kwargs):
"""Nach der Gnadenfrist: Jobs, an denen niemand arbeitet, ehrlich auf
'failed' setzen (Details und Sicherungen in zombies.py).
Läuft im Hintergrund-Thread — der Worker soll sofort Aufgaben annehmen und
nicht zwei Minuten auf die Aufräumung warten.
"""
import threading
import time
def spaeter():
time.sleep(zombies.GNADENFRIST_SEKUNDEN)
try:
db.init_db()
bericht = zombies.raeume_zombies_auf(celery_app, db)
print(f"Zombie-Erkennung: {bericht}")
except Exception as e: # darf den Worker nie mitnehmen
print(f"Zombie-Erkennung fehlgeschlagen: {e}")
threading.Thread(target=spaeter, daemon=True, name="zombie-erkennung").start()
@worker_ready.connect
def melde_faehigkeiten(**kwargs):
"""Beim Start + minütlich: Encoder-Fähigkeiten und Lebenszeichen melden.
last_seen ist die Basis der Erreichbarkeits-Anzeige im UI — ein Worker,
der sich 2 Minuten nicht meldet, gilt als offline.
"""
import threading
import time
# Anzeigename: WORKER_NAME (compose/.env) macht die Maschine im UI
# zuordenbar UND stabil — der Container-Hostname wechselt bei jedem
# Rebuild und hinterließ Offline-Leichen in der Worker-Liste.
anzeigename = os.getenv("WORKER_NAME") or socket.gethostname()
def herzschlag():
while True:
try:
db.init_db()
db.save_worker(
anzeigename,
caps.erkenne_encoder(),
caps.werkzeug_versionen(),
)
except Exception as e: # DB weg → weiterversuchen, nicht sterben
print(f"Fähigkeiten-Meldung fehlgeschlagen: {e}")
time.sleep(60)
threading.Thread(target=herzschlag, daemon=True, name="worker-herzschlag").start()