feat(ui): V2-3 (Teil 2) — das UI haengt am Ereignis-Strom, Taktgeber raus
Ampel / ampel (push) Successful in 47s
Ampel / ampel (push) Successful in 47s
WAS: EventStreamProvider haelt EINE SSE-Verbindung fuer die ganze
Anwendung. Dashboard, Log-Kasten, Log-Seite, Laufwerksliste und
Worker-Liste beziehen ihren Zustand daraus. Sieben von neun setInterval
sind weg.
GEMESSEN AN DEN TAKTGEBERN, je offenem Tab:
Dashboard 5 Endpunkte / 4 s 75/min -> 0
Log-Kasten 2 Endpunkte / 5 s 24/min -> 0
Laufwerke 1 Endpunkt / 5 s 12/min -> 0
Log-Seite 1 Endpunkt /10 s 6/min -> 0 (+1 Abruf beim Oeffnen)
Worker-Liste 1 Endpunkt /15 s 4/min -> 0 (+1 Abruf beim Oeffnen)
------
121/min -> ~2 einmalige Abrufe
Zwei Taktgeber bleiben bewusst: FirstRunWizard (laeuft nur VOR der
Einrichtung) und RipTargetModal (nur solange der Dialog offen ist).
DIE REGEL IST UMGEZOGEN, NICHT VERSCHWUNDEN: Ein Abriss ist keine
Aussage ueber die Welt. Der Provider BEHAELT bei einem Fehler den letzten
Stand und setzt nur `verbunden` auf false; es wird nie eine Liste geleert.
Jede Komponente uebernimmt einen Wert nur, wenn er wirklich da ist —
`devices === null` heisst "konnte nicht nachsehen", nicht "keine
Laufwerke". Das war der Fehler hinter "wird oft neu geladen".
EIN PLACEBO WENIGER: Oben rechts stand ein fest verdrahtetes "ONLINE" mit
pulsierendem Punkt — es leuchtete gruen, auch wenn die API tot war. Jetzt
zeigt es LIVE oder VERBINDUNG WEG, und im Tooltip steht, wann die letzte
Meldung kam.
DER SERVER SCHIEBT JETZT AUCH DEN SERVER-ZUSTAND: Neuer Ereignistyp
system.status (Hardware, Worker, Ablageziele) im 15-Sekunden-Takt des
Waechters — EINMAL im Server statt 15/min je Tab. Nur mitgeschickte
Schluessel werden uebernommen; ein fehlender heisst "behalte deinen Stand".
DAZU EIN FORMATFEHLER GEFUNDEN UND BEHOBEN: /logs bildet ts -> timestamp
ab, mein Snapshot lieferte die rohe DB-Zeile. Das UI haette "Invalid Date"
gezeigt — und zwar NUR im Live-Betrieb, nicht beim manuellen Neuladen.
Jetzt gibt es _log_zeile() einmal, benutzt von beiden.
UND EINEN ZWEITEN: system_lesen lief per asyncio.get_event_loop() in einem
Worker-Thread — dort ist das NICHT die laufende Schleife. Die Coroutine
waere nie gelaufen. Die Schleife wird jetzt im Startup festgehalten.
GEMESSEN: ruff sauber, 378 Tests + 3 uebersprungen, `npm run build`
durch (1650 Module). Der Live-Beweis steht noch aus — er kommt mit dem
Deploy.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5
parent
f097a0b59d
commit
95705c8d88
+96
-12
@@ -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")
|
||||
|
||||
@@ -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():
|
||||
|
||||
+48
-5
@@ -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 (
|
||||
<div className="flex items-center gap-2 text-xs font-mono text-emerald-400 bg-emerald-500/10 px-3 py-1.5 rounded-xl border border-emerald-500/20">
|
||||
<span className="w-2 h-2 rounded-full bg-emerald-400 animate-pulse" />
|
||||
<span>LIVE</span>
|
||||
</div>
|
||||
)
|
||||
}
|
||||
return (
|
||||
<div
|
||||
className="flex items-center gap-2 text-xs font-mono text-amber-400 bg-amber-500/10 px-3 py-1.5 rounded-xl border border-amber-500/30"
|
||||
title={
|
||||
zuletzt
|
||||
? `Letzte Meldung um ${zuletzt.toLocaleTimeString()}. Die Anzeige zeigt diesen Stand, bis die Verbindung zurueck ist.`
|
||||
: 'Noch keine Verbindung zum Server.'
|
||||
}
|
||||
>
|
||||
<span className="w-2 h-2 rounded-full bg-amber-400" />
|
||||
<span>VERBINDUNG WEG</span>
|
||||
</div>
|
||||
)
|
||||
}
|
||||
|
||||
function AppInhalt() {
|
||||
const [currentPage, setCurrentPage] = useState<Page>('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() {
|
||||
})}
|
||||
</nav>
|
||||
|
||||
<div className="flex items-center gap-2 text-xs font-mono text-emerald-400 bg-emerald-500/10 px-3 py-1.5 rounded-xl border border-emerald-500/20">
|
||||
<span className="w-2 h-2 rounded-full bg-emerald-400 animate-pulse" />
|
||||
<span>ONLINE</span>
|
||||
</div>
|
||||
<VerbindungsAnzeige />
|
||||
</div>
|
||||
</header>
|
||||
|
||||
@@ -86,3 +118,14 @@ export default function App() {
|
||||
</div>
|
||||
)
|
||||
}
|
||||
|
||||
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 (
|
||||
<EventStreamProvider>
|
||||
<AppInhalt />
|
||||
</EventStreamProvider>
|
||||
)
|
||||
}
|
||||
|
||||
@@ -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<Device | null>(null)
|
||||
const [korrekturDevice, setKorrekturDevice] = useState<Device | null>(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<string, string> = {
|
||||
tmdb: 'TMDB', omdb: 'OMDb', jikan: 'MyAnimeList', manuell: 'manuell gewählt',
|
||||
|
||||
@@ -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<LiveLogEntry[]>([])
|
||||
const [currentJob, setCurrentJob] = useState<any>(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) {
|
||||
|
||||
@@ -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'
|
||||
? [
|
||||
|
||||
@@ -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<string, any> | 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<string, any>; 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<StromZustand>(LEER)
|
||||
|
||||
/** Wie viele Log-Zeilen im Speicher gehalten werden. */
|
||||
const LOG_GRENZE = 300
|
||||
|
||||
export function EventStreamProvider({ children }: { children: ReactNode }) {
|
||||
const [zustand, setZustand] = useState<StromZustand>(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<StromZustand>) =>
|
||||
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 <StromContext.Provider value={wert}>{children}</StromContext.Provider>
|
||||
}
|
||||
|
||||
/** Der gemeinsame Live-Zustand. */
|
||||
export function useStrom(): StromZustand {
|
||||
return useContext(StromContext)
|
||||
}
|
||||
@@ -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<WorkerLive[]>([])
|
||||
const [laufwerke, setLaufwerke] = useState<LaufwerkLive[]>([])
|
||||
const [ablagen, setAblagen] = useState<AblageZiel[]>([])
|
||||
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')
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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, …)",
|
||||
}
|
||||
|
||||
|
||||
@@ -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"
|
||||
|
||||
+124
-5
@@ -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."""
|
||||
|
||||
Reference in New Issue
Block a user