satellite_manager.py 42 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181
  1. """
  2. Satellite Manager Implementation
  3. This module implements the SatelliteManager class that serves as the central management
  4. system for all registered satellites in the Trixy application. It provides advanced
  5. access patterns, bulk operations, and comprehensive satellite lifecycle management
  6. as specified in CLAUDE.md.
  7. Key Features:
  8. - Direct index access: satellite_manager[0]
  9. - Query-based access: satellite_manager["status=connected,room=kitchen"]
  10. - Bulk operations: disconnect_all(), reconnect_all(), say_all()
  11. - Thread-safe operations for multi-satellite support
  12. - Registration file management and blacklist support
  13. - Integration with network system for socket management
  14. - Integration with event system for satellite events
  15. - Connection state management (registered, not connected, connected)
  16. - Comprehensive logging and error handling
  17. """
  18. import json
  19. import threading
  20. import time
  21. import weakref
  22. from typing import List, Dict, Any, Optional, Union, Iterator, Tuple, Set, Callable
  23. from dataclasses import dataclass, field
  24. from datetime import datetime, timezone
  25. from enum import Enum
  26. from pathlib import Path
  27. import uuid
  28. import os
  29. from .satellite import (
  30. Satellite,
  31. SatelliteStatus,
  32. SatelliteInfo,
  33. SatelliteCapability,
  34. AudioPortInfo,
  35. create_satellite,
  36. create_satellite_info,
  37. SatelliteError,
  38. SatelliteConnectionError
  39. )
  40. from .query_parser import (
  41. QueryParser,
  42. QueryCondition,
  43. QueryError,
  44. parse_query,
  45. filter_satellites
  46. )
  47. def pprint(message: str) -> None:
  48. """Satellite manager logging function."""
  49. print(f"[SATELLITE_MANAGER] {message}")
  50. class ConnectionState(Enum):
  51. """Connection states for satellites."""
  52. REGISTERED = "registered"
  53. NOT_CONNECTED = "not_connected"
  54. CONNECTED = "connected"
  55. ERROR = "error"
  56. BLACKLISTED = "blacklisted"
  57. @dataclass
  58. class SatelliteStats:
  59. """Statistics for satellite management."""
  60. total_satellites: int = 0
  61. connected_satellites: int = 0
  62. registered_satellites: int = 0
  63. blacklisted_satellites: int = 0
  64. error_satellites: int = 0
  65. total_connections: int = 0
  66. total_messages: int = 0
  67. total_errors: int = 0
  68. uptime_seconds: float = 0.0
  69. last_updated: Optional[datetime] = None
  70. def to_dict(self) -> Dict[str, Any]:
  71. """Convert to dictionary representation."""
  72. return {
  73. "total_satellites": self.total_satellites,
  74. "connected_satellites": self.connected_satellites,
  75. "registered_satellites": self.registered_satellites,
  76. "blacklisted_satellites": self.blacklisted_satellites,
  77. "error_satellites": self.error_satellites,
  78. "total_connections": self.total_connections,
  79. "total_messages": self.total_messages,
  80. "total_errors": self.total_errors,
  81. "uptime_seconds": self.uptime_seconds,
  82. "last_updated": self.last_updated.isoformat() if self.last_updated else None,
  83. }
  84. class SatelliteManagerError(Exception):
  85. """Base exception for satellite manager errors."""
  86. pass
  87. class SatelliteNotFoundError(SatelliteManagerError):
  88. """Raised when a requested satellite is not found."""
  89. pass
  90. class SatelliteRegistrationError(SatelliteManagerError):
  91. """Raised when satellite registration fails."""
  92. pass
  93. class SatelliteConnectionError(SatelliteManagerError):
  94. """Raised when satellite connection operations fail."""
  95. pass
  96. class SatelliteManager:
  97. """
  98. Central management system for all registered satellites.
  99. This class implements the interface specified in CLAUDE.md, providing:
  100. - Advanced access patterns (index, query-based, bulk operations)
  101. - Thread-safe operations for multi-satellite environments
  102. - Registration and blacklist management
  103. - Integration with network and event systems
  104. - Comprehensive statistics and monitoring
  105. Usage Examples:
  106. # Direct access
  107. satellite = satellite_manager[0]
  108. # Query-based access
  109. kitchen_satellites = satellite_manager["room=kitchen"]
  110. connected_satellites = satellite_manager["status=connected"]
  111. # Bulk operations
  112. satellite_manager.disconnect_all("room=living_room")
  113. satellite_manager.say_all("Hello everyone!", "status=connected")
  114. """
  115. def __init__(
  116. self,
  117. application,
  118. registration_dir: Optional[str] = None,
  119. blacklist_file: Optional[str] = None,
  120. max_satellites: int = 50,
  121. enable_auto_registration: bool = False,
  122. registration_timeout: float = 60.0,
  123. connection_timeout: float = 30.0,
  124. heartbeat_interval: float = 30.0
  125. ):
  126. """
  127. Initialize the satellite manager.
  128. Args:
  129. application: Application container reference
  130. registration_dir: Directory for registration files
  131. blacklist_file: Path to blacklist file
  132. max_satellites: Maximum number of satellites
  133. enable_auto_registration: Enable automatic registration
  134. registration_timeout: Timeout for registration process
  135. connection_timeout: Timeout for connections
  136. heartbeat_interval: Heartbeat check interval
  137. """
  138. self._application = application
  139. self._application_ref = weakref.ref(application) if application else None
  140. # Configuration
  141. self._max_satellites = max_satellites
  142. self._enable_auto_registration = enable_auto_registration
  143. self._registration_timeout = registration_timeout
  144. self._connection_timeout = connection_timeout
  145. self._heartbeat_interval = heartbeat_interval
  146. # File paths
  147. self._registration_dir = Path(registration_dir or "config/satellites")
  148. self._blacklist_file = Path(blacklist_file or "config/satellites/blacklist.json")
  149. # Ensure directories exist
  150. self._registration_dir.mkdir(parents=True, exist_ok=True)
  151. self._blacklist_file.parent.mkdir(parents=True, exist_ok=True)
  152. # Storage
  153. self._satellites: Dict[str, Satellite] = {} # MAC address -> Satellite
  154. self._satellites_by_id: Dict[str, Satellite] = {} # satellite_id -> Satellite
  155. self._blacklisted_macs: Set[str] = set()
  156. # Thread safety
  157. self._lock = threading.RLock()
  158. # Query parser
  159. self._query_parser = QueryParser(case_sensitive=False)
  160. # Statistics
  161. self._stats = SatelliteStats()
  162. self._start_time = datetime.now(timezone.utc)
  163. # Registration mode
  164. self._registration_mode = False
  165. self._registration_timer = None
  166. # Event integration
  167. self._event_handler = None
  168. self._network_manager = None
  169. # Load existing data
  170. self._load_blacklist()
  171. self._load_registered_satellites()
  172. pprint(f"SatelliteManager initialized (max: {max_satellites}, auto_reg: {enable_auto_registration})")
  173. self._update_stats()
  174. # Advanced Access Patterns Implementation
  175. def __getitem__(self, key: Union[int, str]) -> Union[Satellite, List[Satellite]]:
  176. """
  177. Advanced access pattern implementation.
  178. Supports:
  179. - Direct index access: satellite_manager[0]
  180. - Query-based access: satellite_manager["status=connected,room=kitchen"]
  181. Args:
  182. key: Index (int) or query string (str)
  183. Returns:
  184. Satellite or List[Satellite] depending on access pattern
  185. Raises:
  186. SatelliteNotFoundError: If satellite not found
  187. QueryError: If query is invalid
  188. """
  189. with self._lock:
  190. if isinstance(key, int):
  191. # Direct index access
  192. satellites = list(self._satellites.values())
  193. if key < 0 or key >= len(satellites):
  194. raise SatelliteNotFoundError(f"Satellite index {key} out of range (0-{len(satellites)-1})")
  195. return satellites[key]
  196. elif isinstance(key, str):
  197. # Query-based access
  198. try:
  199. satellites = list(self._satellites.values())
  200. filtered = self._query_parser.filter_satellites(satellites, key)
  201. return filtered
  202. except QueryError as e:
  203. raise SatelliteManagerError(f"Invalid query '{key}': {e}")
  204. else:
  205. raise SatelliteManagerError(f"Unsupported key type: {type(key)}")
  206. def __len__(self) -> int:
  207. """Get the number of registered satellites."""
  208. return len(self._satellites)
  209. def __iter__(self) -> Iterator[Satellite]:
  210. """Iterate over all satellites."""
  211. with self._lock:
  212. return iter(list(self._satellites.values()))
  213. def __contains__(self, item: Union[str, Satellite]) -> bool:
  214. """Check if satellite exists by MAC address or satellite object."""
  215. with self._lock:
  216. if isinstance(item, str):
  217. return item in self._satellites
  218. elif isinstance(item, Satellite):
  219. return item.mac_address in self._satellites
  220. return False
  221. # Core Management Methods
  222. def add_satellite(self, satellite_info: SatelliteInfo) -> Satellite:
  223. """
  224. Add a new satellite to the manager.
  225. Args:
  226. satellite_info: Complete satellite information
  227. Returns:
  228. Satellite: The created satellite instance
  229. Raises:
  230. SatelliteRegistrationError: If registration fails
  231. """
  232. with self._lock:
  233. if len(self._satellites) >= self._max_satellites:
  234. raise SatelliteRegistrationError(
  235. f"Maximum number of satellites ({self._max_satellites}) reached"
  236. )
  237. mac_address = satellite_info.mac_address.lower()
  238. # Check blacklist
  239. if mac_address in self._blacklisted_macs:
  240. raise SatelliteRegistrationError(f"MAC address {mac_address} is blacklisted")
  241. # Check if already exists
  242. if mac_address in self._satellites:
  243. pprint(f"Satellite {mac_address} already registered, updating info")
  244. existing = self._satellites[mac_address]
  245. # Update satellite info
  246. existing._info = satellite_info
  247. return existing
  248. # Create new satellite
  249. satellite = create_satellite(
  250. satellite_info,
  251. self._application,
  252. self._network_manager
  253. )
  254. # Add to storage
  255. self._satellites[mac_address] = satellite
  256. self._satellites_by_id[satellite_info.satellite_id] = satellite
  257. # Save registration
  258. self._save_satellite_registration(satellite_info)
  259. # Update statistics
  260. self._update_stats()
  261. # Trigger event
  262. self._trigger_satellite_registered_event(satellite)
  263. pprint(f"Satellite added: {satellite_info.alias_name} ({mac_address})")
  264. return satellite
  265. def remove_satellite(self, mac_address: str, reason: str = "manual_removal") -> bool:
  266. """
  267. Remove a satellite from the manager.
  268. Args:
  269. mac_address: MAC address of satellite to remove
  270. reason: Reason for removal
  271. Returns:
  272. bool: True if satellite was removed
  273. """
  274. with self._lock:
  275. mac_address = mac_address.lower()
  276. if mac_address not in self._satellites:
  277. return False
  278. satellite = self._satellites[mac_address]
  279. # Disconnect if connected
  280. if satellite.is_connected:
  281. satellite.disconnect(reason=f"removal: {reason}")
  282. # Remove from storage
  283. del self._satellites[mac_address]
  284. if satellite.satellite_id in self._satellites_by_id:
  285. del self._satellites_by_id[satellite.satellite_id]
  286. # Remove registration file
  287. self._remove_satellite_registration(mac_address)
  288. # Update statistics
  289. self._update_stats()
  290. pprint(f"Satellite removed: {satellite.alias_name} ({mac_address})")
  291. return True
  292. def get_satellite_by_mac(self, mac_address: str) -> Optional[Satellite]:
  293. """Get satellite by MAC address."""
  294. with self._lock:
  295. return self._satellites.get(mac_address.lower())
  296. def get_satellite_by_id(self, satellite_id: str) -> Optional[Satellite]:
  297. """Get satellite by satellite ID."""
  298. with self._lock:
  299. return self._satellites_by_id.get(satellite_id)
  300. def get_satellite_by_alias(self, alias_name: str) -> Optional[Satellite]:
  301. """Get satellite by alias name."""
  302. with self._lock:
  303. for satellite in self._satellites.values():
  304. if satellite.alias_name.lower() == alias_name.lower():
  305. return satellite
  306. return None
  307. def get_satellites_by_room(self, room_id: str) -> List[Satellite]:
  308. """Get all satellites in a specific room."""
  309. with self._lock:
  310. return [s for s in self._satellites.values() if s.room_id.lower() == room_id.lower()]
  311. def get_connected_satellites(self) -> List[Satellite]:
  312. """Get all currently connected satellites."""
  313. with self._lock:
  314. return [s for s in self._satellites.values() if s.is_connected]
  315. def get_registered_satellites(self) -> List[Satellite]:
  316. """Get all registered satellites."""
  317. with self._lock:
  318. return list(self._satellites.values())
  319. # Bulk Operations
  320. def disconnect_all(self, query: Optional[str] = None, reason: str = "bulk_disconnect") -> int:
  321. """
  322. Disconnect all satellites or satellites matching a query.
  323. Args:
  324. query: Optional query to filter satellites
  325. reason: Reason for disconnection
  326. Returns:
  327. int: Number of satellites disconnected
  328. """
  329. with self._lock:
  330. if query:
  331. satellites = self[query]
  332. if not isinstance(satellites, list):
  333. satellites = [satellites]
  334. else:
  335. satellites = list(self._satellites.values())
  336. disconnected = 0
  337. for satellite in satellites:
  338. if satellite.is_connected:
  339. if satellite.disconnect(reason=reason):
  340. disconnected += 1
  341. pprint(f"Bulk disconnect: {disconnected} satellites disconnected")
  342. return disconnected
  343. def reconnect_all(self, query: Optional[str] = None) -> int:
  344. """
  345. Reconnect all satellites or satellites matching a query.
  346. Args:
  347. query: Optional query to filter satellites
  348. Returns:
  349. int: Number of satellites reconnected
  350. """
  351. with self._lock:
  352. if query:
  353. satellites = self[query]
  354. if not isinstance(satellites, list):
  355. satellites = [satellites]
  356. else:
  357. satellites = list(self._satellites.values())
  358. reconnected = 0
  359. for satellite in satellites:
  360. if not satellite.is_connected and satellite.status != SatelliteStatus.BLACKLISTED:
  361. if satellite.reconnect():
  362. reconnected += 1
  363. pprint(f"Bulk reconnect: {reconnected} satellites reconnected")
  364. return reconnected
  365. def say_all(
  366. self,
  367. text: str,
  368. query: Optional[str] = None,
  369. voice_settings: Optional[Dict[str, Any]] = None
  370. ) -> int:
  371. """
  372. Send TTS message to all satellites or satellites matching a query.
  373. Args:
  374. text: Text to speak
  375. query: Optional query to filter satellites
  376. voice_settings: Optional voice configuration
  377. Returns:
  378. int: Number of satellites that received the message
  379. """
  380. with self._lock:
  381. if query:
  382. satellites = self[query]
  383. if not isinstance(satellites, list):
  384. satellites = [satellites]
  385. else:
  386. satellites = self.get_connected_satellites()
  387. sent = 0
  388. for satellite in satellites:
  389. if satellite.is_connected:
  390. if satellite.say(text, voice_settings):
  391. sent += 1
  392. pprint(f"Bulk TTS: '{text[:50]}...' sent to {sent} satellites")
  393. return sent
  394. def update_all_capabilities(
  395. self,
  396. capability: SatelliteCapability,
  397. add: bool = True,
  398. query: Optional[str] = None
  399. ) -> int:
  400. """
  401. Add or remove a capability from satellites.
  402. Args:
  403. capability: Capability to add/remove
  404. add: True to add, False to remove
  405. query: Optional query to filter satellites
  406. Returns:
  407. int: Number of satellites updated
  408. """
  409. with self._lock:
  410. if query:
  411. satellites = self[query]
  412. if not isinstance(satellites, list):
  413. satellites = [satellites]
  414. else:
  415. satellites = list(self._satellites.values())
  416. updated = 0
  417. for satellite in satellites:
  418. if add:
  419. if not satellite.has_capability(capability):
  420. satellite.add_capability(capability)
  421. updated += 1
  422. else:
  423. if satellite.has_capability(capability):
  424. satellite.remove_capability(capability)
  425. updated += 1
  426. action = "added" if add else "removed"
  427. pprint(f"Capability {capability.value} {action} for {updated} satellites")
  428. return updated
  429. # Registration Management
  430. def enter_registration_mode(self, timeout: Optional[float] = None) -> bool:
  431. """
  432. Enter registration mode to allow new satellites to register.
  433. Args:
  434. timeout: Optional timeout in seconds (uses default if None)
  435. Returns:
  436. bool: True if registration mode was activated
  437. """
  438. with self._lock:
  439. if self._registration_mode:
  440. pprint("Registration mode is already active")
  441. return True
  442. timeout = timeout or self._registration_timeout
  443. self._registration_mode = True
  444. # Set up automatic exit timer
  445. if self._registration_timer:
  446. self._registration_timer.cancel()
  447. self._registration_timer = threading.Timer(
  448. timeout,
  449. self._exit_registration_mode_auto
  450. )
  451. self._registration_timer.start()
  452. pprint(f"Registration mode activated for {timeout} seconds")
  453. # Trigger event
  454. self._trigger_registration_mode_event(True, timeout)
  455. return True
  456. def exit_registration_mode(self) -> bool:
  457. """
  458. Exit registration mode.
  459. Returns:
  460. bool: True if registration mode was deactivated
  461. """
  462. with self._lock:
  463. if not self._registration_mode:
  464. return True
  465. self._registration_mode = False
  466. if self._registration_timer:
  467. self._registration_timer.cancel()
  468. self._registration_timer = None
  469. pprint("Registration mode deactivated")
  470. # Trigger event
  471. self._trigger_registration_mode_event(False, 0)
  472. return True
  473. def is_registration_mode_active(self) -> bool:
  474. """Check if registration mode is active."""
  475. return self._registration_mode
  476. def can_register_satellite(self, mac_address: str) -> bool:
  477. """
  478. Check if a satellite can be registered.
  479. Args:
  480. mac_address: MAC address to check
  481. Returns:
  482. bool: True if registration is allowed
  483. """
  484. with self._lock:
  485. mac_address = mac_address.lower()
  486. # Check if at capacity
  487. if len(self._satellites) >= self._max_satellites:
  488. return False
  489. # Check if blacklisted
  490. if mac_address in self._blacklisted_macs:
  491. return False
  492. # Check if already registered
  493. if mac_address in self._satellites:
  494. return True # Can update existing registration
  495. # Check if registration mode is active or auto-registration is enabled
  496. return self._registration_mode or self._enable_auto_registration
  497. # Blacklist Management
  498. def add_to_blacklist(self, mac_address: str, reason: str = "manual_blacklist") -> bool:
  499. """
  500. Add a MAC address to the blacklist.
  501. Args:
  502. mac_address: MAC address to blacklist
  503. reason: Reason for blacklisting
  504. Returns:
  505. bool: True if added to blacklist
  506. """
  507. with self._lock:
  508. mac_address = mac_address.lower()
  509. if mac_address in self._blacklisted_macs:
  510. return True
  511. self._blacklisted_macs.add(mac_address)
  512. # Disconnect if currently connected
  513. if mac_address in self._satellites:
  514. satellite = self._satellites[mac_address]
  515. if satellite.is_connected:
  516. satellite.disconnect(reason=f"blacklisted: {reason}")
  517. satellite.set_error(f"Blacklisted: {reason}")
  518. # Save blacklist
  519. self._save_blacklist()
  520. pprint(f"MAC address blacklisted: {mac_address} (reason: {reason})")
  521. return True
  522. def remove_from_blacklist(self, mac_address: str) -> bool:
  523. """
  524. Remove a MAC address from the blacklist.
  525. Args:
  526. mac_address: MAC address to remove from blacklist
  527. Returns:
  528. bool: True if removed from blacklist
  529. """
  530. with self._lock:
  531. mac_address = mac_address.lower()
  532. if mac_address not in self._blacklisted_macs:
  533. return True
  534. self._blacklisted_macs.remove(mac_address)
  535. # Save blacklist
  536. self._save_blacklist()
  537. pprint(f"MAC address removed from blacklist: {mac_address}")
  538. return True
  539. def is_blacklisted(self, mac_address: str) -> bool:
  540. """Check if a MAC address is blacklisted."""
  541. return mac_address.lower() in self._blacklisted_macs
  542. def get_blacklisted_macs(self) -> List[str]:
  543. """Get all blacklisted MAC addresses."""
  544. with self._lock:
  545. return list(self._blacklisted_macs)
  546. # Connection Management
  547. def handle_satellite_connection(
  548. self,
  549. mac_address: str,
  550. ip_address: str,
  551. room_id: str,
  552. alias_name: str,
  553. version: str = "1.0.0",
  554. capabilities: Optional[List[str]] = None,
  555. audio_ports: Optional[Dict[str, int]] = None
  556. ) -> Tuple[bool, str, Optional[Satellite]]:
  557. """
  558. Handle a satellite connection attempt.
  559. Args:
  560. mac_address: Satellite MAC address
  561. ip_address: Satellite IP address
  562. room_id: Room identifier
  563. alias_name: Human-readable name
  564. version: Software version
  565. capabilities: List of capability names
  566. audio_ports: Audio port information
  567. Returns:
  568. Tuple[bool, str, Optional[Satellite]]: (success, message, satellite)
  569. """
  570. with self._lock:
  571. mac_address = mac_address.lower()
  572. # Check if can register/connect
  573. if not self.can_register_satellite(mac_address):
  574. if mac_address in self._blacklisted_macs:
  575. return False, "MAC address is blacklisted", None
  576. elif len(self._satellites) >= self._max_satellites:
  577. return False, "Maximum satellites limit reached", None
  578. else:
  579. return False, "Registration mode is not active", None
  580. # Parse capabilities
  581. parsed_capabilities = []
  582. if capabilities:
  583. for cap_name in capabilities:
  584. try:
  585. cap = SatelliteCapability(cap_name.lower())
  586. parsed_capabilities.append(cap)
  587. except ValueError:
  588. pprint(f"Unknown capability: {cap_name}")
  589. # Parse audio ports
  590. audio_port_info = AudioPortInfo()
  591. if audio_ports:
  592. audio_port_info = AudioPortInfo.from_dict(audio_ports)
  593. # Get or create satellite
  594. satellite = self._satellites.get(mac_address)
  595. if satellite is None:
  596. # Create new satellite
  597. satellite_info = create_satellite_info(
  598. satellite_id=str(uuid.uuid4()),
  599. mac_address=mac_address,
  600. room_id=room_id,
  601. alias_name=alias_name,
  602. ip_address=ip_address,
  603. version=version,
  604. capabilities=parsed_capabilities
  605. )
  606. satellite = self.add_satellite(satellite_info)
  607. else:
  608. # Update existing satellite
  609. satellite._info.room_id = room_id
  610. satellite._info.alias_name = alias_name
  611. satellite._info.version = version
  612. satellite._info.capabilities = parsed_capabilities
  613. satellite.update_ip_address(ip_address)
  614. # Mark as connected
  615. satellite.set_connected(ip_address, audio_port_info)
  616. satellite.update_last_seen()
  617. # Update statistics
  618. self._update_stats()
  619. pprint(f"Satellite connected: {alias_name} ({mac_address}) from {ip_address}")
  620. return True, "Connection accepted", satellite
  621. def handle_satellite_disconnection(
  622. self,
  623. mac_address: str,
  624. reason: str = "normal_disconnect"
  625. ) -> bool:
  626. """
  627. Handle a satellite disconnection.
  628. Args:
  629. mac_address: Satellite MAC address
  630. reason: Reason for disconnection
  631. Returns:
  632. bool: True if handled successfully
  633. """
  634. with self._lock:
  635. mac_address = mac_address.lower()
  636. satellite = self._satellites.get(mac_address)
  637. if satellite is None:
  638. return True
  639. satellite.set_disconnected(reason)
  640. # Update statistics
  641. self._update_stats()
  642. pprint(f"Satellite disconnected: {satellite.alias_name} ({mac_address}) - {reason}")
  643. return True
  644. # Statistics and Monitoring
  645. def get_stats(self) -> SatelliteStats:
  646. """Get comprehensive satellite statistics."""
  647. with self._lock:
  648. self._update_stats()
  649. return self._stats
  650. def get_status(self) -> Dict[str, Any]:
  651. """Get comprehensive status information."""
  652. with self._lock:
  653. stats = self.get_stats()
  654. return {
  655. "statistics": stats.to_dict(),
  656. "configuration": {
  657. "max_satellites": self._max_satellites,
  658. "auto_registration": self._enable_auto_registration,
  659. "registration_timeout": self._registration_timeout,
  660. "connection_timeout": self._connection_timeout,
  661. "heartbeat_interval": self._heartbeat_interval,
  662. "registration_dir": str(self._registration_dir),
  663. "blacklist_file": str(self._blacklist_file),
  664. },
  665. "registration": {
  666. "mode_active": self._registration_mode,
  667. "can_register": len(self._satellites) < self._max_satellites,
  668. },
  669. "satellites_by_status": {
  670. status.value: len([
  671. s for s in self._satellites.values()
  672. if s.status == status
  673. ]) for status in SatelliteStatus
  674. },
  675. "satellites_by_room": self._get_satellites_by_room_stats(),
  676. "blacklist": {
  677. "count": len(self._blacklisted_macs),
  678. "mac_addresses": list(self._blacklisted_macs),
  679. },
  680. "capabilities": self._get_capabilities_stats(),
  681. "timestamp": datetime.now(timezone.utc).isoformat(),
  682. }
  683. def _update_stats(self) -> None:
  684. """Update internal statistics."""
  685. now = datetime.now(timezone.utc)
  686. total_satellites = len(self._satellites)
  687. connected_count = 0
  688. registered_count = 0
  689. blacklisted_count = 0
  690. error_count = 0
  691. total_connections = 0
  692. total_messages = 0
  693. total_errors = 0
  694. for satellite in self._satellites.values():
  695. if satellite.status == SatelliteStatus.CONNECTED:
  696. connected_count += 1
  697. elif satellite.status == SatelliteStatus.REGISTERED:
  698. registered_count += 1
  699. elif satellite.status == SatelliteStatus.BLACKLISTED:
  700. blacklisted_count += 1
  701. elif satellite.status == SatelliteStatus.ERROR:
  702. error_count += 1
  703. # Get satellite stats
  704. sat_status = satellite.get_status()
  705. stats = sat_status["statistics"]
  706. total_connections += stats.get("connection_count", 0)
  707. total_messages += stats.get("message_count", 0)
  708. total_errors += stats.get("error_count", 0)
  709. blacklisted_count += len(self._blacklisted_macs)
  710. uptime = (now - self._start_time).total_seconds()
  711. self._stats = SatelliteStats(
  712. total_satellites=total_satellites,
  713. connected_satellites=connected_count,
  714. registered_satellites=registered_count,
  715. blacklisted_satellites=blacklisted_count,
  716. error_satellites=error_count,
  717. total_connections=total_connections,
  718. total_messages=total_messages,
  719. total_errors=total_errors,
  720. uptime_seconds=uptime,
  721. last_updated=now
  722. )
  723. def _get_satellites_by_room_stats(self) -> Dict[str, int]:
  724. """Get satellite count by room."""
  725. room_stats = {}
  726. for satellite in self._satellites.values():
  727. room = satellite.room_id
  728. room_stats[room] = room_stats.get(room, 0) + 1
  729. return room_stats
  730. def _get_capabilities_stats(self) -> Dict[str, int]:
  731. """Get capability usage statistics."""
  732. capability_stats = {}
  733. for satellite in self._satellites.values():
  734. for capability in satellite.capabilities:
  735. cap_name = capability.value
  736. capability_stats[cap_name] = capability_stats.get(cap_name, 0) + 1
  737. return capability_stats
  738. # File Management
  739. def _load_registered_satellites(self) -> None:
  740. """Load registered satellites from files."""
  741. if not self._registration_dir.exists():
  742. return
  743. loaded_count = 0
  744. for file_path in self._registration_dir.glob("*.json"):
  745. try:
  746. with open(file_path, 'r') as f:
  747. data = json.load(f)
  748. satellite_info = SatelliteInfo.from_dict(data)
  749. satellite = create_satellite(
  750. satellite_info,
  751. self._application,
  752. self._network_manager
  753. )
  754. mac_address = satellite_info.mac_address.lower()
  755. self._satellites[mac_address] = satellite
  756. self._satellites_by_id[satellite_info.satellite_id] = satellite
  757. loaded_count += 1
  758. except Exception as e:
  759. pprint(f"Error loading satellite from {file_path}: {e}")
  760. if loaded_count > 0:
  761. pprint(f"Loaded {loaded_count} registered satellites")
  762. def _save_satellite_registration(self, satellite_info: SatelliteInfo) -> None:
  763. """Save satellite registration to file."""
  764. filename = f"{satellite_info.mac_address.replace(':', '_')}.json"
  765. file_path = self._registration_dir / filename
  766. try:
  767. with open(file_path, 'w') as f:
  768. json.dump(satellite_info.to_dict(), f, indent=2)
  769. pprint(f"Satellite registration saved: {filename}")
  770. except Exception as e:
  771. pprint(f"Error saving satellite registration: {e}")
  772. def _remove_satellite_registration(self, mac_address: str) -> None:
  773. """Remove satellite registration file."""
  774. filename = f"{mac_address.replace(':', '_')}.json"
  775. file_path = self._registration_dir / filename
  776. try:
  777. if file_path.exists():
  778. file_path.unlink()
  779. pprint(f"Satellite registration removed: {filename}")
  780. except Exception as e:
  781. pprint(f"Error removing satellite registration: {e}")
  782. def _load_blacklist(self) -> None:
  783. """Load blacklisted MAC addresses."""
  784. if not self._blacklist_file.exists():
  785. return
  786. try:
  787. with open(self._blacklist_file, 'r') as f:
  788. data = json.load(f)
  789. blacklist = data.get("blacklisted_macs", [])
  790. self._blacklisted_macs = set(mac.lower() for mac in blacklist)
  791. if self._blacklisted_macs:
  792. pprint(f"Loaded {len(self._blacklisted_macs)} blacklisted MAC addresses")
  793. except Exception as e:
  794. pprint(f"Error loading blacklist: {e}")
  795. def _save_blacklist(self) -> None:
  796. """Save blacklisted MAC addresses."""
  797. try:
  798. data = {
  799. "blacklisted_macs": list(self._blacklisted_macs),
  800. "updated_at": datetime.now(timezone.utc).isoformat()
  801. }
  802. with open(self._blacklist_file, 'w') as f:
  803. json.dump(data, f, indent=2)
  804. except Exception as e:
  805. pprint(f"Error saving blacklist: {e}")
  806. # Event Integration
  807. def set_event_handler(self, event_handler) -> None:
  808. """Set the event handler for satellite events."""
  809. self._event_handler = event_handler
  810. def set_network_manager(self, network_manager) -> None:
  811. """Set the network manager for satellite communication."""
  812. self._network_manager = network_manager
  813. # Update all satellites with network manager
  814. with self._lock:
  815. for satellite in self._satellites.values():
  816. satellite._network_manager = network_manager
  817. def _trigger_satellite_registered_event(self, satellite: Satellite) -> None:
  818. """Trigger satellite registered event."""
  819. try:
  820. if self._application_ref and self._application_ref():
  821. app = self._application_ref()
  822. event_handler = app.get_event_handler()
  823. from ..events import EventDataFactory
  824. register_data = EventDataFactory.create_event_data(
  825. "satellite_registered",
  826. satellite_id=satellite.satellite_id,
  827. mac_address=satellite.mac_address,
  828. room_id=satellite.room_id,
  829. alias_name=satellite.alias_name,
  830. registration_time=datetime.now(timezone.utc).isoformat()
  831. )
  832. event_handler.trigger_event("satellite_registered", register_data)
  833. except Exception as e:
  834. pprint(f"Error triggering satellite registered event: {e}")
  835. def _trigger_registration_mode_event(self, active: bool, timeout: float) -> None:
  836. """Trigger registration mode change event."""
  837. try:
  838. if self._application_ref and self._application_ref():
  839. app = self._application_ref()
  840. event_handler = app.get_event_handler()
  841. from ..events import EventDataFactory
  842. mode_data = EventDataFactory.create_event_data(
  843. "registration_mode_changed",
  844. active=active,
  845. timeout=timeout,
  846. timestamp=datetime.now(timezone.utc).isoformat()
  847. )
  848. event_handler.trigger_event("registration_mode_changed", mode_data)
  849. except Exception as e:
  850. pprint(f"Error triggering registration mode event: {e}")
  851. def _exit_registration_mode_auto(self) -> None:
  852. """Automatically exit registration mode (timer callback)."""
  853. self.exit_registration_mode()
  854. pprint("Registration mode automatically deactivated (timeout)")
  855. # Cleanup and Maintenance
  856. def cleanup_disconnected_satellites(self, max_age_hours: float = 24.0) -> int:
  857. """
  858. Remove satellites that have been disconnected for too long.
  859. Args:
  860. max_age_hours: Maximum age in hours for disconnected satellites
  861. Returns:
  862. int: Number of satellites cleaned up
  863. """
  864. with self._lock:
  865. cutoff_time = datetime.now(timezone.utc) - timedelta(hours=max_age_hours)
  866. to_remove = []
  867. for mac_address, satellite in self._satellites.items():
  868. if (not satellite.is_connected and
  869. satellite._info.last_seen and
  870. satellite._info.last_seen < cutoff_time):
  871. to_remove.append(mac_address)
  872. for mac_address in to_remove:
  873. self.remove_satellite(mac_address, "cleanup_old")
  874. if to_remove:
  875. pprint(f"Cleaned up {len(to_remove)} old disconnected satellites")
  876. return len(to_remove)
  877. def validate_satellite_states(self) -> int:
  878. """
  879. Validate and fix inconsistent satellite states.
  880. Returns:
  881. int: Number of satellites with fixed states
  882. """
  883. with self._lock:
  884. fixed_count = 0
  885. for satellite in self._satellites.values():
  886. # Check for inconsistent states
  887. if satellite.status == SatelliteStatus.CONNECTED and not satellite._sockets:
  888. satellite.set_disconnected("state_validation")
  889. fixed_count += 1
  890. elif satellite.status != SatelliteStatus.CONNECTED and satellite._sockets:
  891. satellite._sockets.clear()
  892. fixed_count += 1
  893. if fixed_count > 0:
  894. pprint(f"Fixed {fixed_count} inconsistent satellite states")
  895. return fixed_count
  896. def shutdown(self) -> None:
  897. """Shutdown the satellite manager gracefully."""
  898. with self._lock:
  899. pprint("Shutting down satellite manager...")
  900. # Exit registration mode
  901. self.exit_registration_mode()
  902. # Disconnect all satellites
  903. disconnected = self.disconnect_all(reason="system_shutdown")
  904. # Save final state
  905. self._save_blacklist()
  906. for satellite in self._satellites.values():
  907. self._save_satellite_registration(satellite._info)
  908. pprint(f"Satellite manager shutdown complete ({disconnected} satellites disconnected)")
  909. def __str__(self) -> str:
  910. """String representation of the satellite manager."""
  911. stats = self.get_stats()
  912. return (f"SatelliteManager(total={stats.total_satellites}, "
  913. f"connected={stats.connected_satellites}, "
  914. f"reg_mode={self._registration_mode})")
  915. def __repr__(self) -> str:
  916. """Detailed representation of the satellite manager."""
  917. return (f"SatelliteManager(satellites={len(self._satellites)}, "
  918. f"max={self._max_satellites}, "
  919. f"auto_reg={self._enable_auto_registration}, "
  920. f"reg_mode={self._registration_mode})")
  921. # Factory function
  922. def create_satellite_manager(
  923. application,
  924. **kwargs
  925. ) -> SatelliteManager:
  926. """
  927. Create a SatelliteManager instance with default configuration.
  928. Args:
  929. application: Application container reference
  930. **kwargs: Additional configuration parameters
  931. Returns:
  932. SatelliteManager: Configured satellite manager
  933. """
  934. return SatelliteManager(application, **kwargs)
  935. # Module exports
  936. __all__ = [
  937. "SatelliteManager",
  938. "SatelliteManagerError",
  939. "SatelliteNotFoundError",
  940. "SatelliteRegistrationError",
  941. "SatelliteConnectionError",
  942. "ConnectionState",
  943. "SatelliteStats",
  944. "create_satellite_manager",
  945. ]