| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436 |
- # -*- coding: utf-8 -*-
- """
- Network Retry-Mechanismen.
- Bietet robuste Retry-Strategien für Netzwerkoperationen mit
- exponential backoff, Jitter und konfigurierbaren Policies.
- """
- from __future__ import annotations
- import asyncio
- import random
- import logging
- from dataclasses import dataclass, field
- from datetime import datetime, timedelta
- from enum import Enum, auto
- from typing import Any, Awaitable, Callable, TypeVar, Generic
- T = TypeVar("T")
- class RetryStrategy(Enum):
- """Verfügbare Retry-Strategien."""
- EXPONENTIAL = auto() # Exponentielles Backoff
- LINEAR = auto() # Lineares Backoff
- CONSTANT = auto() # Konstante Wartezeit
- FIBONACCI = auto() # Fibonacci-Sequenz
- DECORRELATED = auto() # Decorrelated Jitter (AWS-Stil)
- class RetryableError(Exception):
- """Exception die einen Retry auslösen soll."""
- pass
- class NonRetryableError(Exception):
- """Exception die keinen Retry auslösen soll."""
- pass
- @dataclass
- class RetryConfig:
- """Konfiguration für Retry-Verhalten."""
- max_attempts: int = 3
- strategy: RetryStrategy = RetryStrategy.EXPONENTIAL
- base_delay: float = 1.0 # Basis-Wartezeit in Sekunden
- max_delay: float = 60.0 # Maximale Wartezeit
- multiplier: float = 2.0 # Multiplikator für exponentielles Backoff
- jitter: float = 0.1 # Zufällige Variation (0.0-1.0)
- retryable_exceptions: tuple = (Exception,) # Exceptions die retry auslösen
- non_retryable_exceptions: tuple = () # Exceptions die NICHT retry auslösen
- on_retry: Callable[[int, Exception, float], None] | None = None # Callback bei Retry
- @dataclass
- class RetryStats:
- """Statistiken über Retry-Versuche."""
- total_attempts: int = 0
- successful_attempts: int = 0
- failed_attempts: int = 0
- total_delay: float = 0.0
- last_error: Exception | None = None
- last_attempt_at: datetime | None = None
- history: list[dict[str, Any]] = field(default_factory=list)
- class RetryContext:
- """
- Kontext für einen Retry-Durchlauf.
- Enthält Informationen über den aktuellen Versuch
- und ermöglicht Manipulation des Retry-Verhaltens.
- """
- def __init__(
- self,
- attempt: int,
- max_attempts: int,
- elapsed_time: float = 0.0,
- ) -> None:
- self.attempt = attempt
- self.max_attempts = max_attempts
- self.elapsed_time = elapsed_time
- self.should_retry = True
- self.next_delay: float | None = None
- self.metadata: dict[str, Any] = {}
- @property
- def is_last_attempt(self) -> bool:
- """Prüft ob dies der letzte Versuch ist."""
- return self.attempt >= self.max_attempts
- @property
- def remaining_attempts(self) -> int:
- """Anzahl verbleibender Versuche."""
- return max(0, self.max_attempts - self.attempt)
- def abort(self) -> None:
- """Bricht Retry ab - keine weiteren Versuche."""
- self.should_retry = False
- def set_next_delay(self, delay: float) -> None:
- """Setzt benutzerdefinierte Verzögerung für nächsten Versuch."""
- self.next_delay = delay
- class NetworkRetry(Generic[T]):
- """
- Retry-Handler für Netzwerkoperationen.
- Bietet verschiedene Backoff-Strategien und detaillierte
- Statistiken über Retry-Versuche.
- Beispiel:
- retry = NetworkRetry(
- config=RetryConfig(
- max_attempts=5,
- strategy=RetryStrategy.EXPONENTIAL,
- base_delay=1.0,
- )
- )
- result = await retry.execute(async_operation)
- """
- def __init__(
- self,
- config: RetryConfig | None = None,
- logger: logging.Logger | None = None,
- ) -> None:
- """
- Initialisiert den Retry-Handler.
- Args:
- config: Retry-Konfiguration
- logger: Optional: Logger für Retry-Logs
- """
- self.config = config or RetryConfig()
- self.logger = logger or logging.getLogger(__name__)
- self.stats = RetryStats()
- self._fibonacci_cache: list[int] = [0, 1]
- def _get_fibonacci(self, n: int) -> int:
- """Berechnet Fibonacci-Zahl (gecached)."""
- while len(self._fibonacci_cache) <= n:
- self._fibonacci_cache.append(
- self._fibonacci_cache[-1] + self._fibonacci_cache[-2]
- )
- return self._fibonacci_cache[n]
- def _calculate_delay(self, attempt: int, context: RetryContext) -> float:
- """Berechnet Wartezeit basierend auf Strategie."""
- # Benutzerdefinierte Verzögerung?
- if context.next_delay is not None:
- base_delay = context.next_delay
- else:
- strategy = self.config.strategy
- base = self.config.base_delay
- if strategy == RetryStrategy.EXPONENTIAL:
- base_delay = base * (self.config.multiplier ** (attempt - 1))
- elif strategy == RetryStrategy.LINEAR:
- base_delay = base * attempt
- elif strategy == RetryStrategy.CONSTANT:
- base_delay = base
- elif strategy == RetryStrategy.FIBONACCI:
- base_delay = base * self._get_fibonacci(attempt)
- elif strategy == RetryStrategy.DECORRELATED:
- # AWS-Stil decorrelated Jitter
- prev_delay = context.metadata.get("prev_delay", base)
- base_delay = random.uniform(base, prev_delay * 3)
- context.metadata["prev_delay"] = base_delay
- else:
- base_delay = base
- # Maximale Verzögerung
- base_delay = min(base_delay, self.config.max_delay)
- # Jitter hinzufügen
- if self.config.jitter > 0:
- jitter_range = base_delay * self.config.jitter
- jitter = random.uniform(-jitter_range, jitter_range)
- base_delay = max(0, base_delay + jitter)
- return base_delay
- def _should_retry(
- self,
- exception: Exception,
- context: RetryContext,
- ) -> bool:
- """Prüft ob bei dieser Exception ein Retry stattfinden soll."""
- if not context.should_retry:
- return False
- if context.is_last_attempt:
- return False
- # Non-retryable Exceptions
- if isinstance(exception, self.config.non_retryable_exceptions):
- return False
- if isinstance(exception, NonRetryableError):
- return False
- # Retryable Exceptions
- if isinstance(exception, RetryableError):
- return True
- if isinstance(exception, self.config.retryable_exceptions):
- return True
- return False
- async def execute(
- self,
- operation: Callable[[], Awaitable[T]],
- on_attempt: Callable[[RetryContext], None] | None = None,
- ) -> T:
- """
- Führt Operation mit Retry aus.
- Args:
- operation: Async-Funktion die ausgeführt werden soll
- on_attempt: Optional: Callback vor jedem Versuch
- Returns:
- Ergebnis der Operation
- Raises:
- Exception: Letzte Exception nach allen Versuchen
- """
- start_time = datetime.now()
- last_exception: Exception | None = None
- for attempt in range(1, self.config.max_attempts + 1):
- elapsed = (datetime.now() - start_time).total_seconds()
- context = RetryContext(
- attempt=attempt,
- max_attempts=self.config.max_attempts,
- elapsed_time=elapsed,
- )
- # Callback vor Versuch
- if on_attempt:
- on_attempt(context)
- # Statistik aktualisieren
- self.stats.total_attempts += 1
- self.stats.last_attempt_at = datetime.now()
- try:
- result = await operation()
- self.stats.successful_attempts += 1
- self.stats.history.append({
- "attempt": attempt,
- "success": True,
- "timestamp": datetime.now().isoformat(),
- })
- return result
- except Exception as e:
- last_exception = e
- self.stats.last_error = e
- self.stats.history.append({
- "attempt": attempt,
- "success": False,
- "error": str(e),
- "error_type": type(e).__name__,
- "timestamp": datetime.now().isoformat(),
- })
- if not self._should_retry(e, context):
- self.stats.failed_attempts += 1
- self.logger.debug(
- f"Retry abgebrochen nach Versuch {attempt}: {e}"
- )
- raise
- # Verzögerung berechnen
- delay = self._calculate_delay(attempt, context)
- self.stats.total_delay += delay
- self.logger.debug(
- f"Versuch {attempt}/{self.config.max_attempts} fehlgeschlagen, "
- f"nächster Versuch in {delay:.2f}s: {e}"
- )
- # Callback bei Retry
- if self.config.on_retry:
- self.config.on_retry(attempt, e, delay)
- # Warten
- await asyncio.sleep(delay)
- self.stats.failed_attempts += 1
- if last_exception:
- raise last_exception
- raise RuntimeError("Retry fehlgeschlagen ohne Exception")
- def reset_stats(self) -> None:
- """Setzt Statistiken zurück."""
- self.stats = RetryStats()
- def with_retry(
- max_attempts: int = 3,
- strategy: RetryStrategy = RetryStrategy.EXPONENTIAL,
- base_delay: float = 1.0,
- max_delay: float = 60.0,
- jitter: float = 0.1,
- retryable_exceptions: tuple = (Exception,),
- non_retryable_exceptions: tuple = (),
- ) -> Callable:
- """
- Decorator für automatisches Retry von Async-Funktionen.
- Args:
- max_attempts: Maximale Anzahl Versuche
- strategy: Backoff-Strategie
- base_delay: Basis-Wartezeit
- max_delay: Maximale Wartezeit
- jitter: Zufällige Variation
- retryable_exceptions: Exceptions die Retry auslösen
- non_retryable_exceptions: Exceptions ohne Retry
- Beispiel:
- @with_retry(max_attempts=5, base_delay=2.0)
- async def fetch_data():
- return await client.get("/api/data")
- """
- config = RetryConfig(
- max_attempts=max_attempts,
- strategy=strategy,
- base_delay=base_delay,
- max_delay=max_delay,
- jitter=jitter,
- retryable_exceptions=retryable_exceptions,
- non_retryable_exceptions=non_retryable_exceptions,
- )
- def decorator(func: Callable[..., Awaitable[T]]) -> Callable[..., Awaitable[T]]:
- async def wrapper(*args: Any, **kwargs: Any) -> T:
- retry = NetworkRetry[T](config=config)
- return await retry.execute(lambda: func(*args, **kwargs))
- wrapper.__name__ = func.__name__
- wrapper.__doc__ = func.__doc__
- return wrapper
- return decorator
- class RetryPolicy:
- """
- Wiederverwendbare Retry-Policy.
- Erlaubt Definition von Retry-Verhalten das auf mehrere
- Operationen angewendet werden kann.
- """
- def __init__(
- self,
- name: str,
- config: RetryConfig,
- ) -> None:
- """
- Initialisiert die Policy.
- Args:
- name: Name der Policy
- config: Retry-Konfiguration
- """
- self.name = name
- self.config = config
- self._stats: dict[str, RetryStats] = {}
- async def execute(
- self,
- operation: Callable[[], Awaitable[T]],
- operation_name: str = "default",
- ) -> T:
- """
- Führt Operation mit dieser Policy aus.
- Args:
- operation: Auszuführende Operation
- operation_name: Name für Statistik-Tracking
- Returns:
- Ergebnis der Operation
- """
- retry = NetworkRetry[T](config=self.config)
- result = await retry.execute(operation)
- # Statistik speichern
- self._stats[operation_name] = retry.stats
- return result
- def get_stats(self, operation_name: str | None = None) -> RetryStats | dict[str, RetryStats]:
- """Holt Statistiken für Operation(en)."""
- if operation_name:
- return self._stats.get(operation_name, RetryStats())
- return self._stats.copy()
- # Vordefinierte Policies
- AGGRESSIVE_RETRY = RetryConfig(
- max_attempts=10,
- strategy=RetryStrategy.EXPONENTIAL,
- base_delay=0.5,
- max_delay=30.0,
- jitter=0.2,
- )
- CONSERVATIVE_RETRY = RetryConfig(
- max_attempts=3,
- strategy=RetryStrategy.LINEAR,
- base_delay=5.0,
- max_delay=60.0,
- jitter=0.1,
- )
- FAST_FAIL_RETRY = RetryConfig(
- max_attempts=2,
- strategy=RetryStrategy.CONSTANT,
- base_delay=1.0,
- max_delay=1.0,
- jitter=0.0,
- )
|