From b2610e35e5d00675e1fd8944422b12af8d29da49 Mon Sep 17 00:00:00 2001 From: HMS MediaEngine Agent Date: Fri, 11 Sep 2026 01:33:47 +0200 Subject: [PATCH] =?UTF-8?q?Phase=202=20abgeschlossen:=20Renderer-State-Syn?= =?UTF-8?q?c=20=C3=BCber=20echte=20IPC=20(=C2=A76.2,=20=C2=A76.4)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - RemoteStateMirror (Renderer): Snapshot-Pflicht nach (Re-)Connect, Deltas mit neu/geaendert/geloescht, Deltas vor erstem Snapshot werden abgelehnt (Re-Sync), Duplikate/Veraltete idempotent - RendererStateLink (Control Core): connect_and_sync mit Pflicht-Snapshot, sync_if_changed versendet Deltas seit letzter gesendeter Revision, Reconnect loest immer neuen Snapshot aus (kein halber Zustand) - End-to-End-Integrationstests ueber echtes TCP-Loopback ohne Mocks: Snapshot->Delta-Kette, Reconnect-Re-Sync, Delta-vor-Snapshot-Ablehnung, Duplikat-No-Op - Gesamtsuite 297 gruen, Ruff gruen --- STATUS.md | 16 +- .../hms_control_server/renderer_link.py | 90 +++++++++ apps/renderer/hms_renderer/__init__.py | 8 +- apps/renderer/hms_renderer/state_mirror.py | 96 ++++++++++ tests/integration/test_state_sync_chain.py | 178 ++++++++++++++++++ 5 files changed, 382 insertions(+), 6 deletions(-) create mode 100644 apps/control_server/hms_control_server/renderer_link.py create mode 100644 apps/renderer/hms_renderer/state_mirror.py create mode 100644 tests/integration/test_state_sync_chain.py diff --git a/STATUS.md b/STATUS.md index a2975f1..2ca3276 100644 --- a/STATUS.md +++ b/STATUS.md @@ -83,13 +83,21 @@ Phase 2 (im Bau): ab 1 s Fenster, Showzeit→lokale-Zeit-Abbildung (§6.4) - [x] zeitgestempelte Preset-Aktivierung: Vorlauf 200 ms, Arm/Execute/ Ack, FAILED bei fehlender Arm-Bestätigung (§6.5) -- [ ] Renderer-Anbindung: IPC-Deltas aus dem State-Store in den Renderer +- [x] Renderer-Anbindung: RemoteStateMirror (Renderer) + RendererStateLink + (Control Core) über IPC; Pflicht-Snapshot nach (Re-)Connect, danach + Deltas (neu/geändert/gelöscht), Deltas vor Snapshot abgelehnt, + Duplikate idempotent; End-to-End-Integrationstest über echtes + TCP-Loopback ohne Mocks (§6.2, §6.4) ## Nächste drei Aufgaben -1. Renderer-Anbindung: IPC-Deltas aus dem State-Store (§6.4) -2. Phase-3-Vorbereitung: Plugin-Loader/Lifecycle und Backend-Adapter (§14) -3. Echtes mDNS mit zeroconf auf Zielsystemen (Gate-1-LAN-Test) +1. Phase 3: Plugin-Loader/Lifecycle (discovered→…→active/quarantined) und + Registry mit Backend-Adapter-Vertrag (§14) +2. Pflicht-Generatoren: hms.generator.solid, gradient, checker_grid, + stripes_chaser, noise_clouds, plasma, wave_bars, shapes, drops_ripples, + starfield (§15.1) +3. Pflicht-Filter: hms.fx.transform2d … feedback_trails mit HLSL/GLSL/GLES + (§15.2) ## Ausstehende Hardware-Validierung (ADR-0008 Testplan) diff --git a/apps/control_server/hms_control_server/renderer_link.py b/apps/control_server/hms_control_server/renderer_link.py new file mode 100644 index 0000000..3e6503e --- /dev/null +++ b/apps/control_server/hms_control_server/renderer_link.py @@ -0,0 +1,90 @@ +"""RendererStateLink auf der Control-Core-Seite (PLAN.md §6.2, §6.4). + +Verbindet ProjectStateStore (autoritativ) mit dem Renderer über IPC: +- connect_and_sync(): Handshake + vollständiger Snapshot (§6.2 Pflicht) +- sync_if_changed(): Delta seit letzter gesendeter Revision +- Reconnect: immer neuer Snapshot, danach erst wieder Deltas (§6.2) + +Der Link kennt den letzten beim Renderer angekommenen Zustand und +berechnet Deltas daraus – keine halben Zustände (§11.4, §6.4). +""" + +from __future__ import annotations + +from hms_persistence.state_store import ProjectStateStore +from hms_protocol import Envelope, IpcClient, MessageType + + +class RendererStateLink: + """Synchronisiert den Showzustand Control Core → Renderer über IPC.""" + + def __init__(self, client: IpcClient, store: ProjectStateStore) -> None: + self._client = client + self._store = store + self._last_sent_revision = 0 + self._renderer_known_state: dict[str, float] = {} + + @property + def last_sent_revision(self) -> int: + return self._last_sent_revision + + @property + def in_sync(self) -> bool: + """True, wenn der Renderer die aktuelle Revision besitzt.""" + return self._last_sent_revision == self._store.state_revision + + # ---------- Verbindung (§6.2) ---------- + + async def connect_and_sync(self) -> None: + """Verbindung aufbauen und vollständigen Snapshot senden. + + Nach jedem (Re-)Connect wird immer zuerst der vollständige Snapshot + übertragen; Deltas folgen erst danach (§6.2: „Re-Sync nach + Reconnect", „Deltas erst nach erfolgreichem Re-Sync akzeptiert"). + """ + await self._client.connect() + await self.send_full_snapshot() + + async def send_full_snapshot(self) -> int: + """Sendet den vollständigen Zustand; liefert die gesendete Revision.""" + snap = self._store.snapshot() + self._last_sent_revision = snap["state_revision"] + self._renderer_known_state = dict(snap["values"]) + envelope = Envelope( + type=MessageType.SNAPSHOT, + revision=snap["state_revision"], + payload=snap, + ) + await self._client.send(envelope) + return snap["state_revision"] + + # ---------- Delta-Versand (§6.4) ---------- + + async def sync_if_changed(self) -> bool: + """Sendet ein Delta, falls sich die Revision seit dem letzten Versand + geändert hat. Rückgabe: True, wenn etwas gesendet wurde. + + Das Delta wird aus dem zuletzt bekannten Renderer-Zustand berechnet + (neu/geändert/gelöscht); die Semantik folgt StateDelta (§6.4). + """ + if self._store.state_revision == self._last_sent_revision: + return False # nichts Neues + delta = self._store.delta_since( + self._last_sent_revision, self._renderer_known_state + ) + if delta is None: + return False + # Buchhaltung: was weiß der Renderer ab jetzt? + for path, value in delta.changes.items(): + if value is None: + self._renderer_known_state.pop(path, None) + else: + self._renderer_known_state[path] = value + self._last_sent_revision = delta.state_revision + envelope = Envelope( + type=MessageType.EVENT, + revision=delta.state_revision, + payload=delta.to_dict(), + ) + await self._client.send(envelope) + return True diff --git a/apps/renderer/hms_renderer/__init__.py b/apps/renderer/hms_renderer/__init__.py index e214e63..6fee5ae 100644 --- a/apps/renderer/hms_renderer/__init__.py +++ b/apps/renderer/hms_renderer/__init__.py @@ -1,8 +1,9 @@ -"""hms_renderer – Render-Worker-Spike (PLAN.md §6.1C, §12, §36 Nr. 4–6). +"""hms_renderer – Render-Worker (PLAN.md §6.1C, §12, §36 Nr. 4–6). Python orchestriert native GStreamer-Komponenten; keine Pixelverarbeitung in Python (§2.1, §33). Pipelines werden als gst-launch-Strings definiert -und auf dem Zielsystem ausgeführt/messbar. +und auf dem Zielsystem ausgeführt/messbar. Der RemoteStateMirror hält den +über IPC übermittelten Showzustand (§6.2, §6.4). """ from hms_renderer.pipelines import ( @@ -12,6 +13,7 @@ from hms_renderer.pipelines import ( build_single_video_pipeline, gst_available, ) +from hms_renderer.state_mirror import RemoteStateMirror, apply_envelope __all__ = [ "D3D11Pipeline", @@ -19,4 +21,6 @@ __all__ = [ "build_single_video_pipeline", "build_compositor_pipeline", "gst_available", + "RemoteStateMirror", + "apply_envelope", ] diff --git a/apps/renderer/hms_renderer/state_mirror.py b/apps/renderer/hms_renderer/state_mirror.py new file mode 100644 index 0000000..606caab --- /dev/null +++ b/apps/renderer/hms_renderer/state_mirror.py @@ -0,0 +1,96 @@ +"""Remote-State-Mirror auf der Renderer-Seite (PLAN.md §6.2, §6.4). + +Der Renderer hält den zuletzt übermittelten Showzustand als lokales +Spiegelbild: +- nach Verbindung: vollständiger Snapshot (§6.2 Pflicht) +- danach inkrementelle Deltas mit monotoner Revision +- Deltas werden erst nach erfolgtem Snapshot akzeptiert (§6.2: + „Deltas erst nach erfolgreichem Re-Sync") +- veraltete/doppelte Deltas sind idempotente No-Ops (§6.5) + +Der Mirror ist reine Zustandslogik ohne Pixelbezug (§33); der Rendergraph +liest Werte über get_value() je Frame. +""" + +from __future__ import annotations + + +class RemoteStateMirror: + """Spiegel des autoritativen Showzustands im Renderer.""" + + def __init__(self) -> None: + self._revision = 0 + self._values: dict[str, float] = {} + self._has_snapshot = False + + @property + def revision(self) -> int: + return self._revision + + @property + def has_snapshot(self) -> bool: + """True, nachdem ein vollständiger Snapshot empfangen wurde.""" + return self._has_snapshot + + def get_value(self, path: str, default: float | None = None) -> float | None: + """Wirksamer Wert für einen Parameterpfad (§10.2).""" + return self._values.get(path, default) + + def all_values(self) -> dict[str, float]: + return dict(self._values) + + # ---------- Snapshot / Delta (§6.2, §6.4) ---------- + + def apply_snapshot(self, payload: dict) -> None: + """Übernimmt einen vollständigen Snapshot (nach Verbindung/Reconnect).""" + values = payload.get("values", {}) + if not isinstance(values, dict): + raise ValueError("Snapshot ohne Werte-Objekt") + self._revision = int(payload["state_revision"]) + self._values = {str(k): float(v) for k, v in values.items()} + self._has_snapshot = True + + def apply_delta(self, payload: dict) -> bool: + """Wendet ein Delta an. + + Rückgabe: + - True: Delta angewendet ODER als veraltetes Duplikat ignoriert + - False: Re-Sync nötig (noch kein Snapshot empfangen, §6.2) + + Gelöschte Pfade sind als None kodiert (StateDelta-Vertrag). + """ + if not self._has_snapshot: + return False # §6.2: Deltas erst nach Snapshot + delta_revision = int(payload["state_revision"]) + if delta_revision <= self._revision: + return True # idempotent: Duplikat/veraltet, kein Handlungsbedarf + changes = payload.get("changes", {}) + if not isinstance(changes, dict): + raise ValueError("Delta ohne Changes-Objekt") + for path, value in changes.items(): + if value is None: + self._values.pop(str(path), None) # gelöscht + else: + self._values[str(path)] = float(value) + self._revision = delta_revision + return True + + +def apply_envelope(mirror: RemoteStateMirror, envelope) -> bool: + """Verarbeitet ein IPC-Envelope in den Mirror. + + - SNAPSHOT: vollständige Übernahme + - EVENT mit changes: Delta-Anwendung + - alles andere (Heartbeat, Ack, …): keine Zustandswirkung + + Rückgabe: True, wenn der Zustand dadurch (re-)synchronisiert wurde; + False, wenn ein Re-Sync (neuer Snapshot) angefordert werden muss. + """ + from hms_protocol import MessageType + + if envelope.type is MessageType.SNAPSHOT: + mirror.apply_snapshot(envelope.payload) + return True + if envelope.type is MessageType.EVENT and "changes" in envelope.payload: + return mirror.apply_delta(envelope.payload) + return True # keine Zustandsnachricht: nichts zu tun diff --git a/tests/integration/test_state_sync_chain.py b/tests/integration/test_state_sync_chain.py new file mode 100644 index 0000000..424425e --- /dev/null +++ b/tests/integration/test_state_sync_chain.py @@ -0,0 +1,178 @@ +"""End-to-End State-Sync über echte IPC (PLAN.md §6.4, §29.2). + +Kette ohne Mocks: ProjectStateStore → RendererStateLink → IpcClient → +TCP-Loopback → IpcServer → RemoteStateMirror. Prüft Snapshot-Pflicht, +Delta-Anwendung (neu/geändert/gelöscht), Reconnect-Re-Sync und die +Ablehnung von Deltas vor dem ersten Snapshot (§6.2). +""" + +from __future__ import annotations + +import asyncio + +from hms_control_server.renderer_link import RendererStateLink +from hms_persistence.state_store import ProjectStateStore +from hms_protocol import IpcClient, IpcServer, MessageType +from hms_renderer import RemoteStateMirror, apply_envelope + + +def _run(coro): + return asyncio.run(coro) + + +async def _setup(): + """Server + Client + Store + Link + Mirror, verbunden und gesynct.""" + server = IpcServer() + port = await server.start() + client = IpcClient(port) + store = ProjectStateStore() + link = RendererStateLink(client, store) + mirror = RemoteStateMirror() + return server, client, store, link, mirror + + +async def _drain_into_mirror(server: IpcServer, mirror: RemoteStateMirror, count: int): + """Liest Nachrichten bis `count` STATE-Nachrichten im Mirror ankamen. + + Heartbeats werden übersprungen, ohne den Zähler zu verbrauchen (das + Heartbeat-Intervall von 500 ms kann otherwise mit Testnachrichten + interleaven).""" + applied = [] + state_seen = 0 + while state_seen < count: + env = await asyncio.wait_for(server.receive(), timeout=2.0) + assert env is not None, "Verbindung vorzeitig geschlossen" + if env.type in (MessageType.SNAPSHOT, MessageType.EVENT): + ok = apply_envelope(mirror, env) + applied.append((env.type, env.revision, ok)) + state_seen += 1 + # Heartbeats und andere: überspringen, nicht mitzählen + return applied + + +def test_full_sync_chain_snapshot_then_deltas() -> None: + """Verbindung → Snapshot → Änderungen → Delta → Mirror stimmt überein.""" + + async def impl() -> None: + server, client, store, link, mirror = await _setup() + try: + # Ausgangszustand im Store + store.set_value("master/intensity", 1.0) + store.set_value("composition/x/layer/y/opacity", 0.4) + + # 1) Verbindung + Pflicht-Snapshot (§6.2) + await link.connect_and_sync() + assert link.last_sent_revision == 2 + applied = await _drain_into_mirror(server, mirror, 1) + assert applied[0][0] is MessageType.SNAPSHOT + assert mirror.has_snapshot + assert mirror.revision == 2 + assert mirror.get_value("master/intensity") == 1.0 + assert mirror.get_value("composition/x/layer/y/opacity") == 0.4 + + # 2) Änderungen im Store → Delta + store.set_value("master/intensity", 0.7) # ändern + store.set_value("master/speed", 1.25) # neu + store.clear_value("composition/x/layer/y/opacity") # löschen + assert await link.sync_if_changed() is True + await _drain_into_mirror(server, mirror, 1) + assert mirror.revision == 5 + assert mirror.get_value("master/intensity") == 0.7 + assert mirror.get_value("master/speed") == 1.25 + assert mirror.get_value("composition/x/layer/y/opacity") is None + assert link.in_sync + + # 3) Keine Änderung → kein Delta + assert await link.sync_if_changed() is False + finally: + await client.disconnect() + await server.stop() + + _run(impl()) + + +def test_reconnect_sends_fresh_snapshot() -> None: + """Nach Reconnect: neuer Snapshot, danach erst wieder Deltas (§6.2).""" + + async def impl() -> None: + server, client, store, link, mirror = await _setup() + try: + store.set_value("master/intensity", 0.9) + await link.connect_and_sync() + await _drain_into_mirror(server, mirror, 1) + assert mirror.revision == 1 + + # Trennung simulieren und erneut verbinden + await client.disconnect() + store.set_value("master/intensity", 0.3) + await link.connect_and_sync() # Re-Sync: Snapshot statt Delta + applied = await _drain_into_mirror(server, mirror, 1) + assert applied[0][0] is MessageType.SNAPSHOT # kein halber Zustand + assert mirror.revision == 2 + assert mirror.get_value("master/intensity") == 0.3 + assert link.in_sync + finally: + await client.disconnect() + await server.stop() + + _run(impl()) + + +def test_delta_before_snapshot_requires_resync() -> None: + """§6.2: Deltas vor dem ersten Snapshot werden abgelehnt (Re-Sync nötig).""" + + async def impl() -> None: + server, client, store, link, mirror = await _setup() + try: + # Delta-artige Nachricht ohne vorherigen Snapshot direkt senden + from hms_protocol import Envelope + + fake_delta = Envelope( + type=MessageType.EVENT, + revision=5, + payload={ + "state_revision": 5, + "project_revision": 1, + "changes": {"master/intensity": 0.5}, + "monotonic_ns": 0, + }, + ) + await client.connect() + await client.send(fake_delta) + env = await asyncio.wait_for(server.receive(), timeout=2.0) + assert env is not None + assert apply_envelope(mirror, env) is False # Re-Sync erforderlich + assert not mirror.has_snapshot + assert mirror.revision == 0 # unverändert + finally: + await client.disconnect() + await server.stop() + + _run(impl()) + + +def test_stale_duplicate_delta_is_idempotent_noop() -> None: + """Veraltete/doppelte Deltas ändern nichts (§6.5 Idempotenz).""" + + async def impl() -> None: + server, client, store, link, mirror = await _setup() + try: + store.set_value("master/intensity", 1.0) + await link.connect_and_sync() + await _drain_into_mirror(server, mirror, 1) + + # Delta triggern und zweimal dieselbe Nachricht zustellen + store.set_value("master/intensity", 0.6) + assert await link.sync_if_changed() is True + env = await asyncio.wait_for(server.receive(), timeout=2.0) + assert env is not None + assert apply_envelope(mirror, env) is True + # Duplikat: dieselbe Revision erneut → No-Op, kein Fehler + assert apply_envelope(mirror, env) is True + assert mirror.revision == 2 # bleibt + assert mirror.get_value("master/intensity") == 0.6 + finally: + await client.disconnect() + await server.stop() + + _run(impl())