throttle.py 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502
  1. # -*- coding: utf-8 -*-
  2. """
  3. Event-Throttling für Rate-Limiting.
  4. Begrenzt die Anzahl der Events pro Zeiteinheit.
  5. """
  6. from __future__ import annotations
  7. import asyncio
  8. import time
  9. from collections import defaultdict
  10. from dataclasses import dataclass, field
  11. from datetime import datetime
  12. from enum import Enum, auto
  13. from typing import TYPE_CHECKING, Any, Callable, Coroutine
  14. if TYPE_CHECKING:
  15. from trixy_core.events.event_data.base import EventData
  16. class ThrottleStrategy(Enum):
  17. """
  18. Strategien für Throttling.
  19. """
  20. DROP = auto()
  21. """Überschreitende Events verwerfen."""
  22. QUEUE = auto()
  23. """Überschreitende Events in Queue stellen."""
  24. DELAY = auto()
  25. """Überschreitende Events verzögern."""
  26. SAMPLE = auto()
  27. """Nur jedes n-te Event durchlassen."""
  28. @dataclass
  29. class ThrottleResult:
  30. """
  31. Ergebnis einer Throttle-Prüfung.
  32. """
  33. allowed: bool
  34. """Ob das Event durchgelassen wird."""
  35. event_name: str
  36. """Name des Events."""
  37. wait_time_ms: float = 0.0
  38. """Wartezeit bis zum nächsten erlaubten Event."""
  39. queued: bool = False
  40. """Ob das Event in die Queue gestellt wurde."""
  41. dropped: bool = False
  42. """Ob das Event verworfen wurde."""
  43. reason: str = ""
  44. """Grund für die Entscheidung."""
  45. @dataclass
  46. class ThrottleConfig:
  47. """
  48. Konfiguration für Event-Throttling.
  49. """
  50. events_per_second: float = 10.0
  51. """Maximale Events pro Sekunde."""
  52. burst_size: int = 5
  53. """Erlaubte Burst-Größe."""
  54. strategy: ThrottleStrategy = ThrottleStrategy.DROP
  55. """Strategie für überschreitende Events."""
  56. queue_size: int = 100
  57. """Maximale Queue-Größe (für QUEUE-Strategie)."""
  58. sample_rate: int = 1
  59. """Sample-Rate (für SAMPLE-Strategie)."""
  60. @dataclass
  61. class ThrottleState:
  62. """
  63. Zustand für einen Event-Typ.
  64. """
  65. tokens: float = 0.0
  66. """Verfügbare Tokens (Token-Bucket)."""
  67. last_update: float = field(default_factory=time.monotonic)
  68. """Zeitpunkt der letzten Aktualisierung."""
  69. event_count: int = 0
  70. """Gesamtzahl der Events."""
  71. dropped_count: int = 0
  72. """Anzahl verworfener Events."""
  73. queued_count: int = 0
  74. """Anzahl in Queue gestellter Events."""
  75. class EventThrottler:
  76. """
  77. Throttler für Event-Rate-Limiting.
  78. Begrenzt die Anzahl der Events pro Zeiteinheit unter
  79. Verwendung des Token-Bucket-Algorithmus.
  80. Example:
  81. throttler = EventThrottler()
  82. # Globale Rate
  83. throttler.set_global_rate(100) # 100 Events/Sekunde
  84. # Pro Event-Typ
  85. throttler.set_event_rate("mouse_move", 10) # 10/Sekunde
  86. throttler.set_event_rate("click", 100)
  87. # Prüfen und verarbeiten
  88. result = await throttler.check("mouse_move", event_data)
  89. if result.allowed:
  90. await process_event(event_data)
  91. """
  92. def __init__(
  93. self,
  94. default_config: ThrottleConfig | None = None,
  95. on_throttle: Callable[[str, "EventData"], Coroutine[Any, Any, None]] | None = None,
  96. ) -> None:
  97. """
  98. Initialisiert den Throttler.
  99. Args:
  100. default_config: Standard-Konfiguration.
  101. on_throttle: Callback wenn ein Event gedrosselt wird.
  102. """
  103. self._default_config = default_config or ThrottleConfig()
  104. self._on_throttle = on_throttle
  105. self._configs: dict[str, ThrottleConfig] = {}
  106. self._states: dict[str, ThrottleState] = defaultdict(ThrottleState)
  107. self._queues: dict[str, asyncio.Queue] = {}
  108. self._lock = asyncio.Lock()
  109. def set_global_rate(
  110. self,
  111. events_per_second: float,
  112. burst_size: int | None = None,
  113. ) -> None:
  114. """
  115. Setzt die globale Rate.
  116. Args:
  117. events_per_second: Maximale Events pro Sekunde.
  118. burst_size: Optionale Burst-Größe.
  119. """
  120. self._default_config.events_per_second = events_per_second
  121. if burst_size is not None:
  122. self._default_config.burst_size = burst_size
  123. def set_event_rate(
  124. self,
  125. event_name: str,
  126. events_per_second: float,
  127. burst_size: int | None = None,
  128. strategy: ThrottleStrategy | None = None,
  129. ) -> None:
  130. """
  131. Setzt die Rate für einen Event-Typ.
  132. Args:
  133. event_name: Name des Events.
  134. events_per_second: Maximale Events pro Sekunde.
  135. burst_size: Optionale Burst-Größe.
  136. strategy: Optionale Strategie.
  137. """
  138. config = ThrottleConfig(
  139. events_per_second=events_per_second,
  140. burst_size=burst_size or self._default_config.burst_size,
  141. strategy=strategy or self._default_config.strategy,
  142. queue_size=self._default_config.queue_size,
  143. sample_rate=self._default_config.sample_rate,
  144. )
  145. self._configs[event_name] = config
  146. def get_config(self, event_name: str) -> ThrottleConfig:
  147. """
  148. Gibt die Konfiguration für ein Event zurück.
  149. Args:
  150. event_name: Name des Events.
  151. Returns:
  152. ThrottleConfig.
  153. """
  154. return self._configs.get(event_name, self._default_config)
  155. def _update_tokens(
  156. self,
  157. state: ThrottleState,
  158. config: ThrottleConfig,
  159. ) -> None:
  160. """
  161. Aktualisiert die Token-Anzahl.
  162. Args:
  163. state: Der Zustand.
  164. config: Die Konfiguration.
  165. """
  166. now = time.monotonic()
  167. elapsed = now - state.last_update
  168. state.last_update = now
  169. # Tokens hinzufügen basierend auf verstrichener Zeit
  170. new_tokens = elapsed * config.events_per_second
  171. state.tokens = min(
  172. state.tokens + new_tokens,
  173. float(config.burst_size),
  174. )
  175. async def check(
  176. self,
  177. event_name: str,
  178. event_data: "EventData | None" = None,
  179. ) -> ThrottleResult:
  180. """
  181. Prüft, ob ein Event erlaubt ist.
  182. Args:
  183. event_name: Name des Events.
  184. event_data: Optionale Event-Daten.
  185. Returns:
  186. ThrottleResult mit der Entscheidung.
  187. """
  188. async with self._lock:
  189. config = self.get_config(event_name)
  190. state = self._states[event_name]
  191. # Tokens aktualisieren
  192. self._update_tokens(state, config)
  193. state.event_count += 1
  194. # Sample-Strategie prüfen
  195. if config.strategy == ThrottleStrategy.SAMPLE:
  196. if state.event_count % config.sample_rate != 0:
  197. return ThrottleResult(
  198. allowed=False,
  199. event_name=event_name,
  200. dropped=True,
  201. reason=f"Sampling (1/{config.sample_rate})",
  202. )
  203. # Token verfügbar?
  204. if state.tokens >= 1.0:
  205. state.tokens -= 1.0
  206. return ThrottleResult(
  207. allowed=True,
  208. event_name=event_name,
  209. )
  210. # Nicht genug Tokens - Strategie anwenden
  211. wait_time_ms = ((1.0 - state.tokens) / config.events_per_second) * 1000
  212. if config.strategy == ThrottleStrategy.DROP:
  213. state.dropped_count += 1
  214. if self._on_throttle and event_data:
  215. await self._on_throttle(event_name, event_data)
  216. return ThrottleResult(
  217. allowed=False,
  218. event_name=event_name,
  219. wait_time_ms=wait_time_ms,
  220. dropped=True,
  221. reason="Rate-Limit überschritten",
  222. )
  223. elif config.strategy == ThrottleStrategy.QUEUE:
  224. queue = self._get_queue(event_name, config.queue_size)
  225. if queue.full():
  226. state.dropped_count += 1
  227. return ThrottleResult(
  228. allowed=False,
  229. event_name=event_name,
  230. wait_time_ms=wait_time_ms,
  231. dropped=True,
  232. reason="Queue voll",
  233. )
  234. if event_data:
  235. await queue.put(event_data)
  236. state.queued_count += 1
  237. return ThrottleResult(
  238. allowed=False,
  239. event_name=event_name,
  240. wait_time_ms=wait_time_ms,
  241. queued=True,
  242. reason="In Queue gestellt",
  243. )
  244. elif config.strategy == ThrottleStrategy.DELAY:
  245. # Warten bis Token verfügbar
  246. await asyncio.sleep(wait_time_ms / 1000)
  247. state.tokens = 0.0 # Token wurde "gewartet"
  248. return ThrottleResult(
  249. allowed=True,
  250. event_name=event_name,
  251. wait_time_ms=wait_time_ms,
  252. reason="Verzögert",
  253. )
  254. return ThrottleResult(
  255. allowed=False,
  256. event_name=event_name,
  257. wait_time_ms=wait_time_ms,
  258. dropped=True,
  259. reason="Unbekannte Strategie",
  260. )
  261. def _get_queue(
  262. self,
  263. event_name: str,
  264. max_size: int,
  265. ) -> asyncio.Queue:
  266. """
  267. Gibt die Queue für ein Event zurück.
  268. Args:
  269. event_name: Name des Events.
  270. max_size: Maximale Queue-Größe.
  271. Returns:
  272. Die Queue.
  273. """
  274. if event_name not in self._queues:
  275. self._queues[event_name] = asyncio.Queue(maxsize=max_size)
  276. return self._queues[event_name]
  277. async def process_queued(
  278. self,
  279. event_name: str,
  280. handler: Callable[["EventData"], Coroutine[Any, Any, None]],
  281. batch_size: int = 10,
  282. ) -> int:
  283. """
  284. Verarbeitet Events aus der Queue.
  285. Args:
  286. event_name: Name des Events.
  287. handler: Handler-Funktion.
  288. batch_size: Maximale Batch-Größe.
  289. Returns:
  290. Anzahl verarbeiteter Events.
  291. """
  292. queue = self._queues.get(event_name)
  293. if not queue:
  294. return 0
  295. processed = 0
  296. while not queue.empty() and processed < batch_size:
  297. try:
  298. event_data = queue.get_nowait()
  299. await handler(event_data)
  300. processed += 1
  301. except asyncio.QueueEmpty:
  302. break
  303. return processed
  304. def get_stats(self, event_name: str) -> dict[str, Any]:
  305. """
  306. Gibt Statistiken für ein Event zurück.
  307. Args:
  308. event_name: Name des Events.
  309. Returns:
  310. Statistik-Dictionary.
  311. """
  312. state = self._states.get(event_name)
  313. if not state:
  314. return {}
  315. config = self.get_config(event_name)
  316. return {
  317. "event_name": event_name,
  318. "total_events": state.event_count,
  319. "dropped_events": state.dropped_count,
  320. "queued_events": state.queued_count,
  321. "drop_rate": (
  322. state.dropped_count / state.event_count
  323. if state.event_count > 0
  324. else 0.0
  325. ),
  326. "current_tokens": state.tokens,
  327. "max_tokens": config.burst_size,
  328. "rate_limit": config.events_per_second,
  329. }
  330. def get_all_stats(self) -> dict[str, dict[str, Any]]:
  331. """
  332. Gibt Statistiken für alle Events zurück.
  333. Returns:
  334. Dictionary von Event-Namen zu Statistiken.
  335. """
  336. return {
  337. name: self.get_stats(name)
  338. for name in self._states.keys()
  339. }
  340. def reset(self, event_name: str | None = None) -> None:
  341. """
  342. Setzt den Zustand zurück.
  343. Args:
  344. event_name: Optionaler Event-Name (None = alle).
  345. """
  346. if event_name:
  347. if event_name in self._states:
  348. del self._states[event_name]
  349. if event_name in self._queues:
  350. del self._queues[event_name]
  351. else:
  352. self._states.clear()
  353. self._queues.clear()
  354. class ThrottlingMiddleware:
  355. """
  356. Middleware-Integration für Event-Throttling.
  357. Kann in die Event-Pipeline integriert werden.
  358. Example:
  359. from trixy_core.events.middleware.base import EventMiddleware
  360. throttler = EventThrottler()
  361. throttler.set_event_rate("frequent_event", 10)
  362. # Als Middleware verwenden
  363. middleware = ThrottlingMiddleware(throttler)
  364. pipeline.add(middleware)
  365. """
  366. def __init__(
  367. self,
  368. throttler: EventThrottler,
  369. ) -> None:
  370. """
  371. Initialisiert die Middleware.
  372. Args:
  373. throttler: Der Throttler.
  374. """
  375. self._throttler = throttler
  376. self._name = "ThrottlingMiddleware"
  377. self._enabled = True
  378. @property
  379. def name(self) -> str:
  380. """Name der Middleware."""
  381. return self._name
  382. @property
  383. def enabled(self) -> bool:
  384. """Ob aktiv."""
  385. return self._enabled
  386. async def before(self, context: Any) -> Any:
  387. """Prüft Throttling vor Event-Verarbeitung."""
  388. from trixy_core.events.middleware.base import MiddlewareResult, MiddlewareAction
  389. result = await self._throttler.check(
  390. context.event_name,
  391. context.event_data,
  392. )
  393. if result.allowed:
  394. context.set_property("throttle_result", result)
  395. return MiddlewareResult()
  396. else:
  397. context.set_property("throttle_result", result)
  398. return MiddlewareResult(
  399. action=MiddlewareAction.SKIP_HANDLERS if result.queued else MiddlewareAction.ABORT,
  400. )
  401. async def after(self, context: Any) -> None:
  402. """Keine After-Verarbeitung."""
  403. pass