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