retry.py 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436
  1. # -*- coding: utf-8 -*-
  2. """
  3. Network Retry-Mechanismen.
  4. Bietet robuste Retry-Strategien für Netzwerkoperationen mit
  5. exponential backoff, Jitter und konfigurierbaren Policies.
  6. """
  7. from __future__ import annotations
  8. import asyncio
  9. import random
  10. import logging
  11. from dataclasses import dataclass, field
  12. from datetime import datetime, timedelta
  13. from enum import Enum, auto
  14. from typing import Any, Awaitable, Callable, TypeVar, Generic
  15. T = TypeVar("T")
  16. class RetryStrategy(Enum):
  17. """Verfügbare Retry-Strategien."""
  18. EXPONENTIAL = auto() # Exponentielles Backoff
  19. LINEAR = auto() # Lineares Backoff
  20. CONSTANT = auto() # Konstante Wartezeit
  21. FIBONACCI = auto() # Fibonacci-Sequenz
  22. DECORRELATED = auto() # Decorrelated Jitter (AWS-Stil)
  23. class RetryableError(Exception):
  24. """Exception die einen Retry auslösen soll."""
  25. pass
  26. class NonRetryableError(Exception):
  27. """Exception die keinen Retry auslösen soll."""
  28. pass
  29. @dataclass
  30. class RetryConfig:
  31. """Konfiguration für Retry-Verhalten."""
  32. max_attempts: int = 3
  33. strategy: RetryStrategy = RetryStrategy.EXPONENTIAL
  34. base_delay: float = 1.0 # Basis-Wartezeit in Sekunden
  35. max_delay: float = 60.0 # Maximale Wartezeit
  36. multiplier: float = 2.0 # Multiplikator für exponentielles Backoff
  37. jitter: float = 0.1 # Zufällige Variation (0.0-1.0)
  38. retryable_exceptions: tuple = (Exception,) # Exceptions die retry auslösen
  39. non_retryable_exceptions: tuple = () # Exceptions die NICHT retry auslösen
  40. on_retry: Callable[[int, Exception, float], None] | None = None # Callback bei Retry
  41. @dataclass
  42. class RetryStats:
  43. """Statistiken über Retry-Versuche."""
  44. total_attempts: int = 0
  45. successful_attempts: int = 0
  46. failed_attempts: int = 0
  47. total_delay: float = 0.0
  48. last_error: Exception | None = None
  49. last_attempt_at: datetime | None = None
  50. history: list[dict[str, Any]] = field(default_factory=list)
  51. class RetryContext:
  52. """
  53. Kontext für einen Retry-Durchlauf.
  54. Enthält Informationen über den aktuellen Versuch
  55. und ermöglicht Manipulation des Retry-Verhaltens.
  56. """
  57. def __init__(
  58. self,
  59. attempt: int,
  60. max_attempts: int,
  61. elapsed_time: float = 0.0,
  62. ) -> None:
  63. self.attempt = attempt
  64. self.max_attempts = max_attempts
  65. self.elapsed_time = elapsed_time
  66. self.should_retry = True
  67. self.next_delay: float | None = None
  68. self.metadata: dict[str, Any] = {}
  69. @property
  70. def is_last_attempt(self) -> bool:
  71. """Prüft ob dies der letzte Versuch ist."""
  72. return self.attempt >= self.max_attempts
  73. @property
  74. def remaining_attempts(self) -> int:
  75. """Anzahl verbleibender Versuche."""
  76. return max(0, self.max_attempts - self.attempt)
  77. def abort(self) -> None:
  78. """Bricht Retry ab - keine weiteren Versuche."""
  79. self.should_retry = False
  80. def set_next_delay(self, delay: float) -> None:
  81. """Setzt benutzerdefinierte Verzögerung für nächsten Versuch."""
  82. self.next_delay = delay
  83. class NetworkRetry(Generic[T]):
  84. """
  85. Retry-Handler für Netzwerkoperationen.
  86. Bietet verschiedene Backoff-Strategien und detaillierte
  87. Statistiken über Retry-Versuche.
  88. Beispiel:
  89. retry = NetworkRetry(
  90. config=RetryConfig(
  91. max_attempts=5,
  92. strategy=RetryStrategy.EXPONENTIAL,
  93. base_delay=1.0,
  94. )
  95. )
  96. result = await retry.execute(async_operation)
  97. """
  98. def __init__(
  99. self,
  100. config: RetryConfig | None = None,
  101. logger: logging.Logger | None = None,
  102. ) -> None:
  103. """
  104. Initialisiert den Retry-Handler.
  105. Args:
  106. config: Retry-Konfiguration
  107. logger: Optional: Logger für Retry-Logs
  108. """
  109. self.config = config or RetryConfig()
  110. self.logger = logger or logging.getLogger(__name__)
  111. self.stats = RetryStats()
  112. self._fibonacci_cache: list[int] = [0, 1]
  113. def _get_fibonacci(self, n: int) -> int:
  114. """Berechnet Fibonacci-Zahl (gecached)."""
  115. while len(self._fibonacci_cache) <= n:
  116. self._fibonacci_cache.append(
  117. self._fibonacci_cache[-1] + self._fibonacci_cache[-2]
  118. )
  119. return self._fibonacci_cache[n]
  120. def _calculate_delay(self, attempt: int, context: RetryContext) -> float:
  121. """Berechnet Wartezeit basierend auf Strategie."""
  122. # Benutzerdefinierte Verzögerung?
  123. if context.next_delay is not None:
  124. base_delay = context.next_delay
  125. else:
  126. strategy = self.config.strategy
  127. base = self.config.base_delay
  128. if strategy == RetryStrategy.EXPONENTIAL:
  129. base_delay = base * (self.config.multiplier ** (attempt - 1))
  130. elif strategy == RetryStrategy.LINEAR:
  131. base_delay = base * attempt
  132. elif strategy == RetryStrategy.CONSTANT:
  133. base_delay = base
  134. elif strategy == RetryStrategy.FIBONACCI:
  135. base_delay = base * self._get_fibonacci(attempt)
  136. elif strategy == RetryStrategy.DECORRELATED:
  137. # AWS-Stil decorrelated Jitter
  138. prev_delay = context.metadata.get("prev_delay", base)
  139. base_delay = random.uniform(base, prev_delay * 3)
  140. context.metadata["prev_delay"] = base_delay
  141. else:
  142. base_delay = base
  143. # Maximale Verzögerung
  144. base_delay = min(base_delay, self.config.max_delay)
  145. # Jitter hinzufügen
  146. if self.config.jitter > 0:
  147. jitter_range = base_delay * self.config.jitter
  148. jitter = random.uniform(-jitter_range, jitter_range)
  149. base_delay = max(0, base_delay + jitter)
  150. return base_delay
  151. def _should_retry(
  152. self,
  153. exception: Exception,
  154. context: RetryContext,
  155. ) -> bool:
  156. """Prüft ob bei dieser Exception ein Retry stattfinden soll."""
  157. if not context.should_retry:
  158. return False
  159. if context.is_last_attempt:
  160. return False
  161. # Non-retryable Exceptions
  162. if isinstance(exception, self.config.non_retryable_exceptions):
  163. return False
  164. if isinstance(exception, NonRetryableError):
  165. return False
  166. # Retryable Exceptions
  167. if isinstance(exception, RetryableError):
  168. return True
  169. if isinstance(exception, self.config.retryable_exceptions):
  170. return True
  171. return False
  172. async def execute(
  173. self,
  174. operation: Callable[[], Awaitable[T]],
  175. on_attempt: Callable[[RetryContext], None] | None = None,
  176. ) -> T:
  177. """
  178. Führt Operation mit Retry aus.
  179. Args:
  180. operation: Async-Funktion die ausgeführt werden soll
  181. on_attempt: Optional: Callback vor jedem Versuch
  182. Returns:
  183. Ergebnis der Operation
  184. Raises:
  185. Exception: Letzte Exception nach allen Versuchen
  186. """
  187. start_time = datetime.now()
  188. last_exception: Exception | None = None
  189. for attempt in range(1, self.config.max_attempts + 1):
  190. elapsed = (datetime.now() - start_time).total_seconds()
  191. context = RetryContext(
  192. attempt=attempt,
  193. max_attempts=self.config.max_attempts,
  194. elapsed_time=elapsed,
  195. )
  196. # Callback vor Versuch
  197. if on_attempt:
  198. on_attempt(context)
  199. # Statistik aktualisieren
  200. self.stats.total_attempts += 1
  201. self.stats.last_attempt_at = datetime.now()
  202. try:
  203. result = await operation()
  204. self.stats.successful_attempts += 1
  205. self.stats.history.append({
  206. "attempt": attempt,
  207. "success": True,
  208. "timestamp": datetime.now().isoformat(),
  209. })
  210. return result
  211. except Exception as e:
  212. last_exception = e
  213. self.stats.last_error = e
  214. self.stats.history.append({
  215. "attempt": attempt,
  216. "success": False,
  217. "error": str(e),
  218. "error_type": type(e).__name__,
  219. "timestamp": datetime.now().isoformat(),
  220. })
  221. if not self._should_retry(e, context):
  222. self.stats.failed_attempts += 1
  223. self.logger.debug(
  224. f"Retry abgebrochen nach Versuch {attempt}: {e}"
  225. )
  226. raise
  227. # Verzögerung berechnen
  228. delay = self._calculate_delay(attempt, context)
  229. self.stats.total_delay += delay
  230. self.logger.debug(
  231. f"Versuch {attempt}/{self.config.max_attempts} fehlgeschlagen, "
  232. f"nächster Versuch in {delay:.2f}s: {e}"
  233. )
  234. # Callback bei Retry
  235. if self.config.on_retry:
  236. self.config.on_retry(attempt, e, delay)
  237. # Warten
  238. await asyncio.sleep(delay)
  239. self.stats.failed_attempts += 1
  240. if last_exception:
  241. raise last_exception
  242. raise RuntimeError("Retry fehlgeschlagen ohne Exception")
  243. def reset_stats(self) -> None:
  244. """Setzt Statistiken zurück."""
  245. self.stats = RetryStats()
  246. def with_retry(
  247. max_attempts: int = 3,
  248. strategy: RetryStrategy = RetryStrategy.EXPONENTIAL,
  249. base_delay: float = 1.0,
  250. max_delay: float = 60.0,
  251. jitter: float = 0.1,
  252. retryable_exceptions: tuple = (Exception,),
  253. non_retryable_exceptions: tuple = (),
  254. ) -> Callable:
  255. """
  256. Decorator für automatisches Retry von Async-Funktionen.
  257. Args:
  258. max_attempts: Maximale Anzahl Versuche
  259. strategy: Backoff-Strategie
  260. base_delay: Basis-Wartezeit
  261. max_delay: Maximale Wartezeit
  262. jitter: Zufällige Variation
  263. retryable_exceptions: Exceptions die Retry auslösen
  264. non_retryable_exceptions: Exceptions ohne Retry
  265. Beispiel:
  266. @with_retry(max_attempts=5, base_delay=2.0)
  267. async def fetch_data():
  268. return await client.get("/api/data")
  269. """
  270. config = RetryConfig(
  271. max_attempts=max_attempts,
  272. strategy=strategy,
  273. base_delay=base_delay,
  274. max_delay=max_delay,
  275. jitter=jitter,
  276. retryable_exceptions=retryable_exceptions,
  277. non_retryable_exceptions=non_retryable_exceptions,
  278. )
  279. def decorator(func: Callable[..., Awaitable[T]]) -> Callable[..., Awaitable[T]]:
  280. async def wrapper(*args: Any, **kwargs: Any) -> T:
  281. retry = NetworkRetry[T](config=config)
  282. return await retry.execute(lambda: func(*args, **kwargs))
  283. wrapper.__name__ = func.__name__
  284. wrapper.__doc__ = func.__doc__
  285. return wrapper
  286. return decorator
  287. class RetryPolicy:
  288. """
  289. Wiederverwendbare Retry-Policy.
  290. Erlaubt Definition von Retry-Verhalten das auf mehrere
  291. Operationen angewendet werden kann.
  292. """
  293. def __init__(
  294. self,
  295. name: str,
  296. config: RetryConfig,
  297. ) -> None:
  298. """
  299. Initialisiert die Policy.
  300. Args:
  301. name: Name der Policy
  302. config: Retry-Konfiguration
  303. """
  304. self.name = name
  305. self.config = config
  306. self._stats: dict[str, RetryStats] = {}
  307. async def execute(
  308. self,
  309. operation: Callable[[], Awaitable[T]],
  310. operation_name: str = "default",
  311. ) -> T:
  312. """
  313. Führt Operation mit dieser Policy aus.
  314. Args:
  315. operation: Auszuführende Operation
  316. operation_name: Name für Statistik-Tracking
  317. Returns:
  318. Ergebnis der Operation
  319. """
  320. retry = NetworkRetry[T](config=self.config)
  321. result = await retry.execute(operation)
  322. # Statistik speichern
  323. self._stats[operation_name] = retry.stats
  324. return result
  325. def get_stats(self, operation_name: str | None = None) -> RetryStats | dict[str, RetryStats]:
  326. """Holt Statistiken für Operation(en)."""
  327. if operation_name:
  328. return self._stats.get(operation_name, RetryStats())
  329. return self._stats.copy()
  330. # Vordefinierte Policies
  331. AGGRESSIVE_RETRY = RetryConfig(
  332. max_attempts=10,
  333. strategy=RetryStrategy.EXPONENTIAL,
  334. base_delay=0.5,
  335. max_delay=30.0,
  336. jitter=0.2,
  337. )
  338. CONSERVATIVE_RETRY = RetryConfig(
  339. max_attempts=3,
  340. strategy=RetryStrategy.LINEAR,
  341. base_delay=5.0,
  342. max_delay=60.0,
  343. jitter=0.1,
  344. )
  345. FAST_FAIL_RETRY = RetryConfig(
  346. max_attempts=2,
  347. strategy=RetryStrategy.CONSTANT,
  348. base_delay=1.0,
  349. max_delay=1.0,
  350. jitter=0.0,
  351. )