Coverage for src/local_deep_research/library/download_management/status_tracker.py: 92%

142 statements  

« prev     ^ index     » next       coverage.py v7.15.1, created at 2026-07-19 23:35 +0000

1""" 

2Resource Status Tracker 

3 

4Tracks download attempts, failures, and cooldowns in the database. 

5Provides persistent storage for failure classifications and retry eligibility. 

6""" 

7 

8from datetime import datetime, timedelta, UTC 

9from typing import Optional, Dict, Any 

10 

11from loguru import logger 

12from sqlalchemy.orm import sessionmaker, Session 

13 

14from .models import ResourceDownloadStatus, Base 

15from .failure_classifier import BaseFailure 

16 

17MAX_TOTAL_RETRIES = 5 

18 

19 

20def compute_retry_cooldown( 

21 attempt: int, default_cooldown: timedelta 

22) -> Optional[timedelta]: 

23 """Return cooldown for this attempt number, or None for permanent failure. 

24 

25 Schedule: 

26 attempt 1: default_cooldown (from failure type) 

27 attempt 2: 1 day 

28 attempt 3-4: 30 days 

29 attempt >= 5: None (permanent failure) 

30 """ 

31 if attempt >= MAX_TOTAL_RETRIES: 

32 return None 

33 if attempt >= 3: 

34 return timedelta(days=30) 

35 if attempt == 2: 

36 return timedelta(days=1) 

37 return default_cooldown 

38 

39 

40class ResourceStatusTracker: 

41 """Track download attempts, failures, and cooldowns in database""" 

42 

43 def __init__(self, username: str, password: Optional[str] = None): 

44 """ 

45 Initialize the status tracker for a user. 

46 

47 Args: 

48 username: Username for database access 

49 password: Optional password for encrypted database 

50 """ 

51 self.username = username 

52 self.password = password 

53 

54 # Use the global db_manager singleton to share cached connections 

55 from ...database.encrypted_db import ( 

56 DatabaseInitializationError, 

57 db_manager, 

58 ) 

59 

60 self.db_manager = db_manager 

61 try: 

62 self.engine = db_manager.open_user_database(username, password) 

63 except DatabaseInitializationError: 

64 # Surface init failures from the schedulers/library-init 

65 # callers as a plain RuntimeError — they all wrap construction 

66 # in try/except already, and propagating the typed exception 

67 # would couple every caller to encrypted_db's internals. 

68 # ``logger.warning`` (no traceback) rather than 

69 # ``logger.exception``: ``password`` is a live local in this 

70 # frame, so rendering a traceback under ``diagnose=True`` 

71 # would dump the plaintext SQLCipher master password 

72 # (unrecoverable — TRUST.md §5). The redacted failure detail 

73 # is already logged at the raise site in 

74 # ``open_user_database`` (#4182). 

75 logger.warning( 

76 f"[STATUS_TRACKER] Database init failed for user: {username}" 

77 ) 

78 raise RuntimeError( 

79 f"Database initialisation failed for user {username}" 

80 ) from None 

81 self.Session = sessionmaker(bind=self.engine) 

82 

83 # Create tables if they don't exist 

84 Base.metadata.create_all(self.engine) 

85 logger.info( 

86 f"[STATUS_TRACKER] Initialized for user: {username} with encrypted database" 

87 ) 

88 

89 def _get_session(self) -> Session: 

90 """Get a database session""" 

91 return self.Session() 

92 

93 def mark_failure( 

94 self, 

95 resource_id: int, 

96 failure: BaseFailure, 

97 session: Optional[Session] = None, 

98 ) -> None: 

99 """ 

100 Mark a resource as failed with classification. 

101 

102 Args: 

103 resource_id: Resource identifier 

104 failure: Classified failure object 

105 session: Optional existing database session to reuse 

106 """ 

107 if session is not None: 107 ↛ 111line 107 didn't jump to line 111 because the condition on line 107 was always true

108 self._apply_failure(session, resource_id, failure) 

109 return 

110 

111 with self._get_session() as session: 

112 self._apply_failure(session, resource_id, failure) 

113 session.commit() 

114 

115 def _apply_failure( 

116 self, session: Session, resource_id: int, failure: BaseFailure 

117 ) -> None: 

118 """Apply failure status updates to a session (does not commit).""" 

119 # Get or create status record 

120 status = ( 

121 session.query(ResourceDownloadStatus) 

122 .filter_by(resource_id=resource_id) 

123 .first() 

124 ) 

125 

126 if not status: 

127 status = ResourceDownloadStatus(resource_id=resource_id) 

128 session.add(status) 

129 

130 # Update status information 

131 if failure.is_permanent(): 

