mirror of
https://github.com/ChuckBuilds/LEDMatrix.git
synced 2026-10-04 14:25:08 +00:00
fix(cache): skip rewriting unchanged data saved through CacheManager.set
DiskCache.set skipped a payload identical to the last one written for the key, but CacheManager.set stamps every record with time.time(), so the payload always differed and the skip never fired: unchanged API data was rewritten to the SD card on every plugin update cycle. Header-first records are now compared without their timestamp (the digest also carries the content length, since a collision is now a missed write). A skipped write moves the file's mtime to the skipped record's timestamp, and a real write sets it to the embedded one, so only a skip moves it forward. DiskCache.get, including the header fast path, treats such a record as fresh from the later of the two and returns that time as the record's timestamp. The skip also checks the file is still the one this process wrote (inode and size), so a file another process replaced is rewritten. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
Vendored
+113
-20
@@ -14,7 +14,7 @@ import tempfile
|
||||
import logging
|
||||
import threading
|
||||
import zlib
|
||||
from typing import Dict, Any, Optional, Protocol
|
||||
from typing import Dict, Any, Optional, Protocol, Tuple
|
||||
from datetime import datetime
|
||||
|
||||
from src.common.path_safety import safe_path_component
|
||||
@@ -111,18 +111,66 @@ _HEAD_RE = re.compile(
|
||||
)
|
||||
|
||||
|
||||
def _stale_from_head(head: bytes, max_age: Optional[int], now: float) -> bool:
|
||||
def _head_timestamp(head: bytes) -> Optional[Tuple[float, int]]:
|
||||
"""A header-first record's timestamp and the offset just past it.
|
||||
|
||||
None when the record does not start with a finite numeric timestamp.
|
||||
"""
|
||||
match = _HEAD_RE.match(head)
|
||||
if not match:
|
||||
return None
|
||||
try:
|
||||
timestamp = float(match.group(1))
|
||||
except ValueError:
|
||||
return None
|
||||
if not math.isfinite(timestamp):
|
||||
return None
|
||||
return timestamp, match.end(1)
|
||||
|
||||
|
||||
# FRESHNESS OF A SKIPPED WRITE
|
||||
# ----------------------------
|
||||
# CacheManager.set stamps every record with time.time(), so re-saving
|
||||
# unchanged data produced a different payload every time and DiskCache.set's
|
||||
# identical-payload skip never fired: every plugin rewrote its unchanged API
|
||||
# data to the SD card every update cycle. set() now compares header-first
|
||||
# records without their timestamp, and on a skip moves the file's mtime to
|
||||
# the timestamp the skipped record carried instead of rewriting it. So the
|
||||
# file's mtime is when its content was last saved, and a header-first
|
||||
# record is as fresh as the later of its embedded timestamp and its mtime.
|
||||
#
|
||||
# A real write sets the mtime to the embedded timestamp too, so mtime is
|
||||
# never later than the timestamp for a record written with an old one on
|
||||
# purpose -- only a skip can move it forward.
|
||||
|
||||
|
||||
def _refreshed_at(timestamp: float, mtime: float) -> float:
|
||||
"""When a header-first record was last saved, embedded time or mtime.
|
||||
|
||||
An mtime within a second of the timestamp is the write that carried it
|
||||
(float rounding, or an older file whose mtime was not set to match),
|
||||
not a skipped rewrite, and leaves the record as written.
|
||||
"""
|
||||
return mtime if mtime > timestamp + 1.0 else timestamp
|
||||
|
||||
|
||||
def _stale_from_head(head: bytes, max_age: Optional[int], now: float,
|
||||
refreshed: 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.
|
||||
``refreshed`` is the file's mtime: a skipped rewrite advances it rather
|
||||
than the embedded timestamp (see "FRESHNESS OF A SKIPPED WRITE").
|
||||
"""
|
||||
match = _HEAD_RE.match(head)
|
||||
if not match:
|
||||
return False
|
||||
try:
|
||||
timestamp = float(match.group(1))
|
||||
if refreshed is not None:
|
||||
timestamp = max(timestamp, refreshed)
|
||||
limit = max_age
|
||||
if match.group(2) is not None:
|
||||
ttl = float(match.group(2))
|
||||
@@ -248,11 +296,13 @@ class DiskCache:
|
||||
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] = {}
|
||||
# key -> ((length, adler32) of the last content written to the
|
||||
# primary cache path, (st_ino, st_size) of the file it left); lets
|
||||
# set() skip rewriting identical data. The file identity catches
|
||||
# another process -- the web interface writes and clears keys too --
|
||||
# having replaced the file since, which would otherwise make the skip
|
||||
# a missed write. Per-process only. Guarded by _lock.
|
||||
self._write_digests: Dict[str, Tuple[Tuple[int, int], Tuple[int, int]]] = {}
|
||||
|
||||
def get_cache_path(self, key: str) -> Optional[str]:
|
||||
"""
|
||||
@@ -306,7 +356,12 @@ class DiskCache:
|
||||
# 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()):
|
||||
head = f.read(_HEAD_BYTES)
|
||||
stamp = _head_timestamp(head)
|
||||
fresh_at = None
|
||||
if stamp is not None:
|
||||
fresh_at = _refreshed_at(stamp[0], os.fstat(f.fileno()).st_mtime)
|
||||
if _stale_from_head(head, max_age, time.time(), fresh_at):
|
||||
return None
|
||||
f.seek(0)
|
||||
record = _loads(f.read())
|
||||
@@ -315,6 +370,11 @@ class DiskCache:
|
||||
record_ts = None
|
||||
if isinstance(record, dict):
|
||||
record_ts = record.get('timestamp')
|
||||
if fresh_at is not None and fresh_at > stamp[0]:
|
||||
# A skipped rewrite refreshed this record (see "FRESHNESS
|
||||
# OF A SKIPPED WRITE"); hand callers the time it was last
|
||||
# saved, as the rewrite would have.
|
||||
record['timestamp'] = record_ts = fresh_at
|
||||
if record_ts is None:
|
||||
try:
|
||||
record_ts = os.path.getmtime(cache_path)
|
||||
@@ -403,24 +463,37 @@ class DiskCache:
|
||||
self.logger.warning("Cache data for key '%s' not serializable: %s", key, e)
|
||||
return
|
||||
|
||||
digest = zlib.adler32(payload)
|
||||
# A header-first record is compared without its timestamp, which
|
||||
# CacheManager.set changes on every call (see "FRESHNESS OF A SKIPPED
|
||||
# WRITE"). The length rides along with adler32, which is weak on its
|
||||
# own for short payloads, and a collision here is a missed write.
|
||||
stamp = _head_timestamp(payload[:_HEAD_BYTES])
|
||||
stamped_at = stamp[0] if stamp is not None else None
|
||||
content = memoryview(payload)[stamp[1]:] if stamp is not None else payload
|
||||
digest = (len(content), zlib.adler32(content))
|
||||
|
||||
try:
|
||||
# Atomic write to avoid partial/corrupt files
|
||||
with self._lock:
|
||||
# Skip the disk entirely when this exact payload was already
|
||||
# Skip the disk entirely when this content 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:
|
||||
# Move the file mtime instead, so the record stays as fresh as
|
||||
# the rewrite would have left it; a metadata touch is
|
||||
# journal-cheap compared to rewriting the data.
|
||||
known = self._write_digests.get(key)
|
||||
if known is not None and known[0] == digest:
|
||||
try:
|
||||
os.utime(cache_path, None)
|
||||
return
|
||||
st = os.stat(cache_path)
|
||||
if (st.st_ino, st.st_size) == known[1]:
|
||||
os.utime(cache_path, None if stamped_at is None
|
||||
else (stamped_at, stamped_at))
|
||||
return
|
||||
except OSError:
|
||||
# File vanished or perms changed — fall through and write
|
||||
self._write_digests.pop(key, None)
|
||||
pass
|
||||
# File vanished, was replaced by another process, or its
|
||||
# times cannot be set — 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
|
||||
@@ -458,7 +531,7 @@ class DiskCache:
|
||||
# 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
|
||||
self._remember_write(key, cache_path, digest, stamped_at)
|
||||
finally:
|
||||
if os.path.exists(tmp_path):
|
||||
try:
|
||||
@@ -471,7 +544,7 @@ class DiskCache:
|
||||
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._remember_write(key, cache_path, digest, stamped_at)
|
||||
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
|
||||
@@ -520,6 +593,26 @@ class DiskCache:
|
||||
)
|
||||
return # Exit gracefully without raising exception
|
||||
|
||||
def _remember_write(self, key: str, cache_path: str,
|
||||
digest: Tuple[int, int], stamped_at: Optional[float]) -> None:
|
||||
"""Record a completed write so an identical set() can skip the disk.
|
||||
|
||||
Caller holds _lock. A header-first record's mtime is set to its
|
||||
timestamp, so only a skipped rewrite ever moves it later (see
|
||||
"FRESHNESS OF A SKIPPED WRITE").
|
||||
"""
|
||||
try:
|
||||
if stamped_at is not None:
|
||||
os.utime(cache_path, (stamped_at, stamped_at))
|
||||
st = os.stat(cache_path)
|
||||
except OSError:
|
||||
# Written but not stamped (another user's file, on the direct
|
||||
# write path): mtime is the write time, which _refreshed_at reads
|
||||
# as the write itself. Remember nothing; the next set() writes.
|
||||
self._write_digests.pop(key, None)
|
||||
return
|
||||
self._write_digests[key] = (digest, (st.st_ino, st.st_size))
|
||||
|
||||
def clear(self, key: Optional[str] = None) -> None:
|
||||
"""
|
||||
Clear cache entry or all entries.
|
||||
|
||||
Reference in New Issue
Block a user