collector.py 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407
  1. # -*- coding: utf-8 -*-
  2. """
  3. Metriken-Collector für Satellites.
  4. Sammelt verschiedene Metriken wie Latenz, Durchsatz,
  5. Fehler und Ressourcennutzung.
  6. """
  7. from __future__ import annotations
  8. import logging
  9. from collections import deque
  10. from dataclasses import dataclass, field
  11. from datetime import datetime, timedelta
  12. from enum import Enum, auto
  13. from typing import Any, Callable
  14. class MetricType(Enum):
  15. """Typ der Metrik."""
  16. COUNTER = auto() # Zähler (nur steigend)
  17. GAUGE = auto() # Momentanwert
  18. HISTOGRAM = auto() # Verteilung
  19. TIMING = auto() # Zeitmessung
  20. @dataclass
  21. class MetricValue:
  22. """Ein einzelner Metrik-Wert."""
  23. name: str
  24. value: float
  25. timestamp: datetime = field(default_factory=datetime.now)
  26. metric_type: MetricType = MetricType.GAUGE
  27. labels: dict[str, str] = field(default_factory=dict)
  28. unit: str = ""
  29. def to_dict(self) -> dict[str, Any]:
  30. """Konvertiert zu Dictionary."""
  31. return {
  32. "name": self.name,
  33. "value": self.value,
  34. "timestamp": self.timestamp.isoformat(),
  35. "type": self.metric_type.name,
  36. "labels": self.labels,
  37. "unit": self.unit,
  38. }
  39. @dataclass
  40. class SatelliteMetrics:
  41. """
  42. Gesammelte Metriken für einen Satellite.
  43. Enthält verschiedene Metrik-Kategorien und
  44. deren historische Werte.
  45. """
  46. satellite_id: str
  47. last_updated: datetime = field(default_factory=datetime.now)
  48. # Netzwerk
  49. latency_ms: float = 0.0 # Aktuelle Latenz
  50. packet_loss_percent: float = 0.0 # Paketverlusr
  51. bandwidth_kbps: float = 0.0 # Bandbreite
  52. messages_sent: int = 0 # Gesendete Nachrichten
  53. messages_received: int = 0 # Empfangene Nachrichten
  54. bytes_sent: int = 0 # Gesendete Bytes
  55. bytes_received: int = 0 # Empfangene Bytes
  56. # Verbindung
  57. connection_uptime_seconds: float = 0.0
  58. reconnect_count: int = 0
  59. last_heartbeat: datetime | None = None
  60. # Audio (falls relevant)
  61. audio_level_db: float = 0.0
  62. audio_samples_processed: int = 0
  63. # Errors
  64. error_count: int = 0
  65. last_error: str = ""
  66. last_error_time: datetime | None = None
  67. # Historische Daten (für Trends)
  68. latency_history: list[tuple[datetime, float]] = field(default_factory=list)
  69. def record_latency(self, latency_ms: float) -> None:
  70. """Zeichnet Latenz-Messung auf."""
  71. self.latency_ms = latency_ms
  72. self.last_updated = datetime.now()
  73. self.latency_history.append((self.last_updated, latency_ms))
  74. # Alte Einträge entfernen (max 1000)
  75. if len(self.latency_history) > 1000:
  76. self.latency_history = self.latency_history[-1000:]
  77. def record_error(self, error: str) -> None:
  78. """Zeichnet Fehler auf."""
  79. self.error_count += 1
  80. self.last_error = error
  81. self.last_error_time = datetime.now()
  82. self.last_updated = datetime.now()
  83. def get_avg_latency(self, seconds: int = 300) -> float:
  84. """Berechnet durchschnittliche Latenz der letzten Sekunden."""
  85. cutoff = datetime.now() - timedelta(seconds=seconds)
  86. recent = [v for t, v in self.latency_history if t >= cutoff]
  87. return sum(recent) / len(recent) if recent else 0.0
  88. def to_dict(self) -> dict[str, Any]:
  89. """Konvertiert zu Dictionary."""
  90. return {
  91. "satellite_id": self.satellite_id,
  92. "last_updated": self.last_updated.isoformat(),
  93. "network": {
  94. "latency_ms": self.latency_ms,
  95. "packet_loss_percent": self.packet_loss_percent,
  96. "bandwidth_kbps": self.bandwidth_kbps,
  97. "messages_sent": self.messages_sent,
  98. "messages_received": self.messages_received,
  99. "bytes_sent": self.bytes_sent,
  100. "bytes_received": self.bytes_received,
  101. },
  102. "connection": {
  103. "uptime_seconds": self.connection_uptime_seconds,
  104. "reconnect_count": self.reconnect_count,
  105. "last_heartbeat": (
  106. self.last_heartbeat.isoformat()
  107. if self.last_heartbeat else None
  108. ),
  109. },
  110. "errors": {
  111. "count": self.error_count,
  112. "last_error": self.last_error,
  113. "last_error_time": (
  114. self.last_error_time.isoformat()
  115. if self.last_error_time else None
  116. ),
  117. },
  118. }
  119. class SatelliteMetricsCollector:
  120. """
  121. Sammelt Metriken von Satellites.
  122. Verwaltet Metriken für mehrere Satellites und bietet
  123. Methoden für Aufzeichnung und Abfrage.
  124. Beispiel:
  125. collector = SatelliteMetricsCollector()
  126. # Latenz aufzeichnen
  127. collector.record_latency("sat-001", 15.5)
  128. # Nachricht aufzeichnen
  129. collector.record_message_sent("sat-001", 256)
  130. # Metriken abrufen
  131. metrics = collector.get_metrics("sat-001")
  132. """
  133. def __init__(
  134. self,
  135. history_size: int = 1000,
  136. logger: logging.Logger | None = None,
  137. ) -> None:
  138. """
  139. Initialisiert den Collector.
  140. Args:
  141. history_size: Maximale Anzahl historischer Werte
  142. logger: Logger-Instanz
  143. """
  144. self._history_size = history_size
  145. self.logger = logger or logging.getLogger(__name__)
  146. self._metrics: dict[str, SatelliteMetrics] = {}
  147. self._custom_metrics: dict[str, deque[MetricValue]] = {}
  148. self._callbacks: list[Callable[[str, MetricValue], None]] = []
  149. @property
  150. def satellite_count(self) -> int:
  151. """Anzahl überwachter Satellites."""
  152. return len(self._metrics)
  153. def _get_or_create(self, satellite_id: str) -> SatelliteMetrics:
  154. """Holt oder erstellt Metriken für Satellite."""
  155. if satellite_id not in self._metrics:
  156. self._metrics[satellite_id] = SatelliteMetrics(
  157. satellite_id=satellite_id
  158. )
  159. return self._metrics[satellite_id]
  160. def _notify(self, satellite_id: str, value: MetricValue) -> None:
  161. """Benachrichtigt registrierte Callbacks."""
  162. for callback in self._callbacks:
  163. try:
  164. callback(satellite_id, value)
  165. except Exception as e:
  166. self.logger.error(f"Callback-Fehler: {e}")
  167. def on_metric(
  168. self,
  169. callback: Callable[[str, MetricValue], None],
  170. ) -> None:
  171. """Registriert Callback für neue Metriken."""
  172. self._callbacks.append(callback)
  173. def record_latency(self, satellite_id: str, latency_ms: float) -> None:
  174. """
  175. Zeichnet Latenz-Messung auf.
  176. Args:
  177. satellite_id: Satellite-ID
  178. latency_ms: Latenz in Millisekunden
  179. """
  180. metrics = self._get_or_create(satellite_id)
  181. metrics.record_latency(latency_ms)
  182. value = MetricValue(
  183. name="latency",
  184. value=latency_ms,
  185. metric_type=MetricType.TIMING,
  186. unit="ms",
  187. )
  188. self._notify(satellite_id, value)
  189. def record_message_sent(
  190. self,
  191. satellite_id: str,
  192. bytes_count: int,
  193. ) -> None:
  194. """Zeichnet gesendete Nachricht auf."""
  195. metrics = self._get_or_create(satellite_id)
  196. metrics.messages_sent += 1
  197. metrics.bytes_sent += bytes_count
  198. metrics.last_updated = datetime.now()
  199. def record_message_received(
  200. self,
  201. satellite_id: str,
  202. bytes_count: int,
  203. ) -> None:
  204. """Zeichnet empfangene Nachricht auf."""
  205. metrics = self._get_or_create(satellite_id)
  206. metrics.messages_received += 1
  207. metrics.bytes_received += bytes_count
  208. metrics.last_updated = datetime.now()
  209. def record_heartbeat(self, satellite_id: str) -> None:
  210. """Zeichnet Heartbeat auf."""
  211. metrics = self._get_or_create(satellite_id)
  212. metrics.last_heartbeat = datetime.now()
  213. metrics.last_updated = datetime.now()
  214. def record_reconnect(self, satellite_id: str) -> None:
  215. """Zeichnet Reconnect auf."""
  216. metrics = self._get_or_create(satellite_id)
  217. metrics.reconnect_count += 1
  218. metrics.connection_uptime_seconds = 0.0
  219. metrics.last_updated = datetime.now()
  220. def record_error(self, satellite_id: str, error: str) -> None:
  221. """Zeichnet Fehler auf."""
  222. metrics = self._get_or_create(satellite_id)
  223. metrics.record_error(error)
  224. value = MetricValue(
  225. name="error",
  226. value=1,
  227. metric_type=MetricType.COUNTER,
  228. labels={"error": error[:100]}, # Begrenzen
  229. )
  230. self._notify(satellite_id, value)
  231. def record_packet_loss(
  232. self,
  233. satellite_id: str,
  234. loss_percent: float,
  235. ) -> None:
  236. """Zeichnet Paketverlust auf."""
  237. metrics = self._get_or_create(satellite_id)
  238. metrics.packet_loss_percent = loss_percent
  239. metrics.last_updated = datetime.now()
  240. def record_bandwidth(
  241. self,
  242. satellite_id: str,
  243. bandwidth_kbps: float,
  244. ) -> None:
  245. """Zeichnet Bandbreite auf."""
  246. metrics = self._get_or_create(satellite_id)
  247. metrics.bandwidth_kbps = bandwidth_kbps
  248. metrics.last_updated = datetime.now()
  249. def record_custom(
  250. self,
  251. satellite_id: str,
  252. name: str,
  253. value: float,
  254. metric_type: MetricType = MetricType.GAUGE,
  255. labels: dict[str, str] | None = None,
  256. unit: str = "",
  257. ) -> None:
  258. """
  259. Zeichnet benutzerdefinierte Metrik auf.
  260. Args:
  261. satellite_id: Satellite-ID
  262. name: Metrik-Name
  263. value: Wert
  264. metric_type: Metrik-Typ
  265. labels: Optionale Labels
  266. unit: Einheit
  267. """
  268. metric_value = MetricValue(
  269. name=name,
  270. value=value,
  271. metric_type=metric_type,
  272. labels=labels or {},
  273. unit=unit,
  274. )
  275. key = f"{satellite_id}:{name}"
  276. if key not in self._custom_metrics:
  277. self._custom_metrics[key] = deque(maxlen=self._history_size)
  278. self._custom_metrics[key].append(metric_value)
  279. self._notify(satellite_id, metric_value)
  280. def get_metrics(self, satellite_id: str) -> SatelliteMetrics | None:
  281. """Holt Metriken für Satellite."""
  282. return self._metrics.get(satellite_id)
  283. def get_all_metrics(self) -> dict[str, SatelliteMetrics]:
  284. """Holt alle Metriken."""
  285. return self._metrics.copy()
  286. def get_custom_metric(
  287. self,
  288. satellite_id: str,
  289. name: str,
  290. ) -> list[MetricValue]:
  291. """Holt benutzerdefinierte Metrik-Historie."""
  292. key = f"{satellite_id}:{name}"
  293. if key in self._custom_metrics:
  294. return list(self._custom_metrics[key])
  295. return []
  296. def remove_satellite(self, satellite_id: str) -> bool:
  297. """Entfernt Metriken für Satellite."""
  298. if satellite_id in self._metrics:
  299. del self._metrics[satellite_id]
  300. # Custom Metriken entfernen
  301. to_remove = [
  302. k for k in self._custom_metrics
  303. if k.startswith(f"{satellite_id}:")
  304. ]
  305. for key in to_remove:
  306. del self._custom_metrics[key]
  307. return True
  308. return False
  309. def get_summary(self) -> dict[str, Any]:
  310. """Liefert Zusammenfassung aller Metriken."""
  311. if not self._metrics:
  312. return {
  313. "satellite_count": 0,
  314. "avg_latency_ms": 0,
  315. "total_messages": 0,
  316. "total_errors": 0,
  317. }
  318. avg_latency = sum(
  319. m.latency_ms for m in self._metrics.values()
  320. ) / len(self._metrics)
  321. total_messages = sum(
  322. m.messages_sent + m.messages_received
  323. for m in self._metrics.values()
  324. )
  325. total_errors = sum(
  326. m.error_count for m in self._metrics.values()
  327. )
  328. return {
  329. "satellite_count": self.satellite_count,
  330. "avg_latency_ms": round(avg_latency, 2),
  331. "total_messages": total_messages,
  332. "total_errors": total_errors,
  333. "satellites": {
  334. sat_id: {
  335. "latency_ms": m.latency_ms,
  336. "messages": m.messages_sent + m.messages_received,
  337. "errors": m.error_count,
  338. }
  339. for sat_id, m in self._metrics.items()
  340. },
  341. }
  342. def clear_all(self) -> None:
  343. """Löscht alle Metriken."""
  344. self._metrics.clear()
  345. self._custom_metrics.clear()