mirror of
https://github.com/ChuckBuilds/LEDMatrix.git
synced 2026-10-04 06:15:09 +00:00
* fix(cache): store keys too long to be a filename The calendar plugin's cache key joins every calendar id the user picked. On hdpi it passed 300 bytes; ext4 caps a filename at 255, so every write (the temp file, the direct-write fallback and the home-directory fallback) failed with ENAMETOOLONG, once an hour, and the final warning said "(permission denied)" whatever the error was. DiskCache.get_cache_path keeps a key of up to 200 UTF-8 bytes as its filename, exactly as before, and turns a longer one into its first 183 bytes (cut on a character boundary) plus a 16-hex-digit hash of the whole key. The temp file adds 15 bytes, so the longest name is 215. The shortened stem is itself short, so the web UI's cache list, which names a key by its filename, deletes the same file. The give-up warning now names the real error. Validated on ledpi's ext4: the old module drops the hdpi-shaped key, the new one writes a 205-byte filename and reads it back. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(cache): judge a memory hit by the record's own timestamp A record loaded from disk went into the memory tier timed from the load, so get(key, max_age=300) could return data close to 600 s old: after a restart, after the memory sweep, or in a second process. A stored ttl was stretched the same way. #728's _fresh_cached works around it for the scoreboard; every other caller was exposed. get_cached_data and load_cache now also check a memory hit against the record's embedded timestamp, with DiskCache.get's rule that a stored ttl wins over max_age. A stale copy is dropped and the read falls through to disk, which returns the other process's newer write if there is one. Records without a timestamp keep the memory tier's own clock. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
970 lines
46 KiB
Python
970 lines
46 KiB
Python
"""
|
|
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 ".<stem>.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":<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*[,}]'
|
|
)
|
|
|
|
|
|
# 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
|
|
".<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
|
|
|