549727f648
Ampel / ampel (push) Successful in 28s
Dreimal in Folge reproduziert, jetzt bei JEDEM Deploy: Nach `docker compose up -d --build` ist die CIFS-Freigabe tot. `mount` meldet Rueckgabewert 0, /proc/mounts zeigt genau eine korrekt aussehende Schicht, die Erreichbarkeits-Probe antwortet direkt nach dem Mount sogar - und Sekunden spaeter laeuft jeder Zugriff in die Zeitgrenze. Derselbe Ablauf ein bis zwei Minuten spaeter stellt sie zuverlaessig her (POST /storage-mounts/rippy/repair, mehrfach belegt). Die Ursache liegt am NAS und ist nicht gefunden. Aber die Wirkung ist teuer: Nach jedem Update war jeder Rip auf die NAS kaputt, ohne dass irgendwo etwas davon zu sehen war - und der Commander haette es jedes Mal von Hand richten muessen. Wenn die Heilung bekannt und billig ist, gehoert sie automatisiert, auch ohne die Ursache zu kennen. Die Wache sieht jede Minute nach und verbindet stumme Freigaben neu. Zwei Dinge sind dabei wichtiger als die Heilung selbst: 1. NIE waehrend ein Job laeuft. Neu verbinden heisst `umount -l`; mitten in einem Rip oder Encode waere das ein Datenverlust. Die Wache steht still, solange irgendein Job nicht durch ist - auch bei einem wartenden, der jeden Moment anlaufen kann (db.hat_arbeit). 2. Gemeldet wird nur der UEBERGANG. Ist das NAS ausgeschaltet, waere ein Log je Minute ein Wasserfall. Der Zustand steht in /health/vorraete, damit man sehen kann, dass die Wache lebt - dieselbe Lehre wie heute Nachmittag beim stillen except. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
285 lines
8.6 KiB
Python
285 lines
8.6 KiB
Python
"""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))
|