""" Disk Cache Handles persistent disk-based caching with atomic writes and error recovery. """ import hashlib import json import math import os import re import stat import time import tempfile import logging import threading import zlib from typing import Dict, Any, Optional, Protocol, Tuple from datetime import datetime from src.common.path_safety import safe_path_component try: # optional: large speedup on the cache write path, see _dumps below import orjson except ImportError: # pragma: no cover - exercised on hosts without the wheel orjson = None # How old an abandoned write's temp file must be before the sweep removes it. # A real write holds its temp file for milliseconds, so an hour is far beyond # any in-flight write while still clearing the same day's debris. Deliberately # not tied to the retention policies: those describe how long data stays # useful, and a half-written file was never useful. _ORPHAN_TEMP_MAX_AGE_SECONDS = 3600 # Longest key, in UTF-8 bytes, used verbatim as a filename stem. ext4 caps a # name at 255 bytes and set()'s temp file is "..json.<8 random>", 15 # bytes longer than the stem, so anything near the cap could never be written: # the calendar plugin's key joins every calendar id and passed 300 bytes on a # real install, failing every write with ENAMETOOLONG. Longer keys keep this # many bytes as a readable prefix and end in a hash of the whole key. _MAX_KEY_FILENAME_BYTES = 200 _KEY_HASH_CHARS = 16 def _filename_stem(key: str) -> str: """The filename stem for a key that is already a safe path component. Short keys are used as they are, so every file already on disk keeps its name. A long one becomes its first bytes plus a hash of the full key: the prefix keeps the stem recognisable (and keeps the data-type words that cleanup's retention lookup reads from it), the hash keeps two keys that share a long prefix apart. The result is itself short, so a stem read back from a filename -- which is how the web UI names a key it deletes -- maps to the same file. """ encoded = key.encode('utf-8') if len(encoded) <= _MAX_KEY_FILENAME_BYTES: return key digest = hashlib.sha256(encoded).hexdigest()[:_KEY_HASH_CHARS] keep = _MAX_KEY_FILENAME_BYTES - _KEY_HASH_CHARS - 1 prefix = encoded[:keep].decode('utf-8', errors='ignore') return f"{prefix}-{digest}" class CacheStrategyProtocol(Protocol): """Protocol for cache strategy objects that categorize cache keys.""" def get_data_type_from_key(self, key: str) -> str: """ Determine the data type from a cache key. Args: key: Cache key Returns: Data type string for strategy lookup """ ... class DateTimeEncoder(json.JSONEncoder): """JSON encoder that handles datetime objects. Retained for the stdlib fallback path and for any caller importing it. """ def default(self, obj: Any) -> Any: if isinstance(obj, datetime): return obj.isoformat() return super().default(obj) def _datetime_default(obj: Any) -> Any: """Serialise datetimes exactly as DateTimeEncoder did.""" if isinstance(obj, datetime): return obj.isoformat() raise TypeError(f"Object of type {type(obj).__name__} is not JSON serializable") def _replace_nonfinite(obj: Any) -> Any: """Non-finite floats -> None, matching what ``orjson.dumps`` writes. Only reached once a strict pass has proved there is something to replace, so the ordinary write path never pays for this walk. """ if isinstance(obj, float): return obj if math.isfinite(obj) else None if isinstance(obj, dict): return {k: _replace_nonfinite(v) for k, v in obj.items()} if isinstance(obj, (list, tuple)): return [_replace_nonfinite(v) for v in obj] return obj # NON-FINITE FLOATS # ----------------- # JSON has no NaN or Infinity. The stdlib emits them anyway as an extension; # orjson refuses to and writes null. That divergence is not acceptable in a # cache whose files outlive the decision of which encoder is installed, so the # policy here is one behaviour on both paths: # # writing non-finite floats become null, whichever encoder is in use # reading files already on disk that carry the stdlib's NaN/Infinity # tokens stay readable, whichever encoder is in use # # Without the write half, installing orjson silently changed cached values. # Without the read half, installing orjson turned every legacy record holding a # NaN into a "corrupted cache file" that DiskCache.get logged as an error and # deleted. Both halves are covered by test/test_cache_nonfinite_floats.py. #: Enough of a record to hold its header: ``{"timestamp":,"ttl":,``. _HEAD_BYTES = 256 #: A record written with its header first (CacheManager.set does). Anything #: else -- older files with "data" first, records from other writers -- does not #: match and is parsed in full, as before. _HEAD_RE = re.compile( rb'\A\s*\{\s*"timestamp"\s*:\s*(-?[0-9][0-9.eE+-]*)\s*' rb'(?:,\s*"ttl"\s*:\s*(-?[0-9][0-9.eE+-]*))?\s*[,}]' ) # UNCHANGED RE-SAVES: THE FILE'S MTIME CARRIES THE NEWER TIMESTAMP # ---------------------------------------------------------------- # Plugins re-save unchanged API data every update cycle, and every one of # those saves was a full rewrite on the SD card. DiskCache.set skips the write # when the payload matches the last one it wrote for the key -- but # CacheManager.set stamps each record with time.time(), so for set() the # payload never matched and the skip never fired. # # The digest now leaves out a header-first record's timestamp, so an unchanged # set() is skipped. What the skip must not do is make the record look older # than it is: the timestamp inside the file is from the last real write, and # a reader in another process (the web interface, with memory_ttl=0) or after # a restart would call fresh data stale. So the newer timestamp goes where it # costs no data write -- the file's mtime -- and readers take a record's age # from the newer of the two. The invariant that makes that safe: # # a file's mtime is the timestamp of the newest record saved for its key # # real write mtime is set to the record's own timestamp, so a record saved # with an old timestamp (data as of some earlier time) cannot # borrow freshness from the moment it hit the disk # skip mtime is set to the skipped record's timestamp -- exactly what # a rewrite would have stored, without the rewrite # # Readers of the on-disk timestamp, all of which go through _effective_timestamp: # DiskCache.get (the header check and the full parse; it also returns the # record with 'timestamp' set to the effective value, so CacheManager.get's # max_age path, the memory tier hydrated from disk, and any plugin reading # record['timestamp'] all see it). Readers that use mtime alone already see the # newer value: the retention sweep below, CacheManager.list_cache_files (the # web UI's cache list). Nothing else opens cache files: web_interface and # scripts reach them only through CacheManager. # # Something other than this class can also move an mtime forward -- a copy # without -p, an rsync without -t, a `touch`. (backup_manager.py does not # back up or restore the cache directory, so the in-tree restore cannot.) That # must not make old data fresh, so the lift is bounded: a reader never takes # the mtime as more than _MAX_TIMESTAMP_LIFT past the embedded timestamp, and # set() rewrites the file for real once a skip would need more than that, so # an honest lift never reaches the bound. A file copied a day after it was # written therefore reads at most an hour fresher than its contents say, and a # 30-second live-score record from yesterday stays stale. CacheManager.set # records written before this change have mtime == write time == embedded # timestamp, give or take the write itself, and read exactly as before; a # file an older version wrote or touched later than its embedded timestamp # says reads at most the same hour fresher, once, until it is next saved. #: Longest a skipped write may stand in for a real one, and so the furthest a #: file's mtime is ever trusted past the record's own timestamp. Unchanged data #: is rewritten at least this often, at most once an hour per key instead of #: once per update cycle. _MAX_TIMESTAMP_LIFT = 3600.0 def _record_timestamp(value: Any) -> Optional[float]: """A record's timestamp as a finite float, or None if it has no usable one.""" if isinstance(value, bool) or not isinstance(value, (int, float)): return None value = float(value) return value if math.isfinite(value) else None def _effective_timestamp(embedded: float, mtime: Optional[float]) -> float: """When a record was last saved: its timestamp, or the file's mtime if a later unchanged save moved that forward -- never by more than _MAX_TIMESTAMP_LIFT. See "UNCHANGED RE-SAVES" above.""" if mtime is None: return embedded return max(embedded, min(mtime, embedded + _MAX_TIMESTAMP_LIFT)) def _stale_from_head(head: bytes, max_age: Optional[int], now: float, mtime: Optional[float] = None) -> bool: """True when a record's header alone shows it has expired. Mirrors the expiry rule in DiskCache.get: a per-entry ttl wins over the caller's max_age, and no limit at all means never stale. False whenever the header cannot be read, so the full parse decides as it always did. ``mtime`` is the file's, which may carry a newer save than the header does. """ match = _HEAD_RE.match(head) if not match: return False try: timestamp = _effective_timestamp(float(match.group(1)), mtime) limit = max_age if match.group(2) is not None: ttl = float(match.group(2)) if ttl >= 0: limit = ttl except ValueError: return False return limit is not None and (now - timestamp) > limit if orjson is not None: # Encoding the cache record dominated the background fetch worker: on a # Pi 4, stdlib json.dumps runs ~12ms per MB and holds the GIL for all of # it, which stalls the render thread mid-scroll. orjson measures ~7x # faster on the same payloads (11.9ms -> 1.6ms for 985KB). Decoding gains # far less (~1.3x on large payloads) because the cost there is building # the Python objects, not scanning the text, but it is still free to take. # # OPT_NON_STR_KEYS: stdlib json coerces int/float dict keys to strings; # orjson raises without this, and cache records do carry numeric keys. # OPT_PASSTHROUGH_DATETIME: orjson would otherwise emit its own RFC 3339 # form for datetimes instead of calling default(). Routing them through # _datetime_default keeps byte-for-byte parity with the records already # on disk. _DUMPS_OPTS = orjson.OPT_NON_STR_KEYS | orjson.OPT_PASSTHROUGH_DATETIME def _dumps(data: Any) -> bytes: return orjson.dumps(data, default=_datetime_default, option=_DUMPS_OPTS) def _loads(raw: bytes) -> Any: try: return orjson.loads(raw) except orjson.JSONDecodeError: # Legacy record written by the stdlib path, carrying NaN or # Infinity. Genuinely malformed files raise again from here, as # json.JSONDecodeError, which is what DiskCache.get expects. return json.loads(raw) else: def _dumps(data: Any) -> bytes: try: return json.dumps(data, cls=DateTimeEncoder, allow_nan=False).encode("utf-8") except ValueError: # allow_nan=False is what detects the non-finite values; the walk # runs only now that we know there is one to replace. return json.dumps(_replace_nonfinite(data), cls=DateTimeEncoder, allow_nan=False).encode("utf-8") def _loads(raw: bytes) -> Any: return json.loads(raw) # SHARING CACHE FILES BETWEEN THE TWO SERVICES # -------------------------------------------- # The display service runs as root and the web interface as the installing # user, and the web interface reads records only the display writes # (display_current_state, display_on_demand_state, plugin_metrics_snapshot). Files are # written 0660, so the web interface can read one only through its group. # # The installers rely on the directory's setgid bit to set that group. That is # not something the cache can count on: systemd's CacheDirectory=, which # ledmatrix-web.service carried until Sept 2026, re-owns the directory and # everything in it to the web user and its primary group whenever the # directory's owner does not match, and the setgid layout never survives that. # From then on every file root creates is root:root 0660, unreadable by the web # interface. Measured on one rig: 365 such files, and the web UI's display # status, on-demand state and plugin health all silently empty. # # So a cache file takes its group from the directory explicitly, whether or # not setgid is set. Only a group-writable directory counts as shared: that # group can already replace any file in it, so reading them grants nothing new. # # Everything here works on an open descriptor, never a path. The directory is # writable by the web user, so between a path check and a path operation that # user could put a symlink in the file's place, and root would then chown and # chmod whatever it points at. _CACHE_FILE_MODE = 0o660 def _shared_group(directory: str) -> Optional[int]: """The group a cache file in ``directory`` should carry, if it is shared.""" try: st = os.stat(directory) except OSError: return None if not st.st_mode & stat.S_IWGRP: return None return st.st_gid def _share_open_file(fd: int, group: Optional[int]) -> None: """Make an open cache file readable by the other service. Best effort.""" fchmod = getattr(os, 'fchmod', None) # absent on Windows before 3.13 if fchmod is not None: try: fchmod(fd, _CACHE_FILE_MODE) except OSError: pass fchown = getattr(os, 'fchown', None) # absent on Windows if fchown is None or group is None: return try: if os.fstat(fd).st_gid != group: fchown(fd, -1, group) except OSError: # Not a member of the directory's group and not root: nothing to do, # and the file keeps the group it was created with. pass class DiskCache: """Manages persistent disk-based cache.""" def __init__(self, cache_dir: Optional[str], logger: Optional[logging.Logger] = None) -> None: """ Initialize disk cache. Args: cache_dir: Directory for cache files (None = disabled) logger: Optional logger instance """ self.cache_dir = cache_dir self.logger = logger or logging.getLogger(__name__) self._lock = threading.Lock() # key -> what set() last put at the primary cache path: the adler32 of # the payload (less a header-first timestamp), the timestamp the file # holds (None for records without one), and the file's inode and size. # Lets set() skip rewriting identical data. Per-process only, and the # inode/size check means another process's write is never mistaken # for ours -- worst case a redundant write, never a missed one. # Guarded by _lock. self._write_digests: Dict[str, Tuple[int, Optional[float], int, int]] = {} def get_cache_path(self, key: str) -> Optional[str]: """ Get the path for a cache file. The key becomes a filename, so it has to be one. Keys reach this method from the web API -- POST /api/v3/cache/delete passes the request body's ``key`` straight through CacheManager.clear_cache to os.remove -- and a key of ``../../../../etc/whatever`` named a file well outside the cache directory. Every real key is the stem of a file already sitting flat in cache_dir (that is how list_cache_files derives them), so rejecting anything with a path component turns away only inputs that could never have been written here. A key too long to be a filename is shortened by _filename_stem. Args: key: Cache key Returns: Path to cache file, or None if cache is disabled or the key is not a usable filename """ if not self.cache_dir: return None safe_key = safe_path_component(key) if safe_key is None: self.logger.warning("Rejected unsafe cache key %r", key) return None return os.path.join(self.cache_dir, f"{_filename_stem(safe_key)}.json") def get(self, key: str, max_age: Optional[int] = 300) -> Optional[Dict[str, Any]]: """ Get data from disk cache. Args: key: Cache key max_age: Maximum age in seconds; None disables age-based expiry (the record never counts as stale). Mirrors MemoryCache.get. Returns: Cached data or None if not found or expired """ cache_path = self.get_cache_path(key) if not cache_path or not os.path.exists(cache_path): return None try: with self._lock: with open(cache_path, 'rb') as f: # The open file's mtime, not the path's: the file a skipped # write touched is the one being read. mtime = os.fstat(f.fileno()).st_mtime # Decide staleness from the header before paying for the # parse. A stale read is the common case for the biggest # records (a season schedule is re-fetched when its cache # expires), and parsing 53MB to throw it away held the GIL # for ~1.8s -- a visible freeze on the panel. if _stale_from_head(f.read(_HEAD_BYTES), max_age, time.time(), mtime): return None f.seek(0) record = _loads(f.read()) # Determine record timestamp: the embedded one, moved forward by a # later unchanged save if there was one (see "UNCHANGED RE-SAVES"), # else the file mtime. record_ts = None if isinstance(record, dict): record_ts = record.get('timestamp') if record_ts is None: record_ts = mtime else: embedded_ts = _record_timestamp(record_ts) if embedded_ts is None: try: record_ts = float(record_ts) except (TypeError, ValueError): record_ts = None else: record_ts = _effective_timestamp(embedded_ts, mtime) if record_ts != embedded_ts: # Hand the record back as a rewrite would have left it, # so callers that age it themselves agree with us. record['timestamp'] = record_ts now = time.time() # An explicit per-entry ttl wins over the caller's max_age. The # caller that wrote the record knows what its data is; max_age is # inferred from substrings in the key ("live", "odds", "stock") and # is only a fallback for records that never said. Until now the ttl # was stored and ignored, so `set(key, data, ttl=...)` did nothing # at all -- 48 plugin call sites and 4 in the core were writing a # number no read path consulted. effective_max_age = max_age if isinstance(record, dict): stored_ttl = record.get('ttl') if isinstance(stored_ttl, (int, float)) and not isinstance(stored_ttl, bool) \ and stored_ttl >= 0: effective_max_age = stored_ttl max_age = effective_max_age # max_age=None means "never expires" (mirrors MemoryCache and the # cache_manager docstring). Guard it explicitly — otherwise the # comparison below raises TypeError and the record is treated as a # miss, which silently breaks callers that persist long-lived state # via get(key, max_age=None) (e.g. plugin health/metrics that must # survive restarts and be read cross-process). if record_ts is None or max_age is None or (now - record_ts) <= max_age: return record else: # Stale on disk; keep file for potential diagnostics but treat as miss return None except json.JSONDecodeError as e: self.logger.error("Error parsing cache file for %s at %s: %s", key, cache_path, e, exc_info=True) # If the file is corrupted, remove it try: os.remove(cache_path) self.logger.info("Removed corrupted cache file: %s", cache_path) except OSError as remove_error: self.logger.warning("Could not remove corrupted cache file %s: %s", cache_path, remove_error) return None except PermissionError as e: # Permission errors are recoverable - cache just won't be available self.logger.warning("Permission denied loading cache for %s from %s: %s. Cache unavailable for this key.", key, cache_path, e) return None except (IOError, OSError) as e: self.logger.error("Error loading cache for %s from %s: %s", key, cache_path, e, exc_info=True) return None except Exception as e: self.logger.error("Unexpected error loading cache for %s from %s: %s", key, cache_path, e, exc_info=True) return None def set(self, key: str, data: Dict[str, Any]) -> None: """ Save data to disk cache with atomic write. This method gracefully handles permission errors. If the cache directory is not writable, it will log a warning and return silently rather than raising an exception. This allows the application to continue functioning even when running as a non-root user without write access to system cache directories. Args: key: Cache key data: Data to cache """ cache_path = self.get_cache_path(key) if not cache_path: return # Serialize once, compact (no indent): the payload is reused by every # write path below, and cache files are machine-read only — indenting # them just multiplied the bytes written to the SD card. try: payload = _dumps(data) except (TypeError, ValueError) as e: self.logger.warning("Cache data for key '%s' not serializable: %s", key, e) return timestamp = _record_timestamp(data.get('timestamp')) if isinstance(data, dict) else None # A header-first record (CacheManager.set's layout) is compared without # its timestamp, which differs on every save; see "UNCHANGED RE-SAVES". # Any other layout is compared whole, as before. head = _HEAD_RE.match(payload) if timestamp is not None else None digest = zlib.adler32(memoryview(payload)[head.end(1):] if head else payload) try: # Atomic write to avoid partial/corrupt files with self._lock: # Skip the disk entirely when this exact payload was already # written for this key (plugins re-save unchanged API data # every update cycle — each write is real SD-card wear). # A metadata touch is journal-cheap compared to rewriting # the data. if self._skip_unchanged(key, cache_path, digest, timestamp): return tmp_dir = os.path.dirname(cache_path) # Try to create temp file in cache directory first # If that fails due to permissions, fall back to direct write tmp_path = None fd = None try: # First try the cache directory if os.access(tmp_dir, os.W_OK): try: fd, tmp_path = tempfile.mkstemp(prefix=f".{os.path.basename(cache_path)}.", dir=tmp_dir) except (IOError, OSError, PermissionError): # If temp file creation fails, try direct write as fallback self.logger.warning("Could not create temp file in %s, using direct write for %s", tmp_dir, key) tmp_path = None fd = None else: # Directory not writable, use direct write self.logger.warning("Cache directory %s not writable, using direct write for %s", tmp_dir, key) tmp_path = None fd = None if tmp_path and fd is not None: # Atomic write with temp file. No fsync: os.replace # already guarantees readers never see a torn file, # and cache data is re-fetchable — forcing a disk # flush per write was the single biggest SD-card # wear source (dozens of fsyncs/min on API-heavy # installs) for data that can be re-downloaded. try: with os.fdopen(fd, 'wb') as tmp_file: tmp_file.write(payload) # Before the rename, not after: mkstemp # creates the file 0600, and a reader that # opened it in between was refused. _share_open_file(tmp_file.fileno(), _shared_group(tmp_dir)) os.replace(tmp_path, cache_path) self._remember_write(key, cache_path, digest, timestamp) finally: if os.path.exists(tmp_path): try: os.remove(tmp_path) except OSError: pass else: # Fallback: direct write (not atomic, but better than failing) try: with open(cache_path, 'wb') as cache_file: cache_file.write(payload) _share_open_file(cache_file.fileno(), _shared_group(tmp_dir)) self._remember_write(key, cache_path, digest, timestamp) self.logger.debug("Wrote cache for %s directly (non-atomic)", key) except (IOError, OSError, PermissionError) as write_error: # If direct write also fails, try fallback location self.logger.warning("Direct write failed for key '%s' to %s: %s", key, cache_path, write_error) raise # Re-raise to trigger fallback logic except (IOError, OSError, PermissionError) as primary_error: # Attempt one-time fallback write to user's home cache directory try: # Try user's home cache directory as fallback home_dir = os.path.expanduser('~') fallback_dir = os.path.join(home_dir, '.ledmatrix_cache') # Ensure fallback directory exists try: os.makedirs(fallback_dir, exist_ok=True) except (OSError, PermissionError): pass if os.path.isdir(fallback_dir) and os.access(fallback_dir, os.W_OK): # NOTE: no digest record here — the fallback file # is a different path, so future sets must keep # retrying the primary location. fallback_path = os.path.join(fallback_dir, os.path.basename(cache_path)) with open(fallback_path, 'wb') as tmp_file: tmp_file.write(payload) _share_open_file(tmp_file.fileno(), _shared_group(fallback_dir)) self.logger.debug("Cache wrote to fallback location: %s", fallback_path) return # Successfully wrote to fallback, exit gracefully except (IOError, OSError, PermissionError) as e2: self.logger.debug("Fallback cache write also failed for key '%s': %s", key, e2) # If all write attempts failed, log warning but don't raise exception # Cache is a performance optimization, not critical for operation. # Name the real error: this used to say "permission denied" # whatever happened, which sent a too-long filename off to # be debugged as a directory-ownership problem. self.logger.warning( "Could not write cache for key '%s' to %s (%s). " "Cache will be unavailable for this key, but application will continue.", key, cache_path, primary_error.strerror or primary_error ) return # Exit gracefully without raising exception except Exception as e: # For any other unexpected errors, log but don't crash self.logger.warning( "Unexpected error saving cache for key '%s' to %s: %s. " "Application will continue without caching for this key.", key, cache_path, e, exc_info=True ) return # Exit gracefully without raising exception def _skip_unchanged(self, key: str, cache_path: str, digest: int, timestamp: Optional[float]) -> bool: """Stand in for a write of an unchanged record by touching the file. True when the file already holds this record bar its timestamp and the touch landed; False means write it. The touch sets mtime to the record's timestamp -- what a rewrite would have stored -- or to now for a record without one, whose age readers already take from mtime. Caller holds _lock. """ last = self._write_digests.get(key) if last is None or last[0] != digest: return False _, written_ts, ino, size = last if (timestamp is None) != (written_ts is None): return False if timestamp is not None and written_ts is not None: # Never backwards (a rewrite would make the record older), and # never further than readers will trust the mtime: past that the # record is rewritten, so its own timestamp catches up. if not written_ts <= timestamp <= written_ts + _MAX_TIMESTAMP_LIFT: return False try: st = os.stat(cache_path) if (st.st_ino, st.st_size) != (ino, size): # Replaced since our write (another process, a restore): # its contents are not the ones the digest describes. self._write_digests.pop(key, None) return False # Setting an explicit time needs the file's owner; a file someone # else wrote fails here and is rewritten (as our own file) instead. os.utime(cache_path, None if timestamp is None else (timestamp, timestamp)) return True except OSError: # File vanished or perms changed — fall through and write self._write_digests.pop(key, None) return False def _remember_write(self, key: str, cache_path: str, digest: int, timestamp: Optional[float]) -> None: """After a real write: pin mtime to the record's timestamp and note what was written, so the next unchanged save can be skipped. Pinning keeps a record saved with an older timestamp from looking as fresh as the moment it was written (see "UNCHANGED RE-SAVES"); for CacheManager.set's records the two differ only by the write itself. A timestamp in the future is left alone, mtime already being older. Never raises: the data is on disk, and anything failing here only costs the next save its skip. Caller holds _lock. """ self._write_digests.pop(key, None) try: if timestamp is not None and timestamp <= time.time(): os.utime(cache_path, (timestamp, timestamp)) st = os.stat(cache_path) except OSError as e: self.logger.debug("Could not pin mtime of %s: %s", cache_path, e) return self._write_digests[key] = (digest, timestamp, st.st_ino, st.st_size) def clear(self, key: Optional[str] = None) -> None: """ Clear cache entry or all entries. Args: key: Specific key to clear, or None to clear all """ if not self.cache_dir: return with self._lock: if key: self._write_digests.pop(key, None) cache_path = self.get_cache_path(key) if cache_path and os.path.exists(cache_path): try: os.remove(cache_path) except OSError as e: self.logger.warning("Could not remove cache file %s: %s", cache_path, e) else: # Clear all cache files self._write_digests.clear() if os.path.exists(self.cache_dir): for filename in os.listdir(self.cache_dir): if filename.endswith('.json'): try: os.remove(os.path.join(self.cache_dir, filename)) except OSError as e: self.logger.warning("Could not remove cache file %s: %s", filename, e) def get_cache_dir(self) -> Optional[str]: """Get the cache directory path.""" return self.cache_dir def share_existing_files(self) -> int: """Give cache files already on disk the group set() now gives new ones. set() fixes every file it writes from now on; this repairs the ones an older version left behind as root:root, which the web interface cannot read until each key happens to be rewritten -- and some, like a plugin's metrics, may not be for a long time. Meant to run once per process, off the startup path. Only this process's own regular files are touched, and each one through a descriptor opened with O_NOFOLLOW and checked for a single link: the directory is writable by the web user, and a root process must not be steered into changing a file outside it. Returns: Number of files whose group or mode was changed. """ fchown = getattr(os, 'fchown', None) geteuid = getattr(os, 'geteuid', None) nofollow = getattr(os, 'O_NOFOLLOW', None) if not self.cache_dir or fchown is None or geteuid is None or nofollow is None: return 0 group = _shared_group(self.cache_dir) if group is None: return 0 euid = geteuid() changed = 0 try: entries = list(os.scandir(self.cache_dir)) except OSError as e: self.logger.debug("Could not scan %s to share cache files: %s", self.cache_dir, e) return 0 for entry in entries: if not entry.name.endswith('.json'): continue try: st = entry.stat(follow_symlinks=False) except OSError: continue if (not stat.S_ISREG(st.st_mode) or st.st_uid != euid or (st.st_gid == group and stat.S_IMODE(st.st_mode) == _CACHE_FILE_MODE)): continue try: fd = os.open(entry.path, os.O_RDONLY | nofollow | getattr(os, 'O_NONBLOCK', 0)) except OSError: continue try: st = os.fstat(fd) if not stat.S_ISREG(st.st_mode) or st.st_uid != euid or st.st_nlink != 1: continue _share_open_file(fd, group) st = os.fstat(fd) if st.st_gid == group and stat.S_IMODE(st.st_mode) == _CACHE_FILE_MODE: changed += 1 except OSError: continue finally: os.close(fd) if changed: self.logger.info( "Made %d cache file(s) in %s readable by the directory's group " "(gid %d) so the web interface can read them", changed, self.cache_dir, group) return changed @staticmethod def _is_orphaned_temp(filename: str) -> bool: """Whether a name is one of set()'s temp files rather than real data. Matches only what this class creates: mkstemp with a prefix of ".." , so ".weather.json.a1b2c3d4". The shape is checked rather than just the leading dot, because this predicate deletes things -- a stray dotfile someone left in the cache directory is not ours to remove, and a completed ".json" never is either. """ if not filename.startswith('.') or filename.endswith('.json'): return False head, sep, suffix = filename.rpartition('.json.') # head is the key (non-empty after the leading dot), suffix is # mkstemp's random component. return bool(sep) and len(head) > 1 and bool(suffix) def cleanup_expired_files(self, cache_strategy: CacheStrategyProtocol, retention_policies: Dict[str, int]) -> Dict[str, Any]: """ Clean up expired cache files based on retention policies. Args: cache_strategy: Object implementing CacheStrategyProtocol for categorizing files retention_policies: Dict mapping data types to retention days Returns: Dictionary with cleanup statistics: - files_scanned: Total files checked - files_deleted: Files removed - space_freed_bytes: Bytes freed - errors: Number of errors encountered """ if not self.cache_dir or not os.path.exists(self.cache_dir): self.logger.warning("Cache directory not available for cleanup") return {'files_scanned': 0, 'files_deleted': 0, 'space_freed_bytes': 0, 'errors': 0} stats = { 'files_scanned': 0, 'files_deleted': 0, 'space_freed_bytes': 0, 'errors': 0 } current_time = time.time() try: # Collect files to process outside the lock to avoid blocking cache operations # Only hold lock during directory listing to get snapshot of files try: with self._lock: # Get snapshot of files while holding lock briefly entries = os.listdir(self.cache_dir) except OSError as list_error: self.logger.error("Error listing cache directory %s: %s", self.cache_dir, list_error, exc_info=True) stats['errors'] += 1 return stats filenames = [f for f in entries if f.endswith('.json')] # Sweep temp files abandoned by a write that never finished. set() # removes its own in a finally, so these are the ones where the # process died between mkstemp and os.replace -- a SIGKILL, a lost # restart race, a power cut. Nothing ever collected them: they are # named "..json.", and the scan above only matches # names ending in .json, so they accumulated indefinitely. Measured # on a live rig: 76 files, 1,050 MB, 81% of the whole cache # directory, the oldest six months old. stats['orphan_temp_files_deleted'] = 0 for filename in (f for f in entries if self._is_orphaned_temp(f)): # Counted as scanned like any other candidate, so files_deleted # can never exceed files_scanned and the summary line reads # honestly ("77/8864", not "77/0"). stats['files_scanned'] += 1 path = os.path.join(self.cache_dir, filename) try: # An in-flight write lives for milliseconds, so anything # this old is certainly abandoned rather than in progress. if (current_time - os.path.getmtime(path)) <= _ORPHAN_TEMP_MAX_AGE_SECONDS: continue with self._lock: size = os.path.getsize(path) os.remove(path) stats['files_deleted'] += 1 stats['orphan_temp_files_deleted'] += 1 stats['space_freed_bytes'] += size except FileNotFoundError: continue # another sweep got there first except OSError as e: stats['errors'] += 1 self.logger.warning("Error deleting orphaned temp file %s: %s", filename, e) if stats['orphan_temp_files_deleted']: self.logger.info( "Removed %d abandoned cache temp file(s)", stats['orphan_temp_files_deleted']) # Process files outside the lock to avoid blocking get/set operations for filename in filenames: stats['files_scanned'] += 1 file_path = os.path.join(self.cache_dir, filename) try: # Get file age (outside lock - stat operations are generally atomic) file_mtime = os.path.getmtime(file_path) file_age_days = (current_time - file_mtime) / 86400 # Convert to days # Extract cache key from filename (remove .json extension) cache_key = filename[:-5] # Determine data type and retention policy data_type = cache_strategy.get_data_type_from_key(cache_key) retention_days = retention_policies.get(data_type, retention_policies.get('default', 30)) # Delete if older than retention period # Only hold lock during actual file deletion to ensure atomicity if file_age_days > retention_days: try: # Hold lock only during delete operation (get size and remove atomically) with self._lock: # Double-check file still exists (may have been deleted by another process) if os.path.exists(file_path): try: file_size = os.path.getsize(file_path) os.remove(file_path) # Only increment stats if removal succeeded stats['files_deleted'] += 1 stats['space_freed_bytes'] += file_size self.logger.debug( "Deleted expired cache file: %s (age: %.1f days, type: %s, retention: %d days)", filename, file_age_days, data_type, retention_days ) except FileNotFoundError: # File was deleted by another process between exists check and remove # This is a benign race condition, silently continue pass else: # File was deleted by another process before lock was acquired # This is a benign race condition, silently continue pass except FileNotFoundError: # File was already deleted by another process, skip it # This is a benign race condition, silently continue continue except OSError as e: # Other file system errors, log but don't fail the entire cleanup stats['errors'] += 1 self.logger.warning("Error deleting cache file %s: %s", filename, e) continue except FileNotFoundError: # File was deleted by another process between listing and processing # This is a benign race condition, silently continue continue except OSError as e: stats['errors'] += 1 self.logger.warning("Error processing cache file %s: %s", filename, e) continue except Exception as e: stats['errors'] += 1 self.logger.error("Unexpected error processing cache file %s: %s", filename, e, exc_info=True) continue except OSError as e: self.logger.error("Error listing cache directory %s: %s", self.cache_dir, e, exc_info=True) stats['errors'] += 1 return stats