Coverage for agentos/tools/retry_queue.py: 0%

121 statements  

« prev     ^ index     » next       coverage.py v7.14.3, created at 2026-07-06 17:01 +0800

1""" 

2RetryQueue — asynchronous retry with exponential backoff, max attempts, and dead letter queue. 

3 

4Supports: 

5 - Exponential backoff with jitter 

6 - Max retry attempts 

7 - Dead letter queue for permanently failed items 

8 - Callback hooks (on_retry, on_failure, on_success) 

9 - Synchronous and async execution modes 

10 - Configurable backoff strategy (exponential, constant, linear) 

11""" 

12 

13from __future__ import annotations 

14 

15import random 

16import threading 

17import time 

18from collections.abc import Callable 

19from dataclasses import dataclass, field 

20from enum import Enum 

21from typing import Any 

22 

23# ============================================================================ 

24# Backoff Strategy 

25# ============================================================================ 

26 

27 

28class BackoffStrategy(Enum): 

29 EXPONENTIAL = "exponential" 

30 CONSTANT = "constant" 

31 LINEAR = "linear" 

32 

33 

34# ============================================================================ 

35# Job 

36# ============================================================================ 

37 

38 

39@dataclass 

40class RetryJob: 

41 id: str 

42 func: Callable[..., Any] 

43 args: tuple = () 

44 kwargs: dict[str, Any] = field(default_factory=dict) 

45 attempts: int = 0 

46 last_error: Exception | None = None 

47 created_at: float = field(default_factory=time.time) 

48 

49 def execute(self) -> Any: 

50 return self.func(*self.args, **self.kwargs) 

51 

52 

53# ============================================================================ 

54# RetryQueue 

55# ============================================================================ 

56 

57 

58class RetryQueue: 

59 """Asynchronous retry queue with exponential backoff. 

60 

61 Usage: 

62 rq = RetryQueue( 

63 max_attempts=3, 

64 base_delay=1.0, 

65 max_delay=30.0, 

66 ) 

67 

68 def risky_call(a, b): 

69 # might fail... 

70 return a / b 

71 

72 # Submit a job — if it fails, it will be retried 

73 result = rq.submit(risky_call, 10, 0) 

74 

75 # Check dead letters 

76 for job, error in rq.dead_letters: 

77 print(f"Job {job.id} permanently failed: {error}") 

78 """ 

79 

80 def __init__( 

81 self, 

82 max_attempts: int = 3, 

83 base_delay: float = 1.0, 

84 max_delay: float = 60.0, 

85 backoff: BackoffStrategy = BackoffStrategy.EXPONENTIAL, 

86 jitter: bool = True, 

87 ): 

88 if max_attempts < 1: 

89 raise ValueError("max_attempts must be at least 1") 

90 self._max_attempts = max_attempts 

91 self._base_delay = base_delay 

92 self._max_delay = max_delay 

93 self._backoff = backoff 

94 self._jitter = jitter 

95 self._dead_letters: list[tuple] = [] 

96 self._lock = threading.RLock() 

97 self._total_submitted: int = 0 

98 self._total_succeeded: int = 0 

99 self._total_failed: int = 0 

100 # Hooks 

101 self._on_retry: list[Callable[[RetryJob, Exception, int], None]] = [] 

102 self._on_failure: list[Callable[[RetryJob, Exception], None]] = [] 

103 self._on_success: list[Callable[[RetryJob, Any], None]] = [] 

104 

105 # ---------- submit ---------- 

106 

107 def submit(self, func: Callable[..., Any], *args: Any, **kwargs: Any) -> Any: 

108 """Submit and execute a job with retry. Raises last error if all attempts fail.""" 

109 import uuid 

110 

111 job = RetryJob( 

112 id=str(uuid.uuid4())[:8], 

113 func=func, 

114 args=args, 

115 kwargs=kwargs, 

116 ) 

117 with self._lock: 

118 self._total_submitted += 1 

119 return self._execute(job) 

120 

121 def _execute(self, job: RetryJob) -> Any: 

122 attempt = 0 

123 while True: 

124 try: 

125 result = job.execute() 

126 self._notify_success(job, result) 

127 with self._lock: 

128 self._total_succeeded += 1 

129 return result 

130 except Exception as e: 

131 job.last_error = e 

132 job.attempts += 1 

133 attempt += 1 

134 

135 if attempt >= self._max_attempts: 

136 self._notify_failure(job, e) 

137 with self._lock: 

138 self._total_failed += 1 

139 self._dead_letters.append((job, e)) 

140 raise 

141 

142 self._notify_retry(job, e, attempt) 

143 delay = self._compute_delay(attempt) 

144 time.sleep(delay) 

145 

146 def _compute_delay(self, attempt: int) -> float: 

147 if self._backoff == BackoffStrategy.CONSTANT: 

148 delay = self._base_delay 

149 elif self._backoff == BackoffStrategy.LINEAR: 

150 delay = self._base_delay * attempt 

151 else: # EXPONENTIAL 

152 delay = self._base_delay * (2 ** (attempt - 1)) 

153 

154 delay = min(delay, self._max_delay) 

155 

156 if self._jitter: 

157 delay = delay * (0.5 + random.random() * 0.5) # 50%-100% of delay 

158 

159 return delay 

160 

161 # ---------- hooks ---------- 

162 

163 def on_retry(self, callback: Callable[[RetryJob, Exception, int], None]) -> None: 

164 self._on_retry.append(callback) 

165 

166 def on_failure(self, callback: Callable[[RetryJob, Exception], None]) -> None: 

167 self._on_failure.append(callback) 

168 

169 def on_success(self, callback: Callable[[RetryJob, Any], None]) -> None: 

170 self._on_success.append(callback) 

171 

172 def _notify_retry(self, job, error, attempt): 

173 for cb in self._on_retry: 

174 try: 

175 cb(job, error, attempt) 

176 except Exception: 

177 pass 

178 

179 def _notify_failure(self, job, error): 

180 for cb in self._on_failure: 

181 try: 

182 cb(job, error) 

183 except Exception: 

184 pass 

185 

186 def _notify_success(self, job, result): 

187 for cb in self._on_success: 

188 try: 

189 cb(job, result) 

190 except Exception: 

191 pass 

192 

193 # ---------- dead letters ---------- 

194 

195 @property 

196 def dead_letters(self) -> list[tuple]: 

197 with self._lock: 

198 return list(self._dead_letters) 

199 

200 def clear_dead_letters(self) -> None: 

201 with self._lock: 

202 self._dead_letters.clear() 

203 

204 def retry_dead_letter(self, index: int) -> Any: 

205 """Re-submit a dead letter job.""" 

206 with self._lock: 

207 if index < 0 or index >= len(self._dead_letters): 

208 raise IndexError("dead letter index out of range") 

209 job, _ = self._dead_letters.pop(index) 

210 job.attempts = 0 

211 job.last_error = None 

212 return self._execute(job) 

213 

214 # ---------- stats ---------- 

215 

216 @property 

217 def stats(self) -> dict[str, Any]: 

218 with self._lock: 

219 return { 

220 "total_submitted": self._total_submitted, 

221 "total_succeeded": self._total_succeeded, 

222 "total_failed": self._total_failed, 

223 "dead_letter_count": len(self._dead_letters), 

224 "max_attempts": self._max_attempts, 

225 "backoff": self._backoff.value, 

226 }