pipeline.py 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361
  1. # -*- coding: utf-8 -*-
  2. """
  3. Audio Processing Pipeline.
  4. Verwaltet eine Kette von Audio-Prozessoren.
  5. """
  6. import logging
  7. import threading
  8. from typing import TYPE_CHECKING
  9. from trixy_core.audio.processing.context import AudioProcessingContext
  10. from trixy_core.audio.processing.processor import AudioProcessor
  11. if TYPE_CHECKING:
  12. from trixy_core.music.track import Track
  13. class AudioProcessingPipeline:
  14. """
  15. Pipeline für Audio-Verarbeitung.
  16. Verwaltet eine sortierte Liste von Audio-Prozessoren,
  17. die nacheinander auf Audio-Chunks angewendet werden.
  18. Features:
  19. - Prozessoren nach Priorität sortiert
  20. - Thread-safe
  21. - Kontext-Management
  22. - Fehlerbehandlung (Chunk wird unverändert weitergegeben)
  23. Usage:
  24. pipeline = AudioProcessingPipeline()
  25. # Prozessoren hinzufügen
  26. pipeline.add(DuckingProcessor())
  27. pipeline.add(CrossfadeProcessor())
  28. # Audio verarbeiten
  29. processed = pipeline.process(chunk, context)
  30. """
  31. def __init__(self) -> None:
  32. """Initialisiert die Pipeline."""
  33. self._processors: list[AudioProcessor] = []
  34. self._lock = threading.RLock()
  35. self._logger = logging.getLogger(__name__)
  36. # Kontext-Zustand (wird von außen gesetzt)
  37. self._is_wakeword_active = False
  38. self._is_conversation_active = False
  39. self._is_tts_playing = False
  40. self._current_track: "Track | None" = None
  41. self._next_track: "Track | None" = None
  42. # ==========================================================================
  43. # Processor Management
  44. # ==========================================================================
  45. def add(self, processor: AudioProcessor) -> None:
  46. """
  47. Fügt einen Prozessor zur Pipeline hinzu.
  48. Prozessoren werden nach Priorität sortiert (niedriger = früher).
  49. Args:
  50. processor: Der hinzuzufügende Prozessor
  51. """
  52. with self._lock:
  53. # Prüfen ob bereits vorhanden
  54. if any(p.id == processor.id for p in self._processors):
  55. self._logger.warning(
  56. f"Prozessor {processor.id} bereits vorhanden, wird ersetzt"
  57. )
  58. self._processors = [
  59. p for p in self._processors if p.id != processor.id
  60. ]
  61. self._processors.append(processor)
  62. self._processors.sort(key=lambda p: p.priority)
  63. self._logger.debug(
  64. f"Prozessor hinzugefügt: {processor.name} (Priorität: {processor.priority})"
  65. )
  66. def remove(self, processor_id: str) -> bool:
  67. """
  68. Entfernt einen Prozessor aus der Pipeline.
  69. Args:
  70. processor_id: ID des zu entfernenden Prozessors
  71. Returns:
  72. True wenn entfernt wurde
  73. """
  74. with self._lock:
  75. before = len(self._processors)
  76. self._processors = [
  77. p for p in self._processors if p.id != processor_id
  78. ]
  79. removed = len(self._processors) < before
  80. if removed:
  81. self._logger.debug(f"Prozessor entfernt: {processor_id}")
  82. return removed
  83. def get(self, processor_id: str) -> AudioProcessor | None:
  84. """
  85. Holt einen Prozessor nach ID.
  86. Args:
  87. processor_id: ID des Prozessors
  88. Returns:
  89. Prozessor oder None
  90. """
  91. with self._lock:
  92. for p in self._processors:
  93. if p.id == processor_id:
  94. return p
  95. return None
  96. def enable(self, processor_id: str) -> bool:
  97. """
  98. Aktiviert einen Prozessor.
  99. Args:
  100. processor_id: ID des Prozessors
  101. Returns:
  102. True wenn gefunden
  103. """
  104. processor = self.get(processor_id)
  105. if processor:
  106. processor.enabled = True
  107. return True
  108. return False
  109. def disable(self, processor_id: str) -> bool:
  110. """
  111. Deaktiviert einen Prozessor.
  112. Args:
  113. processor_id: ID des Prozessors
  114. Returns:
  115. True wenn gefunden
  116. """
  117. processor = self.get(processor_id)
  118. if processor:
  119. processor.enabled = False
  120. return True
  121. return False
  122. def clear(self) -> None:
  123. """Entfernt alle Prozessoren."""
  124. with self._lock:
  125. self._processors.clear()
  126. self._logger.debug("Alle Prozessoren entfernt")
  127. @property
  128. def processors(self) -> list[AudioProcessor]:
  129. """Liste aller Prozessoren (sortiert nach Priorität)."""
  130. with self._lock:
  131. return list(self._processors)
  132. @property
  133. def count(self) -> int:
  134. """Anzahl der Prozessoren."""
  135. return len(self._processors)
  136. # ==========================================================================
  137. # State Management
  138. # ==========================================================================
  139. def set_wakeword_active(self, active: bool) -> None:
  140. """Setzt Wakeword-Status."""
  141. if self._is_wakeword_active != active:
  142. self._is_wakeword_active = active
  143. self._notify_state_change("wakeword_active", active)
  144. def set_conversation_active(self, active: bool) -> None:
  145. """Setzt Conversation-Status."""
  146. if self._is_conversation_active != active:
  147. self._is_conversation_active = active
  148. self._notify_state_change("conversation_active", active)
  149. def set_tts_playing(self, playing: bool) -> None:
  150. """Setzt TTS-Wiedergabe-Status."""
  151. if self._is_tts_playing != playing:
  152. self._is_tts_playing = playing
  153. self._notify_state_change("tts_playing", playing)
  154. def set_current_track(self, track: "Track | None") -> None:
  155. """Setzt den aktuellen Track."""
  156. self._current_track = track
  157. def set_next_track(self, track: "Track | None") -> None:
  158. """Setzt den nächsten Track."""
  159. self._next_track = track
  160. def _notify_state_change(self, state_name: str, value: bool) -> None:
  161. """Benachrichtigt alle Prozessoren über Zustandsänderung."""
  162. with self._lock:
  163. for processor in self._processors:
  164. if processor.enabled:
  165. try:
  166. processor.on_state_change(state_name, value)
  167. except Exception as e:
  168. self._logger.error(
  169. f"Fehler in on_state_change von {processor.id}: {e}"
  170. )
  171. # ==========================================================================
  172. # Audio Processing
  173. # ==========================================================================
  174. def process(
  175. self,
  176. chunk: bytes,
  177. context: AudioProcessingContext | None = None,
  178. ) -> bytes:
  179. """
  180. Verarbeitet einen Audio-Chunk durch alle aktiven Prozessoren.
  181. Args:
  182. chunk: PCM-Audio-Daten
  183. context: Optionaler Kontext (wird automatisch erstellt wenn None)
  184. Returns:
  185. Verarbeiteter Chunk
  186. """
  187. if not chunk:
  188. return chunk
  189. # Kontext erstellen falls nicht vorhanden
  190. if context is None:
  191. context = self._create_default_context()
  192. # Pipeline-Status in Kontext übernehmen
  193. context.is_wakeword_active = self._is_wakeword_active
  194. context.is_conversation_active = self._is_conversation_active
  195. context.is_tts_playing = self._is_tts_playing
  196. context.track = self._current_track
  197. context.next_track = self._next_track
  198. # Durch alle Prozessoren (gefiltert nach enabled + audio_type)
  199. with self._lock:
  200. processors = [p for p in self._processors if p.should_process(context)]
  201. for processor in processors:
  202. try:
  203. result = processor.process(chunk, context)
  204. # Validierung: Chunk-Länge darf sich nicht ändern
  205. if len(result) != len(chunk):
  206. self._logger.warning(
  207. f"Prozessor {processor.id} hat Chunk-Länge geändert "
  208. f"({len(chunk)} -> {len(result)}), ignoriert"
  209. )
  210. continue
  211. chunk = result
  212. except Exception as e:
  213. self._logger.error(
  214. f"Fehler in Prozessor {processor.id}: {e}"
  215. )
  216. # Bei Fehler wird der unveränderte Chunk weitergegeben
  217. return chunk
  218. def _create_default_context(self) -> AudioProcessingContext:
  219. """Erstellt einen Standard-Kontext."""
  220. return AudioProcessingContext(
  221. is_wakeword_active=self._is_wakeword_active,
  222. is_conversation_active=self._is_conversation_active,
  223. is_tts_playing=self._is_tts_playing,
  224. track=self._current_track,
  225. next_track=self._next_track,
  226. )
  227. # ==========================================================================
  228. # Track Events
  229. # ==========================================================================
  230. def notify_track_start(self, context: AudioProcessingContext) -> None:
  231. """Benachrichtigt alle Prozessoren über Track-Start."""
  232. with self._lock:
  233. for processor in self._processors:
  234. if processor.enabled:
  235. try:
  236. processor.on_track_start(context)
  237. except Exception as e:
  238. self._logger.error(
  239. f"Fehler in on_track_start von {processor.id}: {e}"
  240. )
  241. def notify_track_end(self, context: AudioProcessingContext) -> None:
  242. """Benachrichtigt alle Prozessoren über Track-Ende."""
  243. with self._lock:
  244. for processor in self._processors:
  245. if processor.enabled:
  246. try:
  247. processor.on_track_end(context)
  248. except Exception as e:
  249. self._logger.error(
  250. f"Fehler in on_track_end von {processor.id}: {e}"
  251. )
  252. def reset_all(self) -> None:
  253. """Setzt alle Prozessoren zurück."""
  254. with self._lock:
  255. for processor in self._processors:
  256. try:
  257. processor.reset()
  258. except Exception as e:
  259. self._logger.error(
  260. f"Fehler in reset von {processor.id}: {e}"
  261. )
  262. # Zustand zurücksetzen
  263. self._is_wakeword_active = False
  264. self._is_conversation_active = False
  265. self._is_tts_playing = False
  266. self._current_track = None
  267. self._next_track = None
  268. # ==========================================================================
  269. # Status
  270. # ==========================================================================
  271. def get_status(self) -> dict:
  272. """Liefert Pipeline-Status."""
  273. with self._lock:
  274. return {
  275. "processor_count": len(self._processors),
  276. "active_count": sum(1 for p in self._processors if p.enabled),
  277. "is_wakeword_active": self._is_wakeword_active,
  278. "is_conversation_active": self._is_conversation_active,
  279. "is_tts_playing": self._is_tts_playing,
  280. "processors": [
  281. {
  282. "id": p.id,
  283. "name": p.name,
  284. "priority": p.priority,
  285. "enabled": p.enabled,
  286. }
  287. for p in self._processors
  288. ],
  289. }
  290. def __len__(self) -> int:
  291. return len(self._processors)
  292. def __contains__(self, processor_id: str) -> bool:
  293. return any(p.id == processor_id for p in self._processors)
  294. def __repr__(self) -> str:
  295. return f"AudioProcessingPipeline(processors={len(self._processors)})"