| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717 |
- """
- Network Manager for Trixy System
- This module provides the main NetworkManager class that handles:
- - Command socket management (Port 2101)
- - Audio streaming sockets (Ports 2102, 2103, 2104)
- - Multi-client support with thread safety
- - Connection management and monitoring
- - Protocol serialization/deserialization
- - Integration with event system and application container
- The NetworkManager serves as the central hub for all network communication
- in the Trixy system, supporting both server and client modes.
- """
- import socket
- import threading
- import time
- import queue
- import weakref
- from typing import Dict, Any, Optional, List, Callable, Union, Tuple
- from dataclasses import dataclass, field
- from enum import Enum, IntEnum
- from concurrent.futures import ThreadPoolExecutor
- import select
- from .protocol import TrixyProtocol, TrixyMessage, MessageFlags, ProtocolError
- from .cmd import TrixyCommand, CommandResponse
- def pprint(message: str) -> None:
- """Network manager logging function."""
- print(f"[NETWORK_MANAGER] {message}")
- class ConnectionState(Enum):
- """Connection state enumeration."""
- DISCONNECTED = "disconnected"
- CONNECTING = "connecting"
- CONNECTED = "connected"
- AUTHENTICATING = "authenticating"
- AUTHENTICATED = "authenticated"
- DISCONNECTING = "disconnecting"
- ERROR = "error"
- class AudioStreamType(Enum):
- """Audio stream type enumeration."""
- INPUT = "input" # Port 2102 - Raw audio from satellites
- OUTPUT = "output" # Port 2103 - TTS/response audio
- MUSIC = "music" # Port 2104 - Music/media playback
- @dataclass
- class AudioStreamConfig:
- """Configuration for audio streams."""
- sample_rate: int = 16000
- channels: int = 1
- sample_width: int = 2 # bytes
- format: str = "PCM"
- buffer_size: int = 4096
- @dataclass
- class ConnectionInfo:
- """Information about a network connection."""
- connection_id: str
- remote_address: Tuple[str, int]
- local_address: Tuple[str, int]
- state: ConnectionState
- connected_time: float
- last_activity: float
- bytes_sent: int = 0
- bytes_received: int = 0
- commands_sent: int = 0
- commands_received: int = 0
- satellite_id: Optional[str] = None
- session_id: Optional[str] = None
- class NetworkManagerError(Exception):
- """Base exception for network manager errors."""
- pass
- class ConnectionManager:
- """Manages individual connections and their state."""
-
- def __init__(self):
- self._connections: Dict[str, ConnectionInfo] = {}
- self._sockets: Dict[str, socket.socket] = {}
- self._lock = threading.RLock()
-
- def add_connection(self, connection_id: str, sock: socket.socket, remote_addr: Tuple[str, int]) -> ConnectionInfo:
- """Add a new connection."""
- with self._lock:
- local_addr = sock.getsockname()
- connection = ConnectionInfo(
- connection_id=connection_id,
- remote_address=remote_addr,
- local_address=local_addr,
- state=ConnectionState.CONNECTED,
- connected_time=time.time(),
- last_activity=time.time()
- )
-
- self._connections[connection_id] = connection
- self._sockets[connection_id] = sock
-
- pprint(f"Connection added: {connection_id} from {remote_addr}")
- return connection
-
- def remove_connection(self, connection_id: str) -> Optional[ConnectionInfo]:
- """Remove a connection."""
- with self._lock:
- connection = self._connections.pop(connection_id, None)
- sock = self._sockets.pop(connection_id, None)
-
- if sock:
- try:
- sock.close()
- except:
- pass
-
- if connection:
- pprint(f"Connection removed: {connection_id}")
-
- return connection
-
- def get_connection(self, connection_id: str) -> Optional[ConnectionInfo]:
- """Get connection info."""
- with self._lock:
- return self._connections.get(connection_id)
-
- def get_socket(self, connection_id: str) -> Optional[socket.socket]:
- """Get socket for connection."""
- with self._lock:
- return self._sockets.get(connection_id)
-
- def update_activity(self, connection_id: str, bytes_delta: int = 0, command_delta: int = 0):
- """Update connection activity."""
- with self._lock:
- connection = self._connections.get(connection_id)
- if connection:
- connection.last_activity = time.time()
- if bytes_delta > 0:
- connection.bytes_received += bytes_delta
- if command_delta > 0:
- connection.commands_received += command_delta
-
- def get_all_connections(self) -> Dict[str, ConnectionInfo]:
- """Get all connections."""
- with self._lock:
- return self._connections.copy()
- class AudioStreamManager:
- """Manages audio streaming sockets."""
-
- def __init__(self):
- self._streams: Dict[AudioStreamType, Dict[str, Any]] = {}
- self._configs: Dict[AudioStreamType, AudioStreamConfig] = {}
- self._lock = threading.RLock()
-
- # Initialize default configs
- for stream_type in AudioStreamType:
- if stream_type == AudioStreamType.MUSIC:
- self._configs[stream_type] = AudioStreamConfig(
- sample_rate=48000,
- channels=2,
- sample_width=2
- )
- else:
- self._configs[stream_type] = AudioStreamConfig()
-
- def start_stream(self, stream_type: AudioStreamType, port: int, config: Optional[AudioStreamConfig] = None):
- """Start an audio stream."""
- with self._lock:
- if stream_type in self._streams:
- raise NetworkManagerError(f"Audio stream {stream_type.value} already active")
-
- if config:
- self._configs[stream_type] = config
-
- # Create server socket for audio stream
- server_sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
- server_sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
- server_sock.bind(('', port))
- server_sock.listen(10)
-
- self._streams[stream_type] = {
- 'server_socket': server_sock,
- 'port': port,
- 'connections': {},
- 'active': True
- }
-
- pprint(f"Audio stream {stream_type.value} started on port {port}")
-
- def stop_stream(self, stream_type: AudioStreamType):
- """Stop an audio stream."""
- with self._lock:
- stream_info = self._streams.pop(stream_type, None)
- if stream_info:
- stream_info['active'] = False
- server_sock = stream_info.get('server_socket')
- if server_sock:
- server_sock.close()
-
- # Close all connections
- for conn_sock in stream_info.get('connections', {}).values():
- try:
- conn_sock.close()
- except:
- pass
-
- pprint(f"Audio stream {stream_type.value} stopped")
-
- def get_stream_info(self, stream_type: AudioStreamType) -> Optional[Dict[str, Any]]:
- """Get stream information."""
- with self._lock:
- return self._streams.get(stream_type)
- class NetworkManager:
- """
- Main network manager for the Trixy system.
-
- Handles all network communication including command sockets and audio streams.
- Integrates with the event system and application container.
- """
-
- def __init__(self, application_container):
- """
- Initialize the network manager.
-
- Args:
- application_container: The main application container
- """
- self.application = application_container
- self.application_ref = weakref.ref(application_container)
-
- # Core components
- self.protocol = TrixyProtocol()
- self.connection_manager = ConnectionManager()
- self.audio_stream_manager = AudioStreamManager()
-
- # Server state
- self._server_running = False
- self._command_server_socket: Optional[socket.socket] = None
- self._command_port = 2101
-
- # Client state
- self._client_connected = False
- self._client_socket: Optional[socket.socket] = None
- self._server_address: Optional[Tuple[str, int]] = None
-
- # Threading
- self._thread_pool = ThreadPoolExecutor(max_workers=20)
- self._running = True
- self._lock = threading.RLock()
-
- # Message handling
- self._message_queue = queue.Queue()
- self._message_handlers: Dict[str, Callable] = {}
- self._response_handlers: Dict[str, Callable] = {}
-
- # Event system integration
- self._event_handler = None
-
- pprint("NetworkManager initialized")
-
- def start_server(self, command_port: int = 2101, audio_ports: Optional[Dict[str, int]] = None) -> None:
- """
- Start the network manager in server mode.
-
- Args:
- command_port: Port for command socket
- audio_ports: Ports for audio streams
- """
- if self._server_running:
- raise NetworkManagerError("Server already running")
-
- pprint(f"Starting network server on port {command_port}")
-
- self._command_port = command_port
-
- try:
- # Create command server socket
- self._command_server_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
- self._command_server_socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
- self._command_server_socket.bind(('', command_port))
- self._command_server_socket.listen(50)
-
- # Start audio streams
- if audio_ports is None:
- audio_ports = {
- 'input': 2102,
- 'output': 2103,
- 'music': 2104
- }
-
- for stream_name, port in audio_ports.items():
- if stream_name in ['input', 'output', 'music']:
- stream_type = AudioStreamType(stream_name)
- self.audio_stream_manager.start_stream(stream_type, port)
-
- self._server_running = True
-
- # Start accepting connections
- self._thread_pool.submit(self._accept_connections)
- self._thread_pool.submit(self._process_messages)
-
- # Trigger server started event
- self._trigger_event("network_server_started", {
- 'command_port': command_port,
- 'audio_ports': audio_ports
- })
-
- pprint(f"Network server started successfully")
-
- except Exception as e:
- self._cleanup_server()
- raise NetworkManagerError(f"Failed to start server: {e}") from e
-
- def stop_server(self) -> None:
- """Stop the network server."""
- if not self._server_running:
- return
-
- pprint("Stopping network server...")
- self._server_running = False
-
- self._cleanup_server()
-
- self._trigger_event("network_server_stopped", {})
- pprint("Network server stopped")
-
- def connect_to_server(self, host: str, port: int = 2101) -> bool:
- """
- Connect to a Trixy server as a client.
-
- Args:
- host: Server hostname or IP
- port: Server command port
-
- Returns:
- bool: True if connected successfully
- """
- if self._client_connected:
- raise NetworkManagerError("Client already connected")
-
- pprint(f"Connecting to server {host}:{port}")
-
- try:
- self._client_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
- self._client_socket.settimeout(10.0)
- self._client_socket.connect((host, port))
-
- self._server_address = (host, port)
- self._client_connected = True
-
- # Start client message processing
- connection_id = f"client_{int(time.time())}"
- self.connection_manager.add_connection(
- connection_id,
- self._client_socket,
- self._server_address
- )
-
- self._thread_pool.submit(self._handle_client_connection, connection_id)
- self._thread_pool.submit(self._process_messages)
-
- self._trigger_event("network_client_connected", {
- 'server_address': self._server_address
- })
-
- pprint(f"Connected to server successfully")
- return True
-
- except Exception as e:
- self._cleanup_client()
- pprint(f"Failed to connect to server: {e}")
- return False
-
- def disconnect_from_server(self) -> None:
- """Disconnect from server."""
- if not self._client_connected:
- return
-
- pprint("Disconnecting from server...")
- self._client_connected = False
-
- self._cleanup_client()
-
- self._trigger_event("network_client_disconnected", {})
- pprint("Disconnected from server")
-
- def send_command(self, command: TrixyCommand, connection_id: Optional[str] = None) -> bool:
- """
- Send a command over the network.
-
- Args:
- command: Command to send
- connection_id: Target connection (None for client mode)
-
- Returns:
- bool: True if sent successfully
- """
- try:
- # Create protocol message
- message = self.protocol.create_message(
- command.get_command_name(),
- command
- )
-
- # Serialize message
- data = message.serialize()
-
- # Send based on mode
- if self._client_connected:
- # Client mode - send to server
- sock = self._client_socket
- if sock:
- sock.send(data)
- pprint(f"Sent command: {command.get_command_name()}")
- return True
- else:
- # Server mode - send to specific connection
- if connection_id:
- sock = self.connection_manager.get_socket(connection_id)
- if sock:
- sock.send(data)
- pprint(f"Sent command to {connection_id}: {command.get_command_name()}")
- return True
-
- return False
-
- except Exception as e:
- pprint(f"Failed to send command: {e}")
- return False
-
- def broadcast_command(self, command: TrixyCommand) -> int:
- """
- Broadcast a command to all connected clients.
-
- Args:
- command: Command to broadcast
-
- Returns:
- int: Number of clients that received the command
- """
- if not self._server_running:
- return 0
-
- connections = self.connection_manager.get_all_connections()
- sent_count = 0
-
- for connection_id in connections:
- if self.send_command(command, connection_id):
- sent_count += 1
-
- return sent_count
-
- def register_message_handler(self, command_name: str, handler: Callable[[TrixyMessage, str], None]) -> None:
- """
- Register a handler for incoming messages.
-
- Args:
- command_name: Name of command to handle
- handler: Handler function (message, connection_id) -> None
- """
- self._message_handlers[command_name] = handler
- pprint(f"Registered message handler for: {command_name}")
-
- def _accept_connections(self) -> None:
- """Accept incoming connections (server thread)."""
- pprint("Started accepting connections")
-
- while self._server_running and self._command_server_socket:
- try:
- # Use select to check for new connections with timeout
- ready, _, _ = select.select([self._command_server_socket], [], [], 1.0)
-
- if ready:
- client_sock, addr = self._command_server_socket.accept()
- connection_id = f"client_{addr[0]}_{addr[1]}_{int(time.time())}"
-
- self.connection_manager.add_connection(connection_id, client_sock, addr)
-
- # Handle connection in separate thread
- self._thread_pool.submit(self._handle_client_connection, connection_id)
-
- self._trigger_event("network_client_connected", {
- 'connection_id': connection_id,
- 'remote_address': addr
- })
-
- except Exception as e:
- # Only log errors if server is still supposed to be running
- if self._server_running and self._running:
- pprint(f"Error accepting connection: {e}")
-
- def _handle_client_connection(self, connection_id: str) -> None:
- """Handle messages from a client connection."""
- sock = self.connection_manager.get_socket(connection_id)
- if not sock:
- return
-
- pprint(f"Handling connection: {connection_id}")
-
- try:
- while self._running and sock:
- # Receive data with timeout
- sock.settimeout(1.0)
-
- try:
- data = sock.recv(4096)
- if not data:
- break
-
- # Update connection activity
- self.connection_manager.update_activity(connection_id, len(data))
-
- # Process the received data
- self._process_received_data(data, connection_id)
-
- except socket.timeout:
- continue
- except Exception as e:
- pprint(f"Error receiving data from {connection_id}: {e}")
- break
-
- except Exception as e:
- pprint(f"Error handling connection {connection_id}: {e}")
- finally:
- self.connection_manager.remove_connection(connection_id)
- self._trigger_event("network_client_disconnected", {
- 'connection_id': connection_id
- })
-
- def _process_received_data(self, data: bytes, connection_id: str) -> None:
- """Process received data and extract messages."""
- try:
- # Check if it's a hard-coded command
- if self.protocol.is_hardcoded_command(data):
- command, args = self.protocol.parse_hardcoded_command(data)
- pprint(f"Received hard-coded command from {connection_id}: {command}")
-
- # Handle hard-coded commands immediately
- self._handle_hardcoded_command(command, args, connection_id)
- else:
- # Parse as protocol message
- message = self.protocol.deserialize_message(data)
-
- # Queue for processing
- self._message_queue.put((message, connection_id))
-
- except Exception as e:
- pprint(f"Error processing received data: {e}")
-
- def _process_messages(self) -> None:
- """Process queued messages."""
- pprint("Started message processing")
-
- while self._running:
- try:
- # Get message with timeout
- message, connection_id = self._message_queue.get(timeout=1.0)
-
- # Handle the message
- self._handle_message(message, connection_id)
-
- self._message_queue.task_done()
-
- except queue.Empty:
- continue
- except Exception as e:
- pprint(f"Error processing message: {e}")
-
- def _handle_message(self, message: TrixyMessage, connection_id: str) -> None:
- """Handle a received message."""
- try:
- command_name = message.class_name
-
- # Check for registered handler
- handler = self._message_handlers.get(command_name)
- if handler:
- handler(message, connection_id)
- else:
- # Default handling - trigger event
- self._trigger_event("network_message_received", {
- 'message': message,
- 'connection_id': connection_id,
- 'command_name': command_name
- })
-
- pprint(f"Handled message: {command_name} from {connection_id}")
-
- except Exception as e:
- pprint(f"Error handling message: {e}")
-
- def _handle_hardcoded_command(self, command: str, args: str, connection_id: str) -> None:
- """Handle hard-coded commands."""
- try:
- if command == "TRXINOOP":
- # No-op/heartbeat - just update activity
- self.connection_manager.update_activity(connection_id)
-
- elif command == "TRXIPING":
- # Respond with pong
- pong_data = f"TRXIPONG {args}".encode('utf-8')
- sock = self.connection_manager.get_socket(connection_id)
- if sock:
- sock.send(pong_data)
-
- elif command == "TRXIPRNT":
- # Print command
- pprint(f"Print from {connection_id}: {args}")
-
- elif command == "TRXYHELO":
- # Hello command
- pprint(f"Hello from {connection_id}: {args}")
-
- except Exception as e:
- pprint(f"Error handling hard-coded command {command}: {e}")
-
- def _trigger_event(self, event_name: str, event_data: Dict[str, Any]) -> None:
- """Trigger an event through the event system."""
- try:
- # Don't trigger events if shutting down
- if not self._running:
- return
-
- if not self._event_handler:
- app = self.application_ref()
- if app:
- self._event_handler = app.get_event_handler()
-
- if self._event_handler:
- self._event_handler.trigger_event(event_name, event_data)
-
- except Exception as e:
- # Only log errors if still running (avoid shutdown-related errors)
- if self._running:
- pprint(f"Error triggering event {event_name}: {e}")
-
- def _cleanup_server(self) -> None:
- """Clean up server resources."""
- if self._command_server_socket:
- self._command_server_socket.close()
- self._command_server_socket = None
-
- # Stop all audio streams
- for stream_type in AudioStreamType:
- self.audio_stream_manager.stop_stream(stream_type)
-
- # Close all connections
- for connection_id in list(self.connection_manager.get_all_connections().keys()):
- self.connection_manager.remove_connection(connection_id)
-
- def _cleanup_client(self) -> None:
- """Clean up client resources."""
- if self._client_socket:
- self._client_socket.close()
- self._client_socket = None
-
- self._server_address = None
-
- def shutdown(self) -> None:
- """Shutdown the network manager."""
- pprint("Shutting down network manager...")
-
- self._running = False
-
- if self._server_running:
- self.stop_server()
-
- if self._client_connected:
- self.disconnect_from_server()
-
- self._thread_pool.shutdown(wait=True)
-
- pprint("Network manager shutdown complete")
-
- def get_status(self) -> Dict[str, Any]:
- """Get network manager status."""
- connections = self.connection_manager.get_all_connections()
-
- return {
- 'server_running': self._server_running,
- 'client_connected': self._client_connected,
- 'command_port': self._command_port,
- 'server_address': self._server_address,
- 'active_connections': len(connections),
- 'connection_details': {cid: {
- 'remote_address': conn.remote_address,
- 'state': conn.state.value,
- 'connected_time': conn.connected_time,
- 'bytes_sent': conn.bytes_sent,
- 'bytes_received': conn.bytes_received
- } for cid, conn in connections.items()}
- }
- # Module exports
- __all__ = [
- "NetworkManager",
- "NetworkManagerError",
- "ConnectionManager",
- "ConnectionState",
- "ConnectionInfo",
- "AudioStreamManager",
- "AudioStreamType",
- "AudioStreamConfig",
- ]
|