mirror of
https://github.com/ChuckBuilds/LEDMatrix.git
synced 2026-10-04 14:25:08 +00:00
The disk cache's unchanged-payload skip now ignores a CacheManager.set() record's timestamp, so unchanged re-saves are skipped; a skip moves the file's mtime to the new timestamp instead, and readers take a record's age from the newer of the two (never more than an hour past the embedded timestamp). Per-plugin plugin_metrics:<id> records become one plugin_metrics_snapshot written at most once a minute, and CacheManager builds its ConfigManager on first use. On hdpi, cache file writes went from ~37 to 8.6 a minute. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
935 lines
44 KiB
Python
935 lines
44 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, 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
|
|
|
|
|
|
|
|
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.
|
|
|
|
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:
|
|
# 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):
|
|
# 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 _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
|
|
|