| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452 |
- # -*- coding: utf-8 -*-
- """
- Capability Discovery für Satellites.
- Ermöglicht die automatische Erkennung von Fähigkeiten
- auf Satellites durch verschiedene Probe-Methoden.
- """
- from __future__ import annotations
- import asyncio
- import logging
- from abc import ABC, abstractmethod
- from dataclasses import dataclass, field
- from datetime import datetime, timedelta
- from typing import Any, Awaitable, Callable
- from trixy_core.satellite.capability.registry import (
- Capability,
- CapabilityInfo,
- CapabilityRegistry,
- CapabilityState,
- )
- @dataclass
- class DiscoveryConfig:
- """Konfiguration für Capability-Discovery."""
- probe_timeout: float = 10.0 # Timeout für einzelne Probe
- probe_concurrency: int = 5 # Parallele Probes
- retry_failed: bool = True # Fehlgeschlagene erneut prüfen
- retry_delay: float = 60.0 # Verzögerung vor Retry
- rediscover_interval: float = 3600.0 # Interval für erneute Discovery
- auto_discover: bool = True # Automatische Discovery
- @dataclass
- class ProbeResult:
- """Ergebnis einer Capability-Probe."""
- capability_id: str
- success: bool
- state: CapabilityState
- error: str = ""
- parameters: dict[str, Any] = field(default_factory=dict)
- probe_time_ms: float = 0.0
- timestamp: datetime = field(default_factory=datetime.now)
- class CapabilityProbe(ABC):
- """
- Abstrakte Basis für Capability-Probes.
- Eine Probe testet ob eine bestimmte Capability
- auf einem Satellite verfügbar ist.
- """
- def __init__(self, capability_info: CapabilityInfo) -> None:
- """
- Initialisiert die Probe.
- Args:
- capability_info: Capability die geprüft wird
- """
- self.capability_info = capability_info
- @property
- def capability_id(self) -> str:
- """ID der zu prüfenden Capability."""
- return self.capability_info.id
- @abstractmethod
- async def probe(
- self,
- satellite_id: str,
- send_func: Callable[[str, Any], Awaitable[Any]],
- ) -> ProbeResult:
- """
- Führt die Probe durch.
- Args:
- satellite_id: Zu prüfender Satellite
- send_func: Funktion zum Senden von Nachrichten
- Returns:
- Probe-Ergebnis
- """
- pass
- class SimpleProbe(CapabilityProbe):
- """
- Einfache Probe die auf Ping-Antwort prüft.
- Sendet eine Capability-Query und erwartet eine
- Bestätigung vom Satellite.
- """
- def __init__(
- self,
- capability_info: CapabilityInfo,
- query_command: str = "capability_check",
- ) -> None:
- super().__init__(capability_info)
- self._query_command = query_command
- async def probe(
- self,
- satellite_id: str,
- send_func: Callable[[str, Any], Awaitable[Any]],
- ) -> ProbeResult:
- """Prüft Capability durch Query."""
- start = datetime.now()
- try:
- response = await send_func(
- satellite_id,
- {
- "command": self._query_command,
- "capability_id": self.capability_id,
- }
- )
- elapsed_ms = (datetime.now() - start).total_seconds() * 1000
- if response and response.get("available"):
- return ProbeResult(
- capability_id=self.capability_id,
- success=True,
- state=CapabilityState.AVAILABLE,
- parameters=response.get("parameters", {}),
- probe_time_ms=elapsed_ms,
- )
- else:
- return ProbeResult(
- capability_id=self.capability_id,
- success=True,
- state=CapabilityState.UNAVAILABLE,
- error=response.get("error", ""),
- probe_time_ms=elapsed_ms,
- )
- except Exception as e:
- elapsed_ms = (datetime.now() - start).total_seconds() * 1000
- return ProbeResult(
- capability_id=self.capability_id,
- success=False,
- state=CapabilityState.UNKNOWN,
- error=str(e),
- probe_time_ms=elapsed_ms,
- )
- class CapabilityDiscovery:
- """
- Service für Capability-Discovery.
- Verwaltet Probes und führt automatische Erkennung
- von Satellite-Capabilities durch.
- Beispiel:
- discovery = CapabilityDiscovery(
- registry=registry,
- send_func=send_message,
- )
- # Probe registrieren
- discovery.register_probe(AudioInputProbe())
- # Discovery für Satellite durchführen
- results = await discovery.discover_satellite("sat-001")
- # Alle Satellites prüfen
- await discovery.discover_all()
- """
- def __init__(
- self,
- registry: CapabilityRegistry,
- send_func: Callable[[str, Any], Awaitable[Any]],
- config: DiscoveryConfig | None = None,
- logger: logging.Logger | None = None,
- ) -> None:
- """
- Initialisiert den Discovery-Service.
- Args:
- registry: Capability-Registry
- send_func: Funktion zum Senden an Satellites
- config: Discovery-Konfiguration
- logger: Logger-Instanz
- """
- self.registry = registry
- self._send_func = send_func
- self.config = config or DiscoveryConfig()
- self.logger = logger or logging.getLogger(__name__)
- self._probes: dict[str, CapabilityProbe] = {}
- self._last_discovery: dict[str, datetime] = {}
- self._discovery_task: asyncio.Task | None = None
- self._semaphore = asyncio.Semaphore(self.config.probe_concurrency)
- @property
- def probe_count(self) -> int:
- """Anzahl registrierter Probes."""
- return len(self._probes)
- def register_probe(self, probe: CapabilityProbe) -> None:
- """
- Registriert eine Capability-Probe.
- Args:
- probe: Probe-Instanz
- """
- self._probes[probe.capability_id] = probe
- # Capability-Definition registrieren
- self.registry.define_capability(probe.capability_info)
- self.logger.debug(f"Probe registriert: {probe.capability_id}")
- def unregister_probe(self, capability_id: str) -> bool:
- """Entfernt eine Probe."""
- if capability_id in self._probes:
- del self._probes[capability_id]
- return True
- return False
- async def _run_probe(
- self,
- satellite_id: str,
- probe: CapabilityProbe,
- ) -> ProbeResult:
- """Führt einzelne Probe mit Timeout aus."""
- async with self._semaphore:
- try:
- return await asyncio.wait_for(
- probe.probe(satellite_id, self._send_func),
- timeout=self.config.probe_timeout,
- )
- except asyncio.TimeoutError:
- return ProbeResult(
- capability_id=probe.capability_id,
- success=False,
- state=CapabilityState.UNKNOWN,
- error="Timeout",
- probe_time_ms=self.config.probe_timeout * 1000,
- )
- async def discover_capability(
- self,
- satellite_id: str,
- capability_id: str,
- ) -> ProbeResult | None:
- """
- Prüft einzelne Capability auf Satellite.
- Args:
- satellite_id: Satellite-ID
- capability_id: Capability-ID
- Returns:
- Probe-Ergebnis oder None wenn keine Probe
- """
- probe = self._probes.get(capability_id)
- if not probe:
- return None
- result = await self._run_probe(satellite_id, probe)
- # Registry aktualisieren
- if result.success:
- cap = Capability(
- info=probe.capability_info,
- state=result.state,
- parameters=result.parameters,
- )
- self.registry.register(satellite_id, cap)
- return result
- async def discover_satellite(
- self,
- satellite_id: str,
- capabilities: list[str] | None = None,
- ) -> list[ProbeResult]:
- """
- Führt Discovery für einen Satellite durch.
- Args:
- satellite_id: Satellite-ID
- capabilities: Optional: Spezifische Capabilities
- Returns:
- Liste der Probe-Ergebnisse
- """
- probes_to_run = []
- if capabilities:
- probes_to_run = [
- self._probes[cid] for cid in capabilities
- if cid in self._probes
- ]
- else:
- probes_to_run = list(self._probes.values())
- if not probes_to_run:
- return []
- self.logger.info(
- f"Starte Discovery für {satellite_id} "
- f"({len(probes_to_run)} Probes)"
- )
- # Probes parallel ausführen
- tasks = [
- self._run_probe(satellite_id, probe)
- for probe in probes_to_run
- ]
- results = await asyncio.gather(*tasks, return_exceptions=True)
- # Ergebnisse verarbeiten
- probe_results: list[ProbeResult] = []
- for i, result in enumerate(results):
- probe = probes_to_run[i]
- if isinstance(result, Exception):
- probe_result = ProbeResult(
- capability_id=probe.capability_id,
- success=False,
- state=CapabilityState.UNKNOWN,
- error=str(result),
- )
- else:
- probe_result = result
- probe_results.append(probe_result)
- # Registry aktualisieren
- if probe_result.success:
- cap = Capability(
- info=probe.capability_info,
- state=probe_result.state,
- parameters=probe_result.parameters,
- )
- self.registry.register(satellite_id, cap)
- # Discovery-Zeit speichern
- self._last_discovery[satellite_id] = datetime.now()
- # Statistik loggen
- available = sum(
- 1 for r in probe_results
- if r.state == CapabilityState.AVAILABLE
- )
- self.logger.info(
- f"Discovery für {satellite_id} abgeschlossen: "
- f"{available}/{len(probe_results)} verfügbar"
- )
- return probe_results
- async def discover_all(
- self,
- satellite_ids: list[str],
- force: bool = False,
- ) -> dict[str, list[ProbeResult]]:
- """
- Führt Discovery für mehrere Satellites durch.
- Args:
- satellite_ids: Liste von Satellite-IDs
- force: Erneute Discovery auch wenn kürzlich durchgeführt
- Returns:
- Dict mit Ergebnissen pro Satellite
- """
- results: dict[str, list[ProbeResult]] = {}
- for sat_id in satellite_ids:
- # Prüfen ob erneute Discovery nötig
- if not force and sat_id in self._last_discovery:
- elapsed = (
- datetime.now() - self._last_discovery[sat_id]
- ).total_seconds()
- if elapsed < self.config.rediscover_interval:
- continue
- results[sat_id] = await self.discover_satellite(sat_id)
- return results
- async def start_auto_discovery(
- self,
- get_satellites_func: Callable[[], list[str]],
- ) -> None:
- """
- Startet automatische periodische Discovery.
- Args:
- get_satellites_func: Funktion die Satellite-IDs liefert
- """
- if self._discovery_task is not None:
- return
- async def discovery_loop():
- while True:
- try:
- satellites = get_satellites_func()
- await self.discover_all(satellites)
- except asyncio.CancelledError:
- break
- except Exception as e:
- self.logger.error(f"Auto-Discovery Fehler: {e}")
- await asyncio.sleep(self.config.rediscover_interval)
- self._discovery_task = asyncio.create_task(discovery_loop())
- async def stop_auto_discovery(self) -> None:
- """Stoppt automatische Discovery."""
- if self._discovery_task:
- self._discovery_task.cancel()
- try:
- await self._discovery_task
- except asyncio.CancelledError:
- pass
- self._discovery_task = None
- def needs_rediscovery(self, satellite_id: str) -> bool:
- """Prüft ob Satellite erneute Discovery benötigt."""
- if satellite_id not in self._last_discovery:
- return True
- elapsed = (
- datetime.now() - self._last_discovery[satellite_id]
- ).total_seconds()
- return elapsed >= self.config.rediscover_interval
- def get_stats(self) -> dict[str, Any]:
- """Liefert Discovery-Statistiken."""
- return {
- "probe_count": self.probe_count,
- "discovered_satellites": len(self._last_discovery),
- "auto_discovery_running": self._discovery_task is not None,
- "config": {
- "probe_timeout": self.config.probe_timeout,
- "probe_concurrency": self.config.probe_concurrency,
- "rediscover_interval": self.config.rediscover_interval,
- },
- }
|