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

42 statements  

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) 

15 

16from redis.exceptions import ConnectionError, TimeoutError 

17from redis.retry import AbstractRetry 

18 

19T = TypeVar("T") 

20 

21if TYPE_CHECKING: 

22 from redis.backoff import AbstractBackoff 

23 

24 

25class Retry(AbstractRetry[Exception]): 

26 __hash__ = AbstractRetry.__hash__ 

27 

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) 

38 

39 def __eq__(self, other: Any) -> bool: 

40 if not isinstance(other, Retry): 

41 return NotImplemented 

42 

43 return ( 

44 self._backoff == other._backoff 

45 and self._retries == other._retries 

46 and set(self._supported_errors) == set(other._supported_errors) 

47 ) 

48 

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 

76 

77 if with_failure_count: 

78 await fail(error, failures) 

79 else: 

80 await fail(error) 

81 

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) 

87 

88 

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 

95 

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 ) 

103 

104 return Retry( 

105 backoff=retry._backoff, 

106 retries=retry._retries, 

107 supported_errors=retry._supported_errors, 

108 )