# -*- coding: utf-8 -*- """ MessageHandler für die Verarbeitung von Netzwerk-Nachrichten. Zentralisiert die Logik für das Dispatching und die Verarbeitung von eingehenden Protokoll-Nachrichten. """ from __future__ import annotations import asyncio from typing import TYPE_CHECKING, Callable, Awaitable, Any from dataclasses import dataclass, field from trixy_core.network.protocol import ProtocolMessage from trixy_core.utils.debug import pdebug, perror if TYPE_CHECKING: from trixy_core.satellite.satellite import Satellite # Events die nicht geloggt werden sollen (zu häufig/spammend) _EVENT_LOG_BLACKLIST: set[str] = { "music_output_received", "audio_output_received", } # Typ-Alias für Handler-Funktionen MessageHandlerFunc = Callable[["Satellite", ProtocolMessage], Awaitable[None]] MessageFilter = Callable[[ProtocolMessage], bool] @dataclass class HandlerRegistration: """Registrierung eines Nachricht-Handlers.""" handler: MessageHandlerFunc class_names: set[str] = field(default_factory=set) filter_func: MessageFilter | None = None priority: int = 0 class MessageHandler: """ Zentraler Handler für eingehende Protokoll-Nachrichten. Ermöglicht die Registrierung von Handlern für bestimmte Nachrichtentypen oder mit benutzerdefinierten Filtern. Usage: handler = MessageHandler() # Handler für bestimmte Nachrichtentypen handler.register(on_wakeword, ["WakewordDetected", "WakewordSelected"]) # Handler für alle Nachrichten handler.register(log_all) # Handler mit Filter handler.register(on_error, filter_func=lambda m: m.flags & ProtocolFlags.ERROR) # Nachricht verarbeiten await handler.dispatch(satellite, message) """ def __init__(self) -> None: """Initialisiert den MessageHandler.""" self._handlers: list[HandlerRegistration] = [] self._class_index: dict[str, list[HandlerRegistration]] = {} def register( self, handler: MessageHandlerFunc, class_names: list[str] | None = None, filter_func: MessageFilter | None = None, priority: int = 0, ) -> None: """ Registriert einen Handler für Nachrichten. Args: handler: Async-Funktion die (satellite, message) verarbeitet class_names: Liste von Nachrichtenklassen-Namen (None = alle) filter_func: Optionale Filterfunktion priority: Priorität (höher = früher) """ registration = HandlerRegistration( handler=handler, class_names=set(class_names) if class_names else set(), filter_func=filter_func, priority=priority, ) # Nach Priorität sortiert einfügen inserted = False for i, existing in enumerate(self._handlers): if registration.priority > existing.priority: self._handlers.insert(i, registration) inserted = True break if not inserted: self._handlers.append(registration) # Index aktualisieren if class_names: for name in class_names: if name not in self._class_index: self._class_index[name] = [] self._class_index[name].append(registration) pdebug(f"Message-Handler registriert: {handler.__name__} für {class_names or 'alle'}") def unregister(self, handler: MessageHandlerFunc) -> bool: """ Entfernt einen Handler. Args: handler: Der zu entfernende Handler Returns: True wenn entfernt """ for reg in list(self._handlers): if reg.handler == handler: self._handlers.remove(reg) # Aus Index entfernen for name, registrations in list(self._class_index.items()): if reg in registrations: registrations.remove(reg) if not registrations: del self._class_index[name] pdebug(f"Message-Handler entfernt: {handler.__name__}") return True return False async def dispatch( self, satellite: "Satellite", message: ProtocolMessage, stop_on_handled: bool = False, ) -> int: """ Verteilt eine Nachricht an passende Handler. Args: satellite: Der sendende Satellite message: Die empfangene Nachricht stop_on_handled: Bei True nach erstem Handler stoppen Returns: Anzahl aufgerufener Handler """ class_name = message.class_name called = 0 # Schneller Lookup für spezifische Handler specific_handlers = self._class_index.get(class_name, []) for reg in specific_handlers: if await self._call_handler(reg, satellite, message): called += 1 if stop_on_handled: return called # Allgemeine Handler (ohne class_names) for reg in self._handlers: if reg.class_names: continue # Bereits via Index verarbeitet if await self._call_handler(reg, satellite, message): called += 1 if stop_on_handled: return called return called async def _call_handler( self, registration: HandlerRegistration, satellite: "Satellite", message: ProtocolMessage, ) -> bool: """ Ruft einen einzelnen Handler auf. Returns: True wenn erfolgreich aufgerufen """ # Filter prüfen if registration.filter_func: try: if not registration.filter_func(message): return False except Exception as e: perror(f"Filter-Fehler: {e}") return False # Handler aufrufen try: await registration.handler(satellite, message) return True except Exception as e: perror(f"Handler-Fehler in {registration.handler.__name__}: {e}") return False def clear(self) -> None: """Entfernt alle Handler.""" self._handlers.clear() self._class_index.clear() class EventBridge: """ Brücke zwischen Netzwerk-Nachrichten und dem Event-System. Konvertiert eingehende Protokoll-Nachrichten automatisch in Events und löst sie aus. Usage: bridge = EventBridge(event_manager) bridge.map("WakewordDetected", "wakeword_detected") bridge.map("RecordingComplete", "recording_complete") # In NetworkService: await bridge.on_message(satellite, message) """ def __init__(self, event_manager: Any) -> None: """ Initialisiert die EventBridge. Args: event_manager: Der EventManager der Anwendung """ self._events = event_manager self._mappings: dict[str, str] = {} self._transformers: dict[str, Callable[[Any], dict]] = {} def map( self, class_name: str, event_name: str, transformer: Callable[[Any], dict] | None = None, ) -> None: """ Mappt eine Nachrichtenklasse auf einen Event-Namen. Args: class_name: Name der Nachrichtenklasse event_name: Ziel-Event-Name transformer: Optionale Funktion zur Daten-Transformation """ self._mappings[class_name] = event_name if transformer: self._transformers[class_name] = transformer pdebug(f"Event-Mapping: {class_name} -> {event_name}") def unmap(self, class_name: str) -> None: """Entfernt ein Mapping.""" self._mappings.pop(class_name, None) self._transformers.pop(class_name, None) async def on_message( self, satellite: "Satellite", message: ProtocolMessage, ) -> bool: """ Verarbeitet eine Nachricht und löst ggf. ein Event aus. Args: satellite: Der sendende Satellite message: Die empfangene Nachricht Returns: True wenn ein Event ausgelöst wurde """ class_name = message.class_name event_name = self._mappings.get(class_name) if not event_name: return False # Daten vorbereiten data = message.data if isinstance(data, dict): event_data = data.copy() elif hasattr(data, "__dict__"): event_data = data.__dict__.copy() else: event_data = {"data": data} # Satellite-Info hinzufügen event_data["satellite_id"] = satellite.id event_data["satellite_alias"] = satellite.alias event_data["satellite_room"] = satellite.room_id # Transformer anwenden transformer = self._transformers.get(class_name) if transformer: try: event_data = transformer(event_data) except Exception as e: perror(f"Transformer-Fehler für {class_name}: {e}") # Event auslösen (emit() für dict-Daten) try: await self._events.emit(event_name, event_data) # Nicht nochmal loggen - EventManager.trigger() loggt bereits return True except Exception as e: perror(f"Event-Fehler: {e}") return False # Standard-Mappings für häufige Nachrichten DEFAULT_EVENT_MAPPINGS = { "WakewordDetected": "wakeword_detected", "WakewordSelected": "wakeword_selected", "WakewordAbort": "wakeword_abort", "RecordingComplete": "recording_complete", "TranscriptionResult": "transcription_result", "IntentResult": "intent_result", "AssistantResponse": "assistant_response", "ConversationEnd": "conversation_end", "TextInput": "text_input_received", # Media Commands -> Events "MusicPlayPause": "music_play_pause", "MusicNext": "music_next", "MusicPrevious": "music_previous", "MusicStop": "music_stop", "MediaStopAll": "media_stop_all", "MusicVolumeChange": "music_volume_change", "MusicStatus": "music_status_request", # Server-gesteuerte Aufnahme "SatelliteRecordStopped": "satellite_record_stopped", } def setup_default_event_bridge(event_manager: Any) -> EventBridge: """ Erstellt eine EventBridge mit Standard-Mappings. Args: event_manager: Der EventManager der Anwendung Returns: Konfigurierte EventBridge """ bridge = EventBridge(event_manager) for class_name, event_name in DEFAULT_EVENT_MAPPINGS.items(): bridge.map(class_name, event_name) return bridge