132 status.status = "permanently_failed" 

133 status.retry_after_timestamp = None 

134 status.failure_type = failure.error_type 

135 status.failure_message = failure.message 

136 status.permanent_failure_at = datetime.now(UTC) 

137 logger.info( 

138 f"[STATUS_TRACKER] Marked resource {resource_id} as permanently failed: {failure.error_type}" 

139 ) 

140 else: 

141 status.failure_type = failure.error_type 

142 status.failure_message = failure.message 

143 

144 attempt = (status.total_retry_count or 0) + 1 

145 cooldown = compute_retry_cooldown(attempt, failure.retry_after) 

146 

147 if cooldown is None: 

148 status.status = "permanently_failed" 

149 status.permanent_failure_at = datetime.now(UTC) 

150 status.retry_after_timestamp = None 

151 logger.info( 

152 f"[STATUS_TRACKER] Auto-promoted resource {resource_id} to " 

153 f"permanently failed after {attempt} attempts" 

154 ) 

155 else: 

156 status.status = "temporarily_failed" 

157 status.retry_after_timestamp = datetime.now(UTC) + cooldown 

158 logger.info( 

159 f"[STATUS_TRACKER] Marked resource {resource_id} as temporarily failed: " 

160 f"{failure.error_type}, attempt {attempt}, retry after: {cooldown}" 

161 ) 

162 

163 # Update retry statistics 

164 # Ensure total_retry_count is initialized (handle None from legacy data) 

165 if status.total_retry_count is None: 

166 status.total_retry_count = 0 

167 status.total_retry_count += 1 

168 

169 # Check if last attempt was today (before overwriting last_attempt_at) 

170 today = datetime.now(UTC).date() 

171 last_attempt = ( 

172 status.last_attempt_at.date() if status.last_attempt_at else None 

173 ) 

174 status.last_attempt_at = datetime.now(UTC) 

175 if last_attempt == today: 

176 # Ensure today_retry_count is initialized (handle None from legacy data) 

177 if status.today_retry_count is None: 

178 status.today_retry_count = 0 

179 status.today_retry_count += 1 

180 else: 

181 status.today_retry_count = 1 

182 

183 logger.debug( 

184 f"[STATUS_TRACKER] Updated failure status for resource {resource_id}" 

185 ) 

186 

187 def mark_success( 

188 self, resource_id: int, session: Optional[Session] = None 

189 ) -> None: 

190 """ 

191 Mark a resource as successfully downloaded. 

192 

193 Args: 

194 resource_id: Resource identifier 

195 session: Optional existing database session to reuse 

196 """ 

197 if session is not None: 197 ↛ 201line 197 didn't jump to line 201 because the condition on line 197 was always true

198 self._apply_success(session, resource_id) 

199 return 

200 

201 with self._get_session() as session: 

202 self._apply_success(session, resource_id) 

203 session.commit() 

204 

205 def _apply_success(self, session: Session, resource_id: int) -> None: 

206 """Apply success status updates to a session (does not commit).""" 

207 status = ( 

208 session.query(ResourceDownloadStatus) 

209 .filter_by(resource_id=resource_id) 

210 .first() 

211 ) 

212 

213 if status: 

214 status.status = "completed" 

215 status.failure_type = None 

216 status.failure_message = None 

217 status.retry_after_timestamp = None 

218 status.updated_at = datetime.now(UTC) 

219 logger.info( 

220 f"[STATUS_TRACKER] Marked resource {resource_id} as successfully completed" 

221 ) 

222 

223 def can_retry(self, resource_id: int) -> tuple[bool, Optional[str]]: 

224 """ 

225 Check if a resource can be retried right now. 

226 

227 Args: 

228 resource_id: Resource identifier 

229 

230 Returns: 

231 Tuple of (can_retry, reason_if_not) 

232 """ 

233 with self._get_session() as session: 

234 status = ( 

235 session.query(ResourceDownloadStatus) 

236 .filter_by(resource_id=resource_id) 

237 .first() 

238 ) 

239 

240 if not status: 

241 # No status record, can retry 

242 return True, None 

243 

244 # Check if permanently failed 

245 if status.status == "permanently_failed": 

246 return ( 

247 False, 

248 f"Permanently failed: {status.failure_message or status.failure_type}", 

249 ) 

250 

251 # Check if temporarily failed and cooldown not expired 

252 if ( 

253 status.status == "temporarily_failed" 

254 and status.retry_after_timestamp 

255 ): 

256 # Ensure retry_after_timestamp is timezone-aware (handle legacy data) 

257 retry_timestamp = status.retry_after_timestamp 

258 if retry_timestamp.tzinfo is None: 

