"""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())