diff --git a/docker/api/main.py b/docker/api/main.py index bca56f5..2e5419e 100644 --- a/docker/api/main.py +++ b/docker/api/main.py @@ -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"} diff --git a/docker/api/test_api_smoke.py b/docker/api/test_api_smoke.py index cbace9e..89f9df8 100644 --- a/docker/api/test_api_smoke.py +++ b/docker/api/test_api_smoke.py @@ -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) diff --git a/src/rippy/bus/__init__.py b/src/rippy/bus/__init__.py new file mode 100644 index 0000000..4e321cb --- /dev/null +++ b/src/rippy/bus/__init__.py @@ -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"] diff --git a/src/rippy/bus/memory.py b/src/rippy/bus/memory.py new file mode 100644 index 0000000..96854f0 --- /dev/null +++ b/src/rippy/bus/memory.py @@ -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() diff --git a/src/rippy/bus/schema.py b/src/rippy/bus/schema.py new file mode 100644 index 0000000..9e87b86 --- /dev/null +++ b/src/rippy/bus/schema.py @@ -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 {}, + } diff --git a/src/rippy/bus/test_memory.py b/src/rippy/bus/test_memory.py new file mode 100644 index 0000000..5c46d2c --- /dev/null +++ b/src/rippy/bus/test_memory.py @@ -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