mirror of
https://github.com/ChuckBuilds/LEDMatrix.git
synced 2026-10-04 22:35:08 +00:00
* fix(errors): record the exception's own stack trace record_error() called traceback.format_exc(), which only sees an exception while its except block is running. plugin_executor records exceptions caught on a worker thread after that block has ended, so every trace on /errors read "NoneType: None". The trace is now built from the exception's __traceback__. The executor's log call had the same problem with exc_info=True and now passes the exception. record_error() also merged LEDMatrixError context into the caller's dict in place; it now works on a copy. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * docs(wifi): point at configure_wifi_permissions.sh instead of a sudoers list The module docstring told users to grant NOPASSWD sudo on iptables and ip. configure_wifi_permissions.sh refuses those grants on purpose: a wildcard rule for either runs an arbitrary program as root. Point at the script and say why it leaves them out. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(wifi): disconnect finds the saved profile by SSID disconnect_from_network() asked `nmcli -f NAME,802-11-wireless.ssid connection show` for the profile to take down, but nmcli rejects that column for `connection show`, so the lookup always failed and only the device was disconnected. The per-profile lookup _connect_nmcli() already used is now _find_profile_for_ssid(), and both callers share it. It also splits terse output on the last colon and unescapes "\:", so a profile name containing a colon is found. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(wifi): write wifi_config.json atomically and report a failed save _save_config() opened the file for writing in place and swallowed any error, so a wifi_config.json left owned by root made the web toggle for auto-enabling AP mode report success while nothing was saved, and a crash mid-write could truncate the file. It now uses atomic_write_json, which also keeps the file's owner and shared group when root saves it, and returns False on failure. POST /wifi/ap/auto-enable answers 500 in that case. The file is now written with indent=4, like the other config files. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(fonts): resolve plugin:// fonts in the plugin's own directory FontManager looked for a plugin's bundled fonts under Path("plugins") / plugin_id: relative to the process cwd, and not the default install directory (plugin-repos/), so a manifest's plugin:// fonts never loaded. register_plugin_fonts() takes an optional plugin_dir, and PluginManager passes the directory it loaded the plugin from. Callers that omit it get a lookup in the configured plugin_system.plugins_directory, then plugins/, resolved against the install root. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(api-helper): cache responses for the requested cache_ttl APIHelper.get(cache_ttl=...) and set_cache(ttl=...) dropped the ttl on the claim that CacheManager does not support one, but CacheManager.set() takes a ttl, stores it with the entry, and both cache tiers honour it over a reader's max_age. Without it every response expired after the 300-second default read age, whatever the plugin asked for. The ttl is now passed through, and the cache read passes cache_ttl as max_age for entries written without one. The class docstring describes what the helper actually does. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(style): one scale range for the schema, element_scale and LogoHelper The generated Scale field allowed 0.1 to 10, element_style's reader capped at 10 with no floor, and LogoHelper accepted 0.05 to 8 and reset anything else to 1.0. A logo scale of 9, which the form accepts, drew at the shipped size. MIN_ELEMENT_SCALE / MAX_ELEMENT_SCALE (0.1, 10.0) in src.element_style are now the schema bounds and the clamp every reader applies through coerce_scale(): a positive number outside the range is clamped, and anything that is not a finite positive number means the default. That also stops element_scale() passing NaN through, since min(nan, 10.0) is nan. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(logos): placeholder lands at the requested path; empty logos list download_missing_logo() wrote its fallback placeholder to <normalize_abbreviation(abbr)>.png in the logo directory rather than to the logo_path the caller passed, so it could return True while nothing existed where the plugin looks (e.g. "TA&M.png" vs "TAANDM.png"). create_placeholder_logo() takes an optional filepath, and download_missing_logo passes the requested one. download_missing_logo_for_team() only caught KeyError, so a team whose "logos" list is empty raised IndexError; it now treats KeyError, IndexError and TypeError as "no logo URL". The placeholder is drawn with PLACEHOLDER_SIZE / PLACEHOLDER_BG, the constants is_placeholder_logo() recognises it by, instead of repeated literals. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(fonts): resolve bundled font paths against the install root TextHelper's default font_dir, the logo placeholder's font and FontManager's font_overrides.json were all relative to the process cwd, so a process started anywhere but the install root (the plugin safety harness, a manual run, a unit without WorkingDirectory) drew with PIL's default face and read no overrides. They now go through font_layout.resolve_asset_path; the overrides file sits in the install root's config/. The resolver docstrings described an order the code does not follow: resolve_asset_path never consults the cwd, and sports_shared's _resolve_font_path tries the cwd first. Both docstrings now say what the code does, and _resolve_font_path calls resolve_asset_path instead of probing FontManager for it. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(sync): the web UI reads the sync status file the display writes sync_manager writes its status to tempfile.gettempdir(), but GET /api/v3/sync/status read a hardcoded /tmp/led_matrix_sync_status.json and defaulted the port to a literal 5765. Wherever TMPDIR is set (or on any non-/tmp host) the page only ever showed "starting". The endpoint now uses sync_manager.STATUS_FILE and SYNC_PORT. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(http): the rankings resolver sends the project's User-Agent DynamicTeamResolver fetched ESPN rankings with a bare requests.get, so it sent python-requests' default User-Agent, which ESPN rejects; the AP_TOP_N favourites then resolved to nothing. It now sends DEFAULT_HTTP_HEADERS. BaseOddsManager carried its own copy of the User-Agent string and now uses the same shared headers (which also adds Accept-Language). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(backup): record the core release and read the configured plugin dir The manifest's ledmatrix_version came from a VERSION file that does not exist, then from .git/HEAD: a 12-character sha, or "ref: refs/he" when the branch's ref was packed. It is now src.__version__. list_installed_plugins() scanned a hardcoded plugin-repos/, so on an install whose plugin_system.plugins_directory points elsewhere, plugins missing from plugin_state.json were left out of the backup. It now reads the configured directory from config/config.json, defaulting to plugin-repos. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(startup): report a missing display section once A config without a display section produced three errors for the one problem ("Missing required configuration key: display", "Display configuration is missing or empty" and "Display configuration is missing"), and an empty one produced two. _validate_config now reports it once, as a missing key or an empty section, and _validate_display_config leaves it to that. The module docstring said the validator fails fast; nothing in the display service calls raise_on_errors(), so it now says the errors are reported and startup continues. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * refactor(wifi): share the copied blocks and name the AP constants - _parse_nmcli_wifi_list() is the one parser behind _scan_nmcli and _scan_nmcli_cached. - _verify_connected(), _wait_for_device_idle(), _failsafe_ap() and _mark_forced() replace blocks that were pasted two or three times in the connect and enable-AP paths. The device-idle wait now checks before its first one-second sleep instead of after it. - _check_command() calls _find_command_path() instead of repeating it. - AP_IP, PORTAL_PORT, AP_PROFILE_NAME and AP_PROFILE_NAMES name values that were spelled out 14, 12, 8 and 2 times; the two deletion loops now walk the same tuple. The iwconfig status path compares the AP address exactly: startswith() also skipped 192.168.4.10-19. - Dropped a second WIFI.SIGNAL query that repeated the first, a no-op "if ssid: continue", the try/except around _connect_wpa_supplicant's constant return, and a second save of a scan scan_networks already saves. - _ensure_wifi_radio_enabled's docstring says it returns True when the radio state cannot be read at all. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * refactor(config): drop dead branches and history comments in ConfigManager - The module docstring pointed plugin authors at update_plugin_config(), which does not exist; it now names save_config_atomic() and save_raw_file_content(). - load_config's FileNotFoundError handler tested the message for "config_secrets.json", but a missing secrets file is handled where it is read, so only config.json reaches it; the check is gone. - save_raw_file_content's `file_type == "main" or "secrets"` guard was always true (anything else raised earlier). - get_raw_file_content('secrets') already returns {} for a missing file, so the os.path.exists() in front of two calls to it is gone. - Comments that narrated earlier behaviour are rewritten as what the code does now. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * refactor(background-data): present-tense comments, drop unused API - Comments that told the history of each fix (what "used to" happen, "the old per-delivery release") now state the invariant the code keeps. - get_statistics() no longer reports a constant 'queue_size': 0, and the uncalled clear_completed_requests() is gone (_cleanup_completed_requests does that job on every completion). Neither is referenced in core, the web UI or the plugin monorepo. shutdown_background_service() has no production caller either, but it is the only way to tear down the get_background_service() singleton, which the tests rely on, so it stays. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * refactor(odds): drop the unread cache_ttl and merge the odds_data branches BaseOddsManager loaded base_odds_manager.cache_ttl from config and never used it: cached odds live for the update interval (get_odds' ttl=interval). No core or monorepo code reads the attribute, so it is gone along with its log line. The two consecutive `if odds_data:` blocks are one. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * refactor(backup): one table for the single-file sections config, secrets, wifi and ytm_auth were each spelled out in create, preview, validate and restore. _SINGLE_FILE_SECTIONS lists them once, with the RestoreOptions flag that restores each, and all four walk it. Restore error messages keep their wording ("Failed to restore <file name>"). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * refactor(fonts): drop FontManager's write-only state and duplicate logs - fonts_config, font_metadata and font_dependencies were written and never read; the performance_stats keys font_load_times, render_times, total_renders and the per-call "resolve" timings (_record_performance_metric) likewise. get_performance_stats() reads only the counters that remain. Nothing in core or the plugin monorepo references any of them. - A failed BDF load was logged twice, by _load_bdf_font and again by get_font; get_font's line is the one kept. - Removed "NEW:" and commented-out cozette entries, the "Copy font to assets/fonts" comment on code that copies nothing, and local imports of names the module already imports. The deprecated add_font() now resolves assets/fonts against the install root. The @deprecated methods stay. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * refactor(text-helper): cache loaded fonts; drop the pre-textlength fallback TextHelper declared _font_cache, cleared it and reported its size, but never stored anything in it. load_fonts() now keeps each (file, size) it loads there, so clear_font_cache() and get_font_cache_stats() mean what they say and repeated load_fonts() calls reuse the fonts. get_text_width() no longer catches AttributeError for Pillow releases without ImageDraw.textlength; requirements.txt pins Pillow>=12.2. The class docstring describes what the helper does. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * docs(common): fix wrong docstrings in api_helper, permission_utils, snapshot_policy - permission_utils called 0o2775 "sticky bit"; the 2 is setgid, which is what makes new files take the directory's group. - snapshot_policy pointed at web_interface/blueprints/api_v3.py, which is a package now; the health check is in api_v3/misc.py. - APIHelper.clear_cache() lost a history note and a fallback to a clear() method that neither CacheManager nor the testing MockCacheManager has. The session headers are built from DEFAULT_HTTP_HEADERS instead of a copy of them, and the module docstring says what the module offers. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * docs(sports): present-tense comments in the shared scoreboard renderers - sports_scroll and sports_game_renderer comments that referred to "this PR", "the old flat 128px card" or what the renderer "previously" did now describe the current behaviour and its reason. - The block explaining why non-finite settings are rejected sat above _score_reserve_width; it describes _center_gap_width and now lives in it. - unshare_element_fonts wrapped its import of font_layout.load_truetype in an `except ImportError` that cannot fire inside core; the import stays at call time so tests can spy on the pinned loader. - sports_card docstrings that told the history of a fix say what the code does. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * refactor(sports-shared): drop dead code, name the ESPN limit - _get_weeks_data asked for limit=1000, which fetch_espn_scoreboard clamps to ESPN_MAX_LIMIT anyway; it now names that constant. Its unused `immediate_events = []` is gone. - _get_season_schedule_dates() returned ("", "") and has no caller in core or the plugin monorepo. - _should_log keeps its warning_type parameter (part of the inherited signature, though nothing in core or the monorepo calls it) and its docstring says the cooldown is shared across types. - An unused ImageFont import is gone. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * refactor(sync): one follower-mode switch, shared panel defaults - The class docstring said the leader sends PNG frames. Frames go over UDP as raw RGB; PNG is only the Vegas scroll image sent over TCP. It now describes both paths. - _enter_follower_mode() replaces the two copies of "note the leader, switch from standalone to follower, log, write status" in the frame and scroll-position handlers. - The rows/cols fallbacks use DEFAULT_ROWS / DEFAULT_COLS from src.display_geometry, as chain_length already did. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * refactor(style): drop _layout_axis, name the layout group title - ElementStyleResolver._layout_axis() had no caller in core or the plugin monorepo. - _element_block_from_spec checked spec['size'] was a dict again after size_spec already had; it reads size_spec. - The "Layout Offsets" title written into three generated schema blocks is _LAYOUT_TITLE. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * docs(logo-helper): say what the placeholder draws; name the 1.5 box factor - _create_placeholder_logo's docstring said it draws the team abbreviation; it draws an outlined grey box and nothing else. The docstring says so, and the "in a real implementation you'd want text" comments are gone. - The 1.5 x panel default logo box, written out six times, is DEFAULT_LOGO_BOX_FACTOR. - ImageDraw is imported with Image at the top of the module. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * refactor(logos): drop dead code and a duplicate regex in logo_downloader - _SAFE_LEAGUE_CODE_RE was the same pattern as _SAFE_LEAGUE_RE; both checks use the one. - get_logo_filename_variations reassigned the TA&M case to the list it already had; the function returns the two names directly. - _get_team_name_variations() had no caller in core or the plugin monorepo. - fetch_single_team's docstring was copied from fetch_teams_data; a log message read "for{team_id}". Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * refactor: drop the Pillow<9.1 resample shim and a catch-and-reraise - adaptive_images fell back to Image.LANCZOS/NEAREST for Pillow < 9.1; requirements.txt pins Pillow>=12.2. RESAMPLE_LANCZOS and RESAMPLE_NEAREST keep their names (src.common re-exports them). - CacheManager.save_cache caught CacheError only to re-raise it; the disk write is now called directly, with the same result. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * test(api-helper): stop the real CacheManager's cleanup thread The cache-lifetime tests built a CacheManager and left its cleanup thread's class-wide claim on the directory in place, which broke test_cache_cleanup_thread_ownership when it ran later in the session. The fixture now stops the thread on teardown. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * docs(changelog): core-common Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
949 lines
42 KiB
Python
949 lines
42 KiB
Python
"""
|
|
Cache Manager — multi-tier response cache for the LEDMatrix application.
|
|
|
|
:class:`CacheManager` provides a unified caching layer used by all plugins
|
|
to reduce external API calls and survive network outages gracefully.
|
|
|
|
Two storage tiers
|
|
-----------------
|
|
* **Memory tier** (:class:`~src.cache.memory_cache.MemoryCache`): fast LRU
|
|
cache (up to 1 000 entries by default). Hit on this tier before touching
|
|
disk.
|
|
* **Disk tier** (:class:`~src.cache.disk_cache.DiskCache`): filesystem-backed
|
|
persistent store that survives process restarts.
|
|
|
|
Data written to cache is serialised as JSON. :class:`DateTimeEncoder` handles
|
|
``datetime`` objects transparently so callers don't have to pre-serialise them.
|
|
|
|
Typical plugin usage::
|
|
|
|
data = self.cache_manager.get_cached_data('my_key', max_age=300)
|
|
if data is None:
|
|
data = fetch_from_api()
|
|
self.cache_manager.save_cache('my_key', data)
|
|
"""
|
|
|
|
import json
|
|
import os
|
|
import time
|
|
from datetime import datetime
|
|
import pytz
|
|
from typing import Any, Dict, List, Optional
|
|
import logging
|
|
import threading
|
|
import tempfile
|
|
from src.cache.memory_cache import MemoryCache, default_max_size
|
|
from src.cache.disk_cache import DiskCache
|
|
from src.cache.cache_strategy import CacheStrategy
|
|
from src.cache.cache_metrics import CacheMetrics
|
|
from src.logging_config import get_logger
|
|
from src.deprecation import deprecated
|
|
|
|
# Canonical implementation lives in src.cache.disk_cache; re-exported here
|
|
# because this module's docstring documents it and external code may import
|
|
# it from either path.
|
|
from src.cache.disk_cache import DateTimeEncoder # noqa: F401 - deliberate re-export
|
|
|
|
class CacheManager:
|
|
"""Manages caching of API responses to reduce API calls."""
|
|
|
|
# Which cache directories already have a cleanup thread in this process.
|
|
#
|
|
# The sweep is directory-scoped work -- it lists a directory and deletes
|
|
# from it -- so one per directory is the right number no matter how many
|
|
# managers exist. Nothing enforced that before: every instance started its
|
|
# own, and because the loop closes over `self`, a discarded manager could
|
|
# never be collected and its thread woke to re-scan the same directory
|
|
# every 24 hours for the life of the process. Startup validation runs
|
|
# twice and built a throwaway manager each time, so a display process
|
|
# carried three threads for one cache.
|
|
_cleanup_owners: Dict[str, 'CacheManager'] = {}
|
|
_cleanup_owners_lock = threading.Lock()
|
|
|
|
|
|
def __init__(self) -> None:
|
|
# Initialize logger first
|
|
self.logger: logging.Logger = get_logger(__name__)
|
|
|
|
# Determine the most reliable writable directory
|
|
self.cache_dir: Optional[str] = self._get_writable_cache_dir()
|
|
if self.cache_dir:
|
|
self.logger.info(f"Using cache directory: {self.cache_dir}")
|
|
else:
|
|
# This is a critical failure, as caching is essential.
|
|
self.logger.error("Could not find or create a writable cache directory. Caching will be disabled.")
|
|
self.cache_dir = None
|
|
|
|
# Initialize config manager for sport-specific intervals
|
|
try:
|
|
from src.config_manager import ConfigManager
|
|
self.config_manager: Optional[Any] = ConfigManager()
|
|
self.config_manager.load_config()
|
|
except ImportError:
|
|
self.config_manager: Optional[Any] = None
|
|
self.logger.warning("ConfigManager not available, using default cache intervals")
|
|
|
|
# Initialize cache components using composition
|
|
self._memory_cache_component = MemoryCache(
|
|
max_size=default_max_size(), cleanup_interval=300.0
|
|
)
|
|
self._disk_cache_component = DiskCache(cache_dir=self.cache_dir, logger=self.logger)
|
|
self._strategy_component = CacheStrategy(config_manager=self.config_manager, logger=self.logger)
|
|
self._metrics_component = CacheMetrics(logger=self.logger)
|
|
|
|
# Disk cleanup configuration
|
|
self._disk_cleanup_interval_hours = 24 # Run cleanup every 24 hours
|
|
self._disk_cleanup_interval = 3600.0 # Minimum interval between cleanups (1 hour) for throttle
|
|
self._last_disk_cleanup = 0.0 # Timestamp of last disk cleanup
|
|
self._cleanup_thread: Optional[threading.Thread] = None
|
|
self._cleanup_stop_event = threading.Event() # Event to signal thread shutdown
|
|
self._retention_policies = {
|
|
'odds': 2, # Odds data: 2 days (lines move frequently)
|
|
'odds_live': 2, # Live odds: 2 days
|
|
'sports_live': 7, # Live sports: 7 days
|
|
'weather_current': 7, # Current weather: 7 days
|
|
'sports_recent': 7, # Recent games: 7 days
|
|
'news': 14, # News: 14 days
|
|
'sports_upcoming': 60, # Upcoming games: 60 days (schedules stable)
|
|
'sports_schedules': 60, # Schedules: 60 days
|
|
'team_info': 60, # Team info: 60 days
|
|
'stocks': 14, # Stock data: 14 days
|
|
'crypto': 14, # Crypto data: 14 days
|
|
'default': 30 # Default: 30 days
|
|
}
|
|
|
|
# Start background cleanup thread only if disk caching is enabled
|
|
if self.cache_dir:
|
|
self.start_cleanup_thread()
|
|
|
|
def _get_writable_cache_dir(self) -> Optional[str]:
|
|
"""Tries to find or create a writable cache directory, preferring a system path when available."""
|
|
# Attempt 1: System-wide persistent cache directory (preferred for services)
|
|
try:
|
|
system_cache_dir = '/var/cache/ledmatrix'
|
|
if os.path.exists(system_cache_dir):
|
|
test_file = os.path.join(system_cache_dir, '.writetest')
|
|
try:
|
|
with open(test_file, 'w') as f:
|
|
f.write('test')
|
|
os.remove(test_file)
|
|
self.logger.info(f"Using system cache directory: {system_cache_dir}")
|
|
return system_cache_dir
|
|
except (IOError, OSError):
|
|
self.logger.debug(f"System cache directory exists but is not writable: {system_cache_dir}")
|
|
else:
|
|
from pathlib import Path
|
|
from src.common.permission_utils import (
|
|
ensure_directory_permissions,
|
|
get_cache_dir_mode
|
|
)
|
|
try:
|
|
ensure_directory_permissions(Path(system_cache_dir), get_cache_dir_mode())
|
|
if os.access(system_cache_dir, os.W_OK):
|
|
self.logger.info(f"Using system cache directory: {system_cache_dir}")
|
|
return system_cache_dir
|
|
except (OSError, IOError, PermissionError):
|
|
# Permission errors are expected when running as non-root
|
|
self.logger.debug(f"Could not create system cache directory (permission denied): {system_cache_dir}")
|
|
except (OSError, IOError, PermissionError) as e:
|
|
# Permission errors are expected when running as non-root, log at DEBUG level
|
|
self.logger.debug(f"System cache directory not available: {e}")
|
|
|
|
# Attempt 2: User's home directory (handling sudo), but avoid /root preference
|
|
try:
|
|
real_user = os.environ.get('SUDO_USER') or os.environ.get('USER', 'default')
|
|
if real_user and real_user != 'root':
|
|
home_dir = os.path.expanduser(f"~{real_user}")
|
|
else:
|
|
# When running as root and /var/cache/ledmatrix failed, still allow fallback to /root
|
|
home_dir = os.path.expanduser('~')
|
|
user_cache_dir = os.path.join(home_dir, '.ledmatrix_cache')
|
|
from pathlib import Path
|
|
from src.common.permission_utils import (
|
|
ensure_directory_permissions,
|
|
get_cache_dir_mode
|
|
)
|
|
ensure_directory_permissions(Path(user_cache_dir), get_cache_dir_mode())
|
|
test_file = os.path.join(user_cache_dir, '.writetest')
|
|
with open(test_file, 'w') as f:
|
|
f.write('test')
|
|
os.remove(test_file)
|
|
self.logger.info(f"Using user cache directory: {user_cache_dir}")
|
|
return user_cache_dir
|
|
except (OSError, IOError, PermissionError) as e:
|
|
self.logger.warning(f"Could not use user-specific cache directory: {e}")
|
|
|
|
# Attempt 3: /opt/ledmatrix/cache (alternative persistent location)
|
|
try:
|
|
opt_cache_dir = '/opt/ledmatrix/cache'
|
|
|
|
# Check if directory exists and we can write to it
|
|
if os.path.exists(opt_cache_dir):
|
|
# Test if we can write to the existing directory
|
|
test_file = os.path.join(opt_cache_dir, '.writetest')
|
|
try:
|
|
with open(test_file, 'w') as f:
|
|
f.write('test')
|
|
os.remove(test_file)
|
|
return opt_cache_dir
|
|
except (IOError, OSError):
|
|
self.logger.warning(f"Directory exists but is not writable: {opt_cache_dir}")
|
|
else:
|
|
# Try to create the directory
|
|
from pathlib import Path
|
|
from src.common.permission_utils import (
|
|
ensure_directory_permissions,
|
|
get_cache_dir_mode
|
|
)
|
|
ensure_directory_permissions(Path(opt_cache_dir), get_cache_dir_mode())
|
|
if os.access(opt_cache_dir, os.W_OK):
|
|
return opt_cache_dir
|
|
except (OSError, IOError, PermissionError) as e:
|
|
self.logger.warning(f"Could not use /opt/ledmatrix/cache: {e}", exc_info=True)
|
|
|
|
# Attempt 4: System-wide temporary directory (fallback, not persistent)
|
|
try:
|
|
temp_cache_dir = os.path.join(tempfile.gettempdir(), 'ledmatrix_cache')
|
|
from pathlib import Path
|
|
from src.common.permission_utils import (
|
|
ensure_directory_permissions,
|
|
get_cache_dir_mode
|
|
)
|
|
ensure_directory_permissions(Path(temp_cache_dir), get_cache_dir_mode())
|
|
if os.access(temp_cache_dir, os.W_OK):
|
|
self.logger.warning("Using temporary cache directory - cache will NOT persist across restarts")
|
|
return temp_cache_dir
|
|
except (OSError, IOError, PermissionError) as e:
|
|
self.logger.warning(f"Could not use system-wide temporary cache directory: {e}", exc_info=True)
|
|
|
|
# Return None if no directory is writable
|
|
return None
|
|
|
|
def _cleanup_memory_cache(self, force: bool = False) -> int:
|
|
"""Sweep the memory tier: drop entries older than an hour and trim it
|
|
to its size ceiling, at most once per cleanup interval unless forced.
|
|
|
|
Returns:
|
|
Number of entries removed
|
|
"""
|
|
return self._memory_cache_component.cleanup(force=force)
|
|
|
|
def _get_cache_path(self, key: str) -> Optional[str]:
|
|
"""Get the path for a cache file."""
|
|
return self._disk_cache_component.get_cache_path(key)
|
|
|
|
def get_cached_data(self, key: str, max_age: int = 300, memory_ttl: Optional[int] = None) -> Optional[Dict[str, Any]]:
|
|
"""Get data from cache (memory first, then disk) honoring TTLs.
|
|
|
|
- memory_ttl: TTL for in-memory entry; defaults to max_age if not provided
|
|
- max_age: TTL for persisted (on-disk) entry based on the stored timestamp
|
|
"""
|
|
# Periodic cleanup of memory cache
|
|
self._cleanup_memory_cache()
|
|
|
|
in_memory_ttl = memory_ttl if memory_ttl is not None else max_age
|
|
|
|
# 1) Memory cache
|
|
cached = self._memory_cache_component.get(key, max_age=in_memory_ttl)
|
|
if cached is not None:
|
|
return cached
|
|
|
|
# 2) Disk cache
|
|
record = self._disk_cache_component.get(key, max_age=max_age)
|
|
if record is not None:
|
|
# Hydrate memory cache (use current time to start memory TTL window)
|
|
self._memory_cache_component.set(key, record)
|
|
return record
|
|
|
|
# 3) Miss
|
|
return None
|
|
|
|
def save_cache(self, key: str, data: Dict[str, Any]) -> None:
|
|
"""
|
|
Save data to cache.
|
|
Args:
|
|
key: Cache key
|
|
data: Data to cache
|
|
"""
|
|
# Periodic cleanup before adding new entries
|
|
self._cleanup_memory_cache()
|
|
|
|
# Update memory cache first
|
|
self._memory_cache_component.set(key, data)
|
|
|
|
# DiskCache logs a failed write and raises CacheError, which the
|
|
# caller gets as is.
|
|
self._disk_cache_component.set(key, data)
|
|
|
|
def load_cache(self, key: str) -> Optional[Dict[str, Any]]:
|
|
"""Load data from cache with memory caching."""
|
|
# Check memory cache first (1 minute TTL)
|
|
cached = self._memory_cache_component.get(key, max_age=60)
|
|
if cached is not None:
|
|
return cached
|
|
|
|
# Check disk cache
|
|
data = self._disk_cache_component.get(key, max_age=3600) # 1 hour for load_cache
|
|
if data is not None:
|
|
# Update memory cache
|
|
self._memory_cache_component.set(key, data)
|
|
return data
|
|
|
|
return None
|
|
|
|
def clear_cache(self, key: Optional[str] = None) -> None:
|
|
"""Clear cache entries.
|
|
|
|
Pass a non-empty ``key`` to remove a single entry, or pass
|
|
``None`` (the default) to clear every cached entry. An empty
|
|
string is rejected to prevent accidental whole-cache wipes
|
|
from callers that pass through unvalidated input.
|
|
"""
|
|
if key is None:
|
|
# Clear all keys
|
|
memory_count = self._memory_cache_component.size()
|
|
self._memory_cache_component.clear()
|
|
self._disk_cache_component.clear()
|
|
self.logger.info("Cleared all cache: %d memory entries", memory_count)
|
|
return
|
|
|
|
if not isinstance(key, str) or not key:
|
|
raise ValueError(
|
|
"clear_cache(key) requires a non-empty string; "
|
|
"pass key=None to clear all entries"
|
|
)
|
|
|
|
# Clear specific key
|
|
self._memory_cache_component.clear(key)
|
|
self._disk_cache_component.clear(key)
|
|
self.logger.info("Cleared cache for key: %s", key)
|
|
|
|
def delete(self, key: str) -> None:
|
|
"""Remove a single cache entry.
|
|
|
|
Thin wrapper around :meth:`clear_cache` that **requires** a
|
|
non-empty string key — unlike ``clear_cache(None)`` it never
|
|
wipes every entry. Raises ``ValueError`` on ``None`` or an
|
|
empty string.
|
|
"""
|
|
if key is None or not isinstance(key, str) or not key:
|
|
raise ValueError("delete(key) requires a non-empty string key")
|
|
self.clear_cache(key)
|
|
|
|
def list_cache_files(self) -> List[Dict[str, Any]]:
|
|
"""List all cache files with metadata (key, age, size, path).
|
|
|
|
Returns:
|
|
List of dicts with keys: 'key', 'filename', 'age_seconds', 'age_display',
|
|
'size_bytes', 'size_display', 'path', 'modified_time'
|
|
"""
|
|
if not self.cache_dir or not os.path.exists(self.cache_dir):
|
|
return []
|
|
|
|
cache_files = []
|
|
current_time = time.time()
|
|
|
|
try:
|
|
# No lock: this is disk-only work, and the memory-tier lock it used
|
|
# to hold would stall every get/set while thousands of files are
|
|
# stat'd. A file deleted mid-scan is skipped below.
|
|
for filename in os.listdir(self.cache_dir):
|
|
if not filename.endswith('.json'):
|
|
continue
|
|
|
|
# Extract key from filename (remove .json extension)
|
|
key = filename[:-5] # Remove '.json'
|
|
|
|
file_path = os.path.join(self.cache_dir, filename)
|
|
|
|
try:
|
|
# Get file stats
|
|
stat_info = os.stat(file_path)
|
|
size_bytes = stat_info.st_size
|
|
modified_time = stat_info.st_mtime
|
|
age_seconds = current_time - modified_time
|
|
|
|
# Format age display
|
|
if age_seconds < 60:
|
|
age_display = f"{int(age_seconds)}s"
|
|
elif age_seconds < 3600:
|
|
age_display = f"{int(age_seconds / 60)}m"
|
|
elif age_seconds < 86400:
|
|
age_display = f"{int(age_seconds / 3600)}h"
|
|
else:
|
|
age_display = f"{int(age_seconds / 86400)}d"
|
|
|
|
# Format size display
|
|
if size_bytes < 1024:
|
|
size_display = f"{size_bytes}B"
|
|
elif size_bytes < 1024 * 1024:
|
|
size_display = f"{size_bytes / 1024:.1f}KB"
|
|
else:
|
|
size_display = f"{size_bytes / (1024 * 1024):.1f}MB"
|
|
|
|
cache_files.append({
|
|
'key': key,
|
|
'filename': filename,
|
|
'age_seconds': age_seconds,
|
|
'age_display': age_display,
|
|
'size_bytes': size_bytes,
|
|
'size_display': size_display,
|
|
'path': file_path,
|
|
'modified_time': modified_time,
|
|
'modified_datetime': datetime.fromtimestamp(modified_time).isoformat()
|
|
})
|
|
except OSError as e:
|
|
self.logger.warning(f"Error getting stats for cache file {filename} at {file_path}: {e}", exc_info=True)
|
|
continue
|
|
|
|
except OSError as e:
|
|
self.logger.error(f"Error listing cache directory {self.cache_dir}: {e}", exc_info=True)
|
|
return []
|
|
|
|
# Sort by modified time (newest first)
|
|
cache_files.sort(key=lambda x: x['modified_time'], reverse=True)
|
|
return cache_files
|
|
|
|
def get_cache_dir(self) -> Optional[str]:
|
|
"""Get the cache directory path."""
|
|
return self.cache_dir
|
|
|
|
@deprecated("3.7.0")
|
|
def has_data_changed(self, data_type: str, new_data: Dict[str, Any]) -> bool:
|
|
"""Check if data has changed from cached version."""
|
|
cached_data = self.load_cache(data_type)
|
|
if not cached_data:
|
|
return True
|
|
|
|
if data_type == 'weather':
|
|
return self._has_weather_changed(cached_data, new_data)
|
|
elif data_type == 'stocks':
|
|
return self._has_stocks_changed(cached_data, new_data)
|
|
elif data_type == 'stock_news':
|
|
return self._has_news_changed(cached_data, new_data)
|
|
elif data_type == 'nhl':
|
|
return self._has_nhl_changed(cached_data, new_data)
|
|
elif data_type == 'mlb':
|
|
return self._has_mlb_changed(cached_data, new_data)
|
|
|
|
return True
|
|
|
|
def _has_weather_changed(self, cached: Dict[str, Any], new: Dict[str, Any]) -> bool:
|
|
"""Check if weather data has changed."""
|
|
# Handle new cache structure where data is nested under 'data' key
|
|
if 'data' in cached:
|
|
cached = cached['data']
|
|
|
|
# Handle case where cached data might be the weather data directly
|
|
if 'current' in cached:
|
|
# This is the new structure with 'current' and 'forecast' keys
|
|
current_weather = cached.get('current', {})
|
|
if current_weather and 'main' in current_weather and 'weather' in current_weather:
|
|
cached_temp = round(current_weather['main']['temp'])
|
|
cached_condition = current_weather['weather'][0]['main']
|
|
return (cached_temp != new.get('temp') or
|
|
cached_condition != new.get('condition'))
|
|
|
|
# Handle old structure where temp and condition are directly accessible
|
|
return (cached.get('temp') != new.get('temp') or
|
|
cached.get('condition') != new.get('condition'))
|
|
|
|
def _has_stocks_changed(self, cached: Dict[str, Any], new: Dict[str, Any]) -> bool:
|
|
"""Check if stock data has changed."""
|
|
if not self._is_market_open():
|
|
return False
|
|
return cached.get('price') != new.get('price')
|
|
|
|
def _has_news_changed(self, cached: Dict[str, Any], new: Dict[str, Any]) -> bool:
|
|
"""Check if news data has changed."""
|
|
# Handle both dictionary and list formats
|
|
if isinstance(new, list):
|
|
# If new data is a list, cached data should also be a list
|
|
if not isinstance(cached, list):
|
|
return True
|
|
# Compare lengths and content
|
|
if len(cached) != len(new):
|
|
return True
|
|
# Compare titles since they're unique enough for our purposes
|
|
cached_titles = set(item.get('title', '') for item in cached)
|
|
new_titles = set(item.get('title', '') for item in new)
|
|
return cached_titles != new_titles
|
|
else:
|
|
# Original dictionary format handling
|
|
cached_headlines = set(h.get('id') for h in cached.get('headlines', []))
|
|
new_headlines = set(h.get('id') for h in new.get('headlines', []))
|
|
return not cached_headlines.issuperset(new_headlines)
|
|
|
|
def _has_nhl_changed(self, cached: Dict[str, Any], new: Dict[str, Any]) -> bool:
|
|
"""Check if NHL data has changed."""
|
|
return (cached.get('game_status') != new.get('game_status') or
|
|
cached.get('score') != new.get('score'))
|
|
|
|
def _has_mlb_changed(self, cached: Dict[str, Any], new: Dict[str, Any]) -> bool:
|
|
"""Check if MLB game data has changed."""
|
|
if not cached or not new:
|
|
return True
|
|
|
|
# Check if any games have changed status or score
|
|
for game_id, new_game in new.items():
|
|
cached_game = cached.get(game_id)
|
|
if not cached_game:
|
|
return True
|
|
|
|
# Check for score changes
|
|
if (new_game['away_score'] != cached_game['away_score'] or
|
|
new_game['home_score'] != cached_game['home_score']):
|
|
return True
|
|
|
|
# Check for status changes
|
|
if new_game['status'] != cached_game['status']:
|
|
return True
|
|
|
|
# For live games, check inning and count
|
|
if new_game['status'] == 'in':
|
|
if (new_game['inning'] != cached_game['inning'] or
|
|
new_game['inning_half'] != cached_game['inning_half'] or
|
|
new_game['balls'] != cached_game['balls'] or
|
|
new_game['strikes'] != cached_game['strikes'] or
|
|
new_game['bases_occupied'] != cached_game['bases_occupied']):
|
|
return True
|
|
|
|
return False
|
|
|
|
def _is_market_open(self) -> bool:
|
|
"""Check if the US stock market is currently open."""
|
|
return self._strategy_component.is_market_open()
|
|
|
|
@deprecated("3.7.0", "use set()")
|
|
def update_cache(self, data_type: str, data: Dict[str, Any]) -> bool:
|
|
"""Update cache with new data."""
|
|
cache_data = {
|
|
# Header first; see DiskCache's stale check.
|
|
'timestamp': time.time(),
|
|
'data': data,
|
|
}
|
|
return self.save_cache(data_type, cache_data)
|
|
|
|
def get(self, key: str, max_age: Optional[int] = 300,
|
|
memory_ttl: Optional[int] = None) -> Optional[Dict[str, Any]]:
|
|
"""Get data from cache if it exists and is not stale.
|
|
|
|
Args:
|
|
key: Cache key
|
|
max_age: Max age (seconds) for the on-disk entry; None never expires.
|
|
memory_ttl: Max age (seconds) for the in-memory entry. Pass 0 to
|
|
bypass the memory tier and force a fresh read from disk — used by
|
|
cross-process readers that must observe another process's latest
|
|
write rather than a stale first snapshot. Defaults to max_age.
|
|
"""
|
|
cached_data = self.get_cached_data(key, max_age, memory_ttl=memory_ttl)
|
|
if cached_data and 'data' in cached_data:
|
|
return cached_data['data']
|
|
return cached_data
|
|
|
|
def set(self, key: str, data: Dict[str, Any], ttl: Optional[int] = None) -> None:
|
|
"""
|
|
Store data in cache with current timestamp.
|
|
|
|
Args:
|
|
key: Cache key
|
|
data: Data to cache
|
|
ttl: Time-to-live in seconds for this entry. Takes precedence over
|
|
the max_age a reader would otherwise apply, which is inferred
|
|
from the key and is only a fallback for entries that did not
|
|
say. Omit it to keep that inferred behaviour.
|
|
"""
|
|
# timestamp and ttl before data, so they are the first bytes on disk:
|
|
# DiskCache.get reads them from the head of the file and can call a
|
|
# record stale without parsing it. That matters for the big ones -- a
|
|
# whole MLB season is 53MB and ~1.8s of orjson.loads with the GIL held,
|
|
# paid in full only to learn the record had expired.
|
|
cache_data: Dict[str, Any] = {'timestamp': time.time()}
|
|
if ttl is not None:
|
|
cache_data['ttl'] = ttl
|
|
cache_data['data'] = data
|
|
self.save_cache(key, cache_data)
|
|
|
|
@deprecated("3.7.0")
|
|
def setup_persistent_cache(self) -> bool:
|
|
"""
|
|
Set up a persistent cache directory with proper permissions.
|
|
This should be run once with sudo to create the directory.
|
|
"""
|
|
try:
|
|
# Try to create /var/cache/ledmatrix with proper permissions
|
|
from pathlib import Path
|
|
from src.common.permission_utils import (
|
|
ensure_directory_permissions,
|
|
get_cache_dir_mode
|
|
)
|
|
cache_dir = '/var/cache/ledmatrix'
|
|
cache_dir_path = Path(cache_dir)
|
|
ensure_directory_permissions(cache_dir_path, get_cache_dir_mode())
|
|
|
|
# Set ownership to the real user (not root)
|
|
real_user = os.environ.get('SUDO_USER')
|
|
if real_user:
|
|
import pwd
|
|
try:
|
|
uid = pwd.getpwnam(real_user).pw_uid
|
|
gid = pwd.getpwnam(real_user).pw_gid
|
|
os.chown(cache_dir, uid, gid)
|
|
self.logger.info(f"Set ownership of {cache_dir} to {real_user}")
|
|
except (OSError, KeyError) as e:
|
|
self.logger.warning(f"Could not set ownership for {cache_dir}: {e}", exc_info=True)
|
|
|
|
self.logger.info(f"Successfully set up persistent cache directory: {cache_dir}")
|
|
return True
|
|
|
|
except (OSError, IOError, PermissionError) as e:
|
|
self.logger.error(f"Failed to set up persistent cache directory {cache_dir}: {e}", exc_info=True)
|
|
return False
|
|
|
|
def cleanup_disk_cache(self, force: bool = False) -> Dict[str, Any]:
|
|
"""
|
|
Clean up expired disk cache files based on retention policies.
|
|
|
|
Args:
|
|
force: If True, run cleanup regardless of last cleanup time
|
|
|
|
Returns:
|
|
Dictionary with cleanup statistics
|
|
"""
|
|
now = time.time()
|
|
|
|
# Check if cleanup is needed (throttle to prevent too-frequent cleanups)
|
|
if not force and (now - self._last_disk_cleanup) < self._disk_cleanup_interval:
|
|
return {
|
|
'files_scanned': 0,
|
|
'files_deleted': 0,
|
|
'space_freed_mb': 0.0,
|
|
'errors': 0,
|
|
'duration_sec': 0.0
|
|
}
|
|
|
|
start_time = time.time()
|
|
|
|
try:
|
|
# Perform cleanup
|
|
stats = self._disk_cache_component.cleanup_expired_files(
|
|
cache_strategy=self._strategy_component,
|
|
retention_policies=self._retention_policies
|
|
)
|
|
|
|
duration = time.time() - start_time
|
|
space_freed_mb = stats['space_freed_bytes'] / (1024 * 1024)
|
|
|
|
# Record metrics
|
|
self._metrics_component.record_disk_cleanup(
|
|
files_cleaned=stats['files_deleted'],
|
|
space_freed_mb=space_freed_mb,
|
|
duration_sec=duration
|
|
)
|
|
|
|
# Log summary
|
|
if stats['files_deleted'] > 0:
|
|
self.logger.info(
|
|
"Disk cache cleanup completed: %d/%d files deleted, %.2f MB freed, %d errors, took %.2fs",
|
|
stats['files_deleted'], stats['files_scanned'], space_freed_mb,
|
|
stats['errors'], duration
|
|
)
|
|
else:
|
|
self.logger.debug(
|
|
"Disk cache cleanup completed: no files to delete (%d files scanned)",
|
|
stats['files_scanned']
|
|
)
|
|
|
|
# Update last cleanup time
|
|
self._last_disk_cleanup = time.time()
|
|
|
|
return {
|
|
'files_scanned': stats['files_scanned'],
|
|
'files_deleted': stats['files_deleted'],
|
|
'space_freed_mb': space_freed_mb,
|
|
'errors': stats['errors'],
|
|
'duration_sec': duration
|
|
}
|
|
|
|
except Exception as e:
|
|
self.logger.error("Error during disk cache cleanup: %s", e, exc_info=True)
|
|
return {
|
|
'files_scanned': 0,
|
|
'files_deleted': 0,
|
|
'space_freed_mb': 0.0,
|
|
'errors': 1,
|
|
'duration_sec': time.time() - start_time
|
|
}
|
|
|
|
def start_cleanup_thread(self) -> None:
|
|
"""Start background thread for periodic disk cache cleanup.
|
|
|
|
At most one thread per cache directory per process: the sweep is
|
|
directory-scoped, so a second one only duplicates the scan.
|
|
"""
|
|
if self._cleanup_thread and self._cleanup_thread.is_alive():
|
|
self.logger.debug("Cleanup thread already running")
|
|
return
|
|
|
|
with CacheManager._cleanup_owners_lock:
|
|
owner = CacheManager._cleanup_owners.get(self.cache_dir)
|
|
if owner is not None and owner is not self:
|
|
thread = owner._cleanup_thread
|
|
if thread is not None and thread.is_alive():
|
|
self.logger.debug(
|
|
"Cleanup thread for %s already owned by another cache "
|
|
"manager in this process; not starting a second",
|
|
self.cache_dir)
|
|
return
|
|
# The owner's thread died or was stopped -- take over.
|
|
CacheManager._cleanup_owners[self.cache_dir] = self
|
|
|
|
|
|
def cleanup_loop():
|
|
"""Background loop that runs cleanup periodically."""
|
|
self.logger.info("Disk cache cleanup thread started (interval: %d hours)",
|
|
self._disk_cleanup_interval_hours)
|
|
|
|
# Repair files an older version wrote unreadable by the web
|
|
# interface (see disk_cache.py, "SHARING CACHE FILES"). Once per
|
|
# directory per process, which is what this thread already is.
|
|
try:
|
|
self._disk_cache_component.share_existing_files()
|
|
except Exception as e:
|
|
self.logger.error("Error sharing existing cache files: %s", e, exc_info=True)
|
|
|
|
# Run initial cleanup on startup (deferred from __init__ to avoid blocking)
|
|
try:
|
|
self.logger.debug("Running initial disk cache cleanup")
|
|
self.cleanup_disk_cache()
|
|
except Exception as e:
|
|
self.logger.error("Error in initial cleanup: %s", e, exc_info=True)
|
|
|
|
# Main cleanup loop
|
|
while not self._cleanup_stop_event.is_set():
|
|
try:
|
|
# Sleep for the configured interval (interruptible)
|
|
sleep_seconds = self._disk_cleanup_interval_hours * 3600
|
|
if self._cleanup_stop_event.wait(timeout=sleep_seconds):
|
|
# Event was set, exit loop
|
|
break
|
|
|
|
# Run cleanup if not stopped
|
|
if not self._cleanup_stop_event.is_set():
|
|
self.logger.debug("Running scheduled disk cache cleanup")
|
|
self.cleanup_disk_cache()
|
|
|
|
except Exception as e:
|
|
self.logger.error("Error in cleanup thread: %s", e, exc_info=True)
|
|
# Continue running despite errors, but use interruptible sleep
|
|
if self._cleanup_stop_event.wait(timeout=60):
|
|
# Event was set during error recovery sleep, exit loop
|
|
break
|
|
|
|
self.logger.info("Disk cache cleanup thread stopped")
|
|
|
|
self._cleanup_stop_event.clear() # Reset event before starting thread
|
|
self._cleanup_thread = threading.Thread(target=cleanup_loop, daemon=True, name="DiskCacheCleanup")
|
|
self._cleanup_thread.start()
|
|
self.logger.info("Started disk cache cleanup background thread")
|
|
|
|
def stop_cleanup_thread(self) -> None:
|
|
"""
|
|
Stop the background cleanup thread gracefully.
|
|
|
|
Signals the thread to stop and waits for it to finish (with timeout).
|
|
This allows for clean shutdown during testing or application termination.
|
|
"""
|
|
# Release ownership first and unconditionally, so a manager that never
|
|
# started a thread (or whose thread already exited) cannot keep the
|
|
# directory claimed and block a live manager from sweeping it.
|
|
with CacheManager._cleanup_owners_lock:
|
|
if CacheManager._cleanup_owners.get(self.cache_dir) is self:
|
|
del CacheManager._cleanup_owners[self.cache_dir]
|
|
|
|
if not self._cleanup_thread or not self._cleanup_thread.is_alive():
|
|
self.logger.debug("Cleanup thread not running")
|
|
return
|
|
|
|
self.logger.info("Stopping disk cache cleanup thread...")
|
|
self._cleanup_stop_event.set() # Signal thread to stop
|
|
|
|
# Wait for thread to finish (with timeout to avoid hanging)
|
|
self._cleanup_thread.join(timeout=5.0)
|
|
|
|
if self._cleanup_thread.is_alive():
|
|
self.logger.warning("Cleanup thread did not stop within timeout, thread may still be running")
|
|
else:
|
|
self.logger.info("Disk cache cleanup thread stopped successfully")
|
|
|
|
@deprecated("3.7.0")
|
|
def get_sport_live_interval(self, sport_key: str) -> int:
|
|
"""
|
|
Get the live_update_interval for a specific sport from config.
|
|
Falls back to default values if config is not available.
|
|
"""
|
|
return self._strategy_component.get_sport_live_interval(sport_key)
|
|
|
|
def get_cache_strategy(self, data_type: str, sport_key: Optional[str] = None) -> Dict[str, Any]:
|
|
"""
|
|
Get cache strategy for different data types.
|
|
Now respects sport-specific live_update_interval configurations.
|
|
"""
|
|
return self._strategy_component.get_cache_strategy(data_type, sport_key)
|
|
|
|
def get_data_type_from_key(self, key: str) -> str:
|
|
"""
|
|
Determine the appropriate cache strategy based on the cache key.
|
|
This helps automatically select the right cache duration.
|
|
"""
|
|
return self._strategy_component.get_data_type_from_key(key)
|
|
|
|
@deprecated("3.7.0")
|
|
def get_sport_key_from_cache_key(self, key: str) -> Optional[str]:
|
|
"""
|
|
Extract sport key from cache key to determine appropriate live_update_interval.
|
|
"""
|
|
return self._strategy_component.get_sport_key_from_cache_key(key)
|
|
|
|
def get_cached_data_with_strategy(self, key: str, data_type: str = 'default') -> Optional[Dict[str, Any]]:
|
|
"""
|
|
Get data from cache using data-type-specific strategy.
|
|
Now respects sport-specific live_update_interval configurations.
|
|
"""
|
|
# Extract sport key for live sports data
|
|
sport_key = None
|
|
if data_type in ['sports_live', 'live_scores']:
|
|
sport_key = self._strategy_component.get_sport_key_from_cache_key(key)
|
|
|
|
strategy = self._strategy_component.get_cache_strategy(data_type, sport_key)
|
|
max_age = strategy['max_age']
|
|
memory_ttl = strategy.get('memory_ttl', max_age)
|
|
|
|
# For market data, check if market is open
|
|
if strategy.get('market_hours_only', False) and not self._strategy_component.is_market_open():
|
|
# During off-hours, extend cache duration
|
|
max_age *= 4 # 4x longer cache during off-hours
|
|
|
|
record = self.get_cached_data(key, max_age, memory_ttl)
|
|
# Unwrap if stored in { 'data': ..., 'timestamp': ... }
|
|
if isinstance(record, dict) and 'data' in record:
|
|
return record['data']
|
|
return record
|
|
|
|
def get_with_auto_strategy(self, key: str) -> Optional[Dict[str, Any]]:
|
|
"""
|
|
Get cached data using automatically determined strategy.
|
|
Now respects sport-specific live_update_interval configurations.
|
|
"""
|
|
data_type = self.get_data_type_from_key(key)
|
|
return self.get_cached_data_with_strategy(key, data_type)
|
|
|
|
@deprecated("3.7.0", "use get()")
|
|
def get_background_cached_data(self, key: str, sport_key: Optional[str] = None) -> Optional[Dict[str, Any]]:
|
|
"""
|
|
Get data from background service cache with appropriate strategy.
|
|
This method is specifically designed for Recent/Upcoming managers
|
|
to use data cached by the background service.
|
|
|
|
Args:
|
|
key: Cache key to retrieve
|
|
sport_key: Sport key for determining appropriate cache strategy
|
|
|
|
Returns:
|
|
Cached data if available and fresh, None otherwise
|
|
"""
|
|
# Determine the appropriate cache strategy
|
|
data_type = self.get_data_type_from_key(key)
|
|
strategy = self.get_cache_strategy(data_type, sport_key)
|
|
|
|
# For Recent/Upcoming managers, we want to use the background service cache
|
|
# which should have longer TTLs than the individual manager caches
|
|
max_age = strategy['max_age']
|
|
memory_ttl = strategy.get('memory_ttl', max_age)
|
|
|
|
# Get the cached data
|
|
cached_data = self.get_cached_data(key, max_age, memory_ttl)
|
|
|
|
if cached_data:
|
|
# Record cache hit for performance monitoring
|
|
self.record_cache_hit('background')
|
|
# Unwrap if stored in { 'data': ..., 'timestamp': ... } format
|
|
if isinstance(cached_data, dict) and 'data' in cached_data:
|
|
return cached_data['data']
|
|
return cached_data
|
|
|
|
# Record cache miss for performance monitoring
|
|
self.record_cache_miss('background')
|
|
return None
|
|
|
|
@deprecated("3.7.0", "use get()")
|
|
def is_background_data_available(self, key: str, sport_key: Optional[str] = None) -> bool:
|
|
"""
|
|
Check if background service has fresh data available.
|
|
This helps Recent/Upcoming managers determine if they should
|
|
wait for background data or fetch immediately.
|
|
"""
|
|
data_type = self.get_data_type_from_key(key)
|
|
strategy = self.get_cache_strategy(data_type, sport_key)
|
|
|
|
# Check if we have data that's still fresh according to background service TTL
|
|
cached_data = self.get_cached_data(key, strategy['max_age'])
|
|
return cached_data is not None
|
|
|
|
def generate_sport_cache_key(self, sport: str, date_str: Optional[str] = None) -> str:
|
|
"""
|
|
Centralized cache key generation for sports data.
|
|
This ensures consistent cache keys across background service and managers.
|
|
|
|
Args:
|
|
sport: Sport identifier (e.g., 'nba', 'nfl', 'ncaa_fb')
|
|
date_str: Date string in YYYYMMDD format. If None, uses current UTC date.
|
|
|
|
Returns:
|
|
Cache key in format: {sport}_{date}
|
|
"""
|
|
if date_str is None:
|
|
date_str = datetime.now(pytz.utc).strftime('%Y%m%d')
|
|
return f"{sport}_{date_str}"
|
|
|
|
@deprecated("3.7.0")
|
|
def record_cache_hit(self, cache_type: str = 'regular') -> None:
|
|
"""Record a cache hit for performance monitoring."""
|
|
self._metrics_component.record_hit(cache_type)
|
|
|
|
@deprecated("3.7.0")
|
|
def record_cache_miss(self, cache_type: str = 'regular') -> None:
|
|
"""Record a cache miss for performance monitoring."""
|
|
self._metrics_component.record_miss(cache_type)
|
|
|
|
@deprecated("3.7.0")
|
|
def record_fetch_time(self, duration: float) -> None:
|
|
"""Record fetch operation duration for performance monitoring."""
|
|
self._metrics_component.record_fetch_time(duration)
|
|
|
|
@deprecated("3.7.0")
|
|
def get_cache_metrics(self) -> Dict[str, Any]:
|
|
"""Get current cache performance metrics."""
|
|
return self._metrics_component.get_metrics()
|
|
|
|
@deprecated("3.7.0")
|
|
def log_cache_metrics(self) -> None:
|
|
"""Log current cache performance metrics."""
|
|
self._metrics_component.log_metrics()
|
|
|
|
@deprecated("3.7.0")
|
|
def get_memory_cache_stats(self) -> Dict[str, Any]:
|
|
"""
|
|
Get statistics about the memory cache.
|
|
|
|
Returns:
|
|
Dictionary with memory cache statistics
|
|
"""
|
|
return self._memory_cache_component.get_stats()
|
|
|
|
def log_memory_cache_stats(self) -> None:
|
|
"""Log current memory cache statistics."""
|
|
stats = self.get_memory_cache_stats()
|
|
self.logger.info(f"Memory Cache - Size: {stats['size']}/{stats['max_size']} "
|
|
f"({stats['usage_percent']:.1f}%), "
|
|
f"Last cleanup: {time.time() - stats['last_cleanup']:.1f}s ago") |