discovery.py 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452
  1. # -*- coding: utf-8 -*-
  2. """
  3. Capability Discovery für Satellites.
  4. Ermöglicht die automatische Erkennung von Fähigkeiten
  5. auf Satellites durch verschiedene Probe-Methoden.
  6. """
  7. from __future__ import annotations
  8. import asyncio
  9. import logging
  10. from abc import ABC, abstractmethod
  11. from dataclasses import dataclass, field
  12. from datetime import datetime, timedelta
  13. from typing import Any, Awaitable, Callable
  14. from trixy_core.satellite.capability.registry import (
  15. Capability,
  16. CapabilityInfo,
  17. CapabilityRegistry,
  18. CapabilityState,
  19. )
  20. @dataclass
  21. class DiscoveryConfig:
  22. """Konfiguration für Capability-Discovery."""
  23. probe_timeout: float = 10.0 # Timeout für einzelne Probe
  24. probe_concurrency: int = 5 # Parallele Probes
  25. retry_failed: bool = True # Fehlgeschlagene erneut prüfen
  26. retry_delay: float = 60.0 # Verzögerung vor Retry
  27. rediscover_interval: float = 3600.0 # Interval für erneute Discovery
  28. auto_discover: bool = True # Automatische Discovery
  29. @dataclass
  30. class ProbeResult:
  31. """Ergebnis einer Capability-Probe."""
  32. capability_id: str
  33. success: bool
  34. state: CapabilityState
  35. error: str = ""
  36. parameters: dict[str, Any] = field(default_factory=dict)
  37. probe_time_ms: float = 0.0
  38. timestamp: datetime = field(default_factory=datetime.now)
  39. class CapabilityProbe(ABC):
  40. """
  41. Abstrakte Basis für Capability-Probes.
  42. Eine Probe testet ob eine bestimmte Capability
  43. auf einem Satellite verfügbar ist.
  44. """
  45. def __init__(self, capability_info: CapabilityInfo) -> None:
  46. """
  47. Initialisiert die Probe.
  48. Args:
  49. capability_info: Capability die geprüft wird
  50. """
  51. self.capability_info = capability_info
  52. @property
  53. def capability_id(self) -> str:
  54. """ID der zu prüfenden Capability."""
  55. return self.capability_info.id
  56. @abstractmethod
  57. async def probe(
  58. self,
  59. satellite_id: str,
  60. send_func: Callable[[str, Any], Awaitable[Any]],
  61. ) -> ProbeResult:
  62. """
  63. Führt die Probe durch.
  64. Args:
  65. satellite_id: Zu prüfender Satellite
  66. send_func: Funktion zum Senden von Nachrichten
  67. Returns:
  68. Probe-Ergebnis
  69. """
  70. pass
  71. class SimpleProbe(CapabilityProbe):
  72. """
  73. Einfache Probe die auf Ping-Antwort prüft.
  74. Sendet eine Capability-Query und erwartet eine
  75. Bestätigung vom Satellite.
  76. """
  77. def __init__(
  78. self,
  79. capability_info: CapabilityInfo,
  80. query_command: str = "capability_check",
  81. ) -> None:
  82. super().__init__(capability_info)
  83. self._query_command = query_command
  84. async def probe(
  85. self,
  86. satellite_id: str,
  87. send_func: Callable[[str, Any], Awaitable[Any]],
  88. ) -> ProbeResult:
  89. """Prüft Capability durch Query."""
  90. start = datetime.now()
  91. try:
  92. response = await send_func(
  93. satellite_id,
  94. {
  95. "command": self._query_command,
  96. "capability_id": self.capability_id,
  97. }
  98. )
  99. elapsed_ms = (datetime.now() - start).total_seconds() * 1000
  100. if response and response.get("available"):
  101. return ProbeResult(
  102. capability_id=self.capability_id,
  103. success=True,
  104. state=CapabilityState.AVAILABLE,
  105. parameters=response.get("parameters", {}),
  106. probe_time_ms=elapsed_ms,
  107. )
  108. else:
  109. return ProbeResult(
  110. capability_id=self.capability_id,
  111. success=True,
  112. state=CapabilityState.UNAVAILABLE,
  113. error=response.get("error", ""),
  114. probe_time_ms=elapsed_ms,
  115. )
  116. except Exception as e:
  117. elapsed_ms = (datetime.now() - start).total_seconds() * 1000
  118. return ProbeResult(
  119. capability_id=self.capability_id,
  120. success=False,
  121. state=CapabilityState.UNKNOWN,
  122. error=str(e),
  123. probe_time_ms=elapsed_ms,
  124. )
  125. class CapabilityDiscovery:
  126. """
  127. Service für Capability-Discovery.
  128. Verwaltet Probes und führt automatische Erkennung
  129. von Satellite-Capabilities durch.
  130. Beispiel:
  131. discovery = CapabilityDiscovery(
  132. registry=registry,
  133. send_func=send_message,
  134. )
  135. # Probe registrieren
  136. discovery.register_probe(AudioInputProbe())
  137. # Discovery für Satellite durchführen
  138. results = await discovery.discover_satellite("sat-001")
  139. # Alle Satellites prüfen
  140. await discovery.discover_all()
  141. """
  142. def __init__(
  143. self,
  144. registry: CapabilityRegistry,
  145. send_func: Callable[[str, Any], Awaitable[Any]],
  146. config: DiscoveryConfig | None = None,
  147. logger: logging.Logger | None = None,
  148. ) -> None:
  149. """
  150. Initialisiert den Discovery-Service.
  151. Args:
  152. registry: Capability-Registry
  153. send_func: Funktion zum Senden an Satellites
  154. config: Discovery-Konfiguration
  155. logger: Logger-Instanz
  156. """
  157. self.registry = registry
  158. self._send_func = send_func
  159. self.config = config or DiscoveryConfig()
  160. self.logger = logger or logging.getLogger(__name__)
  161. self._probes: dict[str, CapabilityProbe] = {}
  162. self._last_discovery: dict[str, datetime] = {}
  163. self._discovery_task: asyncio.Task | None = None
  164. self._semaphore = asyncio.Semaphore(self.config.probe_concurrency)
  165. @property
  166. def probe_count(self) -> int:
  167. """Anzahl registrierter Probes."""
  168. return len(self._probes)
  169. def register_probe(self, probe: CapabilityProbe) -> None:
  170. """
  171. Registriert eine Capability-Probe.
  172. Args:
  173. probe: Probe-Instanz
  174. """
  175. self._probes[probe.capability_id] = probe
  176. # Capability-Definition registrieren
  177. self.registry.define_capability(probe.capability_info)
  178. self.logger.debug(f"Probe registriert: {probe.capability_id}")
  179. def unregister_probe(self, capability_id: str) -> bool:
  180. """Entfernt eine Probe."""
  181. if capability_id in self._probes:
  182. del self._probes[capability_id]
  183. return True
  184. return False
  185. async def _run_probe(
  186. self,
  187. satellite_id: str,
  188. probe: CapabilityProbe,
  189. ) -> ProbeResult:
  190. """Führt einzelne Probe mit Timeout aus."""
  191. async with self._semaphore:
  192. try:
  193. return await asyncio.wait_for(
  194. probe.probe(satellite_id, self._send_func),
  195. timeout=self.config.probe_timeout,
  196. )
  197. except asyncio.TimeoutError:
  198. return ProbeResult(
  199. capability_id=probe.capability_id,
  200. success=False,
  201. state=CapabilityState.UNKNOWN,
  202. error="Timeout",
  203. probe_time_ms=self.config.probe_timeout * 1000,
  204. )
  205. async def discover_capability(
  206. self,
  207. satellite_id: str,
  208. capability_id: str,
  209. ) -> ProbeResult | None:
  210. """
  211. Prüft einzelne Capability auf Satellite.
  212. Args:
  213. satellite_id: Satellite-ID
  214. capability_id: Capability-ID
  215. Returns:
  216. Probe-Ergebnis oder None wenn keine Probe
  217. """
  218. probe = self._probes.get(capability_id)
  219. if not probe:
  220. return None
  221. result = await self._run_probe(satellite_id, probe)
  222. # Registry aktualisieren
  223. if result.success:
  224. cap = Capability(
  225. info=probe.capability_info,
  226. state=result.state,
  227. parameters=result.parameters,
  228. )
  229. self.registry.register(satellite_id, cap)
  230. return result
  231. async def discover_satellite(
  232. self,
  233. satellite_id: str,
  234. capabilities: list[str] | None = None,
  235. ) -> list[ProbeResult]:
  236. """
  237. Führt Discovery für einen Satellite durch.
  238. Args:
  239. satellite_id: Satellite-ID
  240. capabilities: Optional: Spezifische Capabilities
  241. Returns:
  242. Liste der Probe-Ergebnisse
  243. """
  244. probes_to_run = []
  245. if capabilities:
  246. probes_to_run = [
  247. self._probes[cid] for cid in capabilities
  248. if cid in self._probes
  249. ]
  250. else:
  251. probes_to_run = list(self._probes.values())
  252. if not probes_to_run:
  253. return []
  254. self.logger.info(
  255. f"Starte Discovery für {satellite_id} "
  256. f"({len(probes_to_run)} Probes)"
  257. )
  258. # Probes parallel ausführen
  259. tasks = [
  260. self._run_probe(satellite_id, probe)
  261. for probe in probes_to_run
  262. ]
  263. results = await asyncio.gather(*tasks, return_exceptions=True)
  264. # Ergebnisse verarbeiten
  265. probe_results: list[ProbeResult] = []
  266. for i, result in enumerate(results):
  267. probe = probes_to_run[i]
  268. if isinstance(result, Exception):
  269. probe_result = ProbeResult(
  270. capability_id=probe.capability_id,
  271. success=False,
  272. state=CapabilityState.UNKNOWN,
  273. error=str(result),
  274. )
  275. else:
  276. probe_result = result
  277. probe_results.append(probe_result)
  278. # Registry aktualisieren
  279. if probe_result.success:
  280. cap = Capability(
  281. info=probe.capability_info,
  282. state=probe_result.state,
  283. parameters=probe_result.parameters,
  284. )
  285. self.registry.register(satellite_id, cap)
  286. # Discovery-Zeit speichern
  287. self._last_discovery[satellite_id] = datetime.now()
  288. # Statistik loggen
  289. available = sum(
  290. 1 for r in probe_results
  291. if r.state == CapabilityState.AVAILABLE
  292. )
  293. self.logger.info(
  294. f"Discovery für {satellite_id} abgeschlossen: "
  295. f"{available}/{len(probe_results)} verfügbar"
  296. )
  297. return probe_results
  298. async def discover_all(
  299. self,
  300. satellite_ids: list[str],
  301. force: bool = False,
  302. ) -> dict[str, list[ProbeResult]]:
  303. """
  304. Führt Discovery für mehrere Satellites durch.
  305. Args:
  306. satellite_ids: Liste von Satellite-IDs
  307. force: Erneute Discovery auch wenn kürzlich durchgeführt
  308. Returns:
  309. Dict mit Ergebnissen pro Satellite
  310. """
  311. results: dict[str, list[ProbeResult]] = {}
  312. for sat_id in satellite_ids:
  313. # Prüfen ob erneute Discovery nötig
  314. if not force and sat_id in self._last_discovery:
  315. elapsed = (
  316. datetime.now() - self._last_discovery[sat_id]
  317. ).total_seconds()
  318. if elapsed < self.config.rediscover_interval:
  319. continue
  320. results[sat_id] = await self.discover_satellite(sat_id)
  321. return results
  322. async def start_auto_discovery(
  323. self,
  324. get_satellites_func: Callable[[], list[str]],
  325. ) -> None:
  326. """
  327. Startet automatische periodische Discovery.
  328. Args:
  329. get_satellites_func: Funktion die Satellite-IDs liefert
  330. """
  331. if self._discovery_task is not None:
  332. return
  333. async def discovery_loop():
  334. while True:
  335. try:
  336. satellites = get_satellites_func()
  337. await self.discover_all(satellites)
  338. except asyncio.CancelledError:
  339. break
  340. except Exception as e:
  341. self.logger.error(f"Auto-Discovery Fehler: {e}")
  342. await asyncio.sleep(self.config.rediscover_interval)
  343. self._discovery_task = asyncio.create_task(discovery_loop())
  344. async def stop_auto_discovery(self) -> None:
  345. """Stoppt automatische Discovery."""
  346. if self._discovery_task:
  347. self._discovery_task.cancel()
  348. try:
  349. await self._discovery_task
  350. except asyncio.CancelledError:
  351. pass
  352. self._discovery_task = None
  353. def needs_rediscovery(self, satellite_id: str) -> bool:
  354. """Prüft ob Satellite erneute Discovery benötigt."""
  355. if satellite_id not in self._last_discovery:
  356. return True
  357. elapsed = (
  358. datetime.now() - self._last_discovery[satellite_id]
  359. ).total_seconds()
  360. return elapsed >= self.config.rediscover_interval
  361. def get_stats(self) -> dict[str, Any]:
  362. """Liefert Discovery-Statistiken."""
  363. return {
  364. "probe_count": self.probe_count,
  365. "discovered_satellites": len(self._last_discovery),
  366. "auto_discovery_running": self._discovery_task is not None,
  367. "config": {
  368. "probe_timeout": self.config.probe_timeout,
  369. "probe_concurrency": self.config.probe_concurrency,
  370. "rediscover_interval": self.config.rediscover_interval,
  371. },
  372. }