Coverage for /pythoncovmergedfiles/medio/medio/usr/local/lib/python3.11/site-packages/redis/asyncio/retry.py: 36%
Shortcuts on this page
r m x toggle line displays
j k next/prev highlighted chunk
0 (zero) top of page
1 (one) first highlighted chunk
Shortcuts on this page
r m x toggle line displays
j k next/prev highlighted chunk
0 (zero) top of page
1 (one) first highlighted chunk
1import warnings
2from asyncio import sleep
3from inspect import iscoroutinefunction
4from typing import (
5 TYPE_CHECKING,
6 Any,
7 Awaitable,
8 Callable,
9 Optional,
10 Tuple,
11 Type,
12 TypeVar,
13 Union,
14)
16from redis.exceptions import ConnectionError, TimeoutError
17from redis.retry import AbstractRetry
19T = TypeVar("T")
21if TYPE_CHECKING:
22 from redis.backoff import AbstractBackoff
25class Retry(AbstractRetry[Exception]):
26 __hash__ = AbstractRetry.__hash__
28 def __init__(
29 self,
30 backoff: "AbstractBackoff",
31 retries: int,
32 supported_errors: Tuple[Type[Exception], ...] = (
33 ConnectionError,
34 TimeoutError,
35 ),
36 ):
37 super().__init__(backoff, retries, supported_errors)
39 def __eq__(self, other: Any) -> bool:
40 if not isinstance(other, Retry):
41 return NotImplemented
43 return (
44 self._backoff == other._backoff
45 and self._retries == other._retries
46 and set(self._supported_errors) == set(other._supported_errors)
47 )
49 async def call_with_retry(
50 self,
51 do: Callable[[], Awaitable[T]],
52 fail: Union[
53 Callable[[Exception], Any],
54 Callable[[Exception, int], Any],
55 ],
56 is_retryable: Optional[Callable[[Exception], bool]] = None,
57 with_failure_count: bool = False,
58 ) -> T:
59 """
60 Execute an operation that might fail and returns its result, or
61 raise the exception that was thrown depending on the `Backoff` object.
62 `do`: the operation to call. Expects no argument.
63 `fail`: the failure handler, expects the last error that was thrown
64 ``is_retryable``: optional function to determine if an error is retryable
65 ``with_failure_count``: if True, the failure count is passed to the failure handler
66 """
67 self._backoff.reset()
68 failures = 0
69 while True:
70 try:
71 return await do()
72 except self._supported_errors as error:
73 if is_retryable and not is_retryable(error):
74 raise
75 failures += 1
77 if with_failure_count:
78 await fail(error, failures)
79 else:
80 await fail(error)
82 if self._retries >= 0 and failures > self._retries:
83 raise error
84 backoff = self._backoff.compute(failures)
85 if backoff > 0:
86 await sleep(backoff)
89def _to_async_retry(retry: Any) -> Any:
90 """Convert a synchronous retry policy while preserving async-shaped ones."""
91 if iscoroutinefunction(getattr(type(retry), "call_with_retry", None)):
92 return retry
93 if not isinstance(retry, AbstractRetry):
94 return retry
96 warnings.warn(
97 "A synchronous redis.retry.Retry was passed to an asyncio client and has "
98 "been converted to redis.asyncio.retry.Retry. A custom call_with_retry "
99 "implementation is not preserved - please use redis.asyncio.retry.Retry.",
100 UserWarning,
101 stacklevel=2,
102 )
104 return Retry(
105 backoff=retry._backoff,
106 retries=retry._retries,
107 supported_errors=retry._supported_errors,
108 )