diff --git a/docker/api/main.py b/docker/api/main.py index 2e5419e..abe6d7c 100644 --- a/docker/api/main.py +++ b/docker/api/main.py @@ -20,6 +20,7 @@ import mounts as mount_verwaltung from rippy.core import notify from rippy.bus import schema as bus_schema from rippy.bus.memory import bus as ereignis_bus +from rippy.bus.waechter import Waechter import phasen import presets as preset_auswahl import rohdaten @@ -53,6 +54,12 @@ app = FastAPI( ) +# Der Ereignis-Waechter dieses Prozesses. Im Startup gesetzt; von +# /health/vorraete abgefragt, damit ein still gestorbener Waechter +# sichtbar ist statt nur als "es passiert nichts mehr". +ereignis_waechter = None + + @app.on_event("startup") async def startup_event(): """Initialisiere Cache + Datenbank, validiere Konfiguration, starte Disc-Wache.""" @@ -64,6 +71,55 @@ async def startup_event(): except ConfigValidationError as e: print(f"⚠️ Konfigurations-Warnung: {e}") + # Ereignis-Waechter (V2-3): macht Aenderungen des Workers zu Ereignissen. + # Der Worker ist ein eigener Container — ohne diese Bruecke wuesste die API + # nichts von seinem Fortschritt, und der SSE-Strom bliebe nach dem ersten + # Snapshot stumm. Warum ueber die Datenbank statt per Rueckruf oder Redis: + # siehe den Modul-Docstring von rippy/bus/waechter.py. + global ereignis_waechter + async def system_stand(): + """Server-Zustand fuers Dashboard — im langsamen Takt des Waechters. + + Ohne das haette das Dashboard weiterhin einen eigenen 12-Sekunden-Takt + fuer /system/info, /capabilities und /storage-targets gebraucht — also + je offenem Tab 15 Anfragen pro Minute. So laeuft es EINMAL im Server, + egal wie viele Tabs offen sind. + """ + info, caps, ablagen = await asyncio.gather( + system_info(), capabilities(), storage_targets(), + return_exceptions=True, + ) + stand = {} + if not isinstance(info, Exception): + stand["info"] = info + if not isinstance(caps, Exception): + stand["workers"] = (caps or {}).get("workers") or [] + if not isinstance(ablagen, Exception): + stand["ablagen"] = ablagen + # Nur schicken, was wirklich gelesen wurde. Ein fehlender Schluessel + # heisst fuer das UI „behalte deinen Stand" — ein leerer Wert hiesse + # „es gibt nichts", und das waere wieder eine falsche Aussage. + return stand or None + + # Die Ereignisschleife HIER festhalten. `einmal()` des Waechters laeuft in + # einem Worker-Thread (asyncio.to_thread) — dort gaebe + # `asyncio.get_event_loop()` nicht diese Schleife zurueck, sondern + # erzeugte eine neue oder wuerfe. Die Coroutine liefe dann nie. + schleife = asyncio.get_running_loop() + + def system_stand_sync(): + return asyncio.run_coroutine_threadsafe(system_stand(), schleife).result(timeout=30) + + ereignis_waechter = Waechter( + db, ereignis_bus, + laufwerke_lesen=lambda: [ + device_discovery.device_info(p) + for p in device_discovery.list_optical_devices() + ], + system_lesen=system_stand_sync, + ) + asyncio.create_task(ereignis_waechter.schleife()) + # Gespeicherte Netzwerk-Speicherziele wiederherstellen — NICHT-BLOCKIEREND. # Ein zickiger/langsamer Netz-Mount (der CIFS-Schreibtest in mounten() kann im # Kernel haengen, wait_for_response) darf den API-Start NIE blockieren. Vorfall @@ -370,7 +426,7 @@ async def _snapshot() -> dict: return { "jobs": [_job_row_to_model(z).model_dump() for z in db.list_jobs(limit=50)], "workers": db.list_workers(), - "logs": db.list_logs(limit=50), + "logs": [_log_zeile(z) for z in db.list_logs(limit=50)], } zustand = await asyncio.to_thread(sammeln) @@ -403,10 +459,18 @@ def _sse_rahmen(ereignis: dict) -> str: return "\n".join(zeilen) -@app.get("/api/v2/events") +@app.get("/events") async def events(request: Request, last_event_id: str = None): """Live-Strom: erst ein Snapshot, danach nur noch Änderungen. + Pfad bewusst `/events` und nicht `/api/v2/events`: Der nginx im UI-Container + entfernt das Präfix `/api/` (siehe ui/nginx.conf), und alle bestehenden + Routen hier sind entsprechend unpräfigiert (`/jobs`, `/devices`). Aus dem + Browser heißt der Aufruf damit `/api/events`. Die Versionierung `/api/v2/*` + aus KONZEPT-V2.md § 6.1 kommt, wenn die Routen in Router aufgeteilt werden — + sie jetzt für eine einzige Route einzuführen, hätte zwei Konventionen + nebeneinander bedeutet. + Der Browser schickt beim Wiederverbinden von selbst `Last-Event-ID` mit. Passt die Lücke in den Ringpuffer, wird sie nachgeliefert; passt sie nicht, kommt ein neuer Snapshot — AUSDRÜCKLICH, nicht stillschweigend. Ein UI, das @@ -523,6 +587,17 @@ async def health_vorraete(): "stand": _MOUNT_STAND, "intervall": MOUNT_WACHE_INTERVALL_SEKUNDEN, }, + # Der Ereignis-Waechter (V2-3) ist die Bruecke zwischen Worker und + # SSE-Strom. Stirbt er still, steht das UI — und zwar OHNE Fehlermeldung, + # weil ein Abriss dort bewusst als "nichts Neues" gilt und nicht als + # "nichts da". Deshalb muss sein Alter von aussen abfragbar sein. + "ereignisse": { + "alter_sekunden": None if ereignis_waechter is None + else round(ereignis_waechter.lebt_seit_sekunden(), 1), + "gesund": bool(ereignis_waechter and ereignis_waechter.gesund), + "abonnenten": len(getattr(ereignis_bus, "_abonnenten", ())), + "letzte_seq": ereignis_bus.seq, + }, } @@ -2209,20 +2284,29 @@ async def setup_complete(): return {"done": True} +def _log_zeile(z: dict) -> dict: + """DB-Zeile -> UI-Form. EINE Abbildung, zwei Aufrufer. + + Sie stand bis V2-3 nur im /logs-Endpunkt. Der Snapshot des SSE-Stroms + haette daneben die rohen DB-Zeilen geliefert — mit `ts` statt `timestamp` + und einer Zahl statt eines Strings als id. Das UI haette „Invalid Date" + angezeigt, und zwar NUR im Live-Betrieb, nicht beim manuellen Neuladen: + genau die Sorte Fehler, die man lange sucht. + """ + return { + "id": str(z["id"]), + "timestamp": z["ts"].isoformat() if z.get("ts") else "", + "level": z.get("level") or "info", + "source": z.get("source") or "system", + "message": z.get("message") or "", + } + + @app.get("/logs") async def get_logs(limit: int = 200): """Echte Ereignisse aus der Datenbank (Watcher, API, Worker).""" zeilen = await asyncio.to_thread(db.list_logs, min(limit, 1000)) - return [ - { - "id": str(z["id"]), - "timestamp": z["ts"].isoformat() if z.get("ts") else "", - "level": z.get("level") or "info", - "source": z.get("source") or "system", - "message": z.get("message") or "", - } - for z in zeilen - ] + return [_log_zeile(z) for z in zeilen] @app.get("/settings") diff --git a/docker/api/test_api_smoke.py b/docker/api/test_api_smoke.py index 61242a1..801c4b2 100644 --- a/docker/api/test_api_smoke.py +++ b/docker/api/test_api_smoke.py @@ -257,7 +257,7 @@ def test_events_route_ist_verdrahtet(): Neues, ohne Fehlermeldung. Deshalb hier festgenagelt.""" from main import app - assert "/api/v2/events" in {route.path for route in app.routes} + assert "/events" in {route.path for route in app.routes} def test_sse_rahmen_hat_das_format_das_der_browser_erwartet(): diff --git a/docker/ui/src/App.tsx b/docker/ui/src/App.tsx index 4ce0706..f64baf8 100644 --- a/docker/ui/src/App.tsx +++ b/docker/ui/src/App.tsx @@ -6,10 +6,45 @@ import LogsPage from './pages/Logs' import AnleitungPage from './pages/Anleitung' import FirstRunWizard from './components/FirstRunWizard' import { api } from './lib/api' +import { EventStreamProvider, useStrom } from './lib/useEventStream' type Page = 'dashboard' | 'anleitung' | 'logs' | 'settings' -export default function App() { +/** + * Zeigt, ob der Live-Strom steht. + * + * Hier stand bis V2-3 ein fest verdrahtetes „ONLINE" mit pulsierendem Punkt — + * es leuchtete auch dann gruen, wenn die API tot war. Genau die Sorte + * Placebo-Anzeige, die in Etappe 19 schon einmal aufgeraeumt wurde + * („Fortschritt log, Auswurf tat nichts, Encoder wurden behauptet statt + * gemessen"). Jetzt sagt sie die Wahrheit. + */ +function VerbindungsAnzeige() { + const { verbunden, zuletzt } = useStrom() + if (verbunden) { + return ( +
+ + LIVE +
+ ) + } + return ( +
+ + VERBINDUNG WEG +
+ ) +} + +function AppInhalt() { const [currentPage, setCurrentPage] = useState('dashboard') // First-Run: solange der Einrichtungs-Assistent nicht abgeschlossen ist, zeigen wir // ihn statt des Dashboards. null = /setup noch nicht geprüft (kein Aufblitzen). @@ -69,10 +104,7 @@ export default function App() { })} -
- - ONLINE -
+ @@ -86,3 +118,14 @@ export default function App() { ) } + +export default function App() { + // Der Provider haelt GENAU EINE SSE-Verbindung fuer die ganze Anwendung. + // Ein Hook je Komponente waere eine Verbindung je Komponente — also das + // alte Problem in neuer Form. + return ( + + + + ) +} diff --git a/docker/ui/src/components/DeviceDiscovery.tsx b/docker/ui/src/components/DeviceDiscovery.tsx index b5fdc79..f3d1fa6 100644 --- a/docker/ui/src/components/DeviceDiscovery.tsx +++ b/docker/ui/src/components/DeviceDiscovery.tsx @@ -7,6 +7,7 @@ import { Button } from './ui/Button' import { TypeBadge } from './ui/Badge' import RipTargetModal, { RipOptionen } from './RipTargetModal' import MetadataKorrektur from './MetadataKorrektur' +import { useStrom } from '../lib/useEventStream' interface DiscInfo { title: string @@ -63,6 +64,7 @@ export default function DeviceDiscovery() { const [modalDevice, setModalDevice] = useState(null) const [korrekturDevice, setKorrekturDevice] = useState(null) const { toast } = useToast() + const strom = useStrom() const refreshDevices = async (manuell = false) => { if (manuell) setRefreshing(true) @@ -124,11 +126,20 @@ export default function DeviceDiscovery() { } } + // Vorher: alle 5 s ein eigener Abruf von /devices (12 Anfragen pro Minute). + // Jetzt liefert der Server die Aenderung, sobald sie passiert. + // + // WICHTIG — nur uebernehmen, wenn der Server wirklich etwas WEISS: + // `strom.devices === null` heisst „konnte nicht nachsehen" (haengendes + // Laufwerk, Zeitgrenze) und NICHT „es gibt keine Laufwerke". In dem Fall + // bleibt der letzte bekannte Stand stehen. Genau diese Vermischung hat in + // v1 die Listen im Sekundentakt geleert. useEffect(() => { - refreshDevices() - const interval = setInterval(refreshDevices, 5000) - return () => clearInterval(interval) - }, []) + if (strom.devices !== null) { + setDevices(strom.devices as Device[]) + setLoading(false) + } + }, [strom.devices]) const QUELLEN_LABEL: Record = { tmdb: 'TMDB', omdb: 'OMDb', jikan: 'MyAnimeList', manuell: 'manuell gewählt', diff --git a/docker/ui/src/components/LiveLogSection.tsx b/docker/ui/src/components/LiveLogSection.tsx index fc7a21a..5c5406d 100644 --- a/docker/ui/src/components/LiveLogSection.tsx +++ b/docker/ui/src/components/LiveLogSection.tsx @@ -1,7 +1,6 @@ -import { useState, useEffect } from 'react' import { Terminal, Play, CheckCircle, AlertCircle } from 'lucide-react' -import { api } from '../lib/api' import { Card, CardHeader, CardTitle } from './ui/Card' +import { useStrom } from '../lib/useEventStream' interface LiveLogEntry { id: string @@ -12,40 +11,14 @@ interface LiveLogEntry { } export default function LiveLogSection() { - const [logs, setLogs] = useState([]) - const [currentJob, setCurrentJob] = useState(null) - - useEffect(() => { - const fetchJobs = async () => { - try { - const response = await api.get('/jobs') - const jobs = response.data - const current = jobs.find((j: any) => j.status === 'processing' || j.status === 'transcoding') - setCurrentJob(current || jobs.find((j: any) => j.status === 'pending')) - } catch (error) { - console.error('Fehler beim Laden der Jobs:', error) - } - } - - const fetchLogs = async () => { - try { - const response = await api.get('/logs', { params: { limit: 50 } }) - setLogs(response.data) - } catch (error) { - console.error('Fehler beim Laden der Logs:', error) - } - } - - fetchJobs() - fetchLogs() - const jobInterval = setInterval(fetchJobs, 5000) - const logInterval = setInterval(fetchLogs, 5000) - - return () => { - clearInterval(jobInterval) - clearInterval(logInterval) - } - }, []) + // Vorher: zwei eigene Taktgeber (alle 5 s /jobs und /logs) = 24 Anfragen pro + // Minute allein aus diesem Kasten. Jetzt kommt beides aus der einen offenen + // Verbindung, die der Provider haelt. + const { jobs, logs } = useStrom() + const currentJob = + jobs.find((j: any) => j.status === 'processing' || j.status === 'transcoding') || + jobs.find((j: any) => j.status === 'pending') || + null const getStatusColor = (level: string) => { switch (level) { diff --git a/docker/ui/src/components/WorkerVerwaltung.tsx b/docker/ui/src/components/WorkerVerwaltung.tsx index 7ee4823..434d6eb 100644 --- a/docker/ui/src/components/WorkerVerwaltung.tsx +++ b/docker/ui/src/components/WorkerVerwaltung.tsx @@ -6,6 +6,7 @@ import { useToast } from '../context/ToastContext' import { Card, CardHeader, CardTitle, CardContent } from './ui/Card' import { Button } from './ui/Button' import { Input } from './ui/Input' +import { useStrom } from '../lib/useEventStream' interface WorkerInfo { name: string @@ -54,6 +55,7 @@ export default function WorkerVerwaltung() { // eintragen lassen als eine falsche Adresse vorgaukeln. const [rippyHost, setRippyHost] = useState('') const { toast } = useToast() + const strom = useStrom() // Für die Befehle: leeres Feld → sichtbarer Platzhalter (Befehl ist dann // erkennbar unvollständig, statt still auf die falsche Adresse zu zeigen). @@ -76,11 +78,16 @@ export default function WorkerVerwaltung() { .catch(() => { if (manuell) toast('error', 'Worker-Liste konnte nicht geladen werden') }) } + // Die Worker-Liste kommt aus dem Live-Strom. /capabilities wird nur noch + // EINMAL geholt — dort steht ausser der Liste die Server-Version, und die + // aendert sich nur beim Deploy. Vorher lief der Abruf alle 15 s. + useEffect(() => { laden() }, []) + useEffect(() => { - laden() - const interval = setInterval(laden, 15000) - return () => clearInterval(interval) - }, []) + if (strom.workers.length || strom.verbunden) { + setWorkers(strom.workers as WorkerInfo[]) + } + }, [strom.workers, strom.verbunden]) const installBefehl = variante === 'linux' ? [ diff --git a/docker/ui/src/lib/useEventStream.tsx b/docker/ui/src/lib/useEventStream.tsx new file mode 100644 index 0000000..286e733 --- /dev/null +++ b/docker/ui/src/lib/useEventStream.tsx @@ -0,0 +1,225 @@ +/** + * Der Live-Strom des UI — EINE Verbindung für die ganze Anwendung. + * + * ## Was das ablöst + * + * Bis hierher hatte jede Komponente ihren eigenen Taktgeber. Nachgerechnet an + * den `setInterval`-Aufrufen verursachte EIN offener Tab plus ein Worker rund + * 133 Anfragen pro Minute: + * + * Dashboard 5 Endpunkte alle 4 s -> 75/min + * Log-Kasten 2 Endpunkte alle 5 s -> 24/min + * Laufwerke 1 Endpunkt alle 5 s -> 12/min + * Log-Seite 1 Endpunkt alle 10 s -> 6/min + * Worker-Liste 1 Endpunkt alle 15 s -> 4/min + * Windows-Tray /jobs alle 5 s -> 12/min + * + * Das Rate-Limit stand einmal UNTER dieser Zahl (100/min). Etwa jede vierte + * Anfrage bekam 429 — und weil das UI einen Fehlschlag als „es gibt nichts" + * verbuchte, leerte sich die Job-Liste im Sekundentakt. Der Commander meldete + * das als „wird oft neu geladen". + * + * ## DIE Regel dieser Datei + * + * Ein Verbindungsabriss ist KEINE Aussage über die Welt. + * + * Bei einem Fehler wird der letzte bekannte Stand BEHALTEN und `verbunden` + * auf false gesetzt — das UI zeigt dann ein Banner. Es wird niemals eine + * Liste auf leer gesetzt. Genau das war der alte Fehler: fünfmal stand + * `catch(() => [])` im Code, und jeder fehlgeschlagene Abruf hieß damit + * „keine Jobs, keine Laufwerke, keine Ablagen". + * + * ## Warum ein Context und nicht ein Hook je Komponente + * + * `EventSource` öffnet je Aufruf eine eigene HTTP-Verbindung. Ein Hook, den + * fünf Komponenten benutzen, wären fünf Verbindungen und fünf Snapshots — + * also genau das Problem in neuer Form. Der Provider hält eine. + */ +import { + createContext, + useContext, + useEffect, + useMemo, + useRef, + useState, + type ReactNode, +} from 'react' + +export type Job = { + id: string + status: string + progress: number + title?: string | null + type?: string + device?: string + error?: string | null + meta?: Record | null + [k: string]: any +} + +export type Device = { id: string; name: string; type: string; path: string; status: string; [k: string]: any } +export type Worker = { name: string; encoders?: string[]; info?: Record; last_seen?: string | null } +/** + * Dieselbe Form, die auch `GET /logs` liefert (`_log_zeile` in main.py). + * `timestamp`, nicht `ts` — sonst zeigt das Log-Fenster „Invalid Date", und + * zwar NUR im Live-Betrieb. + */ +export type LogZeile = { id: string; timestamp?: string; level?: string; source?: string; message?: string } + +export type StromZustand = { + jobs: Job[] + /** Server-Zustand aus dem langsamen Takt des Waechters. null = noch nichts gehoert. */ + systemInfo: any | null + /** Ablageziele mit freiem Platz. null = noch nichts gehoert. */ + ablagen: any[] | null + /** null heißt ausdrücklich „konnte nicht nachsehen" — NICHT „es gibt keine". */ + devices: Device[] | null + workers: Worker[] + logs: LogZeile[] + /** false = Verbindung weg. Die Daten oben sind dann der letzte bekannte Stand. */ + verbunden: boolean + /** Wie oft die Verbindung schon neu aufgebaut wurde — fürs Log, nicht fürs UI. */ + neuverbindungen: number + /** Zeitpunkt des letzten empfangenen Ereignisses. */ + zuletzt: Date | null +} + +const LEER: StromZustand = { + jobs: [], + systemInfo: null, + ablagen: null, + devices: null, + workers: [], + logs: [], + verbunden: false, + neuverbindungen: 0, + zuletzt: null, +} + +const StromContext = createContext(LEER) + +/** Wie viele Log-Zeilen im Speicher gehalten werden. */ +const LOG_GRENZE = 300 + +export function EventStreamProvider({ children }: { children: ReactNode }) { + const [zustand, setZustand] = useState(LEER) + // In einem Ref, damit die Handler nicht bei jedem Zustandswechsel neu + // gebunden werden müssen — sonst risse die Verbindung ständig ab. + const neuverbindungen = useRef(0) + + useEffect(() => { + const quelle = new EventSource('/api/events') + + const merken = (patch: Partial) => + setZustand((alt) => ({ ...alt, ...patch, zuletzt: new Date() })) + + // ── Snapshot: das ganze Bild ────────────────────────────────────── + quelle.addEventListener('snapshot', (e) => { + const d = JSON.parse((e as MessageEvent).data)?.daten ?? {} + merken({ + jobs: d.jobs ?? [], + // devices === null kommt vom Server, wenn er die Laufwerke nicht + // lesen konnte. Das wird DURCHGEREICHT, nicht zu [] geglättet. + devices: d.devices === undefined ? null : d.devices, + workers: d.workers ?? [], + logs: d.logs ?? [], + verbunden: true, + }) + }) + + // ── Job-Deltas ──────────────────────────────────────────────────── + const jobPatch = (e: Event) => { + const ereignis = JSON.parse((e as MessageEvent).data) + const id = ereignis.entitaet_id + const daten = ereignis.daten ?? {} + setZustand((alt) => { + if (daten.status === 'geloescht') { + return { ...alt, jobs: alt.jobs.filter((j) => j.id !== id), verbunden: true, zuletzt: new Date() } + } + const bekannt = alt.jobs.some((j) => j.id === id) + const jobs = bekannt + ? alt.jobs.map((j) => (j.id === id ? { ...j, ...daten } : j)) + : [{ id, ...daten } as Job, ...alt.jobs] + return { ...alt, jobs, verbunden: true, zuletzt: new Date() } + }) + } + for (const typ of ['job.created', 'job.progress', 'job.phase', 'job.finished']) { + quelle.addEventListener(typ, jobPatch) + } + + // ── Laufwerke ───────────────────────────────────────────────────── + quelle.addEventListener('drive.changed', (e) => { + const ereignis = JSON.parse((e as MessageEvent).data) + setZustand((alt) => { + if (alt.devices === null) return alt // noch kein Stand — auf Snapshot warten + const id = ereignis.entitaet_id + return { + ...alt, + devices: alt.devices.map((d) => (d.id === id ? { ...d, ...ereignis.daten } : d)), + verbunden: true, + zuletzt: new Date(), + } + }) + }) + + // ── Log-Zeilen ──────────────────────────────────────────────────── + quelle.addEventListener('log.line', (e) => { + const ereignis = JSON.parse((e as MessageEvent).data) + const d = ereignis.daten ?? {} + const zeile: LogZeile = { + id: String(ereignis.entitaet_id), + timestamp: ereignis.ts, + level: d.level, + source: d.source, + message: d.message, + } + setZustand((alt) => ({ + ...alt, + logs: [zeile, ...alt.logs].slice(0, LOG_GRENZE), + verbunden: true, + zuletzt: new Date(), + })) + }) + + // ── Server-Zustand (langsamer Takt) ─────────────────────────────── + quelle.addEventListener('system.status', (e) => { + const d = JSON.parse((e as MessageEvent).data)?.daten ?? {} + // Nur uebernehmen, was WIRKLICH mitgeschickt wurde. Ein fehlender + // Schluessel heisst „konnte ich diesmal nicht lesen" — dann bleibt der + // alte Stand stehen, statt auf leer zu fallen. + setZustand((alt) => ({ + ...alt, + systemInfo: d.info !== undefined ? d.info : alt.systemInfo, + workers: d.workers !== undefined ? d.workers : alt.workers, + ablagen: d.ablagen !== undefined ? d.ablagen : alt.ablagen, + verbunden: true, + zuletzt: new Date(), + })) + }) + + quelle.onopen = () => merken({ verbunden: true }) + + // ── Abriss: Stand BEHALTEN, nur die Verbindung melden ───────────── + quelle.onerror = () => { + // Kein Leeren, kein Zurücksetzen. Der Browser verbindet von selbst neu + // und schickt dabei Last-Event-ID mit; der Server liefert die Lücke nach + // oder schickt einen frischen Snapshot. + neuverbindungen.current += 1 + setZustand((alt) => ({ + ...alt, + verbunden: false, + neuverbindungen: neuverbindungen.current, + })) + } + + return () => quelle.close() + }, []) + + const wert = useMemo(() => zustand, [zustand]) + return {children} +} + +/** Der gemeinsame Live-Zustand. */ +export function useStrom(): StromZustand { + return useContext(StromContext) +} diff --git a/docker/ui/src/pages/Dashboard.tsx b/docker/ui/src/pages/Dashboard.tsx index ac07675..5050135 100644 --- a/docker/ui/src/pages/Dashboard.tsx +++ b/docker/ui/src/pages/Dashboard.tsx @@ -10,6 +10,7 @@ import DeviceDiscovery from '../components/DeviceDiscovery' import JobDetailModal from '../components/JobDetailModal' import LiveLogSection from '../components/LiveLogSection' import RetryDialog from '../components/RetryDialog' +import { useStrom } from '../lib/useEventStream' interface JobMeta { year?: number @@ -126,6 +127,7 @@ export default function Dashboard() { const [workersLive, setWorkersLive] = useState([]) const [laufwerke, setLaufwerke] = useState([]) const [ablagen, setAblagen] = useState([]) + const strom = useStrom() const { toast } = useToast() /* @@ -232,64 +234,43 @@ export default function Dashboard() { } /* - * ZWEI TAKTE statt einem. + * KEIN TAKTGEBER MEHR (Etappe V2-3). * - * ⚠️ JEDER Abruf hier gibt bei Fehlschlag `null` — und `null` bedeutet „nichts - * Neues erfahren", nicht „es gibt nichts". Vor dem 26.07.2026 stand überall - * `catch(() => [])`: ein einziger verpasster Abruf leerte Job-Liste, - * Laufwerke und Ablagen, vier Sekunden später war alles wieder da. Genau das - * hat der Commander als „wird oft neu geladen" gemeldet. Siehe `fetchJobs`. + * Hier standen zwei: alle 4 s Jobs und Laufwerke, alle 12 s Hardware, + * Worker und Ablagen — zusammen 75 Anfragen pro Minute JE OFFENEM TAB. Bei + * der damaligen Grenze von 100/min wies die API laufend mit HTTP 429 ab + * (gemessen: 97 Abweisungen im Log), und weil ein Fehlschlag im UI als + * „es gibt nichts" gelesen wurde, leerte sich die Job-Liste im Sekundentakt. + * Genau das hat der Commander als „wird oft neu geladen" gemeldet. * - * Und verpasst wurde reichlich: Fünf Endpunkte alle vier Sekunden sind 75 - * Anfragen pro Minute, dazu Log-Kasten und Laufwerks-Suche — die API wies bei - * ihrer damaligen Grenze von 100/min laufend mit HTTP 429 ab (gemessen: 97 - * Abweisungen im Log). Die Grenze ist jetzt der echten Last angemessen - * (ratelimit.py), aber das rechtfertigt keine Verschwendung: Jobs und - * Laufwerke ändern sich sekündlich, die drei anderen praktisch nie. - * 75/min sind damit 15 + 15/min geworden. + * Jetzt kommt alles aus der EINEN Verbindung, die App.tsx haelt. + * + * ⚠️ DIE REGEL BLEIBT — sie ist nur umgezogen: Ein fehlender Wert im Strom + * heisst „nichts Neues erfahren", NICHT „es gibt nichts". Deshalb wird unten + * jeder Wert nur uebernommen, wenn er wirklich da ist; sonst bleibt der + * letzte bekannte Stand stehen. Das UI zeigt dann oben rechts + * „VERBINDUNG WEG" statt leerer Listen. */ useEffect(() => { - const nichts = () => null + setJobs(strom.jobs as Job[]) + setLoading(false) + }, [strom.jobs]) - // Schnell (4 s): was sich während eines Rips wirklich bewegt. - const schnellLaden = async () => { - try { - const [jobsData, devData] = await Promise.all([ - fetchJobs(), - // Die Laufwerke: Ohne sie behauptete die Server-Status-Karte - // „keine Disc in Arbeit", während oben die erkannte Disc stand. - api.get('/devices').then(r => Array.isArray(r.data) ? r.data : null).catch(nichts), - ]) - if (jobsData) setJobs(jobsData) - if (devData) setLaufwerke(devData) - } finally { - setLoading(false) - } - } + useEffect(() => { + if (strom.devices !== null) setLaufwerke(strom.devices as LaufwerkLive[]) + }, [strom.devices]) - // Langsam (12 s): Hardware, Worker-Liste, Ablagen. Ein Worker, der online - // geht, oder eine Freigabe, die wegbricht, darf drei Takte später auffallen. - const langsamLaden = async () => { - const [sysData, capsData, ablData] = await Promise.all([ - api.get('/system/info').then(r => r.data).catch(nichts), - // Für die ECHTE Online-Zahl: /capabilities kennt den Celery-Ping, - // /system/info nur die Registrierung. - api.get('/capabilities').then(r => r.data?.workers ?? null).catch(nichts), - // Die Ablageziele inkl. eingehängter Freigaben — ohne sie zeigte der - // Server-Status nur die Container-Platte (Commander-Befund). - api.get('/storage-targets').then(r => Array.isArray(r.data) ? r.data : null).catch(nichts), - ]) - if (sysData) setSystemInfo(sysData) - if (capsData) setWorkersLive(capsData) - if (ablData) setAblagen(ablData) - } + useEffect(() => { + if (strom.workers.length || strom.verbunden) setWorkersLive(strom.workers as WorkerLive[]) + }, [strom.workers, strom.verbunden]) - schnellLaden() - langsamLaden() - const schnell = setInterval(schnellLaden, 4000) - const langsam = setInterval(langsamLaden, 12000) - return () => { clearInterval(schnell); clearInterval(langsam) } - }, []) + useEffect(() => { + if (strom.systemInfo) setSystemInfo(strom.systemInfo as SystemInfo) + }, [strom.systemInfo]) + + useEffect(() => { + if (strom.ablagen !== null) setAblagen(strom.ablagen as AblageZiel[]) + }, [strom.ablagen]) const aktiverJob = jobs.find(j => j.status === 'processing' || j.status === 'transcoding') const queueJobs = jobs.filter(j => j.status === 'pending') diff --git a/docker/ui/src/pages/Logs.tsx b/docker/ui/src/pages/Logs.tsx index 0f409a1..964b878 100644 --- a/docker/ui/src/pages/Logs.tsx +++ b/docker/ui/src/pages/Logs.tsx @@ -5,6 +5,7 @@ import { useToast } from '../context/ToastContext' import { PageHeader } from '../components/ui/PageHeader' import { Card, CardContent } from '../components/ui/Card' import { Button } from '../components/ui/Button' +import { useStrom } from '../lib/useEventStream' interface LogEntry { id: string @@ -21,6 +22,7 @@ export default function LogsPage() { const [loading, setLoading] = useState(true) const [refreshing, setRefreshing] = useState(false) const { toast } = useToast() + const strom = useStrom() const refreshLogs = async (manuell = false) => { if (manuell) setRefreshing(true) @@ -37,11 +39,20 @@ export default function LogsPage() { } } + // Einmal die volle Historie holen (200 Zeilen — der Strom liefert im + // Snapshot nur die letzten 50), danach kommt alles Neue von selbst. + // Der Aktualisieren-Knopf bleibt: Er holt die Historie erneut. + useEffect(() => { refreshLogs() }, []) + useEffect(() => { - refreshLogs() - const interval = setInterval(refreshLogs, 10000) - return () => clearInterval(interval) - }, []) + if (!strom.logs.length) return + setLogs((alt) => { + const bekannt = new Set(alt.map((l) => String(l.id))) + const neue = strom.logs.filter((l) => !bekannt.has(String(l.id))) + return neue.length ? ([...neue, ...alt] as LogEntry[]) : alt + }) + setLoading(false) + }, [strom.logs]) const getLevelColor = (level: string) => { switch (level) { diff --git a/src/rippy/bus/schema.py b/src/rippy/bus/schema.py index 9e87b86..4d536ef 100644 --- a/src/rippy/bus/schema.py +++ b/src/rippy/bus/schema.py @@ -45,6 +45,7 @@ EREIGNIS_TYPEN = { "mount.changed": "Ein Speicherziel ist erreichbar geworden oder weggefallen", + "system.status": "Server-Zustand: Hardware, Worker, Ablageziele (langsamer Takt)", "system.notice": "Hinweis an den Nutzer (Plattenplatz, Rate-Limit, …)", } diff --git a/src/rippy/bus/test_waechter.py b/src/rippy/bus/test_waechter.py index f6ab932..008d21b 100644 --- a/src/rippy/bus/test_waechter.py +++ b/src/rippy/bus/test_waechter.py @@ -6,7 +6,7 @@ genau die Eigenschaft, die dem ersten Anlauf der SSE-Tests gefehlt hat (Ampel-Lauf 170 lief rot, weil ein Test eine Datenbank brauchte). """ -from rippy.bus.waechter import Waechter, unterschiede +from rippy.bus.waechter import Waechter, laufwerks_unterschiede, unterschiede class FakeBus: @@ -129,7 +129,7 @@ def test_neue_logzeilen_kommen_aelteste_zuerst(): {"id": 5, "level": "info", "source": "w", "message": "alt"}, ] w.einmal() - texte = [d["text"] for t, _, d in bus.gesendet if t == "log.line"] + texte = [d["message"] for t, _, d in bus.gesendet if t == "log.line"] assert texte == ["erste", "zweite"] @@ -142,7 +142,7 @@ def test_dieselbe_logzeile_kommt_nur_einmal(): {"id": 1, "level": "info", "source": "w", "message": "a"}] w.einmal() w.einmal() - assert [d["text"] for t, _, d in bus.gesendet if t == "log.line"] == ["b"] + assert [d["message"] for t, _, d in bus.gesendet if t == "log.line"] == ["b"] def test_alter_ist_abfragbar(): @@ -157,3 +157,79 @@ def test_alter_ist_abfragbar(): w.einmal() assert 0 <= w.lebt_seit_sekunden() < 1 assert w.gesund is True + + +# ── Laufwerke ─────────────────────────────────────────────────────────── + + +def _lw(pfad, status="empty", typ="unknown"): + return {"path": pfad, "id": pfad.split("/")[-1], "status": status, "type": typ} + + +def test_eingelegte_disc_bekommt_ein_eigenes_ereignis(): + """Das UI reagiert darauf anders als auf eine bloße Zustandsaenderung — + bei einer eingelegten Disc springt der Rip-Dialog auf.""" + vorher = {"/dev/sr0": _lw("/dev/sr0", "empty")} + jetzt = {"/dev/sr0": _lw("/dev/sr0", "ready", "bluray")} + typen = [t for t, _, _ in laufwerks_unterschiede(vorher, jetzt)] + assert "disc.inserted" in typen + assert typen.index("disc.inserted") < typen.index("drive.changed") + + +def test_entnommene_disc_ebenso(): + vorher = {"/dev/sr0": _lw("/dev/sr0", "ready", "bluray")} + jetzt = {"/dev/sr0": _lw("/dev/sr0", "empty")} + typen = [t for t, _, _ in laufwerks_unterschiede(vorher, jetzt)] + assert "disc.removed" in typen + + +def test_unveraendertes_laufwerk_erzeugt_nichts(): + stand = {"/dev/sr0": _lw("/dev/sr0", "ready")} + assert laufwerks_unterschiede(stand, dict(stand)) == [] + + +def test_erstes_auftauchen_ist_kein_disc_ereignis(): + """Beim Start ist jedes Laufwerk „neu". Wuerde das als disc.inserted + zaehlen, spraenge nach jedem API-Neustart der Rip-Dialog auf.""" + typen = [t for t, _, _ in laufwerks_unterschiede({}, {"/dev/sr0": _lw("/dev/sr0", "ready")})] + assert typen == ["drive.changed"] + + +def test_verschwundenes_laufwerk_wird_gemeldet(): + typ, geraet, daten = laufwerks_unterschiede({"/dev/sr0": _lw("/dev/sr0")}, {})[0] + assert (typ, geraet) == ("drive.changed", "/dev/sr0") + assert daten["status"] == "weg" + + +def test_haengendes_laufwerk_haelt_den_waechter_nicht_an(): + """DIE Zusage: Ein defektes Laufwerk darf nicht den Job-Fortschritt + mitnehmen. Sonst waere ein klemmendes Laufwerk gleichbedeutend mit einem + eingefrorenen UI — und niemand saehe, woran es liegt.""" + bus = FakeBus() + store = FakeStore(jobs=[_job("j1", "ripping", 10)]) + + def kaputt(): + raise OSError("Laufwerk haengt") + + w = Waechter(store, bus, laufwerke_lesen=kaputt, laufwerks_takt=0) + w.einmal() + store.jobs = [_job("j1", "ripping", 20)] + assert w.einmal() == 1 # der Job kommt trotzdem durch + assert bus.typen == ["job.progress"] + + +def test_laufwerke_werden_seltener_abgefragt_als_jobs(): + """Ein ioctl kostet mehr als ein SELECT. Ohne den eigenen Takt liefe + jede Sekunde eine Laufwerksabfrage — auf einem klemmenden Laufwerk waere + das ein Dauerproblem.""" + aufrufe = [] + + def lesen(): + aufrufe.append(1) + return [_lw("/dev/sr0")] + + w = Waechter(FakeStore(), FakeBus(), laufwerke_lesen=lesen, laufwerks_takt=999) + w.einmal() + w.einmal() + w.einmal() + assert len(aufrufe) == 1, "Laufwerke wurden mehrfach im selben Takt gelesen" diff --git a/src/rippy/bus/waechter.py b/src/rippy/bus/waechter.py index 9eef37a..d20a296 100644 --- a/src/rippy/bus/waechter.py +++ b/src/rippy/bus/waechter.py @@ -64,6 +64,16 @@ TAKT_SEKUNDEN = 1.0 # Zustand (DB weg) nicht jede Sekunde eine Meldung erzeugt. FEHLER_TAKT_SEKUNDEN = 5.0 +# Laufwerke seltener: ein ioctl kostet mehr als ein SELECT, und eine Disc +# wird nicht mehrmals pro Sekunde gewechselt. Drei Sekunden entsprechen dem +# Takt, den die Disc-Wache in v1 schon hatte. +LAUFWERKS_TAKT_SEKUNDEN = 3.0 + +# Server-Zustand (Hardware, Worker-Liste, Ablageziele) noch seltener: Das +# aendert sich selten, und /storage-targets fasst Netzpfade an. 15 s +# entsprechen dem alten UI-Takt von 12 s, nur eben EINMAL statt je Tab. +SYSTEM_TAKT_SEKUNDEN = 15.0 + # Endzustände: ab hier ist ein Job durch. ENDE = ("completed", "failed", "canceled") @@ -120,14 +130,62 @@ def unterschiede(vorher: dict, jetzt: dict) -> list: return neu + geaendert + verschwunden -class Waechter: - """Sieht in der Datenbank nach und meldet Änderungen an den Bus.""" +def laufwerks_unterschiede(vorher: dict, jetzt: dict) -> list: + """Was hat sich an den Laufwerken geändert? (pure Funktion) - def __init__(self, store, bus, takt: float = TAKT_SEKUNDEN): + Der Disc-Wechsel bekommt eigene Ereignisse (`disc.inserted` / + `disc.removed`), weil das UI darauf anders reagiert als auf eine bloße + Zustandsänderung: Bei einer eingelegten Disc springt der Rip-Dialog auf. + """ + ereignisse = [] + for geraet, stand in jetzt.items(): + alt = vorher.get(geraet) + if alt == stand: + continue + if alt is not None: + leer_vorher = alt.get("status") != "ready" + leer_jetzt = stand.get("status") != "ready" + if leer_vorher and not leer_jetzt: + ereignisse.append(("disc.inserted", geraet, stand)) + elif not leer_vorher and leer_jetzt: + ereignisse.append(("disc.removed", geraet, stand)) + ereignisse.append(("drive.changed", geraet, stand)) + for geraet in vorher: + if geraet not in jetzt: + ereignisse.append(("drive.changed", geraet, {"status": "weg"})) + return ereignisse + + +class Waechter: + """Sieht in der Datenbank nach und meldet Änderungen an den Bus. + + `laufwerke_lesen` ist eine Funktion ohne Argumente, die die Laufwerksliste + liefert (in der API: `device_info` über alle gefundenen Geräte). Sie wird + eingespritzt statt importiert — so ist der Wächter ohne echtes Laufwerk + testbar, und im Standalone-Betrieb kann derselbe Wächter den + Windows-Treiber benutzen. + + Laufwerke werden SELTENER abgefragt als Jobs: Ein `ioctl` auf einem + optischen Laufwerk kostet spürbar mehr als ein SELECT, und eine Disc wird + nicht mehrmals pro Sekunde gewechselt. Drei Sekunden entsprechen dem Takt, + den die alte Disc-Wache hatte. + """ + + def __init__(self, store, bus, takt: float = TAKT_SEKUNDEN, + laufwerke_lesen=None, laufwerks_takt: float = LAUFWERKS_TAKT_SEKUNDEN, + system_lesen=None, system_takt: float = SYSTEM_TAKT_SEKUNDEN): self._store = store self._bus = bus self._takt = takt + self._laufwerke_lesen = laufwerke_lesen + self._laufwerks_takt = laufwerks_takt self._jobs: dict = {} + self._laufwerke: dict = {} + self._laufwerke_geprueft = 0.0 + self._system_lesen = system_lesen + self._system_takt = system_takt + self._system = None + self._system_geprueft = 0.0 self._letzte_log_id = None self._letzter_lauf = 0.0 self._fehler_in_folge = 0 @@ -173,18 +231,79 @@ class Waechter: else: neue = [z for z in zeilen if z["id"] > (self._letzte_log_id or 0)] for zeile in reversed(neue): # älteste zuerst + # Feldnamen wie in der Datenbank (level/source/message) — nicht + # uebersetzt. Wer sie umbenennt, muss das UI und /logs mitziehen, + # und dann stehen zwei Formen nebeneinander. self._bus.senden("log.line", { "level": zeile.get("level"), - "quelle": zeile.get("source"), - "text": zeile.get("message"), + "source": zeile.get("source"), + "message": zeile.get("message"), }, entitaet="log", entitaet_id=str(zeile["id"])) gesendet += 1 if neue: self._letzte_log_id = max(z["id"] for z in neue) + gesendet += self._laufwerke_pruefen(erster_lauf) + gesendet += self._system_pruefen() + self._letzter_lauf = time.monotonic() return gesendet + def _system_pruefen(self) -> int: + """Server-Zustand im langsamen Takt. + + Anders als bei Jobs wird hier NICHT auf Unterschiede geprüft, sondern + der Stand jedes Mal geschickt: Es geht um Messwerte (freier Platz, + Auslastung), die sich praktisch immer ändern — ein Vergleich würde nur + Rechenzeit kosten und trotzdem jedes Mal „geändert" sagen. + + Fehler werden geschluckt (siehe _laufwerke_pruefen): /storage-targets + fasst Netzpfade an, und ein weggebrochenes NAS darf den Job-Fortschritt + nicht mitnehmen. + """ + if self._system_lesen is None: + return 0 + jetzt = time.monotonic() + if jetzt - self._system_geprueft < self._system_takt: + return 0 + self._system_geprueft = jetzt + try: + stand = self._system_lesen() + except Exception: + return 0 + if not stand: + return 0 + self._system = stand + self._bus.senden("system.status", stand, entitaet="system") + return 1 + + def _laufwerke_pruefen(self, erster_lauf: bool) -> int: + """Laufwerke im eigenen, langsameren Takt. + + Fehler werden hier BEWUSST geschluckt und nicht weitergereicht: Ein + hängendes oder defektes Laufwerk darf nicht den ganzen Wächter + anhalten — sonst käme auch kein Job-Fortschritt mehr durch. Das UI + behält dann seinen letzten Laufwerksstand, was richtig ist: „konnte + nicht nachsehen" ist keine Aussage über die Welt. + """ + if self._laufwerke_lesen is None: + return 0 + jetzt = time.monotonic() + if jetzt - self._laufwerke_geprueft < self._laufwerks_takt: + return 0 + self._laufwerke_geprueft = jetzt + try: + stand = {g["path"]: g for g in self._laufwerke_lesen()} + except Exception: + return 0 + gesendet = 0 + if not erster_lauf and self._laufwerke: + for typ, geraet, daten in laufwerks_unterschiede(self._laufwerke, stand): + self._bus.senden(typ, daten, entitaet="drive", entitaet_id=geraet) + gesendet += 1 + self._laufwerke = stand + return gesendet + # ── Die Schleife ──────────────────────────────────────────────────── async def schleife(self) -> None: """Läuft, bis sie abgebrochen wird. Meldet Fehler LAUT."""