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
« prev ^ index » next coverage.py v7.15.1, created at 2026-07-19 23:35 +0000
1"""
2Resource Status Tracker
4Tracks download attempts, failures, and cooldowns in the database.
5Provides persistent storage for failure classifications and retry eligibility.
6"""
8from datetime import datetime, timedelta, UTC
9from typing import Optional, Dict, Any
11from loguru import logger
12from sqlalchemy.orm import sessionmaker, Session
14from .models import ResourceDownloadStatus, Base
15from .failure_classifier import BaseFailure
17MAX_TOTAL_RETRIES = 5
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.
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
40class ResourceStatusTracker:
41 """Track download attempts, failures, and cooldowns in database"""
43 def __init__(self, username: str, password: Optional[str] = None):
44 """
45 Initialize the status tracker for a user.
47 Args:
48 username: Username for database access
49 password: Optional password for encrypted database
50 """
51 self.username = username
52 self.password = password
54 # Use the global db_manager singleton to share cached connections
55 from ...database.encrypted_db import (
56 DatabaseInitializationError,
57 db_manager,
58 )
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)
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 )
89 def _get_session(self) -> Session:
90 """Get a database session"""
91 return self.Session()
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.
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
111 with self._get_session() as session:
112 self._apply_failure(session, resource_id, failure)
113 session.commit()
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 )
126 if not status:
127 status = ResourceDownloadStatus(resource_id=resource_id)
128 session.add(status)
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
144 attempt = (status.total_retry_count or 0) + 1
145 cooldown = compute_retry_cooldown(attempt, failure.retry_after)
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 )
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
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
183 logger.debug(
184 f"[STATUS_TRACKER] Updated failure status for resource {resource_id}"
185 )
187 def mark_success(
188 self, resource_id: int, session: Optional[Session] = None
189 ) -> None:
190 """
191 Mark a resource as successfully downloaded.
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
201 with self._get_session() as session:
202 self._apply_success(session, resource_id)
203 session.commit()
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 )
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 )
223 def can_retry(self, resource_id: int) -> tuple[bool, Optional[str]]:
224 """
225 Check if a resource can be retried right now.
227 Args:
228 resource_id: Resource identifier
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 )
240 if not status:
241 # No status record, can retry
242 return True, None
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 )
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)
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 )
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 )
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 )
293 # Can retry
294 return True, None
296 def get_resource_status(self, resource_id: int) -> Optional[Dict[str, Any]]:
297 """
298 Get the current status of a resource.
300 Args:
301 resource_id: Resource identifier
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 )
313 if not status:
314 return None
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 }
333 def get_failed_resources_count(self) -> Dict[str, int]:
334 """
335 Get counts of resources by failure type.
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 )
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
356 return counts
358 def clear_permanent_failures(self, older_than_days: int = 30) -> int:
359 """
360 Clear permanent failure statuses for old records.
362 Args:
363 older_than_days: Clear failures older than this many days
365 Returns:
366 Number of records cleared
367 """
368 cutoff_date = datetime.now(UTC) - timedelta(days=older_than_days)
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 )
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)
391 session.commit()
392 logger.info(
393 f"[STATUS_TRACKER] Cleared {count} old permanent failure records"
394 )
395 return count