"""MediaEngine: GStreamer-Compositor mit dynamischen Layern (Video + Bild). Layer lassen sich zur Laufzeit hinzufuegen/entfernen/verschieben. Jeder Layer hat: Quelle (Bibliotheksdatei), Alpha, Position (x/y), Groesse (width/height), Z-Order. Die Pipeline wird bei Struktur-Aenderungen neu gebaut, Parameter gehen nicht verloren. """ from __future__ import annotations import threading import time import gi gi.require_version("Gst", "1.0") from gi.repository import Gst Gst.init(None) class Layer: """Ein aktiver Layer im Compositor.""" def __init__(self, layer_id: int, name: str, kind: str): self.id = layer_id self.name = name self.kind = kind # video | image self.alpha = 1.0 self.x = 0 self.y = 0 self.width = 320 self.height = 180 self.z = 0 def to_dict(self) -> dict: return {"id": self.id, "name": self.name, "kind": self.kind, "alpha": round(self.alpha, 3), "x": self.x, "y": self.y, "width": self.width, "height": self.height, "z": self.z, "effective": None} class MediaEngine: """Compositor-Pipeline + Steuerung. Thread-sicher ueber Lock.""" def __init__(self, settings: dict, library): self.settings = settings self.library = library self.pipeline = None self.current_jpeg = b"" self.frame_count = 0 self.running = False self.last_error: str | None = None self.layers: list[Layer] = [] self._next_id = 1 self.master = 1.0 self.blackout = False self.playing = True self.dmx_stats = {"packets": 0, "accepted": 0, "seq_errors": 0, "last_universe": None, "last_seq": None, "last_time": None} self._lock = threading.RLock() self._monitor: threading.Thread | None = None # ---------------- Pipeline ---------------- def _layer_source(self, layer: Layer, path: str) -> str: uri = f"file://{path}" if layer.kind == "video": caps = (f"video/x-raw,width={layer.width}," f"height={layer.height},format=BGRA") return (f"uridecodebin uri={uri} ! queue ! videoconvert ! " f"videoscale ! {caps} ! mix.sink_{layer.id}") # Bild: imagefreeze + EIN kombinierter Caps-Filter # (zwei Caps-Filter in Serie lassen gst-parse fehlschlagen) caps = (f"video/x-raw,width={layer.width}," f"height={layer.height},format=BGRA,framerate=15/1") return (f"uridecodebin uri={uri} ! imagefreeze ! " f"videoconvert ! videoscale ! {caps} ! mix.sink_{layer.id}") def _build_description(self) -> str: pv = self.settings["preview"] pad_defs = " ".join( f"sink_{l.id}::xpos={l.x} sink_{l.id}::ypos={l.y} " f"sink_{l.id}::zorder={l.z}" for l in self.layers) parts = [ f"compositor name=mix background=black {pad_defs} ! " f"video/x-raw,width={pv['width']},height={pv['height']}," f"format=BGRA ! tee name=t", ] for layer in self.layers: f = self.library.file_path(layer.name) parts.append(self._layer_source(layer, str(f.resolve()))) parts.append( f"t. ! queue ! videoconvert ! " f"jpegenc quality={pv['jpeg_quality']} ! " "appsink name=preview emit-signals=true max-buffers=2 drop=true") if self.settings["output"].get("fullscreen"): import sys sink = ("d3d11videosink sync=true" if sys.platform == "win32" else "autovideosink sync=true") parts.append(f"t. ! queue ! videoconvert ! {sink}") return " ".join(parts) def rebuild(self) -> bool: """Baut die Pipeline neu auf. Atomar: bei Fehler laeuft die alte weiter.""" with self._lock: if not self.layers: self.stop() self.running = True # kein Layer, aber "betriebsbereit" return True desc = self._build_description() old = self.pipeline try: new = Gst.parse_launch(desc) except Exception as e: # noqa: BLE001 self.last_error = f"parse_launch: {e}" print(f"[Engine] FEHLER: {self.last_error}") return False sink = new.get_by_name("preview") if sink: sink.connect("new-sample", self._on_frame) if new.set_state( Gst.State.PLAYING) == Gst.StateChangeReturn.FAILURE: self.last_error = "set_state(PLAYING) fehlgeschlagen" print(f"[Engine] FEHLER: {self.last_error}") new.set_state(Gst.State.NULL) return False # Neue Pipeline laeuft -> alte beenden (deren Bus-Monitor # beendet sich beim naechsten Poll selbst, da self.pipeline # nicht mehr die referenzierte Pipeline ist) if old is not None: old.set_state(Gst.State.NULL) self.pipeline = new self.running = True self.last_error = None self.apply_all() # Master/Blackout/Alphas auf neue Pads anwenden self._monitor = threading.Thread( target=self._bus_monitor, daemon=True) self._monitor.start() return True def _bus_monitor(self): pipeline = self.pipeline if pipeline is None: return bus = pipeline.get_bus() while self.running and self.pipeline is pipeline: msg = bus.timed_pop_filtered( 500 * Gst.MSECOND, Gst.MessageType.EOS | Gst.MessageType.ERROR) if msg is None: continue if msg.type == Gst.MessageType.EOS: if self.settings["engine"].get("loop_default", True): pipeline.seek_simple( Gst.Format.TIME, Gst.SeekFlags.FLUSH | Gst.SeekFlags.KEY_UNITS, 0) else: self.playing = False return elif msg.type == Gst.MessageType.ERROR: err, debug = msg.parse_error() self.last_error = f"{err.message} ({debug})" print(f"[Engine] GStreamer-Fehler: {self.last_error}") return def _on_frame(self, sink): sample = sink.emit("pull-sample") if sample is not None: buf = sample.get_buffer() self.current_jpeg = buf.extract_dup(0, buf.get_size()) self.frame_count += 1 return Gst.FlowReturn.OK # ---------------- Layer-Verwaltung ---------------- def add_layer(self, name: str) -> dict | None: with self._lock: max_l = int(self.settings["engine"].get("max_layers", 8)) if len(self.layers) >= max_l: self.last_error = f"max. {max_l} Layer" return None f = self.library.file_path(name) if not f.exists(): self.last_error = f"Datei nicht gefunden: {name}" return None kind = self.library.kind_of(name) if kind is None: self.last_error = f"nicht unterstuetzt: {name}" return None layer = Layer(self._next_id, name, kind) self._next_id += 1 layer.z = len(self.layers) pv = self.settings["preview"] layer.width = pv["width"] layer.height = pv["height"] self.layers.append(layer) if not self.rebuild(): self.layers.remove(layer) # Rollback: Engine unveraendert return None return layer.to_dict() def remove_layer(self, layer_id: int) -> bool: with self._lock: kept = [l for l in self.layers if l.id != layer_id] if len(kept) == len(self.layers): return False removed = [l for l in self.layers if l.id == layer_id] self.layers = kept if not self.rebuild(): self.layers.extend(removed) # Rollback self.layers.sort(key=lambda l: l.z) return False return True def get_layer(self, layer_id: int) -> Layer | None: for l in self.layers: if l.id == layer_id: return l return None def update_layer(self, layer_id: int, patch: dict ) -> Layer | None: with self._lock: layer = self.get_layer(layer_id) if layer is None: return None if "alpha" in patch: layer.alpha = max(0.0, min(1.0, float(patch["alpha"]))) if "x" in patch: layer.x = int(patch["x"]) if "y" in patch: layer.y = int(patch["y"]) if "width" in patch: layer.width = max(16, min(7680, int(patch["width"]))) if "height" in patch: layer.height = max(16, min(4320, int(patch["height"]))) if "z" in patch: layer.z = max(0, min(15, int(patch["z"]))) if "name" in patch and patch["name"] != layer.name: f = self.library.file_path(str(patch["name"])) if not f.exists(): return None layer.name = str(patch["name"]) layer.kind = self.library.kind_of(layer.name) or layer.kind # Pad-Eigenschaften live setzen (ohne Rebuild) self._apply_pad_props(layer) return layer def _apply_pad_props(self, layer: Layer): if self.pipeline is None: return mix = self.pipeline.get_by_name("mix") if mix is None: return pad = mix.get_static_pad(f"sink_{layer.id}") if pad is None: return pad.set_property("xpos", layer.x) pad.set_property("ypos", layer.y) pad.set_property("zorder", layer.z) pad.set_property("alpha", self._effective_alpha(layer)) # ---------------- Steuerung ---------------- def _effective_alpha(self, layer: Layer) -> float: if self.blackout: return 0.0 return layer.alpha * self.master def apply_all(self): with self._lock: for l in self.layers: self._apply_pad_props(l) if self.pipeline is not None and self.layers: target = (Gst.State.PLAYING if self.playing else Gst.State.PAUSED) current = self.pipeline.get_state(0)[1] if current != target: self.pipeline.set_state(target) def set_master(self, value: float): with self._lock: self.master = max(0.0, min(1.0, value)) self.apply_all() def set_blackout(self, on: bool): with self._lock: self.blackout = bool(on) self.apply_all() def set_playback(self, playing: bool): with self._lock: self.playing = bool(playing) self.apply_all() def retrigger(self): with self._lock: if self.pipeline is not None: self.pipeline.seek_simple( Gst.Format.TIME, Gst.SeekFlags.FLUSH | Gst.SeekFlags.KEY_UNITS, 0) # ---------------- DMX (Mapping aus Settings) ---------------- def apply_dmx(self, channels: bytes, universe: int, seq: int): a = self.settings["artnet"] def ch(num) -> int: if num is None: return 0 i = int(num) - 1 return channels[i] if 0 <= i < len(channels) else 0 with self._lock: self.dmx_stats["packets"] += 1 self.dmx_stats["accepted"] += 1 self.dmx_stats["last_universe"] = universe self.dmx_stats["last_seq"] = seq self.dmx_stats["last_time"] = time.time() start = a.get("layer_start_channel") for i, layer in enumerate(self.layers): v = ch(None if start is None else start + i) if v > 0: layer.alpha = v / 255.0 self.master = ch(a.get("master_channel")) / 255.0 or self.master self.apply_all() pb = a.get("playback_channel") if ch(pb) >= 128: self.set_playback(True) bo = a.get("blackout_channel") if bo is not None: self.set_blackout(ch(bo) >= 128) rt = a.get("retrigger_channel") if rt is not None and ch(rt) >= 128: self.retrigger() # ---------------- Status ---------------- def stop(self): self.running = False if self._monitor is not None: self._monitor.join(timeout=1.0) self._monitor = None if self.pipeline is not None: self.pipeline.set_state(Gst.State.NULL) self.pipeline = None def status(self) -> dict: with self._lock: return { "running": self.running, "layers": [self._layer_dict(l) for l in self.layers], "master": round(self.master, 3), "blackout": self.blackout, "playing": self.playing, "frames_rendered": self.frame_count, "dmx": dict(self.dmx_stats), "error": self.last_error, } def _layer_dict(self, l: Layer) -> dict: d = l.to_dict() d["effective"] = round(self._effective_alpha(l), 3) return d