| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878 |
- # -*- coding: utf-8 -*-
- """
- ClientConnection - Client-seitige Verbindung zum Server.
- Verwaltet alle 4 Sockets und den Handshake für Satellites.
- """
- from __future__ import annotations
- import asyncio
- import struct
- import uuid
- from typing import TYPE_CHECKING, Callable, Awaitable, Any
- from dataclasses import dataclass
- 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.cmd import (
- SatelliteConnect,
- SatelliteDisconnect,
- SatelliteAccepted,
- SatelliteDenied,
- )
- from trixy_core.utils.debug import pinfo, pdebug, perror, pwarn
- from trixy_core.utils.version import VERSION_STRING
- if TYPE_CHECKING:
- from trixy_core.client import ClientApplication
- @dataclass
- class ConnectionConfig:
- """Konfiguration für die Verbindung."""
- connect_timeout: float = 10.0
- reconnect_interval: float = 5.0
- max_reconnect_attempts: int = 0 # 0 = unbegrenzt
- ping_interval: float = 30.0
- pong_timeout: float = 10.0
- def parse_host_list(value: str | list[str] | None) -> list[str]:
- """Parst einen Host-String oder eine Liste in eine Liste von Hostnames/IPs.
- Akzeptiert:
- - String: "trixyone"
- - Kommagetrennt: "trixyone, 10.10.10.1"
- - Semikolon-/Whitespace-getrennt: "trixyone;10.10.10.1 trixy.lan"
- - Liste: ["trixyone", "10.10.10.1"]
- Whitespace und leere Eintraege werden entfernt.
- Reihenfolge bleibt erhalten.
- """
- if not value:
- return []
- if isinstance(value, list):
- return [str(h).strip() for h in value if str(h).strip()]
- import re
- parts = re.split(r"[,;\s]+", str(value).strip())
- return [p.strip() for p in parts if p.strip()]
- class ClientConnection:
- """
- Client-seitige Verbindung zum Trixy-Server.
- Verwaltet alle 4 Sockets:
- - Command: Befehle und Steuerung
- - Audio In: Mikrofon-Daten zum Server
- - Audio Out: Sprachausgabe vom Server
- - Music: Musik-Stream vom Server
- Usage:
- connection = ClientConnection(
- host="192.168.1.100",
- port=2101,
- room="wohnzimmer",
- alias="Echo",
- mac_address="AA:BB:CC:DD:EE:FF"
- )
- if await connection.connect():
- # Verbunden
- await connection.send_command(WakewordDetected(...))
- """
- def __init__(
- self,
- host: str | list[str],
- port: int,
- room: str,
- alias: str,
- mac_address: str,
- config: ConnectionConfig | None = None,
- ) -> None:
- """
- Initialisiert die ClientConnection.
- Args:
- host: Server-Host. Kann sein:
- - Ein einzelner String: "trixyone"
- - Kommagetrennt: "trixyone, 10.10.10.1, trixy.lan"
- - Liste: ["trixyone", "10.10.10.1"]
- Beim Verbinden werden alle Hosts der Reihe nach probiert,
- der erste erreichbare gewinnt.
- port: Server-Port (Command-Port, Standard 2101)
- room: Raum-Kennung
- alias: Anzeigename des Satellites
- mac_address: MAC-Adresse des Satellites
- config: Verbindungskonfiguration
- """
- # Host-Liste normalisieren
- self._hosts = parse_host_list(host)
- # _host bleibt fuer Backward-Compatibility — wird beim Connect aktualisiert
- self._host = self._hosts[0] if self._hosts else ""
- self._port = port
- self._room = room
- self._alias = alias
- self._mac_address = mac_address
- self._config = config or ConnectionConfig()
- # Protokoll
- self._protocol = TrixyProtocol()
- self._encryption: TrixyEncryption | None = None
- self._session_key: bytes = b""
- # Sockets
- self._command_reader: asyncio.StreamReader | None = None
- self._command_writer: asyncio.StreamWriter | None = None
- self._audio_in_reader: asyncio.StreamReader | None = None
- self._audio_in_writer: asyncio.StreamWriter | None = None
- self._audio_out_reader: asyncio.StreamReader | None = None
- self._audio_out_writer: asyncio.StreamWriter | None = None
- self._music_reader: asyncio.StreamReader | None = None
- self._music_writer: asyncio.StreamWriter | None = None
- # Zustand
- self._connected = False
- self._satellite_id: str = ""
- self._audio_ports: dict[str, int] = {}
- self._running = False
- self._reconnect_attempts = 0
- # Letzte Denied-Antwort — fuer Reconnect-Steuerung
- self._last_denied_reason: str = ""
- self._last_denied_retry_after: int = 0
- # Audio-Format (vom Server mitgeteilt)
- self._tts_format: dict[str, int] = {
- "sample_rate": 22050,
- "channels": 1,
- "sample_width": 2,
- }
- self._music_format: dict[str, int] = {
- "sample_rate": 44100,
- "channels": 2,
- "sample_width": 2,
- }
- # Tasks
- self._receive_task: asyncio.Task | None = None
- self._keepalive_task: asyncio.Task | None = None
- self._audio_out_task: asyncio.Task | None = None
- self._music_task: asyncio.Task | None = None
- # Callbacks
- self._on_connected: Callable[[], Awaitable[None]] | None = None
- self._on_disconnected: Callable[[str], Awaitable[None]] | None = None
- self._on_message: Callable[[ProtocolMessage], Awaitable[None]] | None = None
- self._on_audio_out: Callable[[bytes], Awaitable[None]] | None = None
- self._on_music: Callable[[bytes], Awaitable[None]] | None = None
- @property
- def is_connected(self) -> bool:
- """Ist die Verbindung aktiv?"""
- return self._connected
- @property
- def satellite_id(self) -> str:
- """Die vom Server zugewiesene Satellite-ID."""
- return self._satellite_id
- @property
- def last_denied_retry_after(self) -> int:
- """Wartezeit in Sekunden nach letztem SatelliteDenied (0 wenn keins)."""
- return self._last_denied_retry_after
- @property
- def last_denied_reason(self) -> str:
- """Grund der letzten Ablehnung (leer wenn keine)."""
- return self._last_denied_reason
- def clear_denied_state(self) -> None:
- """Setzt den Denied-State zurueck (nach erfolgreichem Connect)."""
- self._last_denied_reason = ""
- self._last_denied_retry_after = 0
- @property
- def host(self) -> str:
- """Server-Host."""
- return self._host
- @property
- def port(self) -> int:
- """Server-Port."""
- return self._port
- @property
- def tts_format(self) -> dict[str, int]:
- """TTS-Audio-Format vom Server."""
- return self._tts_format.copy()
- @property
- def music_format(self) -> dict[str, int]:
- """Musik-Audio-Format vom Server."""
- return self._music_format.copy()
- def on_connected(self, callback: Callable[[], Awaitable[None]]) -> None:
- """Registriert Callback für Verbindungsaufbau."""
- self._on_connected = callback
- def on_disconnected(self, callback: Callable[[str], Awaitable[None]]) -> None:
- """Registriert Callback für Verbindungsabbruch."""
- self._on_disconnected = callback
- def on_message(self, callback: Callable[[ProtocolMessage], Awaitable[None]]) -> None:
- """Registriert Callback für eingehende Nachrichten."""
- self._on_message = callback
- def on_audio_out(self, callback: Callable[[bytes], Awaitable[None]]) -> None:
- """Registriert Callback für Audio-Ausgabe."""
- self._on_audio_out = callback
- def on_music(self, callback: Callable[[bytes], Awaitable[None]]) -> None:
- """Registriert Callback für Musik-Stream."""
- self._on_music = callback
- async def connect(self) -> bool:
- """
- Verbindet sich mit dem Server.
- Probiert alle konfigurierten Hosts der Reihe nach durch und nimmt
- den ersten der antwortet. Fuehrt dann den Handshake durch und
- verbindet alle Audio-Ports.
- Returns:
- True bei erfolgreicher Verbindung
- """
- if not self._hosts:
- perror("Keine Server-Hosts konfiguriert")
- return False
- # Alle Hosts der Reihe nach probieren
- last_error = ""
- connected_host: str | None = None
- for candidate in self._hosts:
- pdebug(f"Verbinde zu {candidate}:{self._port}...")
- try:
- self._command_reader, self._command_writer = await asyncio.wait_for(
- asyncio.open_connection(candidate, self._port),
- timeout=self._config.connect_timeout,
- )
- connected_host = candidate
- self._host = candidate # Aktualisiere fuer Logging und Audio-Sockets
- pdebug(f"Command-Socket verbunden zu {candidate}")
- break
- except asyncio.TimeoutError:
- last_error = f"Timeout zu {candidate}:{self._port}"
- pdebug(last_error)
- continue
- except OSError as e:
- last_error = f"Fehler {candidate}: {e}"
- pdebug(last_error)
- continue
- if connected_host is None:
- if len(self._hosts) > 1:
- perror(f"Kein Server erreichbar — probierte: {', '.join(self._hosts)}")
- else:
- perror(f"Verbindungsfehler: {last_error}")
- return False
- try:
- # Handshake durchführen
- if not await self._perform_handshake():
- await self._close_sockets()
- return False
- # Audio-Sockets verbinden
- await self._connect_audio_sockets()
- # Empfangs-Tasks starten
- self._running = True
- self._receive_task = asyncio.create_task(self._receive_loop())
- self._keepalive_task = asyncio.create_task(self._keepalive_loop())
- if self._audio_out_reader:
- self._audio_out_task = asyncio.create_task(self._audio_out_loop())
- if self._music_reader:
- self._music_task = asyncio.create_task(self._music_loop())
- self._connected = True
- self._reconnect_attempts = 0
- self.clear_denied_state()
- pinfo(f"Verbunden mit Server {self._host}:{self._port}")
- # Callback
- if self._on_connected:
- try:
- await self._on_connected()
- except Exception as e:
- perror(f"on_connected Callback Fehler: {e}")
- return True
- except Exception as e:
- perror(f"Verbindungsfehler: {e}")
- await self._close_sockets()
- return False
- async def disconnect(self, reason: str = "Client-Disconnect") -> None:
- """
- Trennt die Verbindung zum Server.
- Args:
- reason: Grund für die Trennung
- """
- if not self._connected:
- return
- pinfo(f"Trenne Verbindung: {reason}")
- self._running = False
- # Disconnect-Nachricht senden
- try:
- if self._command_writer:
- disconnect_msg = SatelliteDisconnect(
- mac_address=self._mac_address,
- reason=reason
- )
- await self._send_message(disconnect_msg)
- except Exception:
- pass
- # Tasks beenden
- for task in [self._receive_task, self._keepalive_task, self._audio_out_task, self._music_task]:
- if task and not task.done():
- task.cancel()
- try:
- await task
- except asyncio.CancelledError:
- pass
- # Sockets schließen
- await self._close_sockets()
- self._connected = False
- self._satellite_id = ""
- # Callback
- if self._on_disconnected:
- try:
- await self._on_disconnected(reason)
- except Exception as e:
- perror(f"on_disconnected Callback Fehler: {e}")
- async def send_command(self, message: object, flags: ProtocolFlags = ProtocolFlags.PICKLE) -> bool:
- """
- Sendet eine Nachricht zum Server.
- Args:
- message: Die zu sendende Nachricht
- flags: Protokoll-Flags
- Returns:
- True bei Erfolg
- """
- if not self._connected or not self._command_writer:
- return False
- return await self._send_message(message, flags)
- async def send_audio(self, audio_data: bytes) -> bool:
- """
- Sendet Audio-Daten zum Server.
- Args:
- audio_data: Audio-Bytes (16KHz, 16-bit, mono)
- Returns:
- True bei Erfolg
- """
- if not self._audio_in_writer:
- return False
- try:
- self._audio_in_writer.write(audio_data)
- await self._audio_in_writer.drain()
- return True
- except Exception as e:
- perror(f"Audio-Sende-Fehler: {e}")
- return False
- async def _perform_handshake(self) -> bool:
- """
- Führt den Handshake mit dem Server durch.
- Returns:
- True bei erfolgreichem Handshake
- """
- pdebug("Starte Handshake...")
- # HELLO senden
- try:
- self._command_writer.write(COMMAND_HELLO)
- await self._command_writer.drain()
- except Exception as e:
- perror(f"Fehler beim Senden von HELLO: {e}")
- return False
- pdebug("HELLO gesendet")
- # SatelliteConnect senden
- import platform
- import socket as _socket
- connect_msg = SatelliteConnect(
- mac_address=self._mac_address,
- room=self._room,
- alias=self._alias,
- ip_address="", # Server ermittelt IP
- version=VERSION_STRING,
- hostname=_socket.gethostname(),
- os_name=platform.system(),
- os_version=platform.release()
- )
- if not await self._send_message(connect_msg):
- perror("Fehler beim Senden von SatelliteConnect")
- return False
- pdebug("SatelliteConnect gesendet")
- # Antwort empfangen
- try:
- response = await asyncio.wait_for(
- self._read_message(),
- timeout=self._config.connect_timeout
- )
- except asyncio.TimeoutError:
- perror("Timeout beim Warten auf Server-Antwort")
- return False
- if response is None:
- perror("Keine Antwort vom Server")
- return False
- # Antwort verarbeiten
- if response.class_name == "SatelliteDenied":
- data = response.data
- if isinstance(data, dict):
- reason = data.get("reason", "Unbekannt")
- retry_after = data.get("retry_after_seconds", 60)
- elif isinstance(data, SatelliteDenied):
- reason = data.reason
- retry_after = data.retry_after_seconds
- else:
- reason = "Unbekannt"
- retry_after = 60
- # Fuer den Client-Reconnect-Loop speichern
- self._last_denied_reason = reason
- self._last_denied_retry_after = int(retry_after)
- pwarn(f"Verbindung abgelehnt: {reason}")
- pwarn(f"Erneuter Versuch in {retry_after} Sekunden moeglich")
- return False
- if response.class_name != "SatelliteAccepted":
- perror(f"Unerwartete Antwort: {response.class_name}")
- return False
- # SatelliteAccepted verarbeiten
- data = response.data
- if isinstance(data, dict):
- self._satellite_id = data.get("satellite_id", "")
- self._session_key = data.get("session_key", b"")
- self._audio_ports = {
- "audio_in": data.get("audio_input_port", 0),
- "audio_out": data.get("audio_output_port", 0),
- "music": data.get("music_port", 0),
- }
- # Audio-Format vom Server
- if data.get("tts_sample_rate"):
- self._tts_format = {
- "sample_rate": data.get("tts_sample_rate", 22050),
- "channels": data.get("tts_channels", 1),
- "sample_width": data.get("tts_sample_width", 2),
- }
- if data.get("music_sample_rate"):
- self._music_format = {
- "sample_rate": data.get("music_sample_rate", 44100),
- "channels": data.get("music_channels", 2),
- "sample_width": data.get("music_sample_width", 2),
- }
- elif isinstance(data, SatelliteAccepted):
- self._satellite_id = data.satellite_id
- self._session_key = data.session_key
- self._audio_ports = {
- "audio_in": data.audio_input_port,
- "audio_out": data.audio_output_port,
- "music": data.music_port,
- }
- # Audio-Format vom Server
- if data.tts_sample_rate:
- self._tts_format = {
- "sample_rate": data.tts_sample_rate,
- "channels": data.tts_channels,
- "sample_width": data.tts_sample_width,
- }
- if data.music_sample_rate:
- self._music_format = {
- "sample_rate": data.music_sample_rate,
- "channels": data.music_channels,
- "sample_width": data.music_sample_width,
- }
- else:
- perror(f"Ungültige SatelliteAccepted-Daten: {type(data)}")
- return False
- # Session-Verschlüsselung einrichten
- if self._session_key:
- try:
- self._encryption = TrixyEncryption(self._session_key)
- self._protocol.set_encryption(self._encryption)
- pdebug("Session-Verschlüsselung aktiviert")
- except Exception as e:
- pwarn(f"Session-Verschlüsselung fehlgeschlagen: {e}")
- pinfo(f"Handshake erfolgreich, Satellite-ID: {self._satellite_id}")
- return True
- async def _connect_audio_sockets(self) -> None:
- """Verbindet die Audio-Sockets."""
- # Audio-Input Socket
- if self._audio_ports.get("audio_in"):
- try:
- self._audio_in_reader, self._audio_in_writer = await asyncio.wait_for(
- asyncio.open_connection(self._host, self._audio_ports["audio_in"]),
- timeout=self._config.connect_timeout
- )
- # Satellite-ID senden
- self._audio_in_writer.write(self._satellite_id.ljust(36).encode("utf-8"))
- await self._audio_in_writer.drain()
- pdebug(f"Audio-Input verbunden auf Port {self._audio_ports['audio_in']}")
- except Exception as e:
- pwarn(f"Audio-Input Verbindung fehlgeschlagen: {e}")
- # Audio-Output Socket
- if self._audio_ports.get("audio_out"):
- try:
- self._audio_out_reader, self._audio_out_writer = await asyncio.wait_for(
- asyncio.open_connection(self._host, self._audio_ports["audio_out"]),
- timeout=self._config.connect_timeout
- )
- # Satellite-ID senden
- self._audio_out_writer.write(self._satellite_id.ljust(36).encode("utf-8"))
- await self._audio_out_writer.drain()
- pdebug(f"Audio-Output verbunden auf Port {self._audio_ports['audio_out']}")
- except Exception as e:
- pwarn(f"Audio-Output Verbindung fehlgeschlagen: {e}")
- # Music Socket
- if self._audio_ports.get("music"):
- try:
- self._music_reader, self._music_writer = await asyncio.wait_for(
- asyncio.open_connection(self._host, self._audio_ports["music"]),
- timeout=self._config.connect_timeout
- )
- # Satellite-ID senden
- self._music_writer.write(self._satellite_id.ljust(36).encode("utf-8"))
- await self._music_writer.drain()
- pdebug(f"Music verbunden auf Port {self._audio_ports['music']}")
- except Exception as e:
- pwarn(f"Music Verbindung fehlgeschlagen: {e}")
- async def _receive_loop(self) -> None:
- """Empfängt Nachrichten vom Server."""
- while self._running:
- try:
- message = await self._read_message()
- if message is None:
- if self._running:
- pwarn("Verbindung zum Server verloren")
- asyncio.create_task(self._handle_disconnect("Verbindung verloren"))
- break
- await self._handle_message(message)
- except asyncio.CancelledError:
- break
- except Exception as e:
- if self._running:
- perror(f"Empfangsfehler: {e}")
- asyncio.create_task(self._handle_disconnect(f"Fehler: {e}"))
- break
- async def _handle_message(self, message: ProtocolMessage) -> None:
- """Verarbeitet eine empfangene Nachricht."""
- class_name = message.class_name
- # PING beantworten
- if class_name == "TRXIPING":
- await self._send_pong()
- return
- # PONG ignorieren (KeepAlive)
- if class_name == "TRXIPONG":
- return
- # An Callback weiterleiten
- if self._on_message:
- try:
- await self._on_message(message)
- except Exception as e:
- perror(f"Message-Callback Fehler: {e}")
- async def _keepalive_loop(self) -> None:
- """Sendet regelmäßige Pings."""
- while self._running:
- try:
- await asyncio.sleep(self._config.ping_interval)
- if not self._running:
- break
- # Ping senden
- await self._send_ping()
- except asyncio.CancelledError:
- break
- except Exception as e:
- pdebug(f"KeepAlive-Fehler: {e}")
- async def _audio_out_loop(self) -> None:
- """Empfängt Audio-Ausgabe vom Server."""
- first_chunk = True
- total_bytes = 0
- chunk_count = 0
- # Frame-Alignment Buffer (TTS: Mono 16-bit = 2 bytes pro Frame)
- frame_size = self._tts_format["channels"] * self._tts_format["sample_width"]
- leftover = b""
- while self._running and self._audio_out_reader:
- try:
- chunk = await self._audio_out_reader.read(32768) # 32KB chunks
- if not chunk:
- break
- # Kombiniere mit Resten vom letzten Chunk
- data = leftover + chunk
- # Nur komplette Frames weitergeben
- aligned_size = (len(data) // frame_size) * frame_size
- if aligned_size > 0:
- aligned_chunk = data[:aligned_size]
- leftover = data[aligned_size:]
- chunk_count += 1
- total_bytes += len(aligned_chunk)
- if first_chunk:
- pdebug(f"[CLIENT] Audio-Out: Erster Chunk empfangen ({len(aligned_chunk)} bytes)")
- first_chunk = False
- if self._on_audio_out:
- await self._on_audio_out(aligned_chunk)
- else:
- leftover = data
- except asyncio.CancelledError:
- break
- except Exception as e:
- pdebug(f"Audio-Out Fehler: {e}")
- break
- if chunk_count > 0:
- pdebug(f"[CLIENT] Audio-Out Stream beendet: {chunk_count} chunks, {total_bytes} bytes total")
- async def _music_loop(self) -> None:
- """Empfängt Music-Stream vom Server."""
- first_chunk = True
- total_bytes = 0
- chunk_count = 0
- # Frame-Alignment Buffer (Music: Stereo 16-bit = 4 bytes pro Frame)
- frame_size = self._music_format["channels"] * self._music_format["sample_width"]
- leftover = b""
- while self._running and self._music_reader:
- try:
- chunk = await self._music_reader.read(32768) # 32KB chunks
- if not chunk:
- break
- # Kombiniere mit Resten vom letzten Chunk
- data = leftover + chunk
- # Nur komplette Frames weitergeben
- aligned_size = (len(data) // frame_size) * frame_size
- if aligned_size > 0:
- aligned_chunk = data[:aligned_size]
- leftover = data[aligned_size:]
- chunk_count += 1
- total_bytes += len(aligned_chunk)
- if first_chunk:
- pdebug(f"[CLIENT] Music: Erster Chunk empfangen ({len(aligned_chunk)} bytes)")
- first_chunk = False
- if self._on_music:
- await self._on_music(aligned_chunk)
- else:
- leftover = data
- except asyncio.CancelledError:
- break
- except Exception as e:
- pdebug(f"Music Fehler: {e}")
- break
- if chunk_count > 0:
- pdebug(f"[CLIENT] Music Stream beendet: {chunk_count} chunks, {total_bytes} bytes total")
- async def _handle_disconnect(self, reason: str) -> None:
- """Behandelt einen Verbindungsabbruch."""
- was_connected = self._connected
- self._connected = False
- self._running = False
- await self._close_sockets()
- if was_connected and self._on_disconnected:
- try:
- await self._on_disconnected(reason)
- except Exception as e:
- perror(f"on_disconnected Callback Fehler: {e}")
- async def _send_message(self, message: object, flags: ProtocolFlags = ProtocolFlags.PICKLE) -> bool:
- """Sendet eine Nachricht über den Command-Socket."""
- if not self._command_writer:
- return False
- try:
- data = self._protocol.serialize(message, flags)
- self._command_writer.write(data)
- await self._command_writer.drain()
- return True
- except Exception as e:
- perror(f"Sende-Fehler: {e}")
- return False
- async def _send_ping(self) -> None:
- """Sendet einen Ping."""
- if self._command_writer:
- try:
- self._command_writer.write(COMMAND_PING)
- await self._command_writer.drain()
- except Exception:
- pass
- async def _send_pong(self) -> None:
- """Sendet einen Pong."""
- if self._command_writer:
- try:
- self._command_writer.write(COMMAND_PONG)
- await self._command_writer.drain()
- except Exception:
- pass
- async def _read_message(self) -> ProtocolMessage | None:
- """Liest eine Nachricht vom Command-Socket."""
- if not self._command_reader:
- return None
- 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._command_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._command_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 Klassennamen-Länge
- class_name_length = struct.unpack(">H", header[34:36])[0]
- # Lese Klassenname
- class_name = await self._command_reader.readexactly(class_name_length)
- # Lese Datenlänge (4 bytes)
- data_length_bytes = await self._command_reader.readexactly(4)
- data_length = struct.unpack(">I", data_length_bytes)[0]
- # Lese Daten
- data = await self._command_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"Lese-Fehler: {e}")
- return None
- async def _close_sockets(self) -> None:
- """Schließt alle Sockets."""
- for writer in [
- self._command_writer,
- self._audio_in_writer,
- self._audio_out_writer,
- self._music_writer
- ]:
- if writer:
- try:
- writer.close()
- await writer.wait_closed()
- except Exception:
- pass
- self._command_reader = None
- self._command_writer = None
- self._audio_in_reader = None
- self._audio_in_writer = None
- self._audio_out_reader = None
- self._audio_out_writer = None
- self._music_reader = None
- self._music_writer = None
- def get_mac_address() -> str:
- """
- Ermittelt die MAC-Adresse des Systems.
- Returns:
- MAC-Adresse im Format "AA:BB:CC:DD:EE:FF"
- """
- import uuid as uuid_module
- mac = uuid_module.getnode()
- return ":".join(f"{(mac >> (8 * i)) & 0xFF:02X}" for i in range(5, -1, -1))
|