service.py 31 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890
  1. # -*- coding: utf-8 -*-
  2. """
  3. NetworkService - TCP-Server für Trixy.
  4. Implementiert die Server-seitige Netzwerk-Transport-Schicht.
  5. Horcht auf 4 Ports:
  6. - 2101: Command Socket (Befehle, verschlüsselt)
  7. - 2102: Audio Input (Satellite → Server)
  8. - 2103: Audio Output (Server → Satellite)
  9. - 2104: Music Stream (Server → Satellite)
  10. """
  11. from __future__ import annotations
  12. import asyncio
  13. import struct
  14. from typing import TYPE_CHECKING, Callable, Awaitable
  15. from trixy_core.service.iservice import IService
  16. from trixy_core.service.enums import ServicePriority, ServiceGroup
  17. from trixy_core.network.protocol import (
  18. TrixyProtocol,
  19. ProtocolFlags,
  20. ProtocolMessage,
  21. MAGIC,
  22. MAGIC_LENGTH,
  23. COMMAND_HELLO,
  24. COMMAND_PING,
  25. COMMAND_PONG,
  26. )
  27. from trixy_core.network.encryption import TrixyEncryption
  28. from trixy_core.network.keepalive import KeepAliveManager, KeepAliveConfig
  29. from trixy_core.network.cmd import (
  30. SatelliteConnect,
  31. SatelliteDisconnect,
  32. SatelliteAccepted,
  33. SatelliteDenied,
  34. )
  35. from trixy_core.satellite.satellite import Satellite, ConnectionState
  36. from trixy_core.events.event_data.basic import SatelliteConnected, SatelliteDisconnected
  37. from trixy_core.utils.debug import pinfo, pdebug, perror, pwarn
  38. if TYPE_CHECKING:
  39. from trixy_core.server import ServerApplication
  40. class ClientHandler:
  41. """
  42. Handler für eine einzelne Client-Verbindung.
  43. Verwaltet den Lebenszyklus einer Verbindung vom Handshake
  44. bis zur Trennung.
  45. """
  46. def __init__(
  47. self,
  48. reader: asyncio.StreamReader,
  49. writer: asyncio.StreamWriter,
  50. protocol: TrixyProtocol,
  51. service: "NetworkService",
  52. ) -> None:
  53. self._reader = reader
  54. self._writer = writer
  55. self._protocol = protocol
  56. self._service = service
  57. self._satellite: Satellite | None = None
  58. self._authenticated = False
  59. self._running = True
  60. self._peer_address = writer.get_extra_info("peername")
  61. @property
  62. def satellite(self) -> Satellite | None:
  63. """Der zugehörige Satellite."""
  64. return self._satellite
  65. @property
  66. def is_authenticated(self) -> bool:
  67. """Ist der Client authentifiziert?"""
  68. return self._authenticated
  69. @property
  70. def peer_address(self) -> tuple:
  71. """Adresse des Clients."""
  72. return self._peer_address
  73. async def run(self) -> None:
  74. """Hauptschleife für die Verbindung."""
  75. try:
  76. # Handshake durchführen
  77. if not await self._perform_handshake():
  78. return
  79. # Nachrichtenloop
  80. await self._message_loop()
  81. except asyncio.CancelledError:
  82. pdebug(f"Verbindung abgebrochen: {self._peer_address}")
  83. except ConnectionResetError:
  84. pdebug(f"Verbindung zurückgesetzt: {self._peer_address}")
  85. except Exception as e:
  86. perror(f"Verbindungsfehler {self._peer_address}: {e}")
  87. finally:
  88. await self._cleanup()
  89. async def _perform_handshake(self) -> bool:
  90. """
  91. Führt den Handshake mit dem Client durch.
  92. Ablauf:
  93. 1. Empfange HELLO
  94. 2. Empfange SatelliteConnect
  95. 3. Prüfe Registrierung
  96. 4. Sende SatelliteAccepted/Denied
  97. Returns:
  98. True bei erfolgreichem Handshake
  99. """
  100. pdebug(f"Handshake gestartet: {self._peer_address}")
  101. # Warte auf HELLO (max 10 Sekunden)
  102. try:
  103. hello = await asyncio.wait_for(
  104. self._read_raw(8),
  105. timeout=10.0
  106. )
  107. except asyncio.TimeoutError:
  108. pwarn(f"Timeout beim Handshake: {self._peer_address}")
  109. return False
  110. if hello != COMMAND_HELLO:
  111. pwarn(f"Ungültiges HELLO von {self._peer_address}: {hello}")
  112. return False
  113. pdebug(f"HELLO empfangen von {self._peer_address}")
  114. # Warte auf SatelliteConnect
  115. try:
  116. message = await asyncio.wait_for(
  117. self._read_message(),
  118. timeout=10.0
  119. )
  120. except asyncio.TimeoutError:
  121. pwarn(f"Timeout beim Empfang von SatelliteConnect: {self._peer_address}")
  122. return False
  123. if message is None or message.class_name != "SatelliteConnect":
  124. pwarn(f"Erwartete SatelliteConnect, erhielt: {message.class_name if message else 'None'}")
  125. return False
  126. connect_data = message.data
  127. if isinstance(connect_data, dict):
  128. connect_msg = SatelliteConnect(**{
  129. k: v for k, v in connect_data.items()
  130. if k in SatelliteConnect.__dataclass_fields__
  131. })
  132. elif isinstance(connect_data, SatelliteConnect):
  133. connect_msg = connect_data
  134. else:
  135. pwarn(f"Ungültige SatelliteConnect-Daten: {type(connect_data)}")
  136. return False
  137. pdebug(f"SatelliteConnect empfangen: {connect_msg.alias} ({connect_msg.mac_address})")
  138. app = self._service.application
  139. # Pruefe ob der Server vollstaendig bereit ist
  140. # (z.B. alle Plugins geladen — TTS, NLP, etc.)
  141. if hasattr(app, "services") and not app.services.is_ready:
  142. pwarn(
  143. f"Server noch nicht bereit — weise {connect_msg.alias} "
  144. f"({connect_msg.mac_address}) ab"
  145. )
  146. denied = SatelliteDenied(
  147. reason="Server wird gestartet — bitte erneut verbinden",
  148. retry_allowed=True,
  149. retry_after_seconds=10,
  150. )
  151. await self._send_message(denied)
  152. return False
  153. # Prüfe Registrierung
  154. is_registered = app.registration.is_registered(connect_msg.mac_address)
  155. if not is_registered:
  156. # Prüfe ob Registrierung erlaubt (Permanent/Zeitfenster/Datei-Trigger)
  157. if app.registration.should_auto_register():
  158. # Starte Registrierung
  159. app.registration.begin_registration(
  160. connect_msg.mac_address,
  161. connect_msg.room,
  162. connect_msg.alias
  163. )
  164. # Vervollständige Registrierung automatisch
  165. satellite = app.registration.complete_registration(connect_msg.mac_address)
  166. if satellite:
  167. app.satellites.add(satellite)
  168. is_registered = True
  169. if not is_registered:
  170. pwarn(f"Unregistrierter Satellite: {connect_msg.mac_address}")
  171. denied = SatelliteDenied(
  172. reason="Satellite nicht registriert",
  173. retry_allowed=True,
  174. retry_after_seconds=60
  175. )
  176. await self._send_message(denied)
  177. return False
  178. # Hole oder erstelle Satellite
  179. satellite = app.satellites.get_by_mac(connect_msg.mac_address)
  180. if satellite is None:
  181. satellite = app.registration.load_registered(connect_msg.mac_address)
  182. if satellite:
  183. app.satellites.add(satellite)
  184. if satellite is None:
  185. perror(f"Satellite konnte nicht geladen werden: {connect_msg.mac_address}")
  186. denied = SatelliteDenied(
  187. reason="Interner Fehler",
  188. retry_allowed=True,
  189. retry_after_seconds=5
  190. )
  191. await self._send_message(denied)
  192. return False
  193. # Aktualisiere Satellite-Informationen
  194. satellite.ip_address = self._peer_address[0] if self._peer_address else ""
  195. satellite.alias = connect_msg.alias or satellite.alias
  196. satellite.room_id = connect_msg.room or satellite.room_id
  197. satellite.version = connect_msg.version
  198. # Metadata mit System-Infos befuellen
  199. if connect_msg.hostname:
  200. satellite.metadata["hostname"] = connect_msg.hostname
  201. if connect_msg.os_name:
  202. satellite.metadata["os_name"] = connect_msg.os_name
  203. if connect_msg.os_version:
  204. satellite.metadata["os_version"] = connect_msg.os_version
  205. # Session-Key generieren
  206. session_key = TrixyEncryption.generate_key()
  207. # Sockets zuweisen
  208. satellite.sockets.command = self._writer
  209. # Als verbunden markieren
  210. satellite.set_connected()
  211. self._satellite = satellite
  212. self._authenticated = True
  213. # Registriere bei KeepAlive
  214. await self._service.keepalive.register(
  215. satellite.id,
  216. lambda: self._send_ping(),
  217. auto_start=True
  218. )
  219. # Sende Bestätigung mit Audio-Format-Info
  220. network_config = app.server_config.network
  221. audio_config = app.server_config.audio
  222. # TTS-Format vom aktiven TTS-Provider holen
  223. tts_sample_rate = 22050 # Default
  224. tts_channels = 1
  225. try:
  226. extensions = getattr(app, "extensions", None)
  227. if extensions:
  228. tts_point = extensions.get_point("conversation.tts")
  229. if tts_point:
  230. tts_provider = tts_point.get_first_instance()
  231. if tts_provider and hasattr(tts_provider, "config"):
  232. rate = tts_provider.config.sample_rate
  233. chans = tts_provider.config.channels
  234. # Nur verwenden wenn es echte int-Werte sind
  235. if isinstance(rate, int) and isinstance(chans, int):
  236. tts_sample_rate = rate
  237. tts_channels = chans
  238. except Exception:
  239. pass # Bei Fehlern Default-Werte verwenden
  240. # Musik-Format aus Server-Config
  241. music_sample_rate = audio_config.music_sample_rate
  242. music_channels = audio_config.music_channels
  243. music_bit_depth = getattr(audio_config, 'music_bit_depth', 16)
  244. accepted = SatelliteAccepted(
  245. satellite_id=satellite.id,
  246. session_key=session_key,
  247. audio_input_port=network_config.audio_input_port,
  248. audio_output_port=network_config.audio_output_port,
  249. music_port=network_config.music_port,
  250. # TTS-Format vom aktiven Plugin
  251. tts_sample_rate=tts_sample_rate,
  252. tts_channels=tts_channels,
  253. tts_sample_width=2, # 16-bit
  254. # Musik-Format aus Config (konfigurierbar für Audio-Device-Kompatibilität)
  255. music_sample_rate=music_sample_rate,
  256. music_channels=music_channels,
  257. music_sample_width=music_bit_depth // 8, # Bytes
  258. )
  259. await self._send_message(accepted)
  260. pdebug(f"Musik-Format an Client: {music_sample_rate}Hz, {music_channels}ch, {music_bit_depth}-bit")
  261. pinfo(f"Satellite verbunden: {satellite.alias} ({satellite.room_id})")
  262. # Event auslösen
  263. await app.events.trigger("satellite_connected", SatelliteConnected(
  264. satellite_id=satellite.id,
  265. mac_address=satellite.mac_address,
  266. room=satellite.room_id,
  267. alias=satellite.alias,
  268. ip_address=satellite.ip_address,
  269. source="network_service"
  270. ))
  271. return True
  272. async def _message_loop(self) -> None:
  273. """Verarbeitet eingehende Nachrichten."""
  274. while self._running:
  275. message = await self._read_message()
  276. if message is None:
  277. break
  278. await self._handle_message(message)
  279. async def _handle_message(self, message: ProtocolMessage) -> None:
  280. """
  281. Verarbeitet eine empfangene Nachricht.
  282. Args:
  283. message: Die empfangene Nachricht
  284. """
  285. class_name = message.class_name
  286. # Hard-coded Befehle
  287. if class_name == "TRXIPONG":
  288. # Pong empfangen - KeepAlive aktualisieren
  289. if self._satellite:
  290. self._satellite.update_heartbeat()
  291. self._service.keepalive.handle_pong(self._satellite.id)
  292. return
  293. if class_name == "TRXIPING":
  294. # Ping empfangen - Pong senden
  295. await self._send_pong()
  296. return
  297. # SatelliteDisconnect
  298. if class_name == "SatelliteDisconnect":
  299. data = message.data
  300. reason = data.get("reason", "Client-Trennung") if isinstance(data, dict) else "Client-Trennung"
  301. pinfo(f"Satellite trennt Verbindung: {reason}")
  302. self._running = False
  303. return
  304. # Pending-Response pruefen (Config-Proxy)
  305. if self._satellite:
  306. msg_id = ""
  307. data = message.data
  308. if hasattr(data, "message_id"):
  309. msg_id = data.message_id
  310. elif isinstance(data, dict):
  311. msg_id = data.get("message_id", "")
  312. if msg_id and self._satellite.resolve_pending(msg_id, message):
  313. return
  314. # Andere Nachrichten an MessageHandler weiterleiten
  315. if self._satellite:
  316. await self._service.dispatch_message(self._satellite, message)
  317. async def _read_raw(self, size: int) -> bytes:
  318. """Liest rohe Bytes vom Socket."""
  319. data = await self._reader.read(size)
  320. if not data:
  321. raise ConnectionResetError("Verbindung geschlossen")
  322. return data
  323. async def _read_message(self) -> ProtocolMessage | None:
  324. """
  325. Liest eine vollständige Protokoll-Nachricht.
  326. Returns:
  327. Die gelesene Nachricht oder None bei Verbindungsende
  328. """
  329. try:
  330. # Lese Header (Magic + Version + Timestamp + Flags + Checksum + ClassNameLength)
  331. header_size = MAGIC_LENGTH + 2 + 8 + 4 + 16 + 2 # 36 bytes
  332. # Prüfe zuerst auf Hard-coded Befehle (8 bytes)
  333. # Verwende readexactly um sicherzustellen, dass wir alle 8 Bytes bekommen
  334. peek_data = await self._reader.readexactly(8)
  335. # Prüfe auf Hard-coded Befehl
  336. if self._protocol.is_hard_command(peek_data):
  337. return self._protocol.deserialize(peek_data)
  338. # Kein Hard-coded Befehl - lese Rest des Headers
  339. remaining_header = await self._reader.readexactly(header_size - 8)
  340. header = peek_data + remaining_header
  341. # Prüfe Magic
  342. if header[:MAGIC_LENGTH] != MAGIC:
  343. perror(f"Ungültige Magic: {header[:MAGIC_LENGTH]}")
  344. return None
  345. # Parse Header um Klassennamen-Länge und Daten-Länge zu bekommen
  346. # Offset nach Magic + Version + Timestamp + Flags + Checksum
  347. class_name_length = struct.unpack(">H", header[34:36])[0]
  348. # Lese Klassenname
  349. class_name = await self._reader.readexactly(class_name_length)
  350. # Lese Datenlänge (4 bytes)
  351. data_length_bytes = await self._reader.readexactly(4)
  352. data_length = struct.unpack(">I", data_length_bytes)[0]
  353. # Lese Daten
  354. data = await self._reader.readexactly(data_length) if data_length > 0 else b""
  355. # Baue vollständige Nachricht zusammen
  356. full_message = header + class_name + data_length_bytes + data
  357. return self._protocol.deserialize(full_message)
  358. except (ConnectionResetError, asyncio.IncompleteReadError):
  359. return None
  360. except Exception as e:
  361. perror(f"Fehler beim Lesen der Nachricht: {e}")
  362. return None
  363. async def _send_message(self, message: object, flags: ProtocolFlags = ProtocolFlags.PICKLE) -> bool:
  364. """
  365. Sendet eine Nachricht zum Client.
  366. Args:
  367. message: Die zu sendende Nachricht
  368. flags: Protokoll-Flags
  369. Returns:
  370. True bei Erfolg
  371. """
  372. try:
  373. data = self._protocol.serialize(message, flags)
  374. self._writer.write(data)
  375. await self._writer.drain()
  376. return True
  377. except Exception as e:
  378. perror(f"Fehler beim Senden: {e}")
  379. return False
  380. async def _send_ping(self) -> None:
  381. """Sendet einen Ping."""
  382. try:
  383. self._writer.write(COMMAND_PING)
  384. await self._writer.drain()
  385. except Exception as e:
  386. pdebug(f"Ping-Fehler: {e}")
  387. async def _send_pong(self) -> None:
  388. """Sendet einen Pong."""
  389. try:
  390. self._writer.write(COMMAND_PONG)
  391. await self._writer.drain()
  392. except Exception as e:
  393. pdebug(f"Pong-Fehler: {e}")
  394. async def _cleanup(self, reason: str = "Verbindung getrennt") -> None:
  395. """Räumt die Verbindung auf."""
  396. self._running = False
  397. if self._satellite:
  398. # KeepAlive entfernen
  399. await self._service.keepalive.unregister(self._satellite.id)
  400. # Event auslösen
  401. await self._service.application.events.trigger("satellite_disconnected", SatelliteDisconnected(
  402. satellite_id=self._satellite.id,
  403. room=self._satellite.room_id,
  404. alias=self._satellite.alias,
  405. reason=reason,
  406. source="network_service"
  407. ))
  408. # Satellite als getrennt markieren
  409. self._satellite.state = ConnectionState.DISCONNECTED
  410. self._satellite.sockets.command = None
  411. pinfo(f"Satellite getrennt: {self._satellite.alias}")
  412. self._satellite = None
  413. # Socket schließen
  414. try:
  415. self._writer.close()
  416. await self._writer.wait_closed()
  417. except Exception:
  418. pass
  419. class AudioStreamHandler:
  420. """
  421. Handler für Audio-Stream-Verbindungen.
  422. Verwaltet Audio-Input (Satellite → Server) und
  423. Audio-Output (Server → Satellite) Streams.
  424. """
  425. def __init__(
  426. self,
  427. reader: asyncio.StreamReader,
  428. writer: asyncio.StreamWriter,
  429. service: "NetworkService",
  430. stream_type: str,
  431. ) -> None:
  432. self._reader = reader
  433. self._writer = writer
  434. self._service = service
  435. self._stream_type = stream_type
  436. self._satellite: Satellite | None = None
  437. self._running = True
  438. self._peer_address = writer.get_extra_info("peername")
  439. async def run(self) -> None:
  440. """Hauptschleife für Audio-Stream."""
  441. try:
  442. # Warte auf Satellite-ID
  443. satellite_id_bytes = await asyncio.wait_for(
  444. self._reader.read(36), # UUID-Länge
  445. timeout=10.0
  446. )
  447. if not satellite_id_bytes:
  448. return
  449. satellite_id = satellite_id_bytes.decode("utf-8").strip()
  450. satellite = self._service.application.satellites.get(satellite_id)
  451. if not satellite or not satellite.is_connected:
  452. pwarn(f"Audio-Stream für unbekannten Satellite: {satellite_id}")
  453. return
  454. self._satellite = satellite
  455. pdebug(f"{self._stream_type} verbunden: {satellite.alias}")
  456. # Socket zuweisen
  457. if self._stream_type == "audio_in":
  458. satellite.sockets.audio_in = self._writer
  459. await self._handle_audio_input()
  460. elif self._stream_type == "audio_out":
  461. satellite.sockets.audio_out = self._writer
  462. await self._handle_audio_output()
  463. elif self._stream_type == "music":
  464. satellite.sockets.music_out = self._writer
  465. await self._handle_music_output()
  466. except asyncio.TimeoutError:
  467. pdebug(f"Audio-Stream Timeout: {self._peer_address}")
  468. except Exception as e:
  469. perror(f"Audio-Stream Fehler: {e}")
  470. finally:
  471. await self._cleanup()
  472. async def _handle_audio_input(self) -> None:
  473. """Empfängt Audio vom Satellite."""
  474. while self._running and self._satellite:
  475. try:
  476. # Lese Audio-Chunks (4096 bytes = 128ms bei 16KHz mono)
  477. chunk = await self._reader.read(4096)
  478. if not chunk:
  479. break
  480. # Event für Audio-Verarbeitung auslösen
  481. await self._service.application.events.emit(
  482. "audio_input_received",
  483. {
  484. "satellite_id": self._satellite.id,
  485. "audio_data": chunk,
  486. "sample_rate": 16000,
  487. "channels": 1,
  488. }
  489. )
  490. except Exception as e:
  491. pdebug(f"Audio-Input Fehler: {e}")
  492. break
  493. async def _handle_audio_output(self) -> None:
  494. """Verwaltet Audio-Output zum Satellite (wartet auf Events)."""
  495. # Audio-Output wird über Events gesteuert
  496. # Hier nur Verbindung aufrecht halten
  497. while self._running:
  498. await asyncio.sleep(1)
  499. async def _handle_music_output(self) -> None:
  500. """Verwaltet Music-Output zum Satellite (wartet auf Events)."""
  501. # Event auslösen, dass Music-Socket verbunden ist
  502. if self._satellite:
  503. from trixy_core.events.event_data.basic import MusicSocketConnected
  504. await self._service.application.events.trigger(
  505. "music_socket_connected",
  506. MusicSocketConnected(
  507. satellite_id=self._satellite.id,
  508. alias=self._satellite.alias,
  509. room=self._satellite.room_id,
  510. )
  511. )
  512. pdebug(f"Music-Socket verbunden: {self._satellite.alias}")
  513. # Music-Output wird über Events gesteuert
  514. while self._running:
  515. await asyncio.sleep(1)
  516. async def _cleanup(self) -> None:
  517. """Räumt den Stream auf."""
  518. self._running = False
  519. if self._satellite:
  520. if self._stream_type == "audio_in":
  521. self._satellite.sockets.audio_in = None
  522. elif self._stream_type == "audio_out":
  523. self._satellite.sockets.audio_out = None
  524. elif self._stream_type == "music":
  525. self._satellite.sockets.music_out = None
  526. try:
  527. self._writer.close()
  528. await self._writer.wait_closed()
  529. except Exception:
  530. pass
  531. class NetworkService(IService):
  532. """
  533. TCP-Server für Trixy-Kommunikation.
  534. Horcht auf 4 Ports und verwaltet alle Verbindungen:
  535. - Port 2101: Command Socket (Befehle, verschlüsselt)
  536. - Port 2102: Audio Input (Satellite → Server)
  537. - Port 2103: Audio Output (Server → Satellite)
  538. - Port 2104: Music Stream (Server → Satellite)
  539. """
  540. PRIORITY = ServicePriority.NETWORK
  541. GROUP = ServiceGroup.NETWORK
  542. NAME = "network"
  543. DEPENDENCIES: list[str] = []
  544. def __init__(self, application: "ServerApplication") -> None:
  545. super().__init__(application)
  546. self._servers: dict[str, asyncio.Server] = {}
  547. self._handlers: dict[str, set[ClientHandler | AudioStreamHandler]] = {
  548. "command": set(),
  549. "audio_in": set(),
  550. "audio_out": set(),
  551. "music": set(),
  552. }
  553. self._protocol = TrixyProtocol()
  554. self._encryption: TrixyEncryption | None = None
  555. self._keepalive = KeepAliveManager(
  556. default_config=KeepAliveConfig(
  557. ping_interval=30.0,
  558. pong_timeout=15.0, # Erhöht von 10s
  559. max_missed_pongs=5, # Erhöht von 3 (= 5*30s = 2.5min Toleranz)
  560. )
  561. )
  562. self._message_handlers: list[Callable[[Satellite, ProtocolMessage], Awaitable[None]]] = []
  563. @property
  564. def application(self) -> "ServerApplication":
  565. """Server-Anwendung."""
  566. return self._application # type: ignore
  567. @property
  568. def keepalive(self) -> KeepAliveManager:
  569. """KeepAlive-Manager."""
  570. return self._keepalive
  571. @property
  572. def protocol(self) -> TrixyProtocol:
  573. """Protokoll-Handler."""
  574. return self._protocol
  575. def register_message_handler(
  576. self,
  577. handler: Callable[[Satellite, ProtocolMessage], Awaitable[None]]
  578. ) -> None:
  579. """
  580. Registriert einen Handler für eingehende Nachrichten.
  581. Args:
  582. handler: Async-Funktion die (satellite, message) verarbeitet
  583. """
  584. self._message_handlers.append(handler)
  585. async def dispatch_message(self, satellite: Satellite, message: ProtocolMessage) -> None:
  586. """
  587. Verteilt eine Nachricht an alle registrierten Handler.
  588. Args:
  589. satellite: Der sendende Satellite
  590. message: Die empfangene Nachricht
  591. """
  592. for handler in self._message_handlers:
  593. try:
  594. await handler(satellite, message)
  595. except Exception as e:
  596. perror(f"Message-Handler Fehler: {e}")
  597. async def send_to_satellite(
  598. self,
  599. satellite_id: str,
  600. message: object,
  601. flags: ProtocolFlags = ProtocolFlags.PICKLE
  602. ) -> bool:
  603. """
  604. Sendet eine Nachricht an einen Satellite.
  605. Args:
  606. satellite_id: Ziel-Satellite-ID
  607. message: Die zu sendende Nachricht
  608. flags: Protokoll-Flags
  609. Returns:
  610. True bei Erfolg
  611. """
  612. satellite = self.application.satellites.get(satellite_id)
  613. if not satellite or not satellite.sockets.command:
  614. return False
  615. try:
  616. writer = satellite.sockets.command
  617. data = self._protocol.serialize(message, flags)
  618. writer.write(data)
  619. await writer.drain()
  620. return True
  621. except Exception as e:
  622. perror(f"Fehler beim Senden an {satellite_id}: {e}")
  623. return False
  624. async def broadcast(
  625. self,
  626. message: object,
  627. room: str | None = None,
  628. flags: ProtocolFlags = ProtocolFlags.PICKLE
  629. ) -> int:
  630. """
  631. Sendet eine Nachricht an alle (oder gefilterte) Satellites.
  632. Args:
  633. message: Die zu sendende Nachricht
  634. room: Optional: Nur an Satellites in diesem Raum
  635. flags: Protokoll-Flags
  636. Returns:
  637. Anzahl erfolgreicher Sendungen
  638. """
  639. satellites = self.application.satellites.get_by_room(room) if room else self.application.satellites.get_connected()
  640. count = 0
  641. for satellite in satellites:
  642. if await self.send_to_satellite(satellite.id, message, flags):
  643. count += 1
  644. return count
  645. async def start(self) -> None:
  646. """Startet alle TCP-Server."""
  647. config = self.application.server_config.network
  648. # Verschlüsselung laden
  649. if self.application.server_config.security.encryption_enabled:
  650. key_path = self.application.server_config.security.encryption_key_path
  651. try:
  652. self._encryption = TrixyEncryption.load_or_create(key_path)
  653. self._protocol.set_encryption(self._encryption)
  654. pdebug("Verschlüsselung aktiviert")
  655. except Exception as e:
  656. pwarn(f"Verschlüsselung konnte nicht geladen werden: {e}")
  657. # Command-Server starten (Port 2101)
  658. self._servers["command"] = await asyncio.start_server(
  659. self._handle_command_connection,
  660. config.bind_address,
  661. config.command_port,
  662. )
  663. pdebug(f"Command-Server gestartet auf {config.bind_address}:{config.command_port}")
  664. # Audio-Input-Server starten (Port 2102)
  665. self._servers["audio_in"] = await asyncio.start_server(
  666. lambda r, w: self._handle_audio_connection(r, w, "audio_in"),
  667. config.bind_address,
  668. config.audio_input_port,
  669. )
  670. pdebug(f"Audio-Input-Server gestartet auf {config.bind_address}:{config.audio_input_port}")
  671. # Audio-Output-Server starten (Port 2103)
  672. self._servers["audio_out"] = await asyncio.start_server(
  673. lambda r, w: self._handle_audio_connection(r, w, "audio_out"),
  674. config.bind_address,
  675. config.audio_output_port,
  676. )
  677. pdebug(f"Audio-Output-Server gestartet auf {config.bind_address}:{config.audio_output_port}")
  678. # Music-Server starten (Port 2104)
  679. self._servers["music"] = await asyncio.start_server(
  680. lambda r, w: self._handle_audio_connection(r, w, "music"),
  681. config.bind_address,
  682. config.music_port,
  683. )
  684. pdebug(f"Music-Server gestartet auf {config.bind_address}:{config.music_port}")
  685. # KeepAlive-Disconnect-Handler registrieren
  686. self._keepalive.on_disconnect(self._on_keepalive_disconnect)
  687. pinfo(f"NetworkService gestartet auf Ports {config.command_port}-{config.music_port}")
  688. async def stop(self) -> None:
  689. """Stoppt alle Server und Verbindungen."""
  690. pinfo("NetworkService wird gestoppt...")
  691. # KeepAlive stoppen
  692. await self._keepalive.stop_all()
  693. # Alle Handler beenden
  694. for handler_type, handlers in self._handlers.items():
  695. for handler in list(handlers):
  696. if hasattr(handler, "_running"):
  697. handler._running = False
  698. # Server schließen
  699. for name, server in self._servers.items():
  700. server.close()
  701. await server.wait_closed()
  702. pdebug(f"{name}-Server geschlossen")
  703. self._servers.clear()
  704. pinfo("NetworkService gestoppt")
  705. async def _handle_command_connection(
  706. self,
  707. reader: asyncio.StreamReader,
  708. writer: asyncio.StreamWriter
  709. ) -> None:
  710. """Verarbeitet eine neue Command-Verbindung."""
  711. addr = writer.get_extra_info("peername")
  712. pdebug(f"Neue Command-Verbindung von {addr}")
  713. handler = ClientHandler(reader, writer, self._protocol, self)
  714. self._handlers["command"].add(handler)
  715. try:
  716. await handler.run()
  717. finally:
  718. self._handlers["command"].discard(handler)
  719. async def _handle_audio_connection(
  720. self,
  721. reader: asyncio.StreamReader,
  722. writer: asyncio.StreamWriter,
  723. stream_type: str
  724. ) -> None:
  725. """Verarbeitet eine neue Audio-Verbindung."""
  726. addr = writer.get_extra_info("peername")
  727. pdebug(f"Neue {stream_type}-Verbindung von {addr}")
  728. handler = AudioStreamHandler(reader, writer, self, stream_type)
  729. self._handlers[stream_type].add(handler)
  730. try:
  731. await handler.run()
  732. finally:
  733. self._handlers[stream_type].discard(handler)
  734. async def _on_keepalive_disconnect(self, connection_id: str) -> None:
  735. """Handler für KeepAlive-Disconnect."""
  736. satellite = self.application.satellites.get(connection_id)
  737. if satellite:
  738. pwarn(f"KeepAlive-Timeout für {satellite.alias}")
  739. await satellite.disconnect("KeepAlive-Timeout")
  740. async def health_check(self) -> bool:
  741. """Prüft den Gesundheitszustand des Services."""
  742. # Prüfe ob Server existieren und laufen
  743. if not self._servers:
  744. return False
  745. for name, server in self._servers.items():
  746. if not server.is_serving():
  747. return False
  748. return True