mirror of
https://github.com/ChuckBuilds/LEDMatrix.git
synced 2026-10-04 22:35:08 +00:00
The sports plugins cache whole season schedules: 53MB for MLB, 18MB for NHL, 17MB for NCAA baseball. On a Pi 4, orjson.loads of the MLB file takes ~1.8s with the GIL held, and every thread in the display service waits -- the stall watchdog caught the render thread frozen 0.5-1.3s with the interpreter itself blocked, right on these reads. When a season record expired, DiskCache.get paid that whole parse only to find the timestamp too old and throw the result away. CacheManager.set now writes timestamp and ttl ahead of the data, and DiskCache.get reads them from the first 256 bytes of the file, applying the same rule as before (a per-entry ttl wins over max_age; no limit means never stale). A record that is stale is refused without being parsed. Files in the old layout, and records from other writers, don't match the header and are parsed in full as before. Also: ESPN responses in the background data service and espn_dates are parsed with orjson when it is installed (src/common/json_body.py). The stdlib parser behind response.json() takes 3.1s on the MLB season against orjson's 1.8s, both with the GIL held. espn_dates imports it with a fallback, since plugins bundle copies of that module for older cores. Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
791 lines
36 KiB
Python
791 lines
36 KiB
Python
"""
|
|
Disk Cache
|
|
|
|
Handles persistent disk-based caching with atomic writes and error recovery.
|
|
"""
|
|
|
|
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
|
|
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
|
|
|
|
|
|
|
|
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":<float>,"ttl":<n>,``.
|
|
_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*[,}]'
|
|
)
|
|
|
|
|
|
def _stale_from_head(head: bytes, max_age: Optional[int], now: float) -> 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.
|
|
"""
|
|
match = _HEAD_RE.match(head)
|
|
if not match:
|
|
return False
|
|
try:
|
|
timestamp = float(match.group(1))
|
|
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:*). 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 -> adler32 of the last payload successfully written to the
|
|
# primary cache path; lets set() skip rewriting identical data
|
|
# (per-process only — worst case another process rewrites, never
|
|
# a missed write). Guarded by _lock.
|
|
self._write_digests: Dict[str, 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.
|
|
|
|
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"{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:
|
|
# 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()):
|
|
return None
|
|
f.seek(0)
|
|
record = _loads(f.read())
|
|
|
|
# Determine record timestamp (prefer embedded, else file mtime)
|
|
record_ts = None
|
|
if isinstance(record, dict):
|
|
record_ts = record.get('timestamp')
|
|
if record_ts is None:
|
|
try:
|
|
record_ts = os.path.getmtime(cache_path)
|
|
except OSError:
|
|
record_ts = None
|
|
|
|
if record_ts is not None:
|
|
try:
|
|
record_ts = float(record_ts)
|
|
except (TypeError, ValueError):
|
|
record_ts = None
|
|
|
|
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
|
|
|
|
digest = zlib.adler32(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).
|
|
# Refresh the file mtime so records that rely on it for TTL
|
|
# (no embedded 'timestamp') don't expire early; a metadata
|
|
# touch is journal-cheap compared to rewriting the data.
|
|
if self._write_digests.get(key) == digest:
|
|
try:
|
|
os.utime(cache_path, None)
|
|
return
|
|
except OSError:
|
|
# File vanished or perms changed — fall through and write
|
|
self._write_digests.pop(key, None)
|
|
|
|
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._write_digests[key] = digest
|
|
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._write_digests[key] = digest
|
|
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):
|
|
# 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
|
|
self.logger.warning(
|
|
"Could not write cache for key '%s' to %s (permission denied). "
|
|
"Cache will be unavailable for this key, but application will continue.",
|
|
key, cache_path
|
|
)
|
|
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 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
|
|
".<cache filename>." , 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 ".<key>.json.<random>", 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
|
|
|