diff --git a/apps/renderer/hms_renderer/__init__.py b/apps/renderer/hms_renderer/__init__.py index 6fee5ae..a36370e 100644 --- a/apps/renderer/hms_renderer/__init__.py +++ b/apps/renderer/hms_renderer/__init__.py @@ -1,11 +1,16 @@ """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. Der RemoteStateMirror hält den -über IPC übermittelten Showzustand (§6.2, §6.4). -""" +in Python (§2.1, §33). Der RemoteStateMirror hält den über IPC übermittelten +Showzustand (§6.2, §6.4); die RenderEngine verbindet Mirror, Playback und +Preload zu Frame-Snapshots (§11.4).""" +from hms_renderer.engine import ( + FrameSnapshot, + RenderEngine, + RenderTelemetry, + SourceHandle, +) from hms_renderer.pipelines import ( D3D11Pipeline, DevGLPipeline, @@ -23,4 +28,8 @@ __all__ = [ "gst_available", "RemoteStateMirror", "apply_envelope", + "RenderEngine", + "SourceHandle", + "FrameSnapshot", + "RenderTelemetry", ] diff --git a/apps/renderer/hms_renderer/engine.py b/apps/renderer/hms_renderer/engine.py new file mode 100644 index 0000000..cc38581 --- /dev/null +++ b/apps/renderer/hms_renderer/engine.py @@ -0,0 +1,268 @@ +"""Renderer-Engine: verbindet StateMirror, Playback und Preload (§12, §6.4). + +Der Rendergraph läuft in nativem GStreamer/Rust (ADR-0004); diese Engine +orchestriert: +- liest Showparameter aus dem RemoteStateMirror (§6.4) +- berechnet Playback-Positionen je Quelle (§12.5) +- verwaltet Preload-Slots für atomaren Clipwechsel (§12.2) +- bildet pro Frame einen unveränderlichen Parameter-Snapshot (§11.4) +- erzeugt Rendertelemetrie (§28.2) + +Keine Pixelverarbeitung in Python (§33); die GPU-Pipeline wird über +Pipeline-Definitionen an GStreamer übergeben. +""" + +from __future__ import annotations + +import time +from dataclasses import dataclass + +from hms_domain.model import TransportState +from hms_media import PlaybackController, PreloadSlot + +from hms_renderer.state_mirror import RemoteStateMirror + + +@dataclass(frozen=True) +class FrameSnapshot: + """Unveränderlicher Parameter-Snapshot für genau einen Frame (§11.4). + + Der Rendergraph (Rust/GStreamer) erhält genau dieses Objekt pro Frame; + halbe Zustände sind ausgeschlossen. + """ + + frame_index: int + monotonic_ns: int + state_revision: int + parameters: dict[str, float] + source_positions: dict[str, float] # source_id → normalisierte Position + source_states: dict[str, str] # source_id → TransportState + active_asset_ids: dict[str, str | None] # layer_key → asset_id (nach Commit) + + +@dataclass +class RenderTelemetry: + """Rendertelemetrie je Frame (§28.2, §25.4).""" + + frame_index: int = 0 + fps: float = 0.0 + frame_time_ms: float = 0.0 + dropped_frames: int = 0 + active_layers: int = 0 + preload_pending: int = 0 + source_events: int = 0 + + +class SourceHandle: + """Verwaltet eine Medienquelle im Renderer (§12.5, §12.2). + + Verbindet PlaybackController mit Parameterpfaden und einem PreloadSlot: + - Transport-Befehle aus dem StateMirror steuern den Controller + - Der PreloadSlot wechselt Clips atomar bei Commit + """ + + def __init__( + self, + source_id: str, + parameter_prefix: str, + fps: float = 60.0, + ) -> None: + self.source_id = source_id + self.parameter_prefix = parameter_prefix + self.controller = PlaybackController(source_id=source_id) + self.preload = PreloadSlot() + self._active_asset_id: str | None = None + self._fps = fps + self._last_events: list = [] + + @property + def active_asset_id(self) -> str | None: + return self._active_asset_id + + @property + def position(self) -> float: + return self.controller.position + + @property + def state(self) -> TransportState: + return self.controller.state + + def sync_from_mirror(self, mirror: RemoteStateMirror) -> list: + """Liest Transport-Befehle und Parameter aus dem StateMirror. + + Rückgabe: PlaybackEvents dieses Sync-Schritts (§12.5). + """ + events: list = [] + + # Transport-State aus dem Mirror lesen (§11.2) + state_value = mirror.get_value(f"{self.parameter_prefix}/source/state") + if state_value is not None: + state_map = { + 0: TransportState.STOPPED, + 1: TransportState.PLAYING, + 2: TransportState.PAUSED, + } + target_state = state_map.get(int(state_value)) + if target_state is not None: + is_playing = target_state is TransportState.PLAYING + is_paused = target_state is TransportState.PAUSED + is_stopped = target_state is TransportState.STOPPED + if is_playing and self.controller.state is not TransportState.PLAYING: + self.controller.play() + elif is_paused and self.controller.state is TransportState.PLAYING: + self.controller.pause() + elif is_stopped and self.controller.state is not TransportState.STOPPED: + self.controller.stop() + + # Playback-Parameter aus dem Mirror (§10.2) + speed = mirror.get_value(f"{self.parameter_prefix}/source/speed") + if speed is not None and speed != 0: + self.controller.speed = max(-4.0, min(4.0, speed)) + + in_point = mirror.get_value(f"{self.parameter_prefix}/source/in_point") + if in_point is not None: + self.controller.in_point = max(0.0, min(0.99, in_point)) + + out_point = mirror.get_value(f"{self.parameter_prefix}/source/out_point") + if out_point is not None: + self.controller.out_point = max(self.controller.in_point + 0.01, min(1.0, out_point)) + + # Retrigger: Flankenwert im Mirror (1.0 = triggern, danach zurück auf 0) + retrigger = mirror.get_value(f"{self.parameter_prefix}/source/retrigger") + if retrigger is not None and retrigger > 0.5: + self.controller.retrigger() + + # Clip-Auswahl über PreloadSlot (§12.2, §16.5 Load/Commit) + pending_asset = mirror.get_value(f"{self.parameter_prefix}/source/asset_id_pending") + if pending_asset is not None and pending_asset > 0: + # Asset-ID ist als Hash/Integer im Mirror; real: UUID-String aus Registry + # Hier: Asset-Wechsel nur über Commit (PreloadSlot) + self.preload.preload(str(int(pending_asset))) + + commit = mirror.get_value(f"{self.parameter_prefix}/source/commit") + if commit is not None and commit > 0.5: + self.preload.mark_ready(str(int(mirror.get_value( + f"{self.parameter_prefix}/source/asset_id_pending", 0 + )))) + committed = self.preload.commit() + if committed: + self._active_asset_id = committed # atomar gewechselt (§12.2) + + self._last_events = events + return events + + def advance(self, dt_s: float, now_ns: int) -> list: + """Advancement des PlaybackControllers; liefert Events.""" + return self.controller.advance(dt_s, now_ns) + + +class RenderEngine: + """Zentrale Renderer-Engine: orchestriert alle Quellen pro Frame. + + Ablauf pro Frame (§12.2: feste Master-Bildrate): + 1. sync_from_mirror: Transport-Parameter aus dem StateMirror lesen + 2. advance: Playback-Positionen um dt weiterschieben + 3. snapshot: unveränderlichen Frame-Snapshot bilden + 4. telemetry: Frame-Statistiken aktualisieren + """ + + def __init__(self, mirror: RemoteStateMirror, fps: float = 60.0) -> None: + self._mirror = mirror + self._fps = fps + self._frame_duration_s = 1.0 / fps + self._sources: dict[str, SourceHandle] = {} + self._frame_index = 0 + self._last_frame_ns = 0 + self.telemetry = RenderTelemetry() + + @property + def mirror(self) -> RemoteStateMirror: + return self._mirror + + @property + def fps(self) -> float: + return self._fps + + @property + def frame_index(self) -> int: + return self._frame_index + + def add_source(self, source_id: str, parameter_prefix: str) -> SourceHandle: + """Registriert eine Medienquelle im Renderer.""" + handle = SourceHandle(source_id, parameter_prefix, self._fps) + self._sources[source_id] = handle + return handle + + def remove_source(self, source_id: str) -> None: + self._sources.pop(source_id, None) + + def get_source(self, source_id: str) -> SourceHandle | None: + return self._sources.get(source_id) + + def sources(self) -> list[SourceHandle]: + return list(self._sources.values()) + + def tick(self, now_ns: int | None = None) -> FrameSnapshot: + """Ein Frame: Sync → Advance → Snapshot. + + Wird vom Renderer-Loop mit fester Master-Bildrate aufgerufen (§12.2). + """ + now = now_ns if now_ns is not None else time.monotonic_ns() + if self._last_frame_ns == 0: + self._last_frame_ns = now + dt_s = (now - self._last_frame_ns) / 1e9 + self._last_frame_ns = now + + # 1. Sync: Parameter aus dem Mirror lesen + all_events: list = [] + for handle in self._sources.values(): + handle.sync_from_mirror(self._mirror) + + # 2. Advance: Playback weiterschieben + for handle in self._sources.values(): + events = handle.advance(dt_s, now) + all_events.extend(events) + + # 3. Snapshot: unveränderlicher Frame-Zustand (§11.4) + self._frame_index += 1 + parameters = self._mirror.all_values() + source_positions: dict[str, float] = {} + source_states: dict[str, str] = {} + active_assets: dict[str, str | None] = {} + + for source_id, handle in self._sources.items(): + source_positions[source_id] = handle.position + source_states[source_id] = handle.state.value + active_assets[f"{handle.parameter_prefix}/source"] = handle.active_asset_id + + snapshot = FrameSnapshot( + frame_index=self._frame_index, + monotonic_ns=now, + state_revision=self._mirror.revision, + parameters=parameters, + source_positions=source_positions, + source_states=source_states, + active_asset_ids=active_assets, + ) + + # 4. Telemetry (§28.2) + frame_time_ms = dt_s * 1000.0 + active_count = sum( + 1 for h in self._sources.values() + if h.state is TransportState.PLAYING + ) + preload_count = sum( + 1 for h in self._sources.values() + if h.preload.pending is not None + ) + self.telemetry = RenderTelemetry( + frame_index=self._frame_index, + fps=1000.0 / max(frame_time_ms, 0.01), + frame_time_ms=frame_time_ms, + dropped_frames=0, # zählt der native Rendergraph + active_layers=active_count, + preload_pending=preload_count, + source_events=len(all_events), + ) + + return snapshot diff --git a/tests/unit/test_render_engine.py b/tests/unit/test_render_engine.py new file mode 100644 index 0000000..3465e1e --- /dev/null +++ b/tests/unit/test_render_engine.py @@ -0,0 +1,205 @@ +"""Tests Renderer-Engine: Sync→Advance→Snapshot-Kette (§12, §11.4, §29.2).""" + +from __future__ import annotations + +import pytest +from hms_domain import TransportState +from hms_renderer import FrameSnapshot, RemoteStateMirror, RenderEngine + + +def _mirror_with_params(params: dict[str, float]) -> RemoteStateMirror: + """Mirror mit Snapshot, der die gegebenen Parameter enthält.""" + mirror = RemoteStateMirror() + mirror.apply_snapshot({ + "state_revision": 1, + "project_revision": 1, + "values": params, + "monotonic_ns": 0, + }) + return mirror + + +def test_engine_creates_frame_snapshot_per_tick() -> None: + """tick() erzeugt FrameSnapshots mit steigendem Index (§11.4).""" + mirror = _mirror_with_params({"master/intensity": 1.0}) + engine = RenderEngine(mirror, fps=60.0) + engine.add_source("src-1", "composition/a/layer/b") + + t0 = 1_000_000_000 + t1 = t0 + 16_666_667 # ~60 fps (16.67 ms) + snap1 = engine.tick(t0) + snap2 = engine.tick(t1) + + assert isinstance(snap1, FrameSnapshot) + assert snap1.frame_index == 1 + assert snap2.frame_index == 2 + assert snap2.parameters["master/intensity"] == 1.0 + + +def test_engine_syncs_play_command_from_mirror() -> None: + """Mirror-Parameter 'source/state' steuert den Transport (§12.5).""" + # TransportState.PLAYING = 1.0 im Mirror + mirror = _mirror_with_params({ + "layer/x/source/state": 1.0, + "layer/x/source/speed": 1.0, + }) + engine = RenderEngine(mirror, fps=60.0) + handle = engine.add_source("src-1", "layer/x") + + assert handle.state is TransportState.STOPPED + engine.tick(1_000_000_000) + assert handle.state is TransportState.PLAYING + + # zweiter Frame: Position schreitet voran + engine.tick(1_016_666_667) + assert handle.position > 0.0 + + +def test_engine_advances_position_over_time() -> None: + """Playback-Position wächst mit der Frame-Zeit (§12.5).""" + mirror = _mirror_with_params({ + "layer/x/source/state": 1.0, + "layer/x/source/speed": 1.0, + }) + engine = RenderEngine(mirror, fps=60.0) + engine.add_source("src-1", "layer/x") + + t0 = 1_000_000_000 + engine.tick(t0) + # 10 Frames später à 16.67 ms = ~0.167 s Advancement + for i in range(10): + engine.tick(t0 + (i + 1) * 16_666_667) + handle = engine.get_source("src-1") + assert handle is not None + assert 0.1 < handle.position < 0.25 # ca. 10/60 = 0.167 + + +def test_engine_preload_commit_switches_asset_atomically() -> None: + """Clip-Wechsel über Preload: pending → ready → commit (§12.2).""" + mirror = _mirror_with_params({}) + engine = RenderEngine(mirror, fps=60.0) + handle = engine.add_source("src-1", "layer/x") + + # Asset auswählen (pending) + mirror.apply_delta({ + "state_revision": 2, + "project_revision": 1, + "changes": {"layer/x/source/asset_id_pending": 42.0}, + "monotonic_ns": 0, + }) + engine.tick(1_000_000_000) + assert handle.preload.pending == "42" + assert handle.active_asset_id is None # noch nicht gewechselt + + # Commit (ready + Flanke) + mirror.apply_delta({ + "state_revision": 3, + "project_revision": 1, + "changes": {"layer/x/source/commit": 1.0}, + "monotonic_ns": 0, + }) + snap = engine.tick(1_016_666_667) + assert handle.active_asset_id == "42" # atomar gewechselt + assert snap.active_asset_ids["layer/x/source"] == "42" + + +def test_engine_retrigger_resets_position() -> None: + """Retrigger-Flanke startet die Quelle neu (§12.5).""" + mirror = _mirror_with_params({ + "layer/x/source/state": 1.0, + "layer/x/source/speed": 1.0, + }) + engine = RenderEngine(mirror, fps=60.0) + engine.add_source("src-1", "layer/x") + + # 5 Frames weit spielen + t0 = 1_000_000_000 + for i in range(5): + engine.tick(t0 + i * 16_666_667) + handle = engine.get_source("src-1") + assert handle is not None and handle.position > 0.05 + + # Retrigger + mirror.apply_delta({ + "state_revision": 2, + "project_revision": 1, + "changes": {"layer/x/source/retrigger": 1.0}, + "monotonic_ns": 0, + }) + pos_before = handle.position + engine.tick(t0 + 5 * 16_666_667) + # Nach Retrigger: Position deutlich kleiner als davor (Rücksetzung auf + # In-Point im selben Frame; ein Frame dt advancement ist erlaubt) + assert handle.position < pos_before * 0.5 + + +def test_engine_snapshot_is_immutable_per_frame() -> None: + """§11.4: Frame-Snapshot bleibt stabil, auch wenn sich der Mirror ändert.""" + mirror = _mirror_with_params({"master/intensity": 0.5}) + engine = RenderEngine(mirror, fps=60.0) + engine.add_source("src-1", "layer/x") + + snap = engine.tick(1_000_000_000) + intensity_at_snap = snap.parameters["master/intensity"] + + # Mirror ändert sich nach dem Snapshot + mirror.apply_delta({ + "state_revision": 2, + "project_revision": 1, + "changes": {"master/intensity": 0.9}, + "monotonic_ns": 0, + }) + # alter Snapshot unverändert + assert snap.parameters["master/intensity"] == pytest.approx(intensity_at_snap) + # neuer Frame hat den neuen Wert + snap2 = engine.tick(1_016_666_667) + assert snap2.parameters["master/intensity"] == pytest.approx(0.9) + + +def test_engine_telemetry_reports_frame_stats() -> None: + """Telemetry enthält Frame-Index, FPS, aktive Layer (§28.2).""" + mirror = _mirror_with_params({"layer/x/source/state": 1.0}) + engine = RenderEngine(mirror, fps=60.0) + engine.add_source("src-1", "layer/x") + engine.tick(1_000_000_000) + engine.tick(1_016_666_667) # 16.67 ms → ~60 fps + + assert engine.telemetry.frame_index == 2 + assert 50.0 < engine.telemetry.fps < 70.0 # ~60 fps + assert engine.telemetry.active_layers == 1 # src-1 spielt + + +def test_engine_multiple_sources_advance_independently() -> None: + """Zwei Quellen mit unterschiedlicher Geschwindigkeit (§12.5).""" + mirror = _mirror_with_params({ + "layer/a/source/state": 1.0, + "layer/a/source/speed": 2.0, # doppelt + "layer/b/source/state": 1.0, + "layer/b/source/speed": 0.5, # halb + }) + engine = RenderEngine(mirror, fps=60.0) + engine.add_source("src-a", "layer/a") + engine.add_source("src-b", "layer/b") + + t0 = 1_000_000_000 + for i in range(20): + engine.tick(t0 + i * 16_666_667) + + handle_a = engine.get_source("src-a") + handle_b = engine.get_source("src-b") + assert handle_a is not None and handle_b is not None + # src-a (2x) hat ~4x die Position von src-b (0.5x) + assert handle_a.position > handle_b.position * 3.5 + + +def test_engine_stopped_source_does_not_advance() -> None: + """Gestoppte Quelle bleibt bei In-Point (§12.5).""" + mirror = _mirror_with_params({}) # kein state → STOPPED + engine = RenderEngine(mirror, fps=60.0) + handle = engine.add_source("src-1", "layer/x") + + for i in range(10): + engine.tick(1_000_000_000 + i * 16_666_667) + + assert handle.state is TransportState.STOPPED + assert handle.position == pytest.approx(handle.controller.in_point)