client_connection.py 30 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878
  1. # -*- coding: utf-8 -*-
  2. """
  3. ClientConnection - Client-seitige Verbindung zum Server.
  4. Verwaltet alle 4 Sockets und den Handshake für Satellites.
  5. """
  6. from __future__ import annotations
  7. import asyncio
  8. import struct
  9. import uuid
  10. from typing import TYPE_CHECKING, Callable, Awaitable, Any
  11. from dataclasses import dataclass
  12. from trixy_core.network.protocol import (
  13. TrixyProtocol,
  14. ProtocolFlags,
  15. ProtocolMessage,
  16. MAGIC,
  17. MAGIC_LENGTH,
  18. COMMAND_HELLO,
  19. COMMAND_PING,
  20. COMMAND_PONG,
  21. )
  22. from trixy_core.network.encryption import TrixyEncryption
  23. from trixy_core.network.cmd import (
  24. SatelliteConnect,
  25. SatelliteDisconnect,
  26. SatelliteAccepted,
  27. SatelliteDenied,
  28. )
  29. from trixy_core.utils.debug import pinfo, pdebug, perror, pwarn
  30. from trixy_core.utils.version import VERSION_STRING
  31. if TYPE_CHECKING:
  32. from trixy_core.client import ClientApplication
  33. @dataclass
  34. class ConnectionConfig:
  35. """Konfiguration für die Verbindung."""
  36. connect_timeout: float = 10.0
  37. reconnect_interval: float = 5.0
  38. max_reconnect_attempts: int = 0 # 0 = unbegrenzt
  39. ping_interval: float = 30.0
  40. pong_timeout: float = 10.0
  41. def parse_host_list(value: str | list[str] | None) -> list[str]:
  42. """Parst einen Host-String oder eine Liste in eine Liste von Hostnames/IPs.
  43. Akzeptiert:
  44. - String: "trixyone"
  45. - Kommagetrennt: "trixyone, 10.10.10.1"
  46. - Semikolon-/Whitespace-getrennt: "trixyone;10.10.10.1 trixy.lan"
  47. - Liste: ["trixyone", "10.10.10.1"]
  48. Whitespace und leere Eintraege werden entfernt.
  49. Reihenfolge bleibt erhalten.
  50. """
  51. if not value:
  52. return []
  53. if isinstance(value, list):
  54. return [str(h).strip() for h in value if str(h).strip()]
  55. import re
  56. parts = re.split(r"[,;\s]+", str(value).strip())
  57. return [p.strip() for p in parts if p.strip()]
  58. class ClientConnection:
  59. """
  60. Client-seitige Verbindung zum Trixy-Server.
  61. Verwaltet alle 4 Sockets:
  62. - Command: Befehle und Steuerung
  63. - Audio In: Mikrofon-Daten zum Server
  64. - Audio Out: Sprachausgabe vom Server
  65. - Music: Musik-Stream vom Server
  66. Usage:
  67. connection = ClientConnection(
  68. host="192.168.1.100",
  69. port=2101,
  70. room="wohnzimmer",
  71. alias="Echo",
  72. mac_address="AA:BB:CC:DD:EE:FF"
  73. )
  74. if await connection.connect():
  75. # Verbunden
  76. await connection.send_command(WakewordDetected(...))
  77. """
  78. def __init__(
  79. self,
  80. host: str | list[str],
  81. port: int,
  82. room: str,
  83. alias: str,
  84. mac_address: str,
  85. config: ConnectionConfig | None = None,
  86. ) -> None:
  87. """
  88. Initialisiert die ClientConnection.
  89. Args:
  90. host: Server-Host. Kann sein:
  91. - Ein einzelner String: "trixyone"
  92. - Kommagetrennt: "trixyone, 10.10.10.1, trixy.lan"
  93. - Liste: ["trixyone", "10.10.10.1"]
  94. Beim Verbinden werden alle Hosts der Reihe nach probiert,
  95. der erste erreichbare gewinnt.
  96. port: Server-Port (Command-Port, Standard 2101)
  97. room: Raum-Kennung
  98. alias: Anzeigename des Satellites
  99. mac_address: MAC-Adresse des Satellites
  100. config: Verbindungskonfiguration
  101. """
  102. # Host-Liste normalisieren
  103. self._hosts = parse_host_list(host)
  104. # _host bleibt fuer Backward-Compatibility — wird beim Connect aktualisiert
  105. self._host = self._hosts[0] if self._hosts else ""
  106. self._port = port
  107. self._room = room
  108. self._alias = alias
  109. self._mac_address = mac_address
  110. self._config = config or ConnectionConfig()
  111. # Protokoll
  112. self._protocol = TrixyProtocol()
  113. self._encryption: TrixyEncryption | None = None
  114. self._session_key: bytes = b""
  115. # Sockets
  116. self._command_reader: asyncio.StreamReader | None = None
  117. self._command_writer: asyncio.StreamWriter | None = None
  118. self._audio_in_reader: asyncio.StreamReader | None = None
  119. self._audio_in_writer: asyncio.StreamWriter | None = None
  120. self._audio_out_reader: asyncio.StreamReader | None = None
  121. self._audio_out_writer: asyncio.StreamWriter | None = None
  122. self._music_reader: asyncio.StreamReader | None = None
  123. self._music_writer: asyncio.StreamWriter | None = None
  124. # Zustand
  125. self._connected = False
  126. self._satellite_id: str = ""
  127. self._audio_ports: dict[str, int] = {}
  128. self._running = False
  129. self._reconnect_attempts = 0
  130. # Letzte Denied-Antwort — fuer Reconnect-Steuerung
  131. self._last_denied_reason: str = ""
  132. self._last_denied_retry_after: int = 0
  133. # Audio-Format (vom Server mitgeteilt)
  134. self._tts_format: dict[str, int] = {
  135. "sample_rate": 22050,
  136. "channels": 1,
  137. "sample_width": 2,
  138. }
  139. self._music_format: dict[str, int] = {
  140. "sample_rate": 44100,
  141. "channels": 2,
  142. "sample_width": 2,
  143. }
  144. # Tasks
  145. self._receive_task: asyncio.Task | None = None
  146. self._keepalive_task: asyncio.Task | None = None
  147. self._audio_out_task: asyncio.Task | None = None
  148. self._music_task: asyncio.Task | None = None
  149. # Callbacks
  150. self._on_connected: Callable[[], Awaitable[None]] | None = None
  151. self._on_disconnected: Callable[[str], Awaitable[None]] | None = None
  152. self._on_message: Callable[[ProtocolMessage], Awaitable[None]] | None = None
  153. self._on_audio_out: Callable[[bytes], Awaitable[None]] | None = None
  154. self._on_music: Callable[[bytes], Awaitable[None]] | None = None
  155. @property
  156. def is_connected(self) -> bool:
  157. """Ist die Verbindung aktiv?"""
  158. return self._connected
  159. @property
  160. def satellite_id(self) -> str:
  161. """Die vom Server zugewiesene Satellite-ID."""
  162. return self._satellite_id
  163. @property
  164. def last_denied_retry_after(self) -> int:
  165. """Wartezeit in Sekunden nach letztem SatelliteDenied (0 wenn keins)."""
  166. return self._last_denied_retry_after
  167. @property
  168. def last_denied_reason(self) -> str:
  169. """Grund der letzten Ablehnung (leer wenn keine)."""
  170. return self._last_denied_reason
  171. def clear_denied_state(self) -> None:
  172. """Setzt den Denied-State zurueck (nach erfolgreichem Connect)."""
  173. self._last_denied_reason = ""
  174. self._last_denied_retry_after = 0
  175. @property
  176. def host(self) -> str:
  177. """Server-Host."""
  178. return self._host
  179. @property
  180. def port(self) -> int:
  181. """Server-Port."""
  182. return self._port
  183. @property
  184. def tts_format(self) -> dict[str, int]:
  185. """TTS-Audio-Format vom Server."""
  186. return self._tts_format.copy()
  187. @property
  188. def music_format(self) -> dict[str, int]:
  189. """Musik-Audio-Format vom Server."""
  190. return self._music_format.copy()
  191. def on_connected(self, callback: Callable[[], Awaitable[None]]) -> None:
  192. """Registriert Callback für Verbindungsaufbau."""
  193. self._on_connected = callback
  194. def on_disconnected(self, callback: Callable[[str], Awaitable[None]]) -> None:
  195. """Registriert Callback für Verbindungsabbruch."""
  196. self._on_disconnected = callback
  197. def on_message(self, callback: Callable[[ProtocolMessage], Awaitable[None]]) -> None:
  198. """Registriert Callback für eingehende Nachrichten."""
  199. self._on_message = callback
  200. def on_audio_out(self, callback: Callable[[bytes], Awaitable[None]]) -> None:
  201. """Registriert Callback für Audio-Ausgabe."""
  202. self._on_audio_out = callback
  203. def on_music(self, callback: Callable[[bytes], Awaitable[None]]) -> None:
  204. """Registriert Callback für Musik-Stream."""
  205. self._on_music = callback
  206. async def connect(self) -> bool:
  207. """
  208. Verbindet sich mit dem Server.
  209. Probiert alle konfigurierten Hosts der Reihe nach durch und nimmt
  210. den ersten der antwortet. Fuehrt dann den Handshake durch und
  211. verbindet alle Audio-Ports.
  212. Returns:
  213. True bei erfolgreicher Verbindung
  214. """
  215. if not self._hosts:
  216. perror("Keine Server-Hosts konfiguriert")
  217. return False
  218. # Alle Hosts der Reihe nach probieren
  219. last_error = ""
  220. connected_host: str | None = None
  221. for candidate in self._hosts:
  222. pdebug(f"Verbinde zu {candidate}:{self._port}...")
  223. try:
  224. self._command_reader, self._command_writer = await asyncio.wait_for(
  225. asyncio.open_connection(candidate, self._port),
  226. timeout=self._config.connect_timeout,
  227. )
  228. connected_host = candidate
  229. self._host = candidate # Aktualisiere fuer Logging und Audio-Sockets
  230. pdebug(f"Command-Socket verbunden zu {candidate}")
  231. break
  232. except asyncio.TimeoutError:
  233. last_error = f"Timeout zu {candidate}:{self._port}"
  234. pdebug(last_error)
  235. continue
  236. except OSError as e:
  237. last_error = f"Fehler {candidate}: {e}"
  238. pdebug(last_error)
  239. continue
  240. if connected_host is None:
  241. if len(self._hosts) > 1:
  242. perror(f"Kein Server erreichbar — probierte: {', '.join(self._hosts)}")
  243. else:
  244. perror(f"Verbindungsfehler: {last_error}")
  245. return False
  246. try:
  247. # Handshake durchführen
  248. if not await self._perform_handshake():
  249. await self._close_sockets()
  250. return False
  251. # Audio-Sockets verbinden
  252. await self._connect_audio_sockets()
  253. # Empfangs-Tasks starten
  254. self._running = True
  255. self._receive_task = asyncio.create_task(self._receive_loop())
  256. self._keepalive_task = asyncio.create_task(self._keepalive_loop())
  257. if self._audio_out_reader:
  258. self._audio_out_task = asyncio.create_task(self._audio_out_loop())
  259. if self._music_reader:
  260. self._music_task = asyncio.create_task(self._music_loop())
  261. self._connected = True
  262. self._reconnect_attempts = 0
  263. self.clear_denied_state()
  264. pinfo(f"Verbunden mit Server {self._host}:{self._port}")
  265. # Callback
  266. if self._on_connected:
  267. try:
  268. await self._on_connected()
  269. except Exception as e:
  270. perror(f"on_connected Callback Fehler: {e}")
  271. return True
  272. except Exception as e:
  273. perror(f"Verbindungsfehler: {e}")
  274. await self._close_sockets()
  275. return False
  276. async def disconnect(self, reason: str = "Client-Disconnect") -> None:
  277. """
  278. Trennt die Verbindung zum Server.
  279. Args:
  280. reason: Grund für die Trennung
  281. """
  282. if not self._connected:
  283. return
  284. pinfo(f"Trenne Verbindung: {reason}")
  285. self._running = False
  286. # Disconnect-Nachricht senden
  287. try:
  288. if self._command_writer:
  289. disconnect_msg = SatelliteDisconnect(
  290. mac_address=self._mac_address,
  291. reason=reason
  292. )
  293. await self._send_message(disconnect_msg)
  294. except Exception:
  295. pass
  296. # Tasks beenden
  297. for task in [self._receive_task, self._keepalive_task, self._audio_out_task, self._music_task]:
  298. if task and not task.done():
  299. task.cancel()
  300. try:
  301. await task
  302. except asyncio.CancelledError:
  303. pass
  304. # Sockets schließen
  305. await self._close_sockets()
  306. self._connected = False
  307. self._satellite_id = ""
  308. # Callback
  309. if self._on_disconnected:
  310. try:
  311. await self._on_disconnected(reason)
  312. except Exception as e:
  313. perror(f"on_disconnected Callback Fehler: {e}")
  314. async def send_command(self, message: object, flags: ProtocolFlags = ProtocolFlags.PICKLE) -> bool:
  315. """
  316. Sendet eine Nachricht zum Server.
  317. Args:
  318. message: Die zu sendende Nachricht
  319. flags: Protokoll-Flags
  320. Returns:
  321. True bei Erfolg
  322. """
  323. if not self._connected or not self._command_writer:
  324. return False
  325. return await self._send_message(message, flags)
  326. async def send_audio(self, audio_data: bytes) -> bool:
  327. """
  328. Sendet Audio-Daten zum Server.
  329. Args:
  330. audio_data: Audio-Bytes (16KHz, 16-bit, mono)
  331. Returns:
  332. True bei Erfolg
  333. """
  334. if not self._audio_in_writer:
  335. return False
  336. try:
  337. self._audio_in_writer.write(audio_data)
  338. await self._audio_in_writer.drain()
  339. return True
  340. except Exception as e:
  341. perror(f"Audio-Sende-Fehler: {e}")
  342. return False
  343. async def _perform_handshake(self) -> bool:
  344. """
  345. Führt den Handshake mit dem Server durch.
  346. Returns:
  347. True bei erfolgreichem Handshake
  348. """
  349. pdebug("Starte Handshake...")
  350. # HELLO senden
  351. try:
  352. self._command_writer.write(COMMAND_HELLO)
  353. await self._command_writer.drain()
  354. except Exception as e:
  355. perror(f"Fehler beim Senden von HELLO: {e}")
  356. return False
  357. pdebug("HELLO gesendet")
  358. # SatelliteConnect senden
  359. import platform
  360. import socket as _socket
  361. connect_msg = SatelliteConnect(
  362. mac_address=self._mac_address,
  363. room=self._room,
  364. alias=self._alias,
  365. ip_address="", # Server ermittelt IP
  366. version=VERSION_STRING,
  367. hostname=_socket.gethostname(),
  368. os_name=platform.system(),
  369. os_version=platform.release()
  370. )
  371. if not await self._send_message(connect_msg):
  372. perror("Fehler beim Senden von SatelliteConnect")
  373. return False
  374. pdebug("SatelliteConnect gesendet")
  375. # Antwort empfangen
  376. try:
  377. response = await asyncio.wait_for(
  378. self._read_message(),
  379. timeout=self._config.connect_timeout
  380. )
  381. except asyncio.TimeoutError:
  382. perror("Timeout beim Warten auf Server-Antwort")
  383. return False
  384. if response is None:
  385. perror("Keine Antwort vom Server")
  386. return False
  387. # Antwort verarbeiten
  388. if response.class_name == "SatelliteDenied":
  389. data = response.data
  390. if isinstance(data, dict):
  391. reason = data.get("reason", "Unbekannt")
  392. retry_after = data.get("retry_after_seconds", 60)
  393. elif isinstance(data, SatelliteDenied):
  394. reason = data.reason
  395. retry_after = data.retry_after_seconds
  396. else:
  397. reason = "Unbekannt"
  398. retry_after = 60
  399. # Fuer den Client-Reconnect-Loop speichern
  400. self._last_denied_reason = reason
  401. self._last_denied_retry_after = int(retry_after)
  402. pwarn(f"Verbindung abgelehnt: {reason}")
  403. pwarn(f"Erneuter Versuch in {retry_after} Sekunden moeglich")
  404. return False
  405. if response.class_name != "SatelliteAccepted":
  406. perror(f"Unerwartete Antwort: {response.class_name}")
  407. return False
  408. # SatelliteAccepted verarbeiten
  409. data = response.data
  410. if isinstance(data, dict):
  411. self._satellite_id = data.get("satellite_id", "")
  412. self._session_key = data.get("session_key", b"")
  413. self._audio_ports = {
  414. "audio_in": data.get("audio_input_port", 0),
  415. "audio_out": data.get("audio_output_port", 0),
  416. "music": data.get("music_port", 0),
  417. }
  418. # Audio-Format vom Server
  419. if data.get("tts_sample_rate"):
  420. self._tts_format = {
  421. "sample_rate": data.get("tts_sample_rate", 22050),
  422. "channels": data.get("tts_channels", 1),
  423. "sample_width": data.get("tts_sample_width", 2),
  424. }
  425. if data.get("music_sample_rate"):
  426. self._music_format = {
  427. "sample_rate": data.get("music_sample_rate", 44100),
  428. "channels": data.get("music_channels", 2),
  429. "sample_width": data.get("music_sample_width", 2),
  430. }
  431. elif isinstance(data, SatelliteAccepted):
  432. self._satellite_id = data.satellite_id
  433. self._session_key = data.session_key
  434. self._audio_ports = {
  435. "audio_in": data.audio_input_port,
  436. "audio_out": data.audio_output_port,
  437. "music": data.music_port,
  438. }
  439. # Audio-Format vom Server
  440. if data.tts_sample_rate:
  441. self._tts_format = {
  442. "sample_rate": data.tts_sample_rate,
  443. "channels": data.tts_channels,
  444. "sample_width": data.tts_sample_width,
  445. }
  446. if data.music_sample_rate:
  447. self._music_format = {
  448. "sample_rate": data.music_sample_rate,
  449. "channels": data.music_channels,
  450. "sample_width": data.music_sample_width,
  451. }
  452. else:
  453. perror(f"Ungültige SatelliteAccepted-Daten: {type(data)}")
  454. return False
  455. # Session-Verschlüsselung einrichten
  456. if self._session_key:
  457. try:
  458. self._encryption = TrixyEncryption(self._session_key)
  459. self._protocol.set_encryption(self._encryption)
  460. pdebug("Session-Verschlüsselung aktiviert")
  461. except Exception as e:
  462. pwarn(f"Session-Verschlüsselung fehlgeschlagen: {e}")
  463. pinfo(f"Handshake erfolgreich, Satellite-ID: {self._satellite_id}")
  464. return True
  465. async def _connect_audio_sockets(self) -> None:
  466. """Verbindet die Audio-Sockets."""
  467. # Audio-Input Socket
  468. if self._audio_ports.get("audio_in"):
  469. try:
  470. self._audio_in_reader, self._audio_in_writer = await asyncio.wait_for(
  471. asyncio.open_connection(self._host, self._audio_ports["audio_in"]),
  472. timeout=self._config.connect_timeout
  473. )
  474. # Satellite-ID senden
  475. self._audio_in_writer.write(self._satellite_id.ljust(36).encode("utf-8"))
  476. await self._audio_in_writer.drain()
  477. pdebug(f"Audio-Input verbunden auf Port {self._audio_ports['audio_in']}")
  478. except Exception as e:
  479. pwarn(f"Audio-Input Verbindung fehlgeschlagen: {e}")
  480. # Audio-Output Socket
  481. if self._audio_ports.get("audio_out"):
  482. try:
  483. self._audio_out_reader, self._audio_out_writer = await asyncio.wait_for(
  484. asyncio.open_connection(self._host, self._audio_ports["audio_out"]),
  485. timeout=self._config.connect_timeout
  486. )
  487. # Satellite-ID senden
  488. self._audio_out_writer.write(self._satellite_id.ljust(36).encode("utf-8"))
  489. await self._audio_out_writer.drain()
  490. pdebug(f"Audio-Output verbunden auf Port {self._audio_ports['audio_out']}")
  491. except Exception as e:
  492. pwarn(f"Audio-Output Verbindung fehlgeschlagen: {e}")
  493. # Music Socket
  494. if self._audio_ports.get("music"):
  495. try:
  496. self._music_reader, self._music_writer = await asyncio.wait_for(
  497. asyncio.open_connection(self._host, self._audio_ports["music"]),
  498. timeout=self._config.connect_timeout
  499. )
  500. # Satellite-ID senden
  501. self._music_writer.write(self._satellite_id.ljust(36).encode("utf-8"))
  502. await self._music_writer.drain()
  503. pdebug(f"Music verbunden auf Port {self._audio_ports['music']}")
  504. except Exception as e:
  505. pwarn(f"Music Verbindung fehlgeschlagen: {e}")
  506. async def _receive_loop(self) -> None:
  507. """Empfängt Nachrichten vom Server."""
  508. while self._running:
  509. try:
  510. message = await self._read_message()
  511. if message is None:
  512. if self._running:
  513. pwarn("Verbindung zum Server verloren")
  514. asyncio.create_task(self._handle_disconnect("Verbindung verloren"))
  515. break
  516. await self._handle_message(message)
  517. except asyncio.CancelledError:
  518. break
  519. except Exception as e:
  520. if self._running:
  521. perror(f"Empfangsfehler: {e}")
  522. asyncio.create_task(self._handle_disconnect(f"Fehler: {e}"))
  523. break
  524. async def _handle_message(self, message: ProtocolMessage) -> None:
  525. """Verarbeitet eine empfangene Nachricht."""
  526. class_name = message.class_name
  527. # PING beantworten
  528. if class_name == "TRXIPING":
  529. await self._send_pong()
  530. return
  531. # PONG ignorieren (KeepAlive)
  532. if class_name == "TRXIPONG":
  533. return
  534. # An Callback weiterleiten
  535. if self._on_message:
  536. try:
  537. await self._on_message(message)
  538. except Exception as e:
  539. perror(f"Message-Callback Fehler: {e}")
  540. async def _keepalive_loop(self) -> None:
  541. """Sendet regelmäßige Pings."""
  542. while self._running:
  543. try:
  544. await asyncio.sleep(self._config.ping_interval)
  545. if not self._running:
  546. break
  547. # Ping senden
  548. await self._send_ping()
  549. except asyncio.CancelledError:
  550. break
  551. except Exception as e:
  552. pdebug(f"KeepAlive-Fehler: {e}")
  553. async def _audio_out_loop(self) -> None:
  554. """Empfängt Audio-Ausgabe vom Server."""
  555. first_chunk = True
  556. total_bytes = 0
  557. chunk_count = 0
  558. # Frame-Alignment Buffer (TTS: Mono 16-bit = 2 bytes pro Frame)
  559. frame_size = self._tts_format["channels"] * self._tts_format["sample_width"]
  560. leftover = b""
  561. while self._running and self._audio_out_reader:
  562. try:
  563. chunk = await self._audio_out_reader.read(32768) # 32KB chunks
  564. if not chunk:
  565. break
  566. # Kombiniere mit Resten vom letzten Chunk
  567. data = leftover + chunk
  568. # Nur komplette Frames weitergeben
  569. aligned_size = (len(data) // frame_size) * frame_size
  570. if aligned_size > 0:
  571. aligned_chunk = data[:aligned_size]
  572. leftover = data[aligned_size:]
  573. chunk_count += 1
  574. total_bytes += len(aligned_chunk)
  575. if first_chunk:
  576. pdebug(f"[CLIENT] Audio-Out: Erster Chunk empfangen ({len(aligned_chunk)} bytes)")
  577. first_chunk = False
  578. if self._on_audio_out:
  579. await self._on_audio_out(aligned_chunk)
  580. else:
  581. leftover = data
  582. except asyncio.CancelledError:
  583. break
  584. except Exception as e:
  585. pdebug(f"Audio-Out Fehler: {e}")
  586. break
  587. if chunk_count > 0:
  588. pdebug(f"[CLIENT] Audio-Out Stream beendet: {chunk_count} chunks, {total_bytes} bytes total")
  589. async def _music_loop(self) -> None:
  590. """Empfängt Music-Stream vom Server."""
  591. first_chunk = True
  592. total_bytes = 0
  593. chunk_count = 0
  594. # Frame-Alignment Buffer (Music: Stereo 16-bit = 4 bytes pro Frame)
  595. frame_size = self._music_format["channels"] * self._music_format["sample_width"]
  596. leftover = b""
  597. while self._running and self._music_reader:
  598. try:
  599. chunk = await self._music_reader.read(32768) # 32KB chunks
  600. if not chunk:
  601. break
  602. # Kombiniere mit Resten vom letzten Chunk
  603. data = leftover + chunk
  604. # Nur komplette Frames weitergeben
  605. aligned_size = (len(data) // frame_size) * frame_size
  606. if aligned_size > 0:
  607. aligned_chunk = data[:aligned_size]
  608. leftover = data[aligned_size:]
  609. chunk_count += 1
  610. total_bytes += len(aligned_chunk)
  611. if first_chunk:
  612. pdebug(f"[CLIENT] Music: Erster Chunk empfangen ({len(aligned_chunk)} bytes)")
  613. first_chunk = False
  614. if self._on_music:
  615. await self._on_music(aligned_chunk)
  616. else:
  617. leftover = data
  618. except asyncio.CancelledError:
  619. break
  620. except Exception as e:
  621. pdebug(f"Music Fehler: {e}")
  622. break
  623. if chunk_count > 0:
  624. pdebug(f"[CLIENT] Music Stream beendet: {chunk_count} chunks, {total_bytes} bytes total")
  625. async def _handle_disconnect(self, reason: str) -> None:
  626. """Behandelt einen Verbindungsabbruch."""
  627. was_connected = self._connected
  628. self._connected = False
  629. self._running = False
  630. await self._close_sockets()
  631. if was_connected and self._on_disconnected:
  632. try:
  633. await self._on_disconnected(reason)
  634. except Exception as e:
  635. perror(f"on_disconnected Callback Fehler: {e}")
  636. async def _send_message(self, message: object, flags: ProtocolFlags = ProtocolFlags.PICKLE) -> bool:
  637. """Sendet eine Nachricht über den Command-Socket."""
  638. if not self._command_writer:
  639. return False
  640. try:
  641. data = self._protocol.serialize(message, flags)
  642. self._command_writer.write(data)
  643. await self._command_writer.drain()
  644. return True
  645. except Exception as e:
  646. perror(f"Sende-Fehler: {e}")
  647. return False
  648. async def _send_ping(self) -> None:
  649. """Sendet einen Ping."""
  650. if self._command_writer:
  651. try:
  652. self._command_writer.write(COMMAND_PING)
  653. await self._command_writer.drain()
  654. except Exception:
  655. pass
  656. async def _send_pong(self) -> None:
  657. """Sendet einen Pong."""
  658. if self._command_writer:
  659. try:
  660. self._command_writer.write(COMMAND_PONG)
  661. await self._command_writer.drain()
  662. except Exception:
  663. pass
  664. async def _read_message(self) -> ProtocolMessage | None:
  665. """Liest eine Nachricht vom Command-Socket."""
  666. if not self._command_reader:
  667. return None
  668. try:
  669. # Lese Header (Magic + Version + Timestamp + Flags + Checksum + ClassNameLength)
  670. header_size = MAGIC_LENGTH + 2 + 8 + 4 + 16 + 2 # 36 bytes
  671. # Prüfe zuerst auf Hard-coded Befehle (8 bytes)
  672. # Verwende readexactly um sicherzustellen, dass wir alle 8 Bytes bekommen
  673. peek_data = await self._command_reader.readexactly(8)
  674. # Prüfe auf Hard-coded Befehl
  675. if self._protocol.is_hard_command(peek_data):
  676. return self._protocol.deserialize(peek_data)
  677. # Kein Hard-coded Befehl - lese Rest des Headers
  678. remaining_header = await self._command_reader.readexactly(header_size - 8)
  679. header = peek_data + remaining_header
  680. # Prüfe Magic
  681. if header[:MAGIC_LENGTH] != MAGIC:
  682. perror(f"Ungültige Magic: {header[:MAGIC_LENGTH]}")
  683. return None
  684. # Parse Klassennamen-Länge
  685. class_name_length = struct.unpack(">H", header[34:36])[0]
  686. # Lese Klassenname
  687. class_name = await self._command_reader.readexactly(class_name_length)
  688. # Lese Datenlänge (4 bytes)
  689. data_length_bytes = await self._command_reader.readexactly(4)
  690. data_length = struct.unpack(">I", data_length_bytes)[0]
  691. # Lese Daten
  692. data = await self._command_reader.readexactly(data_length) if data_length > 0 else b""
  693. # Baue vollständige Nachricht zusammen
  694. full_message = header + class_name + data_length_bytes + data
  695. return self._protocol.deserialize(full_message)
  696. except (ConnectionResetError, asyncio.IncompleteReadError):
  697. return None
  698. except Exception as e:
  699. perror(f"Lese-Fehler: {e}")
  700. return None
  701. async def _close_sockets(self) -> None:
  702. """Schließt alle Sockets."""
  703. for writer in [
  704. self._command_writer,
  705. self._audio_in_writer,
  706. self._audio_out_writer,
  707. self._music_writer
  708. ]:
  709. if writer:
  710. try:
  711. writer.close()
  712. await writer.wait_closed()
  713. except Exception:
  714. pass
  715. self._command_reader = None
  716. self._command_writer = None
  717. self._audio_in_reader = None
  718. self._audio_in_writer = None
  719. self._audio_out_reader = None
  720. self._audio_out_writer = None
  721. self._music_reader = None
  722. self._music_writer = None
  723. def get_mac_address() -> str:
  724. """
  725. Ermittelt die MAC-Adresse des Systems.
  726. Returns:
  727. MAC-Adresse im Format "AA:BB:CC:DD:EE:FF"
  728. """
  729. import uuid as uuid_module
  730. mac = uuid_module.getnode()
  731. return ":".join(f"{(mac >> (8 * i)) & 0xFF:02X}" for i in range(5, -1, -1))