| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360 |
- # -*- 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
|