259 # Assume UTC for naive timestamps 

260 retry_timestamp = retry_timestamp.replace(tzinfo=UTC) 

261 

262 if datetime.now(UTC) < retry_timestamp: 

263 return ( 

264 False, 

265 f"Cooldown active, retry available at {retry_timestamp.strftime('%Y-%m-%d %H:%M:%S')}", 

266 ) 

267 

268 # Check daily retry limit (max 3 retries per day) 

269 # today_retry_count is only reset inside _apply_failure, so check 

270 # whether the stored count is actually from today before using it. 

271 today = datetime.now(UTC).date() 

272 last_attempt_date = ( 

273 status.last_attempt_at.date() 

274 if status.last_attempt_at 

275 else None 

276 ) 

277 today_count = ( 

278 status.today_retry_count if last_attempt_date == today else 0 

279 ) 

280 if today_count >= 3: 

281 return ( 

282 False, 

283 f"Daily retry limit exceeded ({today_count}/3). Retry available tomorrow.", 

284 ) 

285 

286 # Check total retry limit (safety net for records not yet auto-promoted) 

287 if (status.total_retry_count or 0) >= MAX_TOTAL_RETRIES: 287 ↛ 288line 287 didn't jump to line 288 because the condition on line 287 was never true

288 return ( 

289 False, 

290 f"Permanently failed after {status.total_retry_count} attempts. Will not retry.", 

291 ) 

292 

293 # Can retry 

294 return True, None 

295 

296 def get_resource_status(self, resource_id: int) -> Optional[Dict[str, Any]]: 

297 """ 

298 Get the current status of a resource. 

299 

300 Args: 

301 resource_id: Resource identifier 

302 

303 Returns: 

304 Status information dictionary or None if not found 

305 """ 

306 with self._get_session() as session: 

307 status = ( 

308 session.query(ResourceDownloadStatus) 

309 .filter_by(resource_id=resource_id) 

310 .first() 

311 ) 

312 

313 if not status: 

314 return None 

315 

316 return { 

317 "resource_id": status.resource_id, 

318 "status": status.status, 

319 "failure_type": status.failure_type, 

320 "failure_message": status.failure_message, 

321 "retry_after_timestamp": status.retry_after_timestamp.isoformat() 

322 if status.retry_after_timestamp 

323 else None, 

324 "last_attempt_at": status.last_attempt_at.isoformat() 

325 if status.last_attempt_at 

326 else None, 

327 "total_retry_count": status.total_retry_count, 

328 "today_retry_count": status.today_retry_count, 

329 "created_at": status.created_at.isoformat(), 

330 "updated_at": status.updated_at.isoformat(), 

331 } 

332 

333 def get_failed_resources_count(self) -> Dict[str, int]: 

334 """ 

335 Get counts of resources by failure type. 

336 

337 Returns: 

338 Dictionary mapping failure types to counts 

339 """ 

340 with self._get_session() as session: 

341 failed_resources = ( 

342 session.query(ResourceDownloadStatus) 

343 .filter( 

344 ResourceDownloadStatus.status.in_( 

345 ["temporarily_failed", "permanently_failed"] 

346 ) 

347 ) 

348 .all() 

349 ) 

350 

351 counts = {} 

352 for resource in failed_resources: 

353 failure_type = resource.failure_type or "unknown" 

354 counts[failure_type] = counts.get(failure_type, 0) + 1 

355 

356 return counts 

357 

358 def clear_permanent_failures(self, older_than_days: int = 30) -> int: 

359 """ 

360 Clear permanent failure statuses for old records. 

361 

362 Args: 

363 older_than_days: Clear failures older than this many days 

364 

365 Returns: 

366 Number of records cleared 

367 """ 

368 cutoff_date = datetime.now(UTC) - timedelta(days=older_than_days) 

369 

370 with self._get_session() as session: 

371 old_failures = ( 

372 session.query(ResourceDownloadStatus) 

373 .filter( 

374 ResourceDownloadStatus.status == "permanently_failed", 

375 ResourceDownloadStatus.created_at < cutoff_date, 

376 ) 

377 .all() 

378 ) 

379 

380 count = len(old_failures) 

381 for failure in old_failures: 

382 failure.status = "available" 

383 failure.failure_type = None 

384 failure.failure_message = None 

385 failure.retry_after_timestamp = None 

386 failure.permanent_failure_at = None 

387 failure.total_retry_count = 0 

388 failure.today_retry_count = 0 

389 failure.updated_at = datetime.now(UTC) 

390 

391 session.commit() 

392 logger.info( 

393 f"[STATUS_TRACKER] Cleared {count} old permanent failure records" 

394 ) 

395 return count