| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361 |
- # -*- coding: utf-8 -*-
- """
- Audio Processing Pipeline.
- Verwaltet eine Kette von Audio-Prozessoren.
- """
- import logging
- import threading
- from typing import TYPE_CHECKING
- from trixy_core.audio.processing.context import AudioProcessingContext
- from trixy_core.audio.processing.processor import AudioProcessor
- if TYPE_CHECKING:
- from trixy_core.music.track import Track
- class AudioProcessingPipeline:
- """
- Pipeline für Audio-Verarbeitung.
- Verwaltet eine sortierte Liste von Audio-Prozessoren,
- die nacheinander auf Audio-Chunks angewendet werden.
- Features:
- - Prozessoren nach Priorität sortiert
- - Thread-safe
- - Kontext-Management
- - Fehlerbehandlung (Chunk wird unverändert weitergegeben)
- Usage:
- pipeline = AudioProcessingPipeline()
- # Prozessoren hinzufügen
- pipeline.add(DuckingProcessor())
- pipeline.add(CrossfadeProcessor())
- # Audio verarbeiten
- processed = pipeline.process(chunk, context)
- """
- def __init__(self) -> None:
- """Initialisiert die Pipeline."""
- self._processors: list[AudioProcessor] = []
- self._lock = threading.RLock()
- self._logger = logging.getLogger(__name__)
- # Kontext-Zustand (wird von außen gesetzt)
- self._is_wakeword_active = False
- self._is_conversation_active = False
- self._is_tts_playing = False
- self._current_track: "Track | None" = None
- self._next_track: "Track | None" = None
- # ==========================================================================
- # Processor Management
- # ==========================================================================
- def add(self, processor: AudioProcessor) -> None:
- """
- Fügt einen Prozessor zur Pipeline hinzu.
- Prozessoren werden nach Priorität sortiert (niedriger = früher).
- Args:
- processor: Der hinzuzufügende Prozessor
- """
- with self._lock:
- # Prüfen ob bereits vorhanden
- if any(p.id == processor.id for p in self._processors):
- self._logger.warning(
- f"Prozessor {processor.id} bereits vorhanden, wird ersetzt"
- )
- self._processors = [
- p for p in self._processors if p.id != processor.id
- ]
- self._processors.append(processor)
- self._processors.sort(key=lambda p: p.priority)
- self._logger.debug(
- f"Prozessor hinzugefügt: {processor.name} (Priorität: {processor.priority})"
- )
- def remove(self, processor_id: str) -> bool:
- """
- Entfernt einen Prozessor aus der Pipeline.
- Args:
- processor_id: ID des zu entfernenden Prozessors
- Returns:
- True wenn entfernt wurde
- """
- with self._lock:
- before = len(self._processors)
- self._processors = [
- p for p in self._processors if p.id != processor_id
- ]
- removed = len(self._processors) < before
- if removed:
- self._logger.debug(f"Prozessor entfernt: {processor_id}")
- return removed
- def get(self, processor_id: str) -> AudioProcessor | None:
- """
- Holt einen Prozessor nach ID.
- Args:
- processor_id: ID des Prozessors
- Returns:
- Prozessor oder None
- """
- with self._lock:
- for p in self._processors:
- if p.id == processor_id:
- return p
- return None
- def enable(self, processor_id: str) -> bool:
- """
- Aktiviert einen Prozessor.
- Args:
- processor_id: ID des Prozessors
- Returns:
- True wenn gefunden
- """
- processor = self.get(processor_id)
- if processor:
- processor.enabled = True
- return True
- return False
- def disable(self, processor_id: str) -> bool:
- """
- Deaktiviert einen Prozessor.
- Args:
- processor_id: ID des Prozessors
- Returns:
- True wenn gefunden
- """
- processor = self.get(processor_id)
- if processor:
- processor.enabled = False
- return True
- return False
- def clear(self) -> None:
- """Entfernt alle Prozessoren."""
- with self._lock:
- self._processors.clear()
- self._logger.debug("Alle Prozessoren entfernt")
- @property
- def processors(self) -> list[AudioProcessor]:
- """Liste aller Prozessoren (sortiert nach Priorität)."""
- with self._lock:
- return list(self._processors)
- @property
- def count(self) -> int:
- """Anzahl der Prozessoren."""
- return len(self._processors)
- # ==========================================================================
- # State Management
- # ==========================================================================
- def set_wakeword_active(self, active: bool) -> None:
- """Setzt Wakeword-Status."""
- if self._is_wakeword_active != active:
- self._is_wakeword_active = active
- self._notify_state_change("wakeword_active", active)
- def set_conversation_active(self, active: bool) -> None:
- """Setzt Conversation-Status."""
- if self._is_conversation_active != active:
- self._is_conversation_active = active
- self._notify_state_change("conversation_active", active)
- def set_tts_playing(self, playing: bool) -> None:
- """Setzt TTS-Wiedergabe-Status."""
- if self._is_tts_playing != playing:
- self._is_tts_playing = playing
- self._notify_state_change("tts_playing", playing)
- def set_current_track(self, track: "Track | None") -> None:
- """Setzt den aktuellen Track."""
- self._current_track = track
- def set_next_track(self, track: "Track | None") -> None:
- """Setzt den nächsten Track."""
- self._next_track = track
- def _notify_state_change(self, state_name: str, value: bool) -> None:
- """Benachrichtigt alle Prozessoren über Zustandsänderung."""
- with self._lock:
- for processor in self._processors:
- if processor.enabled:
- try:
- processor.on_state_change(state_name, value)
- except Exception as e:
- self._logger.error(
- f"Fehler in on_state_change von {processor.id}: {e}"
- )
- # ==========================================================================
- # Audio Processing
- # ==========================================================================
- def process(
- self,
- chunk: bytes,
- context: AudioProcessingContext | None = None,
- ) -> bytes:
- """
- Verarbeitet einen Audio-Chunk durch alle aktiven Prozessoren.
- Args:
- chunk: PCM-Audio-Daten
- context: Optionaler Kontext (wird automatisch erstellt wenn None)
- Returns:
- Verarbeiteter Chunk
- """
- if not chunk:
- return chunk
- # Kontext erstellen falls nicht vorhanden
- if context is None:
- context = self._create_default_context()
- # Pipeline-Status in Kontext übernehmen
- context.is_wakeword_active = self._is_wakeword_active
- context.is_conversation_active = self._is_conversation_active
- context.is_tts_playing = self._is_tts_playing
- context.track = self._current_track
- context.next_track = self._next_track
- # Durch alle Prozessoren (gefiltert nach enabled + audio_type)
- with self._lock:
- processors = [p for p in self._processors if p.should_process(context)]
- for processor in processors:
- try:
- result = processor.process(chunk, context)
- # Validierung: Chunk-Länge darf sich nicht ändern
- if len(result) != len(chunk):
- self._logger.warning(
- f"Prozessor {processor.id} hat Chunk-Länge geändert "
- f"({len(chunk)} -> {len(result)}), ignoriert"
- )
- continue
- chunk = result
- except Exception as e:
- self._logger.error(
- f"Fehler in Prozessor {processor.id}: {e}"
- )
- # Bei Fehler wird der unveränderte Chunk weitergegeben
- return chunk
- def _create_default_context(self) -> AudioProcessingContext:
- """Erstellt einen Standard-Kontext."""
- return AudioProcessingContext(
- is_wakeword_active=self._is_wakeword_active,
- is_conversation_active=self._is_conversation_active,
- is_tts_playing=self._is_tts_playing,
- track=self._current_track,
- next_track=self._next_track,
- )
- # ==========================================================================
- # Track Events
- # ==========================================================================
- def notify_track_start(self, context: AudioProcessingContext) -> None:
- """Benachrichtigt alle Prozessoren über Track-Start."""
- with self._lock:
- for processor in self._processors:
- if processor.enabled:
- try:
- processor.on_track_start(context)
- except Exception as e:
- self._logger.error(
- f"Fehler in on_track_start von {processor.id}: {e}"
- )
- def notify_track_end(self, context: AudioProcessingContext) -> None:
- """Benachrichtigt alle Prozessoren über Track-Ende."""
- with self._lock:
- for processor in self._processors:
- if processor.enabled:
- try:
- processor.on_track_end(context)
- except Exception as e:
- self._logger.error(
- f"Fehler in on_track_end von {processor.id}: {e}"
- )
- def reset_all(self) -> None:
- """Setzt alle Prozessoren zurück."""
- with self._lock:
- for processor in self._processors:
- try:
- processor.reset()
- except Exception as e:
- self._logger.error(
- f"Fehler in reset von {processor.id}: {e}"
- )
- # Zustand zurücksetzen
- self._is_wakeword_active = False
- self._is_conversation_active = False
- self._is_tts_playing = False
- self._current_track = None
- self._next_track = None
- # ==========================================================================
- # Status
- # ==========================================================================
- def get_status(self) -> dict:
- """Liefert Pipeline-Status."""
- with self._lock:
- return {
- "processor_count": len(self._processors),
- "active_count": sum(1 for p in self._processors if p.enabled),
- "is_wakeword_active": self._is_wakeword_active,
- "is_conversation_active": self._is_conversation_active,
- "is_tts_playing": self._is_tts_playing,
- "processors": [
- {
- "id": p.id,
- "name": p.name,
- "priority": p.priority,
- "enabled": p.enabled,
- }
- for p in self._processors
- ],
- }
- def __len__(self) -> int:
- return len(self._processors)
- def __contains__(self, processor_id: str) -> bool:
- return any(p.id == processor_id for p in self._processors)
- def __repr__(self) -> str:
- return f"AudioProcessingPipeline(processors={len(self._processors)})"
|