| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502 |
- # -*- coding: utf-8 -*-
- """
- Event-Throttling für Rate-Limiting.
- Begrenzt die Anzahl der Events pro Zeiteinheit.
- """
- from __future__ import annotations
- import asyncio
- import time
- from collections import defaultdict
- from dataclasses import dataclass, field
- from datetime import datetime
- from enum import Enum, auto
- from typing import TYPE_CHECKING, Any, Callable, Coroutine
- if TYPE_CHECKING:
- from trixy_core.events.event_data.base import EventData
- class ThrottleStrategy(Enum):
- """
- Strategien für Throttling.
- """
- DROP = auto()
- """Überschreitende Events verwerfen."""
- QUEUE = auto()
- """Überschreitende Events in Queue stellen."""
- DELAY = auto()
- """Überschreitende Events verzögern."""
- SAMPLE = auto()
- """Nur jedes n-te Event durchlassen."""
- @dataclass
- class ThrottleResult:
- """
- Ergebnis einer Throttle-Prüfung.
- """
- allowed: bool
- """Ob das Event durchgelassen wird."""
- event_name: str
- """Name des Events."""
- wait_time_ms: float = 0.0
- """Wartezeit bis zum nächsten erlaubten Event."""
- queued: bool = False
- """Ob das Event in die Queue gestellt wurde."""
- dropped: bool = False
- """Ob das Event verworfen wurde."""
- reason: str = ""
- """Grund für die Entscheidung."""
- @dataclass
- class ThrottleConfig:
- """
- Konfiguration für Event-Throttling.
- """
- events_per_second: float = 10.0
- """Maximale Events pro Sekunde."""
- burst_size: int = 5
- """Erlaubte Burst-Größe."""
- strategy: ThrottleStrategy = ThrottleStrategy.DROP
- """Strategie für überschreitende Events."""
- queue_size: int = 100
- """Maximale Queue-Größe (für QUEUE-Strategie)."""
- sample_rate: int = 1
- """Sample-Rate (für SAMPLE-Strategie)."""
- @dataclass
- class ThrottleState:
- """
- Zustand für einen Event-Typ.
- """
- tokens: float = 0.0
- """Verfügbare Tokens (Token-Bucket)."""
- last_update: float = field(default_factory=time.monotonic)
- """Zeitpunkt der letzten Aktualisierung."""
- event_count: int = 0
- """Gesamtzahl der Events."""
- dropped_count: int = 0
- """Anzahl verworfener Events."""
- queued_count: int = 0
- """Anzahl in Queue gestellter Events."""
- class EventThrottler:
- """
- Throttler für Event-Rate-Limiting.
- Begrenzt die Anzahl der Events pro Zeiteinheit unter
- Verwendung des Token-Bucket-Algorithmus.
- Example:
- throttler = EventThrottler()
- # Globale Rate
- throttler.set_global_rate(100) # 100 Events/Sekunde
- # Pro Event-Typ
- throttler.set_event_rate("mouse_move", 10) # 10/Sekunde
- throttler.set_event_rate("click", 100)
- # Prüfen und verarbeiten
- result = await throttler.check("mouse_move", event_data)
- if result.allowed:
- await process_event(event_data)
- """
- def __init__(
- self,
- default_config: ThrottleConfig | None = None,
- on_throttle: Callable[[str, "EventData"], Coroutine[Any, Any, None]] | None = None,
- ) -> None:
- """
- Initialisiert den Throttler.
- Args:
- default_config: Standard-Konfiguration.
- on_throttle: Callback wenn ein Event gedrosselt wird.
- """
- self._default_config = default_config or ThrottleConfig()
- self._on_throttle = on_throttle
- self._configs: dict[str, ThrottleConfig] = {}
- self._states: dict[str, ThrottleState] = defaultdict(ThrottleState)
- self._queues: dict[str, asyncio.Queue] = {}
- self._lock = asyncio.Lock()
- def set_global_rate(
- self,
- events_per_second: float,
- burst_size: int | None = None,
- ) -> None:
- """
- Setzt die globale Rate.
- Args:
- events_per_second: Maximale Events pro Sekunde.
- burst_size: Optionale Burst-Größe.
- """
- self._default_config.events_per_second = events_per_second
- if burst_size is not None:
- self._default_config.burst_size = burst_size
- def set_event_rate(
- self,
- event_name: str,
- events_per_second: float,
- burst_size: int | None = None,
- strategy: ThrottleStrategy | None = None,
- ) -> None:
- """
- Setzt die Rate für einen Event-Typ.
- Args:
- event_name: Name des Events.
- events_per_second: Maximale Events pro Sekunde.
- burst_size: Optionale Burst-Größe.
- strategy: Optionale Strategie.
- """
- config = ThrottleConfig(
- events_per_second=events_per_second,
- burst_size=burst_size or self._default_config.burst_size,
- strategy=strategy or self._default_config.strategy,
- queue_size=self._default_config.queue_size,
- sample_rate=self._default_config.sample_rate,
- )
- self._configs[event_name] = config
- def get_config(self, event_name: str) -> ThrottleConfig:
- """
- Gibt die Konfiguration für ein Event zurück.
- Args:
- event_name: Name des Events.
- Returns:
- ThrottleConfig.
- """
- return self._configs.get(event_name, self._default_config)
- def _update_tokens(
- self,
- state: ThrottleState,
- config: ThrottleConfig,
- ) -> None:
- """
- Aktualisiert die Token-Anzahl.
- Args:
- state: Der Zustand.
- config: Die Konfiguration.
- """
- now = time.monotonic()
- elapsed = now - state.last_update
- state.last_update = now
- # Tokens hinzufügen basierend auf verstrichener Zeit
- new_tokens = elapsed * config.events_per_second
- state.tokens = min(
- state.tokens + new_tokens,
- float(config.burst_size),
- )
- async def check(
- self,
- event_name: str,
- event_data: "EventData | None" = None,
- ) -> ThrottleResult:
- """
- Prüft, ob ein Event erlaubt ist.
- Args:
- event_name: Name des Events.
- event_data: Optionale Event-Daten.
- Returns:
- ThrottleResult mit der Entscheidung.
- """
- async with self._lock:
- config = self.get_config(event_name)
- state = self._states[event_name]
- # Tokens aktualisieren
- self._update_tokens(state, config)
- state.event_count += 1
- # Sample-Strategie prüfen
- if config.strategy == ThrottleStrategy.SAMPLE:
- if state.event_count % config.sample_rate != 0:
- return ThrottleResult(
- allowed=False,
- event_name=event_name,
- dropped=True,
- reason=f"Sampling (1/{config.sample_rate})",
- )
- # Token verfügbar?
- if state.tokens >= 1.0:
- state.tokens -= 1.0
- return ThrottleResult(
- allowed=True,
- event_name=event_name,
- )
- # Nicht genug Tokens - Strategie anwenden
- wait_time_ms = ((1.0 - state.tokens) / config.events_per_second) * 1000
- if config.strategy == ThrottleStrategy.DROP:
- state.dropped_count += 1
- if self._on_throttle and event_data:
- await self._on_throttle(event_name, event_data)
- return ThrottleResult(
- allowed=False,
- event_name=event_name,
- wait_time_ms=wait_time_ms,
- dropped=True,
- reason="Rate-Limit überschritten",
- )
- elif config.strategy == ThrottleStrategy.QUEUE:
- queue = self._get_queue(event_name, config.queue_size)
- if queue.full():
- state.dropped_count += 1
- return ThrottleResult(
- allowed=False,
- event_name=event_name,
- wait_time_ms=wait_time_ms,
- dropped=True,
- reason="Queue voll",
- )
- if event_data:
- await queue.put(event_data)
- state.queued_count += 1
- return ThrottleResult(
- allowed=False,
- event_name=event_name,
- wait_time_ms=wait_time_ms,
- queued=True,
- reason="In Queue gestellt",
- )
- elif config.strategy == ThrottleStrategy.DELAY:
- # Warten bis Token verfügbar
- await asyncio.sleep(wait_time_ms / 1000)
- state.tokens = 0.0 # Token wurde "gewartet"
- return ThrottleResult(
- allowed=True,
- event_name=event_name,
- wait_time_ms=wait_time_ms,
- reason="Verzögert",
- )
- return ThrottleResult(
- allowed=False,
- event_name=event_name,
- wait_time_ms=wait_time_ms,
- dropped=True,
- reason="Unbekannte Strategie",
- )
- def _get_queue(
- self,
- event_name: str,
- max_size: int,
- ) -> asyncio.Queue:
- """
- Gibt die Queue für ein Event zurück.
- Args:
- event_name: Name des Events.
- max_size: Maximale Queue-Größe.
- Returns:
- Die Queue.
- """
- if event_name not in self._queues:
- self._queues[event_name] = asyncio.Queue(maxsize=max_size)
- return self._queues[event_name]
- async def process_queued(
- self,
- event_name: str,
- handler: Callable[["EventData"], Coroutine[Any, Any, None]],
- batch_size: int = 10,
- ) -> int:
- """
- Verarbeitet Events aus der Queue.
- Args:
- event_name: Name des Events.
- handler: Handler-Funktion.
- batch_size: Maximale Batch-Größe.
- Returns:
- Anzahl verarbeiteter Events.
- """
- queue = self._queues.get(event_name)
- if not queue:
- return 0
- processed = 0
- while not queue.empty() and processed < batch_size:
- try:
- event_data = queue.get_nowait()
- await handler(event_data)
- processed += 1
- except asyncio.QueueEmpty:
- break
- return processed
- def get_stats(self, event_name: str) -> dict[str, Any]:
- """
- Gibt Statistiken für ein Event zurück.
- Args:
- event_name: Name des Events.
- Returns:
- Statistik-Dictionary.
- """
- state = self._states.get(event_name)
- if not state:
- return {}
- config = self.get_config(event_name)
- return {
- "event_name": event_name,
- "total_events": state.event_count,
- "dropped_events": state.dropped_count,
- "queued_events": state.queued_count,
- "drop_rate": (
- state.dropped_count / state.event_count
- if state.event_count > 0
- else 0.0
- ),
- "current_tokens": state.tokens,
- "max_tokens": config.burst_size,
- "rate_limit": config.events_per_second,
- }
- def get_all_stats(self) -> dict[str, dict[str, Any]]:
- """
- Gibt Statistiken für alle Events zurück.
- Returns:
- Dictionary von Event-Namen zu Statistiken.
- """
- return {
- name: self.get_stats(name)
- for name in self._states.keys()
- }
- def reset(self, event_name: str | None = None) -> None:
- """
- Setzt den Zustand zurück.
- Args:
- event_name: Optionaler Event-Name (None = alle).
- """
- if event_name:
- if event_name in self._states:
- del self._states[event_name]
- if event_name in self._queues:
- del self._queues[event_name]
- else:
- self._states.clear()
- self._queues.clear()
- class ThrottlingMiddleware:
- """
- Middleware-Integration für Event-Throttling.
- Kann in die Event-Pipeline integriert werden.
- Example:
- from trixy_core.events.middleware.base import EventMiddleware
- throttler = EventThrottler()
- throttler.set_event_rate("frequent_event", 10)
- # Als Middleware verwenden
- middleware = ThrottlingMiddleware(throttler)
- pipeline.add(middleware)
- """
- def __init__(
- self,
- throttler: EventThrottler,
- ) -> None:
- """
- Initialisiert die Middleware.
- Args:
- throttler: Der Throttler.
- """
- self._throttler = throttler
- self._name = "ThrottlingMiddleware"
- self._enabled = True
- @property
- def name(self) -> str:
- """Name der Middleware."""
- return self._name
- @property
- def enabled(self) -> bool:
- """Ob aktiv."""
- return self._enabled
- async def before(self, context: Any) -> Any:
- """Prüft Throttling vor Event-Verarbeitung."""
- from trixy_core.events.middleware.base import MiddlewareResult, MiddlewareAction
- result = await self._throttler.check(
- context.event_name,
- context.event_data,
- )
- if result.allowed:
- context.set_property("throttle_result", result)
- return MiddlewareResult()
- else:
- context.set_property("throttle_result", result)
- return MiddlewareResult(
- action=MiddlewareAction.SKIP_HANDLERS if result.queued else MiddlewareAction.ABORT,
- )
- async def after(self, context: Any) -> None:
- """Keine After-Verarbeitung."""
- pass
|