mirror of
https://github.com/ChuckBuilds/LEDMatrix.git
synced 2026-10-10 09:06:36 +00:00
84043468 also gated the ESPN chunk fetches and the background data
service's workers, for the hourly sports refresh. A burst test on hdpi
(baseball and football refreshing every 5 minutes, 10-minute soaks, G F F G):
G prefetch gated only 0.87%, 0.83% late; 6+ late 20, 16; fetches 0.3-1.8s
F + fetch threads gated 1.23%, 1.05% late; 6+ late 12, 16; fetches 1.6-4.6s
Every parked fetch thread wakes at each swap and has to take the GIL again
just to park at the end of the window, so twenty of them cost more than
they saved, and the fetches ran two to three times as long. The plugins'
own copies of espn_dates (half the burst) were never gated anyway.
espn_dates and the background data service go back to main's versions and
the module-level active gate goes. Kept from that commit: the render thread
is never gated, a live refresh from another thread can't take its place,
and nested blocks keep the outer boundary.
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
875 lines
36 KiB
Python
875 lines
36 KiB
Python
"""
|
|
Background Data Service for LEDMatrix
|
|
|
|
This service provides background threading capabilities for season data fetching
|
|
to prevent blocking the main display loop. It's designed to be used across
|
|
all sport managers for consistent background data management.
|
|
|
|
Key Features:
|
|
- Thread-safe data caching
|
|
- Automatic retry logic with exponential backoff
|
|
- Configurable timeouts and intervals
|
|
- Graceful error handling
|
|
- Progress tracking and logging
|
|
- Memory-efficient data storage
|
|
"""
|
|
|
|
import itertools
|
|
import time
|
|
from datetime import datetime
|
|
import logging
|
|
import threading
|
|
import requests
|
|
from typing import Dict, Any, Optional, Callable, List
|
|
from dataclasses import dataclass, field
|
|
from enum import Enum
|
|
from concurrent.futures import ThreadPoolExecutor
|
|
import pytz
|
|
from src.cache_manager import CacheManager
|
|
from src.common.json_body import response_json
|
|
from src.common.espn_dates import (
|
|
RANGE_RETRY_SECONDS,
|
|
_note_range_rejected,
|
|
_ranges_known_rejected,
|
|
clamp_espn_limit,
|
|
fetch_espn_date_chunks,
|
|
parse_espn_date_range,
|
|
)
|
|
# Configure logging
|
|
logger = logging.getLogger(__name__)
|
|
|
|
class FetchStatus(Enum):
|
|
"""Status of background fetch operations."""
|
|
PENDING = "pending"
|
|
IN_PROGRESS = "in_progress"
|
|
COMPLETED = "completed"
|
|
FAILED = "failed"
|
|
CANCELLED = "cancelled"
|
|
|
|
@dataclass
|
|
class FetchRequest:
|
|
"""Represents a background fetch request."""
|
|
id: str
|
|
sport: str
|
|
year: int
|
|
cache_key: str
|
|
url: str
|
|
params: Dict[str, Any] = field(default_factory=dict)
|
|
headers: Dict[str, str] = field(default_factory=dict)
|
|
timeout: int = 30
|
|
retry_count: int = 0
|
|
max_retries: int = 3
|
|
# Recorded but not acted on: requests go straight to the thread pool in
|
|
# submission order. Kept because plugins pass it through.
|
|
priority: int = 1
|
|
callback: Optional[Callable] = None
|
|
# Callbacks from submitters that JOINED this fetch instead of starting a
|
|
# duplicate one. The primary `callback` above belongs to whoever created
|
|
# the request; these belong to everyone who asked for the same cache_key
|
|
# while it was still in flight.
|
|
extra_callbacks: List[Callable] = field(default_factory=list)
|
|
created_at: float = field(default_factory=time.time)
|
|
status: FetchStatus = FetchStatus.PENDING
|
|
# Set once the worker has decided this response will be cached, while it
|
|
# still holds the lock. From that point cancelling is refused: the write
|
|
# is already authorised, and abandoning it here would put the payload in
|
|
# the cache with the callbacks suppressed -- joiners waiting forever for a
|
|
# fetch that did, in fact, succeed.
|
|
commit_claimed: bool = False
|
|
result: Optional[Any] = None
|
|
error: Optional[str] = None
|
|
|
|
@dataclass
|
|
class FetchResult:
|
|
"""Result of a background fetch operation.
|
|
|
|
``data`` survives on the stored result only for requests submitted without
|
|
a ``callback``, where polling ``get_result()`` is the sole way to collect
|
|
it. When a callback was given, the payload has already been delivered and
|
|
the service releases it -- see :meth:`BackgroundDataService._release_payload`.
|
|
Either way the data remains in the cache under the request's ``cache_key``,
|
|
which is where consumers read it from.
|
|
"""
|
|
request_id: str
|
|
success: bool
|
|
data: Optional[Any] = None
|
|
error: Optional[str] = None
|
|
cached: bool = False
|
|
fetch_time: float = 0.0
|
|
retry_count: int = 0
|
|
completed_at: float = field(default_factory=time.time) # Timestamp when request completed
|
|
# The request's final status, recorded so a finished request can still be
|
|
# reported accurately. Without it a caller can only be told COMPLETED or
|
|
# FAILED, which turns "you cancelled this" into "this errored".
|
|
final_status: Optional[FetchStatus] = None
|
|
|
|
class BackgroundDataService:
|
|
"""
|
|
Background data service for fetching season data without blocking the main thread.
|
|
|
|
This service manages a pool of background threads to fetch data asynchronously,
|
|
with intelligent caching, retry logic, and progress tracking.
|
|
"""
|
|
|
|
# Plugins feature-detect this. A core without it sends season ranges to
|
|
# ESPN as-is and gets 400s since 2026-09-15, so plugins fetch those
|
|
# ranges themselves instead of submitting them here.
|
|
handles_espn_date_ranges = True
|
|
|
|
def __init__(self, cache_manager: CacheManager, max_workers: int = 3, request_timeout: int = 30):
|
|
"""
|
|
Initialize the background data service.
|
|
|
|
Args:
|
|
cache_manager: Cache manager instance for storing fetched data
|
|
max_workers: Maximum number of background threads
|
|
request_timeout: Default timeout for HTTP requests
|
|
"""
|
|
self.cache_manager = cache_manager
|
|
self.max_workers = max_workers
|
|
self.request_timeout = request_timeout
|
|
|
|
# Thread management
|
|
self.executor = ThreadPoolExecutor(max_workers=max_workers, thread_name_prefix="BackgroundData")
|
|
# cache_key -> request_id for fetches currently in flight. Submitting
|
|
# the same key twice used to start two identical fetches: request_id
|
|
# carries a millisecond timestamp, so every submit looked new, and
|
|
# active_requests is keyed by it rather than by what is being fetched.
|
|
# On a real board the season-schedule key is requested by both the
|
|
# Recent and the Upcoming manager, which miss the cache in the same
|
|
# millisecond and each download and parse the same payload.
|
|
self._inflight_by_cache_key: Dict[str, str] = {}
|
|
# request_id was sport_year_milliseconds, which is not unique: two
|
|
# submits inside the same millisecond produced the SAME id, so one
|
|
# silently replaced the other in active_requests and completed_requests.
|
|
# Rare before, but dedupe hands this id back to every joiner as their
|
|
# handle for get_result(), so it has to be unique. A counter is enough.
|
|
self._request_seq = itertools.count()
|
|
self.active_requests: Dict[str, FetchRequest] = {}
|
|
self.completed_requests: Dict[str, FetchResult] = {}
|
|
|
|
# Thread safety
|
|
self._lock = threading.RLock()
|
|
self._shutdown = False
|
|
|
|
# Cleanup tracking
|
|
self._max_completed_requests = 500 # Maximum completed requests to keep
|
|
self._completed_requests_cleanup_interval = 600.0 # Cleanup every 10 minutes
|
|
self._last_completed_requests_cleanup = time.time()
|
|
|
|
# Statistics
|
|
self.stats = {
|
|
'total_requests': 0,
|
|
'completed_requests': 0,
|
|
'failed_requests': 0,
|
|
'cached_hits': 0,
|
|
'cache_misses': 0,
|
|
'total_fetch_time': 0.0,
|
|
'average_fetch_time': 0.0
|
|
}
|
|
|
|
# Session for HTTP requests
|
|
self.session = requests.Session()
|
|
self.session.mount('http://', requests.adapters.HTTPAdapter(max_retries=3))
|
|
self.session.mount('https://', requests.adapters.HTTPAdapter(max_retries=3))
|
|
|
|
# Default headers: core's shared set (real User-Agent, no hand-set
|
|
# Accept-Encoding) -- see src/common/api_helper.py.
|
|
from src.common.api_helper import DEFAULT_HTTP_HEADERS
|
|
self.default_headers = dict(DEFAULT_HTTP_HEADERS)
|
|
|
|
logger.info(f"BackgroundDataService initialized with {max_workers} workers")
|
|
|
|
def get_sport_cache_key(self, sport: str, date_str: str = None) -> str:
|
|
"""
|
|
Generate consistent cache keys for sports data.
|
|
This ensures Recent/Upcoming managers and background service
|
|
use the same cache keys.
|
|
"""
|
|
# Same format as CacheManager.generate_sport_cache_key(). This used to
|
|
# build a whole CacheManager to call it -- config load, cache-dir
|
|
# probing with test writes -- on every submit without a cache_key.
|
|
if date_str is None:
|
|
date_str = datetime.now(pytz.utc).strftime('%Y%m%d')
|
|
return f"{sport}_{date_str}"
|
|
|
|
def submit_fetch_request(self,
|
|
sport: str,
|
|
year: int,
|
|
url: str,
|
|
cache_key: str = None,
|
|
params: Optional[Dict[str, Any]] = None,
|
|
headers: Optional[Dict[str, str]] = None,
|
|
timeout: Optional[int] = None,
|
|
max_retries: int = 3,
|
|
priority: int = 1,
|
|
callback: Optional[Callable] = None) -> str:
|
|
"""
|
|
Submit a background fetch request.
|
|
|
|
Args:
|
|
sport: Sport identifier (e.g., 'nfl', 'ncaafb')
|
|
year: Year to fetch data for
|
|
url: URL to fetch data from
|
|
cache_key: Cache key for storing/retrieving data
|
|
params: URL parameters
|
|
headers: HTTP headers
|
|
timeout: Request timeout
|
|
max_retries: Maximum number of retries
|
|
priority: Accepted for compatibility and ignored; requests run in
|
|
submission order.
|
|
callback: Optional callback function when request completes
|
|
|
|
Returns:
|
|
Request ID for tracking the fetch operation
|
|
"""
|
|
if self._shutdown:
|
|
raise RuntimeError("BackgroundDataService is shutting down")
|
|
|
|
# Generate cache key if not provided
|
|
if cache_key is None:
|
|
cache_key = self.get_sport_cache_key(sport)
|
|
|
|
with self._lock:
|
|
request_id = (f"{sport}_{year}_{int(time.time() * 1000)}"
|
|
f"_{next(self._request_seq)}")
|
|
|
|
# Check cache first
|
|
cached_data = self.cache_manager.get(cache_key)
|
|
if cached_data:
|
|
with self._lock:
|
|
self.stats['cached_hits'] += 1
|
|
result = FetchResult(
|
|
request_id=request_id,
|
|
success=True,
|
|
data=cached_data,
|
|
cached=True,
|
|
fetch_time=0.0
|
|
)
|
|
# Filed before the callback runs, as it always was: a callback
|
|
# that queries get_result()/is_request_complete() for its own
|
|
# request must still find it. Releasing afterwards mutates the
|
|
# same object the dict holds.
|
|
self.completed_requests[request_id] = result
|
|
|
|
if callback:
|
|
try:
|
|
callback(result)
|
|
except Exception as e:
|
|
logger.error(f"Error in callback for request {request_id}: {e}")
|
|
self._release_payload(result)
|
|
|
|
logger.debug(f"Cache hit for {sport} {year} data")
|
|
return request_id
|
|
|
|
# limit above 500 makes an ESPN *scoreboard* return a truncated list
|
|
# (src/common/espn_dates.py). Other endpoints need more: /teams has 762
|
|
# college-football teams, so only scoreboards are clamped.
|
|
if url.split('?', 1)[0].rstrip('/').endswith('/scoreboard'):
|
|
params = clamp_espn_limit(params)
|
|
|
|
# Create fetch request
|
|
request = FetchRequest(
|
|
id=request_id,
|
|
sport=sport,
|
|
year=year,
|
|
cache_key=cache_key,
|
|
url=url,
|
|
params=dict(params or {}),
|
|
headers={**self.default_headers, **(headers or {})},
|
|
timeout=timeout or self.request_timeout,
|
|
max_retries=max_retries,
|
|
priority=priority,
|
|
callback=callback
|
|
)
|
|
|
|
with self._lock:
|
|
existing_id = self._inflight_by_cache_key.get(cache_key)
|
|
existing = self.active_requests.get(existing_id) if existing_id else None
|
|
if existing_id and existing is None:
|
|
# Stranded index entry: the request it names is gone. Drop it and
|
|
# fetch normally. Looking the request up rather than trusting the
|
|
# id is what stops a stale entry wedging a key forever.
|
|
del self._inflight_by_cache_key[cache_key]
|
|
if existing is not None:
|
|
# Someone is already fetching this key. Ride along rather than
|
|
# duplicating the download, the parse and the resident copy.
|
|
if callback:
|
|
existing.extra_callbacks.append(callback)
|
|
self.stats['deduplicated_requests'] = (
|
|
self.stats.get('deduplicated_requests', 0) + 1
|
|
)
|
|
logger.info(
|
|
"Joined in-flight fetch %s for %s (cache_key=%s) instead of "
|
|
"starting a duplicate", existing_id, sport, cache_key
|
|
)
|
|
return existing_id
|
|
|
|
self.active_requests[request_id] = request
|
|
self._inflight_by_cache_key[cache_key] = request_id
|
|
self.stats['total_requests'] += 1
|
|
self.stats['cache_misses'] += 1
|
|
|
|
# Submit to executor
|
|
self.executor.submit(self._fetch_data_worker, request)
|
|
|
|
logger.info(f"Submitted background fetch request {request_id} for {sport} {year}")
|
|
return request_id
|
|
|
|
def _fetch_data_worker(self, request: FetchRequest) -> FetchResult:
|
|
"""
|
|
Worker function that performs the actual data fetching.
|
|
|
|
Args:
|
|
request: Fetch request to process
|
|
|
|
Returns:
|
|
Fetch result with data or error information
|
|
"""
|
|
start_time = time.time()
|
|
result = FetchResult(request_id=request.id, success=False, retry_count=request.retry_count)
|
|
|
|
try:
|
|
with self._lock:
|
|
# A request cancelled while it sat in the executor queue must
|
|
# stay cancelled. Overwriting the status here undid the cancel
|
|
# outright: the worker went on to download, cache and call back
|
|
# for work the caller had already withdrawn.
|
|
if request.status == FetchStatus.CANCELLED:
|
|
cancelled_before_start = True
|
|
else:
|
|
cancelled_before_start = False
|
|
request.status = FetchStatus.IN_PROGRESS
|
|
if cancelled_before_start:
|
|
logger.info(
|
|
"Request %s was cancelled before its worker started; "
|
|
"skipping the fetch entirely", request.id
|
|
)
|
|
# Assign before returning: the finally block stores `result`
|
|
# in completed_requests, so building a fresh one here would
|
|
# file the untouched placeholder instead of this outcome.
|
|
result = FetchResult(
|
|
request_id=request.id,
|
|
success=False,
|
|
error="cancelled",
|
|
fetch_time=time.time() - start_time,
|
|
retry_count=request.retry_count
|
|
)
|
|
return result
|
|
|
|
logger.info(f"Starting background fetch for {request.sport} {request.year}")
|
|
|
|
# ESPN stopped accepting dates=YYYYMMDD-YYYYMMDD on 2026-09-15 and
|
|
# answers 400 for every sport. Re-ask in months and days rather
|
|
# than let a whole season fail. See src/common/espn_dates.py.
|
|
# The "ranges are rejected" memo is shared with
|
|
# fetch_espn_scoreboard(): once either path has seen a range
|
|
# rejected, the other skips the doomed range request too.
|
|
is_range = parse_espn_date_range(request.params.get("dates")) is not None
|
|
data = None
|
|
chunks_tried = False
|
|
if is_range and _ranges_known_rejected():
|
|
data = self._fetch_in_date_chunks(request)
|
|
# Every chunk failed: ask for the range itself below so the
|
|
# failure carries a real HTTP error, without re-spending chunks.
|
|
chunks_tried = data is None
|
|
|
|
if data is None:
|
|
# Perform HTTP request with retry logic
|
|
response = self._make_request_with_retry(request)
|
|
if is_range and response.status_code == 400 and not chunks_tried:
|
|
_note_range_rejected()
|
|
logger.warning(
|
|
"ESPN rejected the date range %s (400); fetching it as "
|
|
"month/day chunks, and fetching ranges that way for the "
|
|
"next %d hours",
|
|
request.params.get("dates"), RANGE_RETRY_SECONDS // 3600,
|
|
)
|
|
data = self._fetch_in_date_chunks(request)
|
|
if data is None:
|
|
response.raise_for_status()
|
|
else:
|
|
response.raise_for_status()
|
|
data = response_json(response)
|
|
|
|
# Validate data structure
|
|
if not isinstance(data, dict):
|
|
raise ValueError(f"Expected dict response, got {type(data)}")
|
|
|
|
if 'events' not in data:
|
|
raise ValueError("Response missing 'events' field")
|
|
|
|
# Validate events structure
|
|
events = data.get('events', [])
|
|
if not isinstance(events, list):
|
|
raise ValueError(f"Expected events to be list, got {type(events)}")
|
|
|
|
# Log data validation
|
|
logger.debug(f"Validated {len(events)} events for {request.sport} {request.year}")
|
|
|
|
# A cancelled request must not commit anything. Cancelling
|
|
# releases the cache_key, so a replacement fetch for the same key
|
|
# may already be in flight or finished -- writing this response to
|
|
# the cache now would overwrite fresher data with the response
|
|
# nobody wanted. The worker has no way to abort the HTTP call, so
|
|
# this is where the work gets discarded.
|
|
with self._lock:
|
|
cancelled = request.status == FetchStatus.CANCELLED
|
|
if not cancelled:
|
|
# Claim the commit in the same critical section that read
|
|
# the status, so a cancel cannot slip in between the check
|
|
# and the cache write below. The write itself stays outside
|
|
# the lock: it serialises a multi-megabyte payload to the
|
|
# SD card, and holding the service lock across that would
|
|
# stall every submit, status query and cancel behind it.
|
|
request.commit_claimed = True
|
|
if cancelled:
|
|
logger.info(
|
|
"Discarding response for cancelled request %s; %s may "
|
|
"already belong to a replacement fetch",
|
|
request.id, request.cache_key
|
|
)
|
|
result = FetchResult(
|
|
request_id=request.id,
|
|
success=False,
|
|
error="cancelled",
|
|
fetch_time=time.time() - start_time,
|
|
retry_count=request.retry_count
|
|
)
|
|
return result
|
|
|
|
# Cache the data
|
|
self.cache_manager.set(request.cache_key, data)
|
|
|
|
# Update request status
|
|
with self._lock:
|
|
request.status = FetchStatus.COMPLETED
|
|
request.result = data
|
|
|
|
# Create successful result
|
|
fetch_time = time.time() - start_time
|
|
result = FetchResult(
|
|
request_id=request.id,
|
|
success=True,
|
|
data=data,
|
|
fetch_time=fetch_time,
|
|
retry_count=request.retry_count
|
|
)
|
|
|
|
logger.info(f"Successfully fetched {request.sport} {request.year} data in {fetch_time:.2f}s")
|
|
|
|
except Exception as e:
|
|
error_msg = str(e)
|
|
logger.error(f"Failed to fetch {request.sport} {request.year} data: {error_msg}")
|
|
|
|
with self._lock:
|
|
# Don't relabel a cancelled request. The callback gate in the
|
|
# finally block only suppresses CANCELLED, so promoting it to
|
|
# FAILED here delivered an error callback for a fetch nobody
|
|
# was waiting on any more.
|
|
if request.status != FetchStatus.CANCELLED:
|
|
request.status = FetchStatus.FAILED
|
|
request.error = error_msg
|
|
|
|
result = FetchResult(
|
|
request_id=request.id,
|
|
success=False,
|
|
error=error_msg,
|
|
fetch_time=time.time() - start_time,
|
|
retry_count=request.retry_count
|
|
)
|
|
|
|
finally:
|
|
# Store result and clean up
|
|
with self._lock:
|
|
result.final_status = request.status
|
|
self.completed_requests[request.id] = result
|
|
if request.id in self.active_requests:
|
|
del self.active_requests[request.id]
|
|
# Stop accepting joiners and take the callback list in the same
|
|
# critical section. A submitter that arrives after this point
|
|
# finds no in-flight entry and either hits the cache (written
|
|
# above, before the result was built) or starts a fresh fetch --
|
|
# what it must never do is join a fetch whose callbacks have
|
|
# already run and then never be called.
|
|
if self._inflight_by_cache_key.get(request.cache_key) == request.id:
|
|
del self._inflight_by_cache_key[request.cache_key]
|
|
# A cancelled request delivers nothing: its joiners were told
|
|
# about a fetch that has been abandoned, and a replacement will
|
|
# call them via its own request.
|
|
if request.status == FetchStatus.CANCELLED:
|
|
callbacks = []
|
|
else:
|
|
callbacks = ([request.callback] if request.callback else [])
|
|
callbacks.extend(request.extra_callbacks)
|
|
|
|
# Update statistics
|
|
if result.success:
|
|
self.stats['completed_requests'] += 1
|
|
else:
|
|
self.stats['failed_requests'] += 1
|
|
|
|
self.stats['total_fetch_time'] += result.fetch_time
|
|
self.stats['average_fetch_time'] = (
|
|
self.stats['total_fetch_time'] /
|
|
(self.stats['completed_requests'] + self.stats['failed_requests'])
|
|
)
|
|
|
|
# Periodic cleanup after storing result
|
|
self._cleanup_completed_requests()
|
|
|
|
# Call every callback: the original submitter's and any that joined
|
|
# this fetch. One raising must not stop the others being delivered.
|
|
for cb in callbacks:
|
|
try:
|
|
cb(result)
|
|
except Exception as e:
|
|
logger.error(f"Error in callback for request {request.id}: {e}")
|
|
|
|
# Released AFTER the loop, not inside it. Every callback here holds
|
|
# the same FetchResult, so releasing per-delivery handed the first
|
|
# one the data and every joiner `result.data is None` -- which is
|
|
# not a quiet degradation: they read `result.data.get('events')` and
|
|
# raise AttributeError, which this very loop catches and logs, so
|
|
# the symptom was one ERROR line and a manager that silently never
|
|
# got its schedule. Deduplication is the normal case, not a corner:
|
|
# a sport's recent, upcoming and live managers all ride one season
|
|
# fetch.
|
|
#
|
|
# Guarded on `callbacks`, because a request submitted without one
|
|
# has no other way to collect its payload than polling get_result().
|
|
# The old per-delivery release got that right by accident: an empty
|
|
# list never entered the loop body.
|
|
if callbacks:
|
|
self._release_payload(result)
|
|
request.result = None
|
|
|
|
return result
|
|
|
|
@staticmethod
|
|
def _release_payload(result: FetchResult) -> None:
|
|
"""Drop a delivered payload, keeping the result's status and timings.
|
|
|
|
Only called once EVERY callback has been handed the data -- callers
|
|
that joined an in-flight fetch share this object, so releasing between
|
|
deliveries strips the payload out from under the ones still queued.
|
|
Consumers read fetched data back from the cache under ``cache_key``;
|
|
the copy carried
|
|
here was pinning a parsed season schedule -- 946 games for NCAA
|
|
football, roughly a tenth of total RAM on a 1GB Pi -- in memory until
|
|
the hourly sweep.
|
|
|
|
The cache-hit path matters most: it runs once per update interval per
|
|
sport, mints a fresh request_id each time, and a memory-tier miss
|
|
re-parses the payload from disk. Those were genuinely separate copies
|
|
accumulating toward the 500-entry cap, not shared references.
|
|
"""
|
|
result.data = None
|
|
|
|
def _fetch_in_date_chunks(self, request: FetchRequest) -> Optional[Dict[str, Any]]:
|
|
"""Re-fetch a rejected ``YYYYMMDD-YYYYMMDD`` range as month/day chunks.
|
|
|
|
None means the request was not a day range, or every chunk failed; the
|
|
caller then re-raises the original 400 instead of caching an empty
|
|
season. See src/common/espn_dates.py.
|
|
"""
|
|
logger.info("Recovering %s %s from a rejected date range", request.sport, request.year)
|
|
return fetch_espn_date_chunks(
|
|
self.session,
|
|
request.url,
|
|
params=request.params,
|
|
headers=request.headers,
|
|
timeout=request.timeout,
|
|
logger=logger,
|
|
)
|
|
|
|
def _make_request_with_retry(self, request: FetchRequest) -> requests.Response:
|
|
"""
|
|
Make HTTP request with retry logic and exponential backoff.
|
|
|
|
Args:
|
|
request: Fetch request containing request details
|
|
|
|
Returns:
|
|
HTTP response
|
|
|
|
Raises:
|
|
requests.RequestException: If all retries fail
|
|
"""
|
|
last_exception = None
|
|
|
|
for attempt in range(request.max_retries + 1):
|
|
try:
|
|
response = self.session.get(
|
|
request.url,
|
|
params=request.params,
|
|
headers=request.headers,
|
|
timeout=request.timeout
|
|
)
|
|
return response
|
|
|
|
except requests.RequestException as e:
|
|
last_exception = e
|
|
request.retry_count = attempt + 1
|
|
|
|
if attempt < request.max_retries:
|
|
# Exponential backoff: 1s, 2s, 4s, 8s...
|
|
delay = 2 ** attempt
|
|
logger.warning(f"Request failed (attempt {attempt + 1}/{request.max_retries + 1}), retrying in {delay}s: {e}")
|
|
time.sleep(delay)
|
|
else:
|
|
logger.error(f"All {request.max_retries + 1} attempts failed for {request.sport} {request.year}")
|
|
|
|
raise last_exception
|
|
|
|
def get_result(self, request_id: str) -> Optional[FetchResult]:
|
|
"""
|
|
Get the result of a fetch request.
|
|
|
|
Args:
|
|
request_id: Request ID to get result for
|
|
|
|
Returns:
|
|
Fetch result if available, None otherwise
|
|
"""
|
|
# Periodic cleanup
|
|
self._cleanup_completed_requests()
|
|
|
|
with self._lock:
|
|
return self.completed_requests.get(request_id)
|
|
|
|
def is_request_complete(self, request_id: str) -> bool:
|
|
"""
|
|
Check if a request has completed.
|
|
|
|
Args:
|
|
request_id: Request ID to check
|
|
|
|
Returns:
|
|
True if request is complete, False otherwise
|
|
"""
|
|
# Periodic cleanup
|
|
self._cleanup_completed_requests()
|
|
|
|
with self._lock:
|
|
return request_id in self.completed_requests
|
|
|
|
def get_request_status(self, request_id: str) -> Optional[FetchStatus]:
|
|
"""
|
|
Get the status of a fetch request.
|
|
|
|
Args:
|
|
request_id: Request ID to get status for
|
|
|
|
Returns:
|
|
Request status if found, None otherwise
|
|
"""
|
|
with self._lock:
|
|
if request_id in self.active_requests:
|
|
return self.active_requests[request_id].status
|
|
elif request_id in self.completed_requests:
|
|
result = self.completed_requests[request_id]
|
|
if result.final_status is not None:
|
|
return result.final_status
|
|
return FetchStatus.COMPLETED if result.success else FetchStatus.FAILED
|
|
return None
|
|
|
|
def cancel_request(self, request_id: str) -> bool:
|
|
"""
|
|
Cancel a pending or in-progress request.
|
|
|
|
Args:
|
|
request_id: Request ID to cancel
|
|
|
|
Returns:
|
|
True if request was cancelled, False if not found or already complete
|
|
"""
|
|
with self._lock:
|
|
if request_id in self.active_requests:
|
|
request = self.active_requests[request_id]
|
|
if request.commit_claimed:
|
|
# Too late: the worker holds an authorised commit. Report
|
|
# the failure rather than half-cancelling a request whose
|
|
# data is about to land in the cache.
|
|
logger.debug(
|
|
"Not cancelling %s: its response is already being "
|
|
"committed", request_id
|
|
)
|
|
return False
|
|
request.status = FetchStatus.CANCELLED
|
|
del self.active_requests[request_id]
|
|
# Cancelling is the other way a request leaves active_requests,
|
|
# so the in-flight index has to be released here too or the key
|
|
# stays pointed at a request that no longer exists.
|
|
if self._inflight_by_cache_key.get(request.cache_key) == request_id:
|
|
del self._inflight_by_cache_key[request.cache_key]
|
|
logger.info(f"Cancelled request {request_id}")
|
|
return True
|
|
return False
|
|
|
|
def get_statistics(self) -> Dict[str, Any]:
|
|
"""
|
|
Get service statistics.
|
|
|
|
Returns:
|
|
Dictionary containing service statistics
|
|
"""
|
|
with self._lock:
|
|
return {
|
|
**self.stats,
|
|
'active_requests': len(self.active_requests),
|
|
'completed_requests_count': len(self.completed_requests),
|
|
'max_completed_requests': self._max_completed_requests,
|
|
'completed_requests_usage_percent': (len(self.completed_requests) / self._max_completed_requests * 100) if self._max_completed_requests > 0 else 0,
|
|
# Nothing is queued outside the executor; kept for callers
|
|
# that read the key.
|
|
'queue_size': 0,
|
|
'last_cleanup': self._last_completed_requests_cleanup,
|
|
'cleanup_interval': self._completed_requests_cleanup_interval
|
|
}
|
|
|
|
def log_memory_stats(self):
|
|
"""Log current memory usage statistics."""
|
|
stats = self.get_statistics()
|
|
logger.info(f"BackgroundDataService Memory - Active: {stats['active_requests']}, "
|
|
f"Completed: {stats['completed_requests_count']}/{stats['max_completed_requests']} "
|
|
f"({stats['completed_requests_usage_percent']:.1f}%), "
|
|
f"Last cleanup: {time.time() - stats['last_cleanup']:.1f}s ago")
|
|
|
|
def _cleanup_completed_requests(self, force: bool = False) -> int:
|
|
"""
|
|
Automatically clean up old completed requests.
|
|
|
|
Args:
|
|
force: If True, perform cleanup regardless of time interval
|
|
|
|
Returns:
|
|
Number of requests removed
|
|
"""
|
|
now = time.time()
|
|
|
|
# Check if cleanup is needed
|
|
if not force and (now - self._last_completed_requests_cleanup) < self._completed_requests_cleanup_interval:
|
|
return 0
|
|
|
|
with self._lock:
|
|
removed_count = 0
|
|
current_time = time.time()
|
|
|
|
# Remove requests older than 1 hour
|
|
cutoff_time = current_time - 3600 # 1 hour
|
|
|
|
to_remove = []
|
|
for request_id, result in self.completed_requests.items():
|
|
# Check if request is old enough to remove
|
|
if result.completed_at < cutoff_time:
|
|
to_remove.append(request_id)
|
|
|
|
# Also enforce size limit if we have too many requests
|
|
if len(self.completed_requests) > self._max_completed_requests:
|
|
# Sort by completion time (oldest first)
|
|
sorted_requests = sorted(
|
|
self.completed_requests.items(),
|
|
key=lambda x: x[1].completed_at
|
|
)
|
|
|
|
# Remove oldest entries until we're under the limit
|
|
excess_count = len(self.completed_requests) - self._max_completed_requests
|
|
for i in range(excess_count):
|
|
if i < len(sorted_requests):
|
|
request_id = sorted_requests[i][0]
|
|
if request_id not in to_remove:
|
|
to_remove.append(request_id)
|
|
|
|
# Remove the requests
|
|
for request_id in to_remove:
|
|
del self.completed_requests[request_id]
|
|
removed_count += 1
|
|
|
|
self._last_completed_requests_cleanup = current_time
|
|
|
|
if removed_count > 0:
|
|
logger.debug(f"Cleaned up {removed_count} old completed requests (remaining: {len(self.completed_requests)})")
|
|
|
|
return removed_count
|
|
|
|
def clear_completed_requests(self, older_than_hours: int = 24):
|
|
"""
|
|
Clear completed requests older than specified time.
|
|
|
|
Args:
|
|
older_than_hours: Clear requests older than this many hours
|
|
"""
|
|
cutoff_time = time.time() - (older_than_hours * 3600)
|
|
|
|
with self._lock:
|
|
to_remove = []
|
|
for request_id, result in self.completed_requests.items():
|
|
if result.completed_at < cutoff_time:
|
|
to_remove.append(request_id)
|
|
|
|
for request_id in to_remove:
|
|
del self.completed_requests[request_id]
|
|
|
|
if to_remove:
|
|
logger.info(f"Cleared {len(to_remove)} old completed requests")
|
|
|
|
def shutdown(self, wait: bool = True):
|
|
"""
|
|
Shutdown the background data service.
|
|
|
|
Args:
|
|
wait: Whether to wait for active requests to complete
|
|
"""
|
|
logger.info("Shutting down BackgroundDataService...")
|
|
|
|
self._shutdown = True
|
|
|
|
# Cancel all active requests
|
|
with self._lock:
|
|
for request_id in list(self.active_requests.keys()):
|
|
self.cancel_request(request_id)
|
|
|
|
self.executor.shutdown(wait=wait)
|
|
|
|
logger.info("BackgroundDataService shutdown complete")
|
|
|
|
def __del__(self):
|
|
"""Cleanup when service is destroyed."""
|
|
if not self._shutdown:
|
|
self.shutdown(wait=False)
|
|
|
|
# Global service instance
|
|
_background_service: Optional[BackgroundDataService] = None
|
|
_service_lock = threading.Lock()
|
|
|
|
def get_background_service(cache_manager=None, max_workers: int = 3) -> BackgroundDataService:
|
|
"""
|
|
Get the global background data service instance.
|
|
|
|
Args:
|
|
cache_manager: Cache manager instance (required for first call)
|
|
max_workers: Maximum number of background threads
|
|
|
|
Returns:
|
|
Background data service instance
|
|
"""
|
|
global _background_service
|
|
|
|
with _service_lock:
|
|
if _background_service is None:
|
|
if cache_manager is None:
|
|
raise ValueError("cache_manager is required for first call to get_background_service")
|
|
_background_service = BackgroundDataService(cache_manager, max_workers)
|
|
|
|
return _background_service
|
|
|
|
def shutdown_background_service():
|
|
"""Shutdown the global background data service."""
|
|
global _background_service
|
|
|
|
with _service_lock:
|
|
if _background_service is not None:
|
|
_background_service.shutdown()
|
|
_background_service = None
|