# -*- coding: utf-8 -*- """ NetworkService - TCP-Server für Trixy. Implementiert die Server-seitige Netzwerk-Transport-Schicht. Horcht auf 4 Ports: - 2101: Command Socket (Befehle, verschlüsselt) - 2102: Audio Input (Satellite → Server) - 2103: Audio Output (Server → Satellite) - 2104: Music Stream (Server → Satellite) """ from __future__ import annotations import asyncio import struct from typing import TYPE_CHECKING, Callable, Awaitable from trixy_core.service.iservice import IService from trixy_core.service.enums import ServicePriority, ServiceGroup from trixy_core.network.protocol import ( TrixyProtocol, ProtocolFlags, ProtocolMessage, MAGIC, MAGIC_LENGTH, COMMAND_HELLO, COMMAND_PING, COMMAND_PONG, ) from trixy_core.network.encryption import TrixyEncryption from trixy_core.network.keepalive import KeepAliveManager, KeepAliveConfig from trixy_core.network.cmd import ( SatelliteConnect, SatelliteDisconnect, SatelliteAccepted, SatelliteDenied, ) from trixy_core.satellite.satellite import Satellite, ConnectionState from trixy_core.events.event_data.basic import SatelliteConnected, SatelliteDisconnected from trixy_core.utils.debug import pinfo, pdebug, perror, pwarn if TYPE_CHECKING: from trixy_core.server import ServerApplication class ClientHandler: """ Handler für eine einzelne Client-Verbindung. Verwaltet den Lebenszyklus einer Verbindung vom Handshake bis zur Trennung. """ def __init__( self, reader: asyncio.StreamReader, writer: asyncio.StreamWriter, protocol: TrixyProtocol, service: "NetworkService", ) -> None: self._reader = reader self._writer = writer self._protocol = protocol self._service = service self._satellite: Satellite | None = None self._authenticated = False self._running = True self._peer_address = writer.get_extra_info("peername") @property def satellite(self) -> Satellite | None: """Der zugehörige Satellite.""" return self._satellite @property def is_authenticated(self) -> bool: """Ist der Client authentifiziert?""" return self._authenticated @property def peer_address(self) -> tuple: """Adresse des Clients.""" return self._peer_address async def run(self) -> None: """Hauptschleife für die Verbindung.""" try: # Handshake durchführen if not await self._perform_handshake(): return # Nachrichtenloop await self._message_loop() except asyncio.CancelledError: pdebug(f"Verbindung abgebrochen: {self._peer_address}") except ConnectionResetError: pdebug(f"Verbindung zurückgesetzt: {self._peer_address}") except Exception as e: perror(f"Verbindungsfehler {self._peer_address}: {e}") finally: await self._cleanup() async def _perform_handshake(self) -> bool: """ Führt den Handshake mit dem Client durch. Ablauf: 1. Empfange HELLO 2. Empfange SatelliteConnect 3. Prüfe Registrierung 4. Sende SatelliteAccepted/Denied Returns: True bei erfolgreichem Handshake """ pdebug(f"Handshake gestartet: {self._peer_address}") # Warte auf HELLO (max 10 Sekunden) try: hello = await asyncio.wait_for( self._read_raw(8), timeout=10.0 ) except asyncio.TimeoutError: pwarn(f"Timeout beim Handshake: {self._peer_address}") return False if hello != COMMAND_HELLO: pwarn(f"Ungültiges HELLO von {self._peer_address}: {hello}") return False pdebug(f"HELLO empfangen von {self._peer_address}") # Warte auf SatelliteConnect try: message = await asyncio.wait_for( self._read_message(), timeout=10.0 ) except asyncio.TimeoutError: pwarn(f"Timeout beim Empfang von SatelliteConnect: {self._peer_address}") return False if message is None or message.class_name != "SatelliteConnect": pwarn(f"Erwartete SatelliteConnect, erhielt: {message.class_name if message else 'None'}") return False connect_data = message.data if isinstance(connect_data, dict): connect_msg = SatelliteConnect(**{ k: v for k, v in connect_data.items() if k in SatelliteConnect.__dataclass_fields__ }) elif isinstance(connect_data, SatelliteConnect): connect_msg = connect_data else: pwarn(f"Ungültige SatelliteConnect-Daten: {type(connect_data)}") return False pdebug(f"SatelliteConnect empfangen: {connect_msg.alias} ({connect_msg.mac_address})") app = self._service.application # Pruefe ob der Server vollstaendig bereit ist # (z.B. alle Plugins geladen — TTS, NLP, etc.) if hasattr(app, "services") and not app.services.is_ready: pwarn( f"Server noch nicht bereit — weise {connect_msg.alias} " f"({connect_msg.mac_address}) ab" ) denied = SatelliteDenied( reason="Server wird gestartet — bitte erneut verbinden", retry_allowed=True, retry_after_seconds=10, ) await self._send_message(denied) return False # Prüfe Registrierung is_registered = app.registration.is_registered(connect_msg.mac_address) if not is_registered: # Prüfe ob Registrierung erlaubt (Permanent/Zeitfenster/Datei-Trigger) if app.registration.should_auto_register(): # Starte Registrierung app.registration.begin_registration( connect_msg.mac_address, connect_msg.room, connect_msg.alias ) # Vervollständige Registrierung automatisch satellite = app.registration.complete_registration(connect_msg.mac_address) if satellite: app.satellites.add(satellite) is_registered = True if not is_registered: pwarn(f"Unregistrierter Satellite: {connect_msg.mac_address}") denied = SatelliteDenied( reason="Satellite nicht registriert", retry_allowed=True, retry_after_seconds=60 ) await self._send_message(denied) return False # Hole oder erstelle Satellite satellite = app.satellites.get_by_mac(connect_msg.mac_address) if satellite is None: satellite = app.registration.load_registered(connect_msg.mac_address) if satellite: app.satellites.add(satellite) if satellite is None: perror(f"Satellite konnte nicht geladen werden: {connect_msg.mac_address}") denied = SatelliteDenied( reason="Interner Fehler", retry_allowed=True, retry_after_seconds=5 ) await self._send_message(denied) return False # Aktualisiere Satellite-Informationen satellite.ip_address = self._peer_address[0] if self._peer_address else "" satellite.alias = connect_msg.alias or satellite.alias satellite.room_id = connect_msg.room or satellite.room_id satellite.version = connect_msg.version # Metadata mit System-Infos befuellen if connect_msg.hostname: satellite.metadata["hostname"] = connect_msg.hostname if connect_msg.os_name: satellite.metadata["os_name"] = connect_msg.os_name if connect_msg.os_version: satellite.metadata["os_version"] = connect_msg.os_version # Session-Key generieren session_key = TrixyEncryption.generate_key() # Sockets zuweisen satellite.sockets.command = self._writer # Als verbunden markieren satellite.set_connected() self._satellite = satellite self._authenticated = True # Registriere bei KeepAlive await self._service.keepalive.register( satellite.id, lambda: self._send_ping(), auto_start=True ) # Sende Bestätigung mit Audio-Format-Info network_config = app.server_config.network audio_config = app.server_config.audio # TTS-Format vom aktiven TTS-Provider holen tts_sample_rate = 22050 # Default tts_channels = 1 try: extensions = getattr(app, "extensions", None) if extensions: tts_point = extensions.get_point("conversation.tts") if tts_point: tts_provider = tts_point.get_first_instance() if tts_provider and hasattr(tts_provider, "config"): rate = tts_provider.config.sample_rate chans = tts_provider.config.channels # Nur verwenden wenn es echte int-Werte sind if isinstance(rate, int) and isinstance(chans, int): tts_sample_rate = rate tts_channels = chans except Exception: pass # Bei Fehlern Default-Werte verwenden # Musik-Format aus Server-Config music_sample_rate = audio_config.music_sample_rate music_channels = audio_config.music_channels music_bit_depth = getattr(audio_config, 'music_bit_depth', 16) accepted = SatelliteAccepted( satellite_id=satellite.id, session_key=session_key, audio_input_port=network_config.audio_input_port, audio_output_port=network_config.audio_output_port, music_port=network_config.music_port, # TTS-Format vom aktiven Plugin tts_sample_rate=tts_sample_rate, tts_channels=tts_channels, tts_sample_width=2, # 16-bit # Musik-Format aus Config (konfigurierbar für Audio-Device-Kompatibilität) music_sample_rate=music_sample_rate, music_channels=music_channels, music_sample_width=music_bit_depth // 8, # Bytes ) await self._send_message(accepted) pdebug(f"Musik-Format an Client: {music_sample_rate}Hz, {music_channels}ch, {music_bit_depth}-bit") pinfo(f"Satellite verbunden: {satellite.alias} ({satellite.room_id})") # Event auslösen await app.events.trigger("satellite_connected", SatelliteConnected( satellite_id=satellite.id, mac_address=satellite.mac_address, room=satellite.room_id, alias=satellite.alias, ip_address=satellite.ip_address, source="network_service" )) return True async def _message_loop(self) -> None: """Verarbeitet eingehende Nachrichten.""" while self._running: message = await self._read_message() if message is None: break await self._handle_message(message) async def _handle_message(self, message: ProtocolMessage) -> None: """ Verarbeitet eine empfangene Nachricht. Args: message: Die empfangene Nachricht """ class_name = message.class_name # Hard-coded Befehle if class_name == "TRXIPONG": # Pong empfangen - KeepAlive aktualisieren if self._satellite: self._satellite.update_heartbeat() self._service.keepalive.handle_pong(self._satellite.id) return if class_name == "TRXIPING": # Ping empfangen - Pong senden await self._send_pong() return # SatelliteDisconnect if class_name == "SatelliteDisconnect": data = message.data reason = data.get("reason", "Client-Trennung") if isinstance(data, dict) else "Client-Trennung" pinfo(f"Satellite trennt Verbindung: {reason}") self._running = False return # Pending-Response pruefen (Config-Proxy) if self._satellite: msg_id = "" data = message.data if hasattr(data, "message_id"): msg_id = data.message_id elif isinstance(data, dict): msg_id = data.get("message_id", "") if msg_id and self._satellite.resolve_pending(msg_id, message): return # Andere Nachrichten an MessageHandler weiterleiten if self._satellite: await self._service.dispatch_message(self._satellite, message) async def _read_raw(self, size: int) -> bytes: """Liest rohe Bytes vom Socket.""" data = await self._reader.read(size) if not data: raise ConnectionResetError("Verbindung geschlossen") return data async def _read_message(self) -> ProtocolMessage | None: """ Liest eine vollständige Protokoll-Nachricht. Returns: Die gelesene Nachricht oder None bei Verbindungsende """ try: # Lese Header (Magic + Version + Timestamp + Flags + Checksum + ClassNameLength) header_size = MAGIC_LENGTH + 2 + 8 + 4 + 16 + 2 # 36 bytes # Prüfe zuerst auf Hard-coded Befehle (8 bytes) # Verwende readexactly um sicherzustellen, dass wir alle 8 Bytes bekommen peek_data = await self._reader.readexactly(8) # Prüfe auf Hard-coded Befehl if self._protocol.is_hard_command(peek_data): return self._protocol.deserialize(peek_data) # Kein Hard-coded Befehl - lese Rest des Headers remaining_header = await self._reader.readexactly(header_size - 8) header = peek_data + remaining_header # Prüfe Magic if header[:MAGIC_LENGTH] != MAGIC: perror(f"Ungültige Magic: {header[:MAGIC_LENGTH]}") return None # Parse Header um Klassennamen-Länge und Daten-Länge zu bekommen # Offset nach Magic + Version + Timestamp + Flags + Checksum class_name_length = struct.unpack(">H", header[34:36])[0] # Lese Klassenname class_name = await self._reader.readexactly(class_name_length) # Lese Datenlänge (4 bytes) data_length_bytes = await self._reader.readexactly(4) data_length = struct.unpack(">I", data_length_bytes)[0] # Lese Daten data = await self._reader.readexactly(data_length) if data_length > 0 else b"" # Baue vollständige Nachricht zusammen full_message = header + class_name + data_length_bytes + data return self._protocol.deserialize(full_message) except (ConnectionResetError, asyncio.IncompleteReadError): return None except Exception as e: perror(f"Fehler beim Lesen der Nachricht: {e}") return None async def _send_message(self, message: object, flags: ProtocolFlags = ProtocolFlags.PICKLE) -> bool: """ Sendet eine Nachricht zum Client. Args: message: Die zu sendende Nachricht flags: Protokoll-Flags Returns: True bei Erfolg """ try: data = self._protocol.serialize(message, flags) self._writer.write(data) await self._writer.drain() return True except Exception as e: perror(f"Fehler beim Senden: {e}") return False async def _send_ping(self) -> None: """Sendet einen Ping.""" try: self._writer.write(COMMAND_PING) await self._writer.drain() except Exception as e: pdebug(f"Ping-Fehler: {e}") async def _send_pong(self) -> None: """Sendet einen Pong.""" try: self._writer.write(COMMAND_PONG) await self._writer.drain() except Exception as e: pdebug(f"Pong-Fehler: {e}") async def _cleanup(self, reason: str = "Verbindung getrennt") -> None: """Räumt die Verbindung auf.""" self._running = False if self._satellite: # KeepAlive entfernen await self._service.keepalive.unregister(self._satellite.id) # Event auslösen await self._service.application.events.trigger("satellite_disconnected", SatelliteDisconnected( satellite_id=self._satellite.id, room=self._satellite.room_id, alias=self._satellite.alias, reason=reason, source="network_service" )) # Satellite als getrennt markieren self._satellite.state = ConnectionState.DISCONNECTED self._satellite.sockets.command = None pinfo(f"Satellite getrennt: {self._satellite.alias}") self._satellite = None # Socket schließen try: self._writer.close() await self._writer.wait_closed() except Exception: pass class AudioStreamHandler: """ Handler für Audio-Stream-Verbindungen. Verwaltet Audio-Input (Satellite → Server) und Audio-Output (Server → Satellite) Streams. """ def __init__( self, reader: asyncio.StreamReader, writer: asyncio.StreamWriter, service: "NetworkService", stream_type: str, ) -> None: self._reader = reader self._writer = writer self._service = service self._stream_type = stream_type self._satellite: Satellite | None = None self._running = True self._peer_address = writer.get_extra_info("peername") async def run(self) -> None: """Hauptschleife für Audio-Stream.""" try: # Warte auf Satellite-ID satellite_id_bytes = await asyncio.wait_for( self._reader.read(36), # UUID-Länge timeout=10.0 ) if not satellite_id_bytes: return satellite_id = satellite_id_bytes.decode("utf-8").strip() satellite = self._service.application.satellites.get(satellite_id) if not satellite or not satellite.is_connected: pwarn(f"Audio-Stream für unbekannten Satellite: {satellite_id}") return self._satellite = satellite pdebug(f"{self._stream_type} verbunden: {satellite.alias}") # Socket zuweisen if self._stream_type == "audio_in": satellite.sockets.audio_in = self._writer await self._handle_audio_input() elif self._stream_type == "audio_out": satellite.sockets.audio_out = self._writer await self._handle_audio_output() elif self._stream_type == "music": satellite.sockets.music_out = self._writer await self._handle_music_output() except asyncio.TimeoutError: pdebug(f"Audio-Stream Timeout: {self._peer_address}") except Exception as e: perror(f"Audio-Stream Fehler: {e}") finally: await self._cleanup() async def _handle_audio_input(self) -> None: """Empfängt Audio vom Satellite.""" while self._running and self._satellite: try: # Lese Audio-Chunks (4096 bytes = 128ms bei 16KHz mono) chunk = await self._reader.read(4096) if not chunk: break # Event für Audio-Verarbeitung auslösen await self._service.application.events.emit( "audio_input_received", { "satellite_id": self._satellite.id, "audio_data": chunk, "sample_rate": 16000, "channels": 1, } ) except Exception as e: pdebug(f"Audio-Input Fehler: {e}") break async def _handle_audio_output(self) -> None: """Verwaltet Audio-Output zum Satellite (wartet auf Events).""" # Audio-Output wird über Events gesteuert # Hier nur Verbindung aufrecht halten while self._running: await asyncio.sleep(1) async def _handle_music_output(self) -> None: """Verwaltet Music-Output zum Satellite (wartet auf Events).""" # Event auslösen, dass Music-Socket verbunden ist if self._satellite: from trixy_core.events.event_data.basic import MusicSocketConnected await self._service.application.events.trigger( "music_socket_connected", MusicSocketConnected( satellite_id=self._satellite.id, alias=self._satellite.alias, room=self._satellite.room_id, ) ) pdebug(f"Music-Socket verbunden: {self._satellite.alias}") # Music-Output wird über Events gesteuert while self._running: await asyncio.sleep(1) async def _cleanup(self) -> None: """Räumt den Stream auf.""" self._running = False if self._satellite: if self._stream_type == "audio_in": self._satellite.sockets.audio_in = None elif self._stream_type == "audio_out": self._satellite.sockets.audio_out = None elif self._stream_type == "music": self._satellite.sockets.music_out = None try: self._writer.close() await self._writer.wait_closed() except Exception: pass class NetworkService(IService): """ TCP-Server für Trixy-Kommunikation. Horcht auf 4 Ports und verwaltet alle Verbindungen: - Port 2101: Command Socket (Befehle, verschlüsselt) - Port 2102: Audio Input (Satellite → Server) - Port 2103: Audio Output (Server → Satellite) - Port 2104: Music Stream (Server → Satellite) """ PRIORITY = ServicePriority.NETWORK GROUP = ServiceGroup.NETWORK NAME = "network" DEPENDENCIES: list[str] = [] def __init__(self, application: "ServerApplication") -> None: super().__init__(application) self._servers: dict[str, asyncio.Server] = {} self._handlers: dict[str, set[ClientHandler | AudioStreamHandler]] = { "command": set(), "audio_in": set(), "audio_out": set(), "music": set(), } self._protocol = TrixyProtocol() self._encryption: TrixyEncryption | None = None self._keepalive = KeepAliveManager( default_config=KeepAliveConfig( ping_interval=30.0, pong_timeout=15.0, # Erhöht von 10s max_missed_pongs=5, # Erhöht von 3 (= 5*30s = 2.5min Toleranz) ) ) self._message_handlers: list[Callable[[Satellite, ProtocolMessage], Awaitable[None]]] = [] @property def application(self) -> "ServerApplication": """Server-Anwendung.""" return self._application # type: ignore @property def keepalive(self) -> KeepAliveManager: """KeepAlive-Manager.""" return self._keepalive @property def protocol(self) -> TrixyProtocol: """Protokoll-Handler.""" return self._protocol def register_message_handler( self, handler: Callable[[Satellite, ProtocolMessage], Awaitable[None]] ) -> None: """ Registriert einen Handler für eingehende Nachrichten. Args: handler: Async-Funktion die (satellite, message) verarbeitet """ self._message_handlers.append(handler) async def dispatch_message(self, satellite: Satellite, message: ProtocolMessage) -> None: """ Verteilt eine Nachricht an alle registrierten Handler. Args: satellite: Der sendende Satellite message: Die empfangene Nachricht """ for handler in self._message_handlers: try: await handler(satellite, message) except Exception as e: perror(f"Message-Handler Fehler: {e}") async def send_to_satellite( self, satellite_id: str, message: object, flags: ProtocolFlags = ProtocolFlags.PICKLE ) -> bool: """ Sendet eine Nachricht an einen Satellite. Args: satellite_id: Ziel-Satellite-ID message: Die zu sendende Nachricht flags: Protokoll-Flags Returns: True bei Erfolg """ satellite = self.application.satellites.get(satellite_id) if not satellite or not satellite.sockets.command: return False try: writer = satellite.sockets.command data = self._protocol.serialize(message, flags) writer.write(data) await writer.drain() return True except Exception as e: perror(f"Fehler beim Senden an {satellite_id}: {e}") return False async def broadcast( self, message: object, room: str | None = None, flags: ProtocolFlags = ProtocolFlags.PICKLE ) -> int: """ Sendet eine Nachricht an alle (oder gefilterte) Satellites. Args: message: Die zu sendende Nachricht room: Optional: Nur an Satellites in diesem Raum flags: Protokoll-Flags Returns: Anzahl erfolgreicher Sendungen """ satellites = self.application.satellites.get_by_room(room) if room else self.application.satellites.get_connected() count = 0 for satellite in satellites: if await self.send_to_satellite(satellite.id, message, flags): count += 1 return count async def start(self) -> None: """Startet alle TCP-Server.""" config = self.application.server_config.network # Verschlüsselung laden if self.application.server_config.security.encryption_enabled: key_path = self.application.server_config.security.encryption_key_path try: self._encryption = TrixyEncryption.load_or_create(key_path) self._protocol.set_encryption(self._encryption) pdebug("Verschlüsselung aktiviert") except Exception as e: pwarn(f"Verschlüsselung konnte nicht geladen werden: {e}") # Command-Server starten (Port 2101) self._servers["command"] = await asyncio.start_server( self._handle_command_connection, config.bind_address, config.command_port, ) pdebug(f"Command-Server gestartet auf {config.bind_address}:{config.command_port}") # Audio-Input-Server starten (Port 2102) self._servers["audio_in"] = await asyncio.start_server( lambda r, w: self._handle_audio_connection(r, w, "audio_in"), config.bind_address, config.audio_input_port, ) pdebug(f"Audio-Input-Server gestartet auf {config.bind_address}:{config.audio_input_port}") # Audio-Output-Server starten (Port 2103) self._servers["audio_out"] = await asyncio.start_server( lambda r, w: self._handle_audio_connection(r, w, "audio_out"), config.bind_address, config.audio_output_port, ) pdebug(f"Audio-Output-Server gestartet auf {config.bind_address}:{config.audio_output_port}") # Music-Server starten (Port 2104) self._servers["music"] = await asyncio.start_server( lambda r, w: self._handle_audio_connection(r, w, "music"), config.bind_address, config.music_port, ) pdebug(f"Music-Server gestartet auf {config.bind_address}:{config.music_port}") # KeepAlive-Disconnect-Handler registrieren self._keepalive.on_disconnect(self._on_keepalive_disconnect) pinfo(f"NetworkService gestartet auf Ports {config.command_port}-{config.music_port}") async def stop(self) -> None: """Stoppt alle Server und Verbindungen.""" pinfo("NetworkService wird gestoppt...") # KeepAlive stoppen await self._keepalive.stop_all() # Alle Handler beenden for handler_type, handlers in self._handlers.items(): for handler in list(handlers): if hasattr(handler, "_running"): handler._running = False # Server schließen for name, server in self._servers.items(): server.close() await server.wait_closed() pdebug(f"{name}-Server geschlossen") self._servers.clear() pinfo("NetworkService gestoppt") async def _handle_command_connection( self, reader: asyncio.StreamReader, writer: asyncio.StreamWriter ) -> None: """Verarbeitet eine neue Command-Verbindung.""" addr = writer.get_extra_info("peername") pdebug(f"Neue Command-Verbindung von {addr}") handler = ClientHandler(reader, writer, self._protocol, self) self._handlers["command"].add(handler) try: await handler.run() finally: self._handlers["command"].discard(handler) async def _handle_audio_connection( self, reader: asyncio.StreamReader, writer: asyncio.StreamWriter, stream_type: str ) -> None: """Verarbeitet eine neue Audio-Verbindung.""" addr = writer.get_extra_info("peername") pdebug(f"Neue {stream_type}-Verbindung von {addr}") handler = AudioStreamHandler(reader, writer, self, stream_type) self._handlers[stream_type].add(handler) try: await handler.run() finally: self._handlers[stream_type].discard(handler) async def _on_keepalive_disconnect(self, connection_id: str) -> None: """Handler für KeepAlive-Disconnect.""" satellite = self.application.satellites.get(connection_id) if satellite: pwarn(f"KeepAlive-Timeout für {satellite.alias}") await satellite.disconnect("KeepAlive-Timeout") async def health_check(self) -> bool: """Prüft den Gesundheitszustand des Services.""" # Prüfe ob Server existieren und laufen if not self._servers: return False for name, server in self._servers.items(): if not server.is_serving(): return False return True