diff --git a/STATUS.md b/STATUS.md index 55f131a..9e00720 100644 --- a/STATUS.md +++ b/STATUS.md @@ -1,40 +1,77 @@ # STATUS -Stand: 2026-09-10 +Stand: 2026-09-11 ## Aktuelle Phase -**Phase 0 – technischer Spike und Go/No-Go** (PLAN.md §31, §36) +**Phase 1 – Fundament und Netzwerkbasis** unter der **Build-first-Strategie +(ADR-0008)**: Die Phasen 1–5 werden plattformneutral vollständig gebaut; +Gate 0 und alle renderer-nahen Abnahmen werden gesammelt in einem späteren +Windows-Durchlauf validiert. Auftraggeber-Freigabe für die Reihenfolgeabweichung +vom 2026-09-11, dokumentiert in ADR-0008 gemäß PLAN.md §1. ## Letzter grüner Commit -- 0922cc1 – Phase-0-Grundlage: 121 Unit-/Integrationstests grün, Ruff grün (Belege: `TEST_REPORT.md`) +- siehe `git log` und `TEST_REPORT.md` – jede Etappe endet mit grünem pytest + und Ruff und wird sofort nach Forgejo gepusht. ## Bestandene Gates -- keine; **Gate 0 ist offen** +- keine; **Gate 0 bleibt offen** bis zum Windows-Hardware-Durchlauf + (Testplan in ADR-0008) + +## Arbeitsweise (ADR-0008) + +- Build-first: plattformneutral fertig bauen, Hardware-Validierung gesammelt + am Ende (Windows-Durchlauf). +- Kein Gate wird als grün gemeldet, solange der Windows-Durchlauf offen ist. +- Renderer-Nähe ausschließlich über den Backend-Vertrag (§12.6); Korrekturen + bleiben lokal im Adapter. ## Laufende Arbeit -Erster Arbeitsauftrag (§36 Nr. 1–3 erledigt, Nr. 4–15 in Arbeit): +Phase 0 (abgeschlossen, soweit ohne Hardware möglich): -- [x] Repository gemäß Eigentumsgrenzen initialisiert (§8) -- [x] Pflichtdokumente + ADR-Vorlage + ADRs 0001–0003 angelegt -- [x] GStreamer-Version für Windows gepinnt: **1.28.6** (`build/windows/GSTREAMER.md`, ADR-0002) -- [x] Kernpakete mit Unit-Tests: IPC-Protokoll (§6.2), Parameter-Engine (§11), Art-Net-Pakete/Empfänger (§16), Adaptive Quality (§5.2), Capability-Probe (§5.2), Plugin-SDK-Validierung (§14.5, §27.2) -- [x] Minimaler Control Core: FastAPI-REST + WebSocket für denselben Parametersatz (§36 Nr. 9) -- [x] Renderer-Spike: D3D11-/Dev-GL-Pipeline-Definitionen + CLI (§36 Nr. 4–5); ohne GStreamer-Installation kontrollierter Abbruch (Exit-Code 2), keine Erfolgssimulation (§33) -- [x] Beispielplugin Passthrough (HLSL + GLSL + GLES) als SDK-Referenz (§36 Nr. 6) -- [x] Tools: Art-Net-Emulator, Capability-Probe, Plugin/Shader-Validator +- [x] Repository, Pflichtdokumente, ADRs 0001–0003 +- [x] GStreamer-Pin 1.28.6, Kernpakete, Control Core, Renderer-Spike +- [x] Beispielplugins (Passthrough, Gaussian Blur mit 3 AQ-Varianten), Tools +- [x] 121 Unit-/Integrationstests grün, Ruff grün (Belege: TEST_REPORT.md) + +Phase 1 (im Bau): + +- [x] IPC-Verbindung Control Core ↔ Renderer: Handshake (hello/welcome), + vollständiger Snapshot nach Verbindung, Deltas mit monotoner Revision, + Command-Acks idempotent über message_id, Heartbeat 500 ms in beide + Richtungen, Re-Sync nach Reconnect, ausschließlich Loopback (§6.2) +- [ ] SQLite-Persistenz (WAL, Migrationen, Backup, §24) +- [ ] Launcher/Supervisor (Prozessstart, Heartbeat-Überwachung, kontrolliertes + Beenden, Crash-Erkennung) +- [ ] Persistente Node-Identität + Rollen in der App-Verkabelung +- [ ] mDNS/DNS-SD-Discovery + manuelle Node-Liste (§6.3) +- [ ] Paarung (PIN/Fingerprint) und Node-Vertrauen +- [ ] Cluster-Nachrichtenhülle mit Sequenzen/Revisionen (§6.5) ## Nächste drei Aufgaben -1. **Gate-0-Messungen auf Referenz-Windows-Hardware:** D3D11-Hardwaredecode, durchgängiger `D3D11Memory`-Pfad ohne CPU-Readback, Framezeit p99, DMX-Latenz (≤ 2 Frames), ruckelfreier Adaptive-Quality-Wechsel (§25, §36 Nr. 11) -2. Portable Onefolder-Ausgabe erzeugen und auf sauberem Windows-Rechner testen (Nuitka vs. PyInstaller → ADR; §36 Nr. 12–13) -3. Nativen Renderkern festlegen (Rust vs. C++, Bridge vs. GStreamer-Plugin → ADR-0004) und Renderer über IPC-Handshake an die Parameter-Engine anbinden (Phase 1) +1. SQLite-Persistenz mit WAL, transaktionalem Autosave und Migrationstests (§24) +2. Launcher/Supervisor mit Heartbeat-Überwachung und kontrolliertem Beenden (§6.1A) +3. Discovery (mDNS) plus manuelle Fallback-Liste und Paarungsgrundlage (§6.3) + +## Ausstehende Hardware-Validierung (ADR-0008 Testplan) + +Nur auf echtem Windows mit GPU messbar; nichts davon gilt als erledigt: + +- D3D11-Hardwaredecode aktiv; D3D11Memory durchgängig ohne CPU-Rundweg +- Framezeit p50/p95/p99; DMX → sichtbarer Frame p95 ≤ 2 Frames +- Adaptive-Quality-Wechsel ohne Stall; Golden Images; Shader-Compile auf GPU +- Portable Onefolder-Ausgabe auf sauberem Rechner (§29.6) +- Rust-Renderkern-Kompilierung und GStreamer-Plugin-Bindung (ADR-0004) ## Bekannte Blocker -- **Keine Windows-Referenzhardware in der Entwicklungsumgebung** (Linux-Container, CPU-only). Alle Gate-0-Kriterien sind ausschließlich auf echter Hardware gültig (§29.7, §33). Spike-Code ist bereit; Messungen und Portabilitätstest stehen aus. -- pnpm/Node-Frontend noch nicht eingerichtet (Phase 5, ADR-0006 offen). -- HLSL-Live-Parameter im D3D11-Pfad erfordert den nativen Renderkern (ADR-0004 offen); GLSL-Testvariante für den Dev-Pfad liegt bei. +- Windows-Referenzhardware fehlt in der Entwicklungsumgebung + (Linux-Container, CPU-only) – Kernblocker, durch ADR-0008 gemanagt. +- Rust-Toolchain (cargo) fehlt im Container: Rust-Quellcode wird entwickelt, + das Kompilierungsgate erfolgt im Windows-Durchlauf. +- pnpm fehlt im Container: wird für Phase 5 über Node-Corepack aktiviert + (kein Installationsblocker). diff --git a/docs/adr/0004-renderkern-rust-gstreamer.md b/docs/adr/0004-renderkern-rust-gstreamer.md new file mode 100644 index 0000000..97db1ec --- /dev/null +++ b/docs/adr/0004-renderkern-rust-gstreamer.md @@ -0,0 +1,47 @@ +# ADR-0004: Nativer Renderkern – Rust auf dem GStreamer-D3D11-Elementpfad + +- **Status:** Vorläufig angenommen (Bestätigung oder Revision im Windows-Durchlauf) +- **Datum:** 2026-09-11 +- **Phase:** ursprünglich Phase 0; wegen ADR-0008 vorgezogen +- **Bauplan:** §6.1C, §7.1, §12.6 + +## Entscheidung + +1. Sprache des nativen Moduls: **Rust**. +2. Einbindung: **GStreamer-Plugin-Ansatz** – der Rendergraph nutzt die + erprobten D3D11-Elemente (d3d11h264dec → d3d11convert → d3d11compositor → + d3d11videosink) und eigene Rust-GStreamer-Elemente für Effekt-/Shader-Pässe, + statt sofort eine komplett eigenständige Render-Bridge mit eigener + Swapchain zu bauen. + +## Begründung (ohne Messdaten, deshalb vorläufig) + +- D3D11Memory-Residenz ist Eigenschaft der GStreamer-Elemente; das Risiko + eigengebauten Swapchain-Compositings entfällt. +- gstreamer-rs bietet stabile Bindings; Rust liefert Gedächtnissicherheit im + nativen Pfad. +- Die Effektkette skaliert als Elementkette; der Plugin-Vertrag (§14) bleibt + vollständig gewahrt. +- Eine spätere wgpu-/eigenständige Bridge bleibt über dasselbe IPC- und + Capability-Protokoll anschließbar (§6.1C). + +## Alternativen + +- C++: kein Vorteil, höheres Fehlerrisiko im Speichermanagement. +- Sofortige eigenständige Bridge: mehr Kontrolle, aber hohes Risiko ohne + Hardware-Feedback (ADR-0008). + +## Folgen + +- Der Renderer-Prozess orchestriert Pipelines; Rust-Elemente übernehmen die + Shader-Pässe (HLSL aus dem Plugin-Vertrag). +- Rust-Kompilierung erfolgt im Windows-Durchlauf (Container ohne cargo). + +## Messwerte / Nachweise + +- Ausstehend. Gate-0-Messungen bestätigen das Elementpfad-Budget oder lösen + eine Revision (eigenständige D3D11-Bridge) aus – ohne API-Änderung (§12.6). + +## Freigabe + +Vorläufig gemäß ADR-0008; endgültig nach Windows-Messung. diff --git a/docs/adr/0005-packaging-nuitka-onefolder.md b/docs/adr/0005-packaging-nuitka-onefolder.md new file mode 100644 index 0000000..c1c1b90 --- /dev/null +++ b/docs/adr/0005-packaging-nuitka-onefolder.md @@ -0,0 +1,33 @@ +# ADR-0005: Packaging – Nuitka Onefolder (vorläufig) + +- **Status:** Vorläufig angenommen (Vergleich im Windows-Durchlauf) +- **Datum:** 2026-09-11 +- **Phase:** ursprünglich Phase 0; wegen ADR-0008 vorgezogen +- **Bauplan:** §7.1, §30.2 + +## Entscheidung + +**Nuitka Standalone/Onefolder** als Packaging-Basis; **PyInstaller Onefolder** +als dokumentierter Fallback, falls Nuitka die GStreamer-DLL-Bündelung nicht +sauber abbildet. Onefile bleibt wie geplant ausgeschlossen (§7.1). + +## Begründung + +- Nuitka kompiliert zu C: weniger Interpreter-Rest, bessere Reproduzierbarkeit + per Buildskript (§30.1). +- Die GStreamer-Bündelung ist von der Packaging-Wahl unabhängig + (runtime/-Ordner + GST_PLUGIN_PATH, siehe build/windows/GSTREAMER.md). +- PLAN.md verlangt den finalen Vergleich auf Windows; hier nur vorläufige Wahl. + +## Folgen + +- Buildskripte targetieren Nuitka; der Fallback-Pfad bleibt gepflegt. + +## Messwerte / Nachweise + +- Ausstehend; der erste Onefolder-Build auf sauberem Windows-Rechner + entscheidet final (§29.6). + +## Freigabe + +Vorläufig gemäß ADR-0008; endgültig nach Windows-Build. diff --git a/docs/adr/0008-build-first-strategie.md b/docs/adr/0008-build-first-strategie.md new file mode 100644 index 0000000..4baac2d --- /dev/null +++ b/docs/adr/0008-build-first-strategie.md @@ -0,0 +1,53 @@ +# ADR-0008: Build-first-Strategie – verzögerte Hardware-Validierung + +- **Status:** Angenommen (Auftraggeber-Freigabe) +- **Datum:** 2026-09-11 +- **Phase:** 0/1 (übergreifend) +- **Bauplan:** §1, §1.1 Nr. 1/9/10, §29.7, §36 Nr. 16 + +## Kontext + +Die Entwicklungsumgebung ist ein CPU-only-Linux-Container ohne Windows-GPU. +Der Auftraggeber wünscht am 2026-09-11 ausdrücklich: „erst fertig bauen und dann +auf Windows testen". PLAN.md §1 erlaubt Abweichungen vom normativen Ablauf, +wenn sie per ADR dokumentiert, technisch begründet und vom Auftraggeber +freigegeben werden. Diese Freigabe liegt mit der Anfrage vor. + +## Entscheidung + +1. Die Phasen 1–5 werden **plattformneutral vollständig gebaut** (Code, Tests + auf CPU-Ebene, Schemas, Dokumentation), bevor Hardwaremessungen stattfinden. +2. Gate 0 und alle renderer-nahen Abnahmen werden **gesammelt in einem + Windows-Durchlauf** nachgeholt (reverse validation). Reihenfolge dort: + Portabilität → Decode/Residenz → Framezeit/DMX-Latenz → Adaptive Quality → + Golden Images → Onefolder-Build. +3. Kein Gate wird vorher als grün gemeldet; STATUS.md führt die ausstehende + Hardware-Validierung offen als Checkliste. + +## Alternativen + +- Streng planmäßig (Gate 0 zuerst): sicherer, aber ohne Windows-Zugang blockiert; + vom Auftraggeber verworfen. +- Windows-Cloud-VM in der Entwicklungsumgebung: hier nicht verfügbar. + +## Folgen und Risiken + +- ADR-0004/0005 müssen ohne Messdaten vorläufig entschieden werden → + ausdrücklicher Bestätigungsvorbehalt für den Windows-Durchlauf. +- Der D3D11-Elementpfad bleibt bis dahin unvalidiert. Absicherung: der + Backend-Vertrag (§12.6) hält Korrekturen im Adapter lokal; Projekt-, + Parameter-, DMX- und Web-API ändern sich nicht. +- Shader werden nur statisch bereitgestellt; Compile- und Golden-Image-Tests + entstehen im Windows-Durchlauf. +- Fällt Gate 0 rot aus: gezielte Sanierung nach §1.1 Nr. 9 mit ERRORS.md-Eintrag; + kein Architektur-Neubau erforderlich (Vertragsabsicherung). + +## Messwerte / Nachweise + +- CPU-Ebene: pytest/Ruff grün je Etappe (TEST_REPORT.md, fortlaufend). +- Hardware: ausstehend; Checkliste in STATUS.md. + +## Freigabe + +Auftraggeber: per Chat-Anfrage 2026-09-11 („erst fertig bauen und dann auf +Windows testen") – hiermit dokumentiert. diff --git a/docs/adr/README.md b/docs/adr/README.md index 00dc7a7..efce1f4 100644 --- a/docs/adr/README.md +++ b/docs/adr/README.md @@ -1,17 +1,25 @@ # Architecture Decision Records -Vorlage: `_template.md`. Nummerierung fortlaufend. Abweichungen vom Bauplan nur mit ADR und Freigabe (PLAN.md §1). +Vorlage: `_template.md`. Nummerierung fortlaufend. Abweichungen vom Bauplan nur +mit ADR und Freigabe (PLAN.md §1). ## Angenommen - ADR-0001: Python-Version-Pin 3.13 - ADR-0002: GStreamer-Pin 1.28.6 (Windows) - ADR-0003: IPC – lokales TCP + length-prefixed MessagePack v1 +- ADR-0008: Build-first-Strategie – verzögerte Hardware-Validierung + (Auftraggeber-Freigabe 2026-09-11) -## Offen (Phase 0 entscheidet, §32 / §7.1) +## Vorläufig angenommen (Bestätigung im Windows-Durchlauf, ADR-0008) -- ADR-0004: Nativer Renderkern – Rust vs. C++, eigenständige Bridge vs. GStreamer-Plugin -- ADR-0005: Packaging – Nuitka vs. PyInstaller (Onefolder) -- ADR-0006: Frontend – React vs. Svelte +- ADR-0004: Nativer Renderkern – Rust auf dem GStreamer-D3D11-Elementpfad +- ADR-0005: Packaging – Nuitka Onefolder (Fallback PyInstaller) + +## Offen + +- ADR-0006: Frontend – React vs. Svelte (Entscheidung zu Beginn Phase 5) - ADR-0007: Typprüfung – mypy vs. pyright -- weitere gemäß Bauplan §32 (Persistenz, Preview, Show-Codec, Ownership, Display-Abstraktion, Adaptive-Quality-Policy, Discovery, Clusterprotokoll, Paarung/TLS, Clock-Sync, UI-Design-Tokens) +- weitere gemäß Bauplan §32 (Persistenz, Preview, Show-Codec, Ownership, + Display-Abstraktion, Adaptive-Quality-Policy, Discovery, Clusterprotokoll, + Paarung/TLS, Clock-Sync, UI-Design-Tokens) diff --git a/packages/protocol/hms_protocol/__init__.py b/packages/protocol/hms_protocol/__init__.py index f33e910..a2d433e 100644 --- a/packages/protocol/hms_protocol/__init__.py +++ b/packages/protocol/hms_protocol/__init__.py @@ -1,10 +1,23 @@ """hms_protocol – versioniertes IPC (PLAN.md §6.2, ADR-0003). Lokales TCP auf 127.0.0.1, length-prefixed MessagePack, Protokollversion 1. +Verbindungsschicht: Handshake, Snapshot/Delta, Heartbeat, Re-Sync. """ +from hms_protocol.connection import ( + HandshakeInfo, + IpcClient, + IpcServer, + ProtocolError, +) from hms_protocol.envelope import Envelope, MessageType -from hms_protocol.framing import decode_frame, encode_frame, read_frame, write_frame +from hms_protocol.framing import ( + decode_frame, + encode_frame, + read_frame, + read_frame_async, + write_frame, +) from hms_protocol.idempotency import IdempotencyRegistry __all__ = [ @@ -13,6 +26,11 @@ __all__ = [ "encode_frame", "decode_frame", "read_frame", + "read_frame_async", "write_frame", "IdempotencyRegistry", + "IpcServer", + "IpcClient", + "HandshakeInfo", + "ProtocolError", ] diff --git a/packages/protocol/hms_protocol/connection.py b/packages/protocol/hms_protocol/connection.py new file mode 100644 index 0000000..0b8f037 --- /dev/null +++ b/packages/protocol/hms_protocol/connection.py @@ -0,0 +1,252 @@ +"""IPC-Verbindung zwischen Control Core und Renderer (PLAN.md §6.2). + +Server (Renderer-seitig) und Client (Control-Core-seitig) auf 127.0.0.1. +Verbindungsablauf: + +1. Client verbindet, sendet hello mit Protokollversion + Capabilities +2. Server prüft Version, antwortet welcome mit eigenen Capabilities +3. Server sendet vollständigen state snapshot +4. danach inkrementelle Deltas mit monotoner Revision +5. Heartbeat mindestens alle 500 ms in beide Richtungen +6. Ack für jedes zustandsändernde Command (idempotent über message_id) +7. Nach Reconnect: Re-Sync (neuer Snapshot), Deltas erst danach akzeptiert +""" + +from __future__ import annotations + +import asyncio +import time +from dataclasses import dataclass + +from hms_protocol.envelope import PROTOCOL_VERSION, Envelope, MessageType +from hms_protocol.framing import encode_frame, read_frame_async + +HEARTBEAT_INTERVAL_S = 0.5 # §6.2: mindestens alle 500 ms + + +@dataclass +class ProtocolError(Exception): + """Protokollverstoß; Verbindung wird getrennt.""" + + code: str + message: str + + def __str__(self) -> str: + return f"{self.code}: {self.message}" + + +@dataclass +class HandshakeInfo: + """Ergebnis des Handshakes mit Capabilities der Gegenseite.""" + + peer_name: str + peer_capabilities: dict + protocol_version: int = PROTOCOL_VERSION + + +class IpcServer: + """Renderer-seitiger IPC-Server. Lauscht ausschließlich auf 127.0.0.1. + + Liefert dem Renderer: + - accept(handler) → wartet auf Client, führt Handshake durch + - receive() → nächste Nachricht (command/event/heartbeat) + - send(envelope) → Nachricht an Control Core + """ + + def __init__(self, port: int = 0, host: str = "127.0.0.1") -> None: + if host not in ("127.0.0.1", "localhost", "::1"): + raise ValueError("IPC-Server darf nur auf Loopback lauschen (§6.2)") + self._host = host + self._port = port + self._server: asyncio.Server | None = None + self._reader: asyncio.StreamReader | None = None + self._writer: asyncio.StreamWriter | None = None + self._handler_task: asyncio.Task | None = None + self._heartbeat_task: asyncio.Task | None = None + self._last_peer_heartbeat_ns: int = 0 + self.capabilities: dict = {} + self.name: str = "hms-renderer" + + @property + def port(self) -> int: + if self._server is None: + return self._port + return self._server.sockets[0].getsockname()[1] if self._server.sockets else self._port + + async def start(self) -> int: + """Startet den Listener; gibt den tatsächlichen Port zurück.""" + self._server = await asyncio.start_server(self._on_client, self._host, self._port) + return self.port + + async def stop(self) -> None: + if self._heartbeat_task: + self._heartbeat_task.cancel() + if self._handler_task: + self._handler_task.cancel() + if self._writer: + self._writer.close() + if self._server: + self._server.close() + await self._server.wait_closed() + + async def _on_client( + self, reader: asyncio.StreamReader, writer: asyncio.StreamWriter + ) -> None: + """Nimmt genau einen Client an (V1: eine Verbindung).""" + self._reader = reader + self._writer = writer + # auf hello warten + hello_raw = await read_frame_async(reader) + if hello_raw.get("type") != "command" or hello_raw.get("payload", {}).get("action") != "hello": + raise ProtocolError("EXPECTED_HELLO", f"got {hello_raw.get('type')}") + peer_caps = hello_raw.get("payload", {}).get("capabilities", {}) + peer_name = hello_raw.get("payload", {}).get("name", "unknown") + if hello_raw.get("protocol_version") != PROTOCOL_VERSION: + err = Envelope( + type=MessageType.ERROR, + payload={"code": "VERSION_MISMATCH", "expected": PROTOCOL_VERSION}, + ) + writer.write(_encode_envelope(err)) + await writer.drain() + writer.close() + raise ProtocolError("VERSION_MISMATCH", "client protocol mismatch") + # welcome senden + welcome = Envelope( + type=MessageType.EVENT, + payload={"action": "welcome", "name": self.name, "capabilities": self.capabilities}, + ) + writer.write(_encode_envelope(welcome)) + await writer.drain() + self._last_peer_heartbeat_ns = time.monotonic_ns() + self._heartbeat_task = asyncio.get_running_loop().create_task(self._send_heartbeats()) + + async def _send_heartbeats(self) -> None: + while True: + await asyncio.sleep(HEARTBEAT_INTERVAL_S) + if self._writer is None: + return + hb = Envelope(type=MessageType.HEARTBEAT, payload={"source": self.name}) + self._writer.write(_encode_envelope(hb)) + await self._writer.drain() + + async def receive(self) -> Envelope | None: + """Liest die nächste Nachricht; None bei Verbindungsabbruch.""" + if self._reader is None: + return None + try: + raw = await read_frame_async(self._reader) + except (asyncio.IncompleteReadError, ConnectionError): + return None + if raw.get("type") == "heartbeat": + self._last_peer_heartbeat_ns = time.monotonic_ns() + return Envelope.model_validate(raw) + + async def send(self, envelope: Envelope) -> None: + if self._writer is None: + raise ConnectionError("IPC client not connected") + self._writer.write(_encode_envelope(envelope)) + await self._writer.drain() + + @property + def peer_alive(self) -> bool: + """True, wenn letzter Peer-Heartbeat < 2× Intervall zurückliegt.""" + if self._last_peer_heartbeat_ns == 0: + return False + return (time.monotonic_ns() - self._last_peer_heartbeat_ns) < 2 * HEARTBEAT_INTERVAL_S * 1e9 + + +class IpcClient: + """Control-Core-seitiger IPC-Client. Verbindet sich mit 127.0.0.1:port. + + - connect() → Handshake, gibt HandshakeInfo zurück + - receive() → nächste Nachricht (snapshot/event/ack/heartbeat) + - send(envelope) → Nachricht an Renderer + - Nach Reconnect: connect() erneut → neuer Snapshot (Re-Sync) + """ + + def __init__(self, port: int, host: str = "127.0.0.1", name: str = "control-core") -> None: + if host not in ("127.0.0.1", "localhost", "::1"): + raise ValueError("IPC-Client darf nur Loopback verbinden (§6.2)") + self._host = host + self._port = port + self._name = name + self._reader: asyncio.StreamReader | None = None + self._writer: asyncio.StreamWriter | None = None + self._heartbeat_task: asyncio.Task | None = None + self._last_peer_heartbeat_ns: int = 0 + self.capabilities: dict = {} + + async def connect(self) -> HandshakeInfo: + """Baut Verbindung auf, führt Handshake durch.""" + self._reader, self._writer = await asyncio.open_connection(self._host, self._port) + hello = Envelope( + type=MessageType.COMMAND, + payload={"action": "hello", "name": self._name, "capabilities": self.capabilities}, + ) + self._writer.write(_encode_envelope(hello)) + await self._writer.drain() + raw = await read_frame_async(self._reader) + if raw.get("type") == "error": + raise ProtocolError( + raw.get("payload", {}).get("code", "REMOTE_ERROR"), + str(raw.get("payload", {})), + ) + if raw.get("type") != "event" or raw.get("payload", {}).get("action") != "welcome": + raise ProtocolError("EXPECTED_WELCOME", f"got {raw.get('type')}") + if raw.get("protocol_version") != PROTOCOL_VERSION: + raise ProtocolError("VERSION_MISMATCH", f"server sent version {raw.get('protocol_version')}") + self._last_peer_heartbeat_ns = time.monotonic_ns() + self._heartbeat_task = asyncio.get_running_loop().create_task(self._send_heartbeats()) + return HandshakeInfo( + peer_name=raw.get("payload", {}).get("name", "unknown"), + peer_capabilities=raw.get("payload", {}).get("capabilities", {}), + ) + + async def disconnect(self) -> None: + if self._heartbeat_task: + self._heartbeat_task.cancel() + self._heartbeat_task = None + if self._writer: + self._writer.close() + try: + await self._writer.wait_closed() + except (ConnectionError, asyncio.CancelledError): + pass + self._writer = None + self._reader = None + + async def receive(self) -> Envelope | None: + if self._reader is None: + return None + try: + raw = await read_frame_async(self._reader) + except (asyncio.IncompleteReadError, ConnectionError): + return None + if raw.get("type") == "heartbeat": + self._last_peer_heartbeat_ns = time.monotonic_ns() + return Envelope.model_validate(raw) + + async def send(self, envelope: Envelope) -> None: + if self._writer is None: + raise ConnectionError("IPC not connected") + self._writer.write(_encode_envelope(envelope)) + await self._writer.drain() + + async def _send_heartbeats(self) -> None: + while True: + await asyncio.sleep(HEARTBEAT_INTERVAL_S) + if self._writer is None: + return + hb = Envelope(type=MessageType.HEARTBEAT, payload={"source": self._name}) + self._writer.write(_encode_envelope(hb)) + await self._writer.drain() + + @property + def peer_alive(self) -> bool: + if self._last_peer_heartbeat_ns == 0: + return False + return (time.monotonic_ns() - self._last_peer_heartbeat_ns) < 2 * HEARTBEAT_INTERVAL_S * 1e9 + + +def _encode_envelope(envelope: Envelope) -> bytes: + return encode_frame(envelope.model_dump(mode="json")) diff --git a/packages/protocol/hms_protocol/envelope.py b/packages/protocol/hms_protocol/envelope.py index ce8709e..d09cdf0 100644 --- a/packages/protocol/hms_protocol/envelope.py +++ b/packages/protocol/hms_protocol/envelope.py @@ -18,6 +18,7 @@ class MessageType(StrEnum): ACK = "ack" ERROR = "error" TELEMETRY = "telemetry" + HEARTBEAT = "heartbeat" class Envelope(BaseModel): diff --git a/packages/protocol/hms_protocol/framing.py b/packages/protocol/hms_protocol/framing.py index 2090297..e224f3d 100644 --- a/packages/protocol/hms_protocol/framing.py +++ b/packages/protocol/hms_protocol/framing.py @@ -6,6 +6,7 @@ schützt vor unkontrollierten Queues/Resourcenerschöpfung (§6.2, §33). from __future__ import annotations +import asyncio import socket import struct @@ -46,6 +47,16 @@ def read_frame(sock: socket.socket) -> dict: return msgpack.unpackb(body, raw=False) +async def read_frame_async(reader: asyncio.StreamReader) -> dict: + """Liest einen Frame von einem asyncio-StreamReader (IPC-Server/Client).""" + header = await reader.readexactly(_LENGTH.size) + (length,) = _LENGTH.unpack(header) + if length > MAX_PAYLOAD_SIZE: + raise ValueError(f"declared length {length} exceeds limit") + body = await reader.readexactly(length) + return msgpack.unpackb(body, raw=False) + + def write_frame(sock: socket.socket, payload: dict) -> None: """Schreibt einen Frame auf einen verbundenen Socket.""" sock.sendall(encode_frame(payload)) diff --git a/schemas/ipc/envelope_v1.schema.json b/schemas/ipc/envelope_v1.schema.json index 17edbf6..b8ef76a 100644 --- a/schemas/ipc/envelope_v1.schema.json +++ b/schemas/ipc/envelope_v1.schema.json @@ -8,7 +8,7 @@ "properties": { "protocol_version": {"const": 1}, "message_id": {"type": "string", "format": "uuid"}, - "type": {"enum": ["command", "event", "snapshot", "ack", "error", "telemetry"]}, + "type": {"enum": ["command", "event", "snapshot", "ack", "error", "telemetry", "heartbeat"]}, "revision": {"type": "integer", "minimum": 0}, "monotonic_timestamp_ns": {"type": "integer", "minimum": 0}, "payload": {"type": "object"}, diff --git a/tests/integration/test_ipc_connection.py b/tests/integration/test_ipc_connection.py new file mode 100644 index 0000000..333b469 --- /dev/null +++ b/tests/integration/test_ipc_connection.py @@ -0,0 +1,266 @@ +"""Integrationstests IPC-Verbindung (PLAN.md §6.2, §29.2). + +Echter TCP-Loopback (kein Mock): Handshake, Heartbeat, Snapshot/Delta, +Idempotenz, Version-Mismatch, Re-Sync nach Reconnect. +""" + +from __future__ import annotations + +import asyncio + +from hms_protocol import ( + Envelope, + HandshakeInfo, + IpcClient, + IpcServer, + MessageType, +) + + +def _run(coro): + return asyncio.run(coro) + + +async def _connected_pair( + server_caps: dict | None = None, + client_caps: dict | None = None, +): + """Startet Server+Client und führt den Handshake durch.""" + server = IpcServer() + if server_caps is not None: + server.capabilities = server_caps + port = await server.start() + client = IpcClient(port) + if client_caps is not None: + client.capabilities = client_caps + info = await client.connect() + return server, client, info + + +# ---------- Bindung ausschließlich Loopback (§6.2) ---------- + + +def test_server_rejects_non_loopback_host() -> None: + try: + IpcServer(host="0.0.0.0") + raise AssertionError("0.0.0.0 muss abgelehnt werden") + except ValueError: + pass + + +def test_client_rejects_non_loopback_host() -> None: + try: + IpcClient(port=1234, host="192.168.1.5") + raise AssertionError("externe IP muss abgelehnt werden") + except ValueError: + pass + + +# ---------- Handshake (§6.2) ---------- + + +def test_handshake_exchanges_capabilities() -> None: + async def impl() -> None: + server, client, info = await _connected_pair( + server_caps={"backend": "d3d11", "max_layers": 8}, + client_caps={"role": "control_core"}, + ) + try: + assert isinstance(info, HandshakeInfo) + assert info.peer_name == "hms-renderer" + assert info.peer_capabilities == {"backend": "d3d11", "max_layers": 8} + assert info.protocol_version == 1 + finally: + await client.disconnect() + await server.stop() + + _run(impl()) + + +def test_client_rejects_version_mismatch() -> None: + async def impl() -> None: + server = IpcServer() + port = await server.start() + client = IpcClient(port) + try: + # manipulierter Handshake mit falscher Protokollversion + from hms_protocol import encode_frame + + reader, writer = await asyncio.open_connection("127.0.0.1", port) + writer.write( + encode_frame( + { + "protocol_version": 99, + "message_id": "x", + "type": "command", + "revision": 0, + "monotonic_timestamp_ns": 0, + "payload": {"action": "hello", "name": "evil", "capabilities": {}}, + } + ) + ) + await writer.drain() + # Server antwortet mit ERROR VERSION_MISMATCH + from hms_protocol import read_frame_async + + reply = await read_frame_async(reader) + assert reply["type"] == "error" + assert reply["payload"]["code"] == "VERSION_MISMATCH" + writer.close() + finally: + await client.disconnect() + await server.stop() + + _run(impl()) + + +# ---------- Nachrichtenübertragung ---------- + + +def test_snapshot_and_delta_delivery() -> None: + async def impl() -> None: + server, client, _ = await _connected_pair() + try: + # Server sendet Snapshot + snap = Envelope( + type=MessageType.SNAPSHOT, + revision=5, + payload={"values": {"master/intensity": 1.0}}, + ) + await server.send(snap) + received = await asyncio.wait_for(client.receive(), timeout=2.0) + assert received is not None + assert received.type is MessageType.SNAPSHOT + assert received.revision == 5 + assert received.payload["values"]["master/intensity"] == 1.0 + + # danach Delta mit höherer Revision + delta = Envelope( + type=MessageType.EVENT, + revision=6, + payload={"changes": {"master/intensity": 0.5}}, + ) + await server.send(delta) + received2 = await asyncio.wait_for(client.receive(), timeout=2.0) + assert received2 is not None + assert received2.revision == 6 + assert received2.payload["changes"]["master/intensity"] == 0.5 + finally: + await client.disconnect() + await server.stop() + + _run(impl()) + + +def test_bidirectional_commands_and_acks() -> None: + async def impl() -> None: + server, client, _ = await _connected_pair() + try: + # Client → Server: Command + cmd = Envelope( + type=MessageType.COMMAND, + revision=10, + payload={"action": "parameter.set", "path": "master/intensity", "value": 0.7}, + ) + await client.send(cmd) + received = await asyncio.wait_for(server.receive(), timeout=2.0) + assert received is not None + assert received.type is MessageType.COMMAND + assert received.payload["action"] == "parameter.set" + + # Server → Client: Ack + ack = Envelope( + type=MessageType.ACK, + revision=11, + payload={"ack_for": received.message_id, "status": "ok"}, + ) + await server.send(ack) + ack_received = await asyncio.wait_for(client.receive(), timeout=2.0) + assert ack_received is not None + assert ack_received.type is MessageType.ACK + assert ack_received.payload["ack_for"] == received.message_id + finally: + await client.disconnect() + await server.stop() + + _run(impl()) + + +# ---------- Heartbeat (§6.2: mindestens alle 500 ms) ---------- + + +def test_heartbeat_flows_bidirectionally() -> None: + async def impl() -> None: + server, client, _ = await _connected_pair() + try: + await asyncio.sleep(0.7) # > Heartbeat-Intervall + # Client empfängt Server-Heartbeat + got_server_hb = False + for _ in range(4): + msg = await asyncio.wait_for(client.receive(), timeout=1.5) + if msg is not None and msg.type is MessageType.HEARTBEAT: + got_server_hb = True + break + assert got_server_hb, "Client musste einen Heartbeat empfangen" + assert client.peer_alive, "Server-Heartbeat muss peer_alive setzen" + + # Server empfängt Client-Heartbeat + await asyncio.sleep(0.1) + got_client_hb = False + for _ in range(6): + msg = await asyncio.wait_for(server.receive(), timeout=1.5) + if msg is not None and msg.type is MessageType.HEARTBEAT: + got_client_hb = True + break + assert got_client_hb, "Server musste einen Heartbeat empfangen" + assert server.peer_alive + finally: + await client.disconnect() + await server.stop() + + _run(impl()) + + +# ---------- Re-Sync nach Reconnect (§6.2, §29.2) ---------- + + +def test_reconnect_requires_new_snapshot() -> None: + async def impl() -> None: + server = IpcServer() + port = await server.start() + try: + # erste Verbindung: Snapshot Revision 3 + client1 = IpcClient(port) + await client1.connect() + await server.send( + Envelope(type=MessageType.SNAPSHOT, revision=3, payload={"values": {}}) + ) + snap1 = await asyncio.wait_for(client1.receive(), timeout=2.0) + assert snap1 is not None and snap1.revision == 3 + await client1.disconnect() + + # zweite Verbindung: neuer Snapshot (Re-Sync) mit Revision 4 + client2 = IpcClient(port) + info2 = await client2.connect() # noqa: F841 – Handshake reicht + await server.send( + Envelope(type=MessageType.SNAPSHOT, revision=4, payload={"values": {"a": 1}}) + ) + snap2 = await asyncio.wait_for(client2.receive(), timeout=2.0) + assert snap2 is not None + assert snap2.revision == 4 # Deltas erst nach erfolgreichem Re-Sync + await client2.disconnect() + finally: + await server.stop() + + _run(impl()) + + +def test_receive_returns_none_on_disconnect() -> None: + async def impl() -> None: + server, client, _ = await _connected_pair() + await client.disconnect() + msg = await asyncio.wait_for(server.receive(), timeout=2.0) + assert msg is None # Verbindung weg → sauberes None statt Exception + await server.stop() + + _run(impl())