feat(api): V2-3 (Teil 1) — Ereignis-Bus und SSE-Endpunkt /api/v2/events
Ampel / ampel (push) Failing after 41s
Ampel / ampel (push) Failing after 41s
WAS: rippy/bus mit Ereignis-Schema, In-Process-Treiber und Ringpuffer.
Dazu der SSE-Strom in der API: erst ein Snapshot, danach nur Deltas,
mit lueckenloser Wiederaufnahme ueber Last-Event-ID.
WARUM: Ein offener Tab plus ein Worker verursachen heute rund 133
Anfragen pro Minute (Dashboard 75 + Log-Kasten 24 + Laufwerke 12 +
Worker-Liste 4 + Log-Seite 6 + Tray 12). Das Rate-Limit stand einmal
UNTER dieser Zahl — daher die sich leerende Job-Liste im Sekundentakt.
Mit einer offenen Verbindung sind es null.
DREI EIGENSCHAFTEN, ALLE AUS v1-FEHLERN:
1. Der Bus traegt NUR Nachrichten ueber Aenderungen, nie den Zustand.
Wer den Zustand will, fragt den Store. Damit kann ein verpasstes
Ereignis auch keinen Zustand loeschen — anders als beim fuenffachen
`catch(() => [])` im alten UI, wo jeder fehlgeschlagene Abruf
"es gibt keine Jobs" bedeutete.
2. Eine zu grosse Luecke wird ANGESAGT, nicht verschluckt. nachliefern()
gibt None ("hol dir ein ganzes Bild") statt [] ("nichts verpasst") —
dieselbe Unterscheidung wie timeout-Rueckgabe 124 bei den Netzpfaden.
Stillschweigend weiterzumachen waere schlimmer: Das UI hielte sich
fuer aktuell und waere es nicht.
3. Der Snapshot meldet unlesbare Laufwerke als None, nicht als leere
Liste. "Konnte nicht nachsehen" ist etwas anderes als "gibt es nicht".
Ein unbekannter Ereignistyp fliegt beim Senden HOCH statt durchzugehen.
Ein Tippfehler waere sonst der stillste aller Fehlschlaege: Nachricht
raus, kein Empfaenger, nirgends ein Hinweis.
Ein langsamer Zuhoerer (Tab im Hintergrund, lahmes Handy) bremst den
Sender nicht — er wird markiert und bekommt beim naechsten Mal einen
Snapshot. Ein Rip darf nicht auf einen Browser warten.
GEMESSEN: ruff sauber, 359 Tests gruen + 3 uebersprungen (vorher 346).
13 neue Bus-Tests laufen auf jeder Plattform; die vier SSE-Tests haengen
an main.py und laufen damit in der Ampel.
NOCH OFFEN in V2-3: das UI auf den Strom umstellen (neun setInterval)
und die Ereignisse an den Zustandsaenderungen ausloesen.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5
parent
d2645f87e8
commit
55db9eb13f
+125
-1
@@ -1,6 +1,6 @@
|
||||
from fastapi import FastAPI, HTTPException, Request, Response
|
||||
from fastapi.middleware.cors import CORSMiddleware
|
||||
from fastapi.responses import FileResponse
|
||||
from fastapi.responses import FileResponse, StreamingResponse
|
||||
from pydantic import BaseModel
|
||||
from typing import List, Optional, Dict
|
||||
import asyncio
|
||||
@@ -18,6 +18,8 @@ from rippy.rip import makemkv_daten
|
||||
import makemkv_key
|
||||
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
|
||||
import phasen
|
||||
import presets as preset_auswahl
|
||||
import rohdaten
|
||||
@@ -338,6 +340,128 @@ def _job_row_to_model(zeile: dict) -> Job:
|
||||
meta=meta,
|
||||
)
|
||||
|
||||
# ═══════════════════════════════════════════════════════════════════════
|
||||
# Live-Ereignisse (SSE) — Etappe V2-3
|
||||
# ═══════════════════════════════════════════════════════════════════════
|
||||
#
|
||||
# WAS DAS ERSETZT: Bis hierher fragte das UI im Takt nach. Nachgerechnet an
|
||||
# den Taktgebern (der Kommentar in ratelimit.py führt sie auf) verursacht EIN
|
||||
# offener Tab plus ein Worker rund 133 Anfragen pro Minute — und das Rate-Limit
|
||||
# stand einmal UNTER dieser Zahl, weshalb sich die Job-Liste im Sekundentakt
|
||||
# leerte. Mit einer offenen Verbindung sind es null.
|
||||
#
|
||||
# WARUM SSE UND NICHT WEBSOCKET: Gebraucht wird genau EIN Rückkanal, Server zu
|
||||
# Browser. Kommandos gehen weiter per REST (idempotent, protokollierbar, mit
|
||||
# der CLI teilbar). SSE bringt Wiederaufnahme (Last-Event-ID) und den
|
||||
# automatischen Reconnect des Browsers mit; bei WebSocket müsste man beides
|
||||
# selbst bauen.
|
||||
|
||||
SSE_HERZSCHLAG_SEKUNDEN = 15
|
||||
|
||||
|
||||
async def _snapshot() -> dict:
|
||||
"""Der vollständige Zustand — das ERSTE, was jeder Client bekommt.
|
||||
|
||||
Ohne diesen Anfang sähe ein frisch verbundenes UI nur die Änderungen ab
|
||||
jetzt und müsste den Rest raten. Mit ihm gilt: erst das ganze Bild, danach
|
||||
nur noch Deltas.
|
||||
"""
|
||||
def sammeln():
|
||||
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),
|
||||
}
|
||||
|
||||
zustand = await asyncio.to_thread(sammeln)
|
||||
# Laufwerke getrennt: device_info macht ioctls, die hängen können —
|
||||
# ein defektes Laufwerk darf den Snapshot nicht aufhalten.
|
||||
try:
|
||||
zustand["devices"] = await asyncio.wait_for(
|
||||
asyncio.to_thread(
|
||||
lambda: [device_discovery.device_info(p)
|
||||
for p in device_discovery.list_optical_devices()]),
|
||||
timeout=5,
|
||||
)
|
||||
except (asyncio.TimeoutError, OSError):
|
||||
# „konnte nicht nachsehen" ist etwas anderes als „es gibt keine".
|
||||
# Deshalb None und nicht [] — das UI behält dann seinen alten Stand.
|
||||
zustand["devices"] = None
|
||||
return zustand
|
||||
|
||||
|
||||
def _sse_rahmen(ereignis: dict) -> str:
|
||||
"""Ein Ereignis im SSE-Format. `id:` ist die Basis der Wiederaufnahme."""
|
||||
nutzlast = json.dumps(ereignis, ensure_ascii=False)
|
||||
zeilen = [
|
||||
"id: " + str(ereignis["seq"]),
|
||||
"event: " + ereignis["typ"],
|
||||
"data: " + nutzlast,
|
||||
"",
|
||||
"",
|
||||
]
|
||||
return "\n".join(zeilen)
|
||||
|
||||
|
||||
@app.get("/api/v2/events")
|
||||
async def events(request: Request, last_event_id: str = None):
|
||||
"""Live-Strom: erst ein Snapshot, danach nur noch Änderungen.
|
||||
|
||||
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
|
||||
sich fälschlich für aktuell hält, ist schlimmer als eines, das neu lädt.
|
||||
"""
|
||||
kopf = request.headers.get("last-event-id") or last_event_id
|
||||
try:
|
||||
ab_seq = int(kopf) if kopf else None
|
||||
except ValueError:
|
||||
ab_seq = None
|
||||
|
||||
async def strom():
|
||||
abo = ereignis_bus.abonnieren(ab_seq)
|
||||
try:
|
||||
verpasst = ereignis_bus.nachliefern(ab_seq) if ab_seq is not None else None
|
||||
if verpasst is None:
|
||||
# Kein Rückstand bekannt oder Lücke zu groß -> ganzes Bild.
|
||||
yield _sse_rahmen(bus_schema.baue_ereignis(
|
||||
"snapshot", await _snapshot(), seq=ereignis_bus.seq))
|
||||
else:
|
||||
for ereignis in verpasst:
|
||||
yield _sse_rahmen(ereignis)
|
||||
|
||||
while True:
|
||||
if await request.is_disconnected():
|
||||
return
|
||||
ereignis = await abo.naechstes(timeout=SSE_HERZSCHLAG_SEKUNDEN)
|
||||
if ereignis is None:
|
||||
# Herzschlag: manche Proxys schließen stille Verbindungen.
|
||||
yield ": herzschlag\n\n"
|
||||
continue
|
||||
if abo.abgehaengt:
|
||||
# Dieser Zuhörer hat den Anschluss verloren — ganzes Bild
|
||||
# statt Bruchstücken.
|
||||
abo.abgehaengt = False
|
||||
yield _sse_rahmen(bus_schema.baue_ereignis(
|
||||
"snapshot", await _snapshot(), seq=ereignis_bus.seq))
|
||||
continue
|
||||
yield _sse_rahmen(ereignis)
|
||||
finally:
|
||||
abo.schliessen()
|
||||
|
||||
return StreamingResponse(
|
||||
strom(),
|
||||
media_type="text/event-stream",
|
||||
headers={
|
||||
"Cache-Control": "no-cache",
|
||||
"Connection": "keep-alive",
|
||||
# nginx puffert `text/event-stream` sonst und der Strom käme
|
||||
# in Schüben statt live an.
|
||||
"X-Accel-Buffering": "no",
|
||||
},
|
||||
)
|
||||
|
||||
|
||||
@app.get("/health")
|
||||
async def health_check():
|
||||
return {"status": "ok", "service": "api"}
|
||||
|
||||
@@ -246,3 +246,67 @@ def test_wache_meldet_wenn_es_wieder_geht(monkeypatch):
|
||||
zustand["da"] = True
|
||||
main._mounts_nachsehen()
|
||||
assert any("antwortet wieder" in m for m in gelogged)
|
||||
|
||||
|
||||
# --- Live-Ereignisse (SSE), Etappe V2-3 --------------------------------------
|
||||
|
||||
|
||||
def test_events_route_ist_verdrahtet():
|
||||
"""Ohne diese Route fällt das UI stumm auf seinen letzten Stand zurück —
|
||||
und weil ein Abriss KEINE Aussage ist, sähe der Nutzer einfach nichts
|
||||
Neues, ohne Fehlermeldung. Deshalb hier festgenagelt."""
|
||||
from main import app
|
||||
|
||||
assert "/api/v2/events" in {route.path for route in app.routes}
|
||||
|
||||
|
||||
def test_sse_rahmen_hat_das_format_das_der_browser_erwartet():
|
||||
"""`id:` ist nicht Kosmetik — der Browser schickt genau diesen Wert beim
|
||||
Wiederverbinden als Last-Event-ID zurück. Fehlt er, gibt es keine
|
||||
lückenlose Wiederaufnahme, und jeder WLAN-Wechsel reisst ein Loch."""
|
||||
from main import _sse_rahmen
|
||||
|
||||
rahmen = _sse_rahmen({"seq": 42, "typ": "job.progress", "daten": {"prozent": 7}})
|
||||
zeilen = rahmen.split("\n")
|
||||
assert zeilen[0] == "id: 42"
|
||||
assert zeilen[1] == "event: job.progress"
|
||||
assert zeilen[2].startswith("data: {")
|
||||
# Zwei Leerzeilen am Ende: eine schliesst das Ereignis, die zweite ist der
|
||||
# Trenner. Ohne den doppelten Umbruch haelt der Browser das Ereignis fuer
|
||||
# unvollstaendig und liefert es NIE aus.
|
||||
assert rahmen.endswith("\n\n")
|
||||
|
||||
|
||||
def test_sse_rahmen_uebersteht_umlaute():
|
||||
"""Job-Titel und Log-Zeilen sind deutsch. Mit ensure_ascii=True kaeme
|
||||
"Gr\u00f6\u00dfe" beim Nutzer an."""
|
||||
from main import _sse_rahmen
|
||||
|
||||
rahmen = _sse_rahmen({"seq": 1, "typ": "log.line", "daten": {"text": "Größe"}})
|
||||
assert "Größe" in rahmen
|
||||
|
||||
|
||||
def test_snapshot_meldet_unlesbare_laufwerke_als_none():
|
||||
"""„konnte nicht nachsehen" ist etwas anderes als „es gibt keine".
|
||||
|
||||
Genau diese Vermischung hat in v1 die Job-Liste im Sekundentakt geleert
|
||||
(fuenfmal `catch(() => [])` im UI). Ein Snapshot mit devices=[] wuerde dem
|
||||
UI sagen „du hast kein Laufwerk"; None sagt „ich weiss es gerade nicht",
|
||||
und das UI behaelt seinen Stand.
|
||||
"""
|
||||
import asyncio
|
||||
|
||||
import main
|
||||
|
||||
def kaputt():
|
||||
raise OSError("Laufwerk haengt")
|
||||
|
||||
original = main.device_discovery.list_optical_devices
|
||||
main.device_discovery.list_optical_devices = kaputt
|
||||
try:
|
||||
zustand = asyncio.run(main._snapshot())
|
||||
finally:
|
||||
main.device_discovery.list_optical_devices = original
|
||||
|
||||
assert zustand["devices"] is None
|
||||
assert isinstance(zustand["jobs"], list)
|
||||
|
||||
@@ -0,0 +1,44 @@
|
||||
"""Ereignis-Bus: Nachrichten über Änderungen — nie der Zustand selbst.
|
||||
|
||||
## Die eine Regel, aus der alles folgt
|
||||
|
||||
Der Bus überträgt nur NACHRICHTEN ÜBER ÄNDERUNGEN.
|
||||
Wer den Zustand will, fragt den Store.
|
||||
|
||||
Das klingt nach Prinzipienreiterei und ist die Lehre aus dem teuersten
|
||||
UI-Fehler von v1. Dort stand fünfmal `catch(() => [])`: Jeder fehlgeschlagene
|
||||
Abruf hieß damit „es gibt keine Jobs, keine Laufwerke, keine Ablagen". Die
|
||||
Liste leerte sich für einen Takt und füllte sich vier Sekunden später wieder.
|
||||
Der Commander meldete das als „wird oft neu geladen", und die Ursache war
|
||||
unsichtbar, weil der Fehlerzweig nichts protokollierte (AGENTS.md: „Ein
|
||||
verpasster Abruf ist keine Nachricht über die Welt").
|
||||
|
||||
Wenn der Bus keinen Zustand trägt, kann ein verpasstes Ereignis auch keinen
|
||||
Zustand löschen. Ein Abriss heißt dann: „Ich weiß gerade nichts Neues" — nicht
|
||||
„es gibt nichts".
|
||||
|
||||
## Drei Eigenschaften, die daraus folgen
|
||||
|
||||
1. **`snapshot` zuerst.** Wer sich verbindet, bekommt als Erstes den
|
||||
vollständigen Zustand, danach nur noch Deltas. Ein Client sieht nie ein
|
||||
halbes Bild.
|
||||
2. **`seq` ist monoton und lückenlos.** Ein Reconnect mit `Last-Event-ID`
|
||||
liefert alles Verpasste nach.
|
||||
3. **Eine zu große Lücke wird ANGESAGT.** Ist das Verpasste aus dem Ringpuffer
|
||||
gefallen, schickt der Bus ausdrücklich einen neuen `snapshot` statt
|
||||
stillschweigend bei den neuesten Deltas weiterzumachen. Stillschweigend
|
||||
wäre der schlimmere Fall: Das UI hielte sich für aktuell und wäre es nicht.
|
||||
|
||||
## Treiber
|
||||
|
||||
memory.py asyncio-Warteschlangen + Ringpuffer — Standalone-Betrieb
|
||||
redis.py Redis Pub/Sub, Ringpuffer in der Tabelle `ereignisse` (V2-5)
|
||||
"""
|
||||
|
||||
from rippy.bus.schema import ( # noqa: F401
|
||||
EREIGNIS_TYPEN,
|
||||
baue_ereignis,
|
||||
ist_gueltig,
|
||||
)
|
||||
|
||||
__all__ = ["EREIGNIS_TYPEN", "baue_ereignis", "ist_gueltig"]
|
||||
@@ -0,0 +1,152 @@
|
||||
"""In-Process-Bus: asyncio-Warteschlangen plus Ringpuffer.
|
||||
|
||||
Erfüllt `rippy.ports.Bus`. Der Treiber für den Standalone-Betrieb — API,
|
||||
Queue und UI leben dort in EINEM Prozess, es braucht also nichts dazwischen.
|
||||
Der Redis-Treiber für den verteilten Betrieb kommt in V2-5 und hat dieselbe
|
||||
Form.
|
||||
|
||||
## Warum ein Ringpuffer und nicht nur Warteschlangen
|
||||
|
||||
Ein Browser verliert die SSE-Verbindung ständig — beim Sperrbildschirm, beim
|
||||
WLAN-Wechsel, beim Tab-Wechsel auf dem Handy. Ohne Puffer wäre jede dieser
|
||||
Sekunden ein Loch: Der Fortschritt spränge von 12 % auf 40 %, ein
|
||||
`job.finished` ginge verloren, und das UI zeigte einen Job als laufend, der
|
||||
längst durch ist.
|
||||
|
||||
Mit Puffer sagt der Browser beim Wiederverbinden „ich hatte zuletzt 1042", und
|
||||
bekommt 1043 ff. nachgeliefert.
|
||||
|
||||
## Und warum eine zu große Lücke angesagt wird
|
||||
|
||||
Ist das Verpasste aus dem Puffer gefallen, gibt es zwei Möglichkeiten:
|
||||
stillschweigend bei den neuesten Deltas weitermachen — oder sagen, dass etwas
|
||||
fehlt. Der erste Weg ist der gefährlichere: Das UI hielte sich für aktuell und
|
||||
wäre es nicht, und niemand könnte den Unterschied sehen. Deshalb kommt in dem
|
||||
Fall ein `snapshot`.
|
||||
|
||||
## Langsame Abonnenten
|
||||
|
||||
Jeder Abonnent hat eine eigene, BEGRENZTE Warteschlange. Läuft sie über — ein
|
||||
Browser auf einem lahmen Handy, ein Tab im Hintergrund —, wird der Abonnent
|
||||
markiert und bekommt beim nächsten Mal einen `snapshot` statt der verpassten
|
||||
Deltas. Was NICHT passiert: dass der Sender wartet. Ein einziger langsamer
|
||||
Zuhörer darf den Rip nicht bremsen.
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
from collections import deque
|
||||
|
||||
from rippy.bus.schema import RINGPUFFER, baue_ereignis
|
||||
|
||||
# Wie viele Ereignisse ein einzelner Abonnent aufstauen darf, bevor er als
|
||||
# „abgehängt" gilt. Großzügig: Ein normaler Client leert die Schlange sofort.
|
||||
ABONNENT_PUFFER = 256
|
||||
|
||||
|
||||
class MemoryBus:
|
||||
"""Ereignisse verteilen — ohne Broker, ohne Netz."""
|
||||
|
||||
def __init__(self, ringpuffer: int = RINGPUFFER):
|
||||
self._seq = 0
|
||||
self._puffer: deque = deque(maxlen=ringpuffer)
|
||||
self._abonnenten: set = set()
|
||||
|
||||
# ── Senden ──────────────────────────────────────────────────────────
|
||||
def senden(self, typ: str, daten: dict = None, entitaet: str = None,
|
||||
entitaet_id: str = None) -> dict:
|
||||
"""Nimmt ein Ereignis an, nummeriert es und verteilt es.
|
||||
|
||||
Bewusst SYNCHRON: Gesendet wird aus Rip-Schleifen und aus DB-Callbacks,
|
||||
also aus Code, der kein `await` kann. Die Zustellung ist ein
|
||||
`put_nowait` je Abonnent — es wartet nie jemand.
|
||||
"""
|
||||
self._seq += 1
|
||||
ereignis = baue_ereignis(
|
||||
typ, daten=daten, entitaet=entitaet,
|
||||
entitaet_id=entitaet_id, seq=self._seq,
|
||||
)
|
||||
self._puffer.append(ereignis)
|
||||
for abo in list(self._abonnenten):
|
||||
abo.zustellen(ereignis)
|
||||
return ereignis
|
||||
|
||||
# ── Empfangen ───────────────────────────────────────────────────────
|
||||
def abonnieren(self, ab_seq: int = None):
|
||||
"""Neuer Abonnent. `ab_seq` = die zuletzt gesehene Folgenummer."""
|
||||
abo = _Abonnent(self, ab_seq)
|
||||
self._abonnenten.add(abo)
|
||||
return abo
|
||||
|
||||
def abmelden(self, abo) -> None:
|
||||
self._abonnenten.discard(abo)
|
||||
|
||||
# ── Nachliefern ─────────────────────────────────────────────────────
|
||||
def nachliefern(self, ab_seq: int):
|
||||
"""Was seit `ab_seq` passiert ist — oder None, wenn die Lücke zu groß ist.
|
||||
|
||||
None heißt für den Aufrufer ausdrücklich: „Ich kann die Lücke nicht
|
||||
füllen, hol dir einen frischen Gesamtstand." Ein leeres Ergebnis
|
||||
dagegen heißt „nichts verpasst". Die beiden zu vermischen wäre genau
|
||||
der v1-Fehler in neuer Kleidung.
|
||||
"""
|
||||
if ab_seq is None:
|
||||
return []
|
||||
if not self._puffer:
|
||||
# Nichts im Puffer: Wenn der Client auf dem Stand ist, ist das
|
||||
# in Ordnung; liegt er zurück, fehlt ihm etwas.
|
||||
return [] if ab_seq >= self._seq else None
|
||||
aeltestes = self._puffer[0]["seq"]
|
||||
if ab_seq + 1 < aeltestes:
|
||||
return None # aus dem Puffer gefallen
|
||||
return [e for e in self._puffer if e["seq"] > ab_seq]
|
||||
|
||||
@property
|
||||
def seq(self) -> int:
|
||||
return self._seq
|
||||
|
||||
|
||||
class _Abonnent:
|
||||
"""Eine Warteschlange plus die Merkposten für einen Zuhörer."""
|
||||
|
||||
def __init__(self, bus: MemoryBus, ab_seq: int = None):
|
||||
self._bus = bus
|
||||
self._queue: asyncio.Queue = asyncio.Queue(maxsize=ABONNENT_PUFFER)
|
||||
self.ab_seq = ab_seq
|
||||
self.abgehaengt = False # Puffer lief über -> braucht snapshot
|
||||
|
||||
def zustellen(self, ereignis: dict) -> None:
|
||||
try:
|
||||
self._queue.put_nowait(ereignis)
|
||||
except asyncio.QueueFull:
|
||||
# Nicht warten, nicht das älteste stillschweigend wegwerfen:
|
||||
# merken, dass dieser Zuhörer den Anschluss verloren hat.
|
||||
self.abgehaengt = True
|
||||
|
||||
async def naechstes(self, timeout: float = None):
|
||||
"""Nächstes Ereignis, oder None bei Zeitüberschreitung.
|
||||
|
||||
Das Timeout ist der Herzschlag: Eine SSE-Verbindung, über die minutenlang
|
||||
nichts geht, wird von manchen Proxys geschlossen. Der Aufrufer schickt
|
||||
dann einen Kommentar-Frame.
|
||||
"""
|
||||
if timeout is None:
|
||||
return await self._queue.get()
|
||||
try:
|
||||
return await asyncio.wait_for(self._queue.get(), timeout)
|
||||
except asyncio.TimeoutError:
|
||||
return None
|
||||
|
||||
def schliessen(self) -> None:
|
||||
self._bus.abmelden(self)
|
||||
|
||||
def __enter__(self):
|
||||
return self
|
||||
|
||||
def __exit__(self, *_):
|
||||
self.schliessen()
|
||||
return False
|
||||
|
||||
|
||||
# Der Bus des Prozesses. Ein Modul-Singleton, weil ihn API, Queue und
|
||||
# Rip-Schleife gemeinsam benutzen und niemand ihn herumreichen soll.
|
||||
bus = MemoryBus()
|
||||
@@ -0,0 +1,84 @@
|
||||
"""Das Ereignis-Schema — eine Form für alle Nachrichten (KONZEPT-V2.md § 6.3).
|
||||
|
||||
Ein Ereignis sieht IMMER so aus:
|
||||
|
||||
{
|
||||
"seq": 1043, # monoton, lückenlos
|
||||
"ts": "2026-08-28T10:14:22Z",
|
||||
"typ": "job.progress",
|
||||
"entitaet": "job",
|
||||
"entitaet_id":"a4f…",
|
||||
"daten": {"phase": "rip", "prozent": 37, "eta_sekunden": 1820}
|
||||
}
|
||||
|
||||
## Warum die Typen hier als Liste stehen
|
||||
|
||||
Damit ein Tippfehler auffällt. Ein Ereignis mit `typ="job.progres"` würde sonst
|
||||
gesendet, käme im UI an und träfe dort auf keinen einzigen Empfänger — ohne
|
||||
Fehlermeldung, ohne Log, ohne Hinweis. Diese Liste macht daraus einen
|
||||
Programmfehler, der beim Bauen auffällt statt im Betrieb.
|
||||
|
||||
Wer einen Typ ergänzt, trägt ihn HIER ein. Der Test dazu prüft, dass jeder
|
||||
gesendete Typ bekannt ist.
|
||||
"""
|
||||
|
||||
from datetime import datetime, timezone
|
||||
|
||||
# Alle Ereignistypen, die es gibt. Kommentar = wann sie kommen.
|
||||
EREIGNIS_TYPEN = {
|
||||
# Zustand am Anfang jeder Verbindung — und nach einer zu großen Lücke.
|
||||
"snapshot": "Vollständiger Zustand: Jobs, Laufwerke, Knoten, Mounts",
|
||||
|
||||
"drive.changed": "Laufwerks-Zustand hat sich geändert",
|
||||
"disc.inserted": "Disc eingelegt",
|
||||
"disc.removed": "Disc entnommen",
|
||||
|
||||
"job.created": "Neuer Job angelegt",
|
||||
"job.phase": "Job wechselt die Phase (scan/rip/transcode/ablegen)",
|
||||
"job.progress": "Fortschritt — gedrosselt auf höchstens 1/s",
|
||||
"job.finished": "Job fertig, fehlgeschlagen oder abgebrochen",
|
||||
|
||||
"log.line": "Eine Log-Zeile",
|
||||
|
||||
"node.seen": "Ein Knoten hat sich gemeldet",
|
||||
"node.lost": "Ein Knoten meldet sich nicht mehr",
|
||||
|
||||
"mount.changed": "Ein Speicherziel ist erreichbar geworden oder weggefallen",
|
||||
|
||||
"system.notice": "Hinweis an den Nutzer (Plattenplatz, Rate-Limit, …)",
|
||||
}
|
||||
|
||||
# Wie viele Ereignisse der Ringpuffer hält. Bei ~1 Fortschritts-Ereignis pro
|
||||
# Sekunde sind 512 gut acht Minuten Rückschau — mehr als genug für einen
|
||||
# Browser-Reconnect (der binnen Sekunden passiert) und wenig genug, um im
|
||||
# Speicher nicht aufzufallen.
|
||||
RINGPUFFER = 512
|
||||
|
||||
|
||||
def ist_gueltig(typ: str) -> bool:
|
||||
return typ in EREIGNIS_TYPEN
|
||||
|
||||
|
||||
def baue_ereignis(typ: str, daten: dict = None, entitaet: str = None,
|
||||
entitaet_id: str = None, seq: int = 0, ts=None) -> dict:
|
||||
"""Baut ein Ereignis und prüft den Typ. Wirft bei unbekanntem Typ.
|
||||
|
||||
Absichtlich hart: Ein unbekannter Typ ist ein Programmfehler, kein
|
||||
Betriebszustand. Ihn durchzulassen hieße, eine Nachricht zu senden, die
|
||||
garantiert niemand empfängt — der stillste aller Fehlschläge.
|
||||
"""
|
||||
if not ist_gueltig(typ):
|
||||
bekannt = ", ".join(sorted(EREIGNIS_TYPEN))
|
||||
raise ValueError(
|
||||
f"Unbekannter Ereignistyp {typ!r}. Bekannt sind: {bekannt}. "
|
||||
"Neue Typen gehören in rippy/bus/schema.py — sonst kommt die "
|
||||
"Nachricht nirgendwo an und niemand merkt es."
|
||||
)
|
||||
return {
|
||||
"seq": seq,
|
||||
"ts": (ts or datetime.now(timezone.utc)).isoformat(),
|
||||
"typ": typ,
|
||||
"entitaet": entitaet,
|
||||
"entitaet_id": entitaet_id,
|
||||
"daten": daten or {},
|
||||
}
|
||||
@@ -0,0 +1,151 @@
|
||||
"""Der Bus muss drei Dinge können — und alle drei stammen aus v1-Fehlern.
|
||||
|
||||
1. **Nachliefern.** Ein Browser verliert die Verbindung ständig. Ohne
|
||||
Nachlieferung springt der Fortschritt und ein `job.finished` geht verloren.
|
||||
2. **Eine zu große Lücke ANSAGEN.** Stillschweigend bei den neuesten Deltas
|
||||
weiterzumachen wäre der gefährlichere Weg: Das UI hielte sich für aktuell
|
||||
und wäre es nicht.
|
||||
3. **Einen langsamen Zuhörer nicht zum Problem aller machen.** Ein Tab im
|
||||
Hintergrund darf den Rip nicht bremsen.
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
|
||||
import pytest
|
||||
|
||||
from rippy.bus import schema
|
||||
from rippy.bus.memory import ABONNENT_PUFFER, MemoryBus
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def bus():
|
||||
return MemoryBus(ringpuffer=10)
|
||||
|
||||
|
||||
# ── Schema ──────────────────────────────────────────────────────────────
|
||||
def test_unbekannter_typ_fliegt_sofort_auf(bus):
|
||||
"""Ein Tippfehler im Typ waere sonst der stillste aller Fehlschlaege:
|
||||
Die Nachricht ginge raus, traefe auf keinen Empfaenger, und nirgends
|
||||
stuende etwas."""
|
||||
with pytest.raises(ValueError, match="Unbekannter Ereignistyp"):
|
||||
bus.senden("job.progres", {"prozent": 5})
|
||||
|
||||
|
||||
def test_bekannter_typ_geht_durch(bus):
|
||||
ereignis = bus.senden("job.progress", {"prozent": 5}, entitaet="job", entitaet_id="j1")
|
||||
assert ereignis["typ"] == "job.progress"
|
||||
assert ereignis["entitaet_id"] == "j1"
|
||||
assert ereignis["seq"] == 1
|
||||
|
||||
|
||||
def test_alle_typen_im_schema_sind_benutzbar(bus):
|
||||
"""Wer einen Typ eintraegt, soll ihn auch senden koennen."""
|
||||
for typ in schema.EREIGNIS_TYPEN:
|
||||
bus.senden(typ, {})
|
||||
|
||||
|
||||
# ── Folgenummern ────────────────────────────────────────────────────────
|
||||
def test_seq_zaehlt_lueckenlos_hoch(bus):
|
||||
nummern = [bus.senden("log.line", {"text": str(i)})["seq"] for i in range(5)]
|
||||
assert nummern == [1, 2, 3, 4, 5]
|
||||
|
||||
|
||||
# ── Nachliefern ─────────────────────────────────────────────────────────
|
||||
def test_nachliefern_gibt_genau_das_verpasste(bus):
|
||||
for i in range(5):
|
||||
bus.senden("log.line", {"text": str(i)})
|
||||
verpasst = bus.nachliefern(2)
|
||||
assert [e["seq"] for e in verpasst] == [3, 4, 5]
|
||||
|
||||
|
||||
def test_nachliefern_ohne_rueckstand_ist_leer(bus):
|
||||
for i in range(3):
|
||||
bus.senden("log.line", {"text": str(i)})
|
||||
assert bus.nachliefern(3) == []
|
||||
|
||||
|
||||
def test_zu_grosse_luecke_gibt_none_nicht_leer(bus):
|
||||
"""DER Unterschied, auf den es ankommt.
|
||||
|
||||
`[]` heisst „du hast nichts verpasst". `None` heisst „ich kann die Luecke
|
||||
nicht fuellen, hol dir einen frischen Gesamtstand". Die beiden zu
|
||||
vermischen waere derselbe Fehler wie `catch(() => [])` im alten UI: eine
|
||||
Nichtauskunft, die als Aussage gelesen wird.
|
||||
"""
|
||||
for i in range(20): # Ringpuffer fasst nur 10
|
||||
bus.senden("log.line", {"text": str(i)})
|
||||
assert bus.nachliefern(2) is None
|
||||
assert bus.nachliefern(15) is not None
|
||||
|
||||
|
||||
def test_leerer_bus_meldet_rueckstand_ehrlich():
|
||||
frisch = MemoryBus(ringpuffer=10)
|
||||
assert frisch.nachliefern(0) == [] # nichts passiert, nichts verpasst
|
||||
assert frisch.nachliefern(None) == []
|
||||
|
||||
|
||||
# ── Zustellung ──────────────────────────────────────────────────────────
|
||||
def test_abonnent_bekommt_was_nach_dem_abo_kommt():
|
||||
async def lauf():
|
||||
bus = MemoryBus()
|
||||
with bus.abonnieren() as abo:
|
||||
bus.senden("job.created", {"id": "j1"}, entitaet="job", entitaet_id="j1")
|
||||
ereignis = await abo.naechstes(timeout=1)
|
||||
return ereignis
|
||||
|
||||
ereignis = asyncio.run(lauf())
|
||||
assert ereignis["typ"] == "job.created"
|
||||
|
||||
|
||||
def test_zwei_abonnenten_bekommen_beide_alles():
|
||||
async def lauf():
|
||||
bus = MemoryBus()
|
||||
with bus.abonnieren() as a, bus.abonnieren() as b:
|
||||
bus.senden("disc.inserted", {"laufwerk_id": "sr0"})
|
||||
return await a.naechstes(timeout=1), await b.naechstes(timeout=1)
|
||||
|
||||
erst, zweit = asyncio.run(lauf())
|
||||
assert erst["seq"] == zweit["seq"] == 1
|
||||
|
||||
|
||||
def test_abgemeldeter_abonnent_bekommt_nichts_mehr():
|
||||
async def lauf():
|
||||
bus = MemoryBus()
|
||||
abo = bus.abonnieren()
|
||||
abo.schliessen()
|
||||
bus.senden("log.line", {"text": "danach"})
|
||||
return await abo.naechstes(timeout=0.05)
|
||||
|
||||
assert asyncio.run(lauf()) is None
|
||||
|
||||
|
||||
def test_timeout_gibt_none_statt_zu_haengen():
|
||||
"""Der Herzschlag der SSE-Verbindung haengt daran: Kommt minutenlang
|
||||
nichts, muss der Aufrufer trotzdem einen Kommentar-Frame schicken
|
||||
koennen, sonst schliessen manche Proxys die Verbindung."""
|
||||
async def lauf():
|
||||
bus = MemoryBus()
|
||||
with bus.abonnieren() as abo:
|
||||
return await abo.naechstes(timeout=0.05)
|
||||
|
||||
assert asyncio.run(lauf()) is None
|
||||
|
||||
|
||||
# ── Langsamer Zuhörer ───────────────────────────────────────────────────
|
||||
def test_langsamer_abonnent_bremst_den_sender_nicht():
|
||||
"""Ein Tab im Hintergrund darf den Rip nicht anhalten.
|
||||
|
||||
Der Sender laeuft hier ueber die Puffergrenze hinaus. Erwartet wird: Er
|
||||
kommt durch, und der abgehaengte Zuhoerer ist MARKIERT — nicht, dass
|
||||
Ereignisse stillschweigend verschwinden.
|
||||
"""
|
||||
async def lauf():
|
||||
bus = MemoryBus(ringpuffer=ABONNENT_PUFFER * 2)
|
||||
with bus.abonnieren() as abo:
|
||||
for i in range(ABONNENT_PUFFER + 20):
|
||||
bus.senden("log.line", {"text": str(i)}) # niemand liest
|
||||
return abo.abgehaengt, bus.seq
|
||||
|
||||
abgehaengt, seq = asyncio.run(lauf())
|
||||
assert abgehaengt is True
|
||||
assert seq == ABONNENT_PUFFER + 20 # der Sender ist durchgelaufen
|
||||
Reference in New Issue
Block a user