diff --git a/CHANGELOG.md b/CHANGELOG.md index b9b0001a..ed9b2095 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -19,6 +19,18 @@ accepts both, but the store flags the old spelling as deprecated ## Unreleased +### Fixes + +- Re-saving unchanged data through `CacheManager.set` no longer rewrites its + cache file. The disk cache already skipped a payload identical to the last + one written, but `set()` stamps every record with the current time, so the + skip never fired and every plugin rewrote its unchanged API data to the SD + card on every update cycle. Records are now compared without that + timestamp, and a skipped write moves the file's mtime to it instead; reads + treat a record as fresh from the later of the two, so it expires exactly + when the rewrite would have. Changed data or a changed `ttl` still writes, + and so does a file another process has replaced since. + ## 3.7.0 Sports consolidation stage 3 (#672). No behaviour change: nothing in core diff --git a/src/cache/disk_cache.py b/src/cache/disk_cache.py index fb776b56..4e092603 100644 --- a/src/cache/disk_cache.py +++ b/src/cache/disk_cache.py @@ -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. diff --git a/test/test_cache_write_dedup.py b/test/test_cache_write_dedup.py new file mode 100644 index 00000000..f80c3f76 --- /dev/null +++ b/test/test_cache_write_dedup.py @@ -0,0 +1,158 @@ +"""Re-saving unchanged data through CacheManager.set does not rewrite the file. + +Regression under test: DiskCache.set skipped the disk when a payload matched +the last one written for the key, but CacheManager.set stamps every record +with time.time(), so no two payloads ever matched and every plugin rewrote its +unchanged API data to the SD card on every update cycle. The skip now ignores +the timestamp, and a skipped write moves the file's mtime instead -- so the +entry must stay exactly as fresh as the rewrite would have left it. +""" + +import os +import time +from unittest.mock import patch + +import pytest + +from src.cache import disk_cache as disk_cache_module +from src.cache.disk_cache import DiskCache +from src.cache_manager import CacheManager + +class Clock: + def __init__(self): + self.now = time.time() + + def __call__(self): + return self.now + + +@pytest.fixture +def clock(monkeypatch): + fake = Clock() + monkeypatch.setattr(time, "time", fake) + return fake + + +@pytest.fixture +def writes(monkeypatch): + """Paths DiskCache.set actually replaced (its atomic write path).""" + replaced = [] + real = os.replace + + def counting(src, dst, *args, **kwargs): + replaced.append(os.path.basename(dst)) + return real(src, dst, *args, **kwargs) + + monkeypatch.setattr(disk_cache_module.os, "replace", counting) + return replaced + + +def _manager(cache_dir): + with patch('src.cache_manager.CacheManager._get_writable_cache_dir', + return_value=str(cache_dir)): + manager = CacheManager() + manager.stop_cleanup_thread() + return manager + + +@pytest.fixture +def cm(tmp_path): + return _manager(tmp_path) + + +DATA = {"events": [{"id": n, "name": "x" * 20} for n in range(50)]} + + +def test_resaving_unchanged_data_does_not_rewrite_the_file(cm, clock, writes): + cm.set("scores", DATA) + path = cm._get_cache_path("scores") + first = os.stat(path) + with open(path, "rb") as f: + first_bytes = f.read() + + for _ in range(5): + clock.now += 60 + cm.set("scores", DATA) + + assert writes == ["scores.json"] + after = os.stat(path) + assert after.st_ino == first.st_ino + with open(path, "rb") as f: + assert f.read() == first_bytes + # The file records when its content was last saved. + assert after.st_mtime == pytest.approx(clock.now, abs=1e-3) + + +def test_skipped_writes_keep_the_entry_fresh(cm, clock, tmp_path): + cm.set("scores", DATA) + for _ in range(4): # re-saved unchanged every 200s + clock.now += 200 + cm.set("scores", DATA) + # 800s past the only real write, well beyond max_age=300 of it. + assert cm.get("scores", max_age=300) == DATA + + # A reader with no memory tier -- the web interface, or this service + # after a restart -- sees the same freshness from disk. + other = _manager(tmp_path) + record = other.get_cached_data("scores", max_age=300) + assert record is not None and record["data"] == DATA + assert record["timestamp"] == pytest.approx(clock.now, abs=1e-3) + assert DiskCache(str(tmp_path)).get("scores", max_age=300) is not None + + # Freshness is the last save, not forever. + clock.now += 301 + assert DiskCache(str(tmp_path)).get("scores", max_age=300) is None + assert _manager(tmp_path).get("scores", max_age=300) is None + + +def test_skipped_writes_keep_a_ttl_entry_fresh(cm, clock, tmp_path): + cm.set("odds", DATA, ttl=120) + for _ in range(3): + clock.now += 100 + cm.set("odds", DATA, ttl=120) + assert _manager(tmp_path).get("odds", max_age=10) == DATA + + +def test_changed_data_writes(cm, clock, writes, tmp_path): + cm.set("scores", DATA) + clock.now += 60 + cm.set("scores", {"events": []}) + assert writes == ["scores.json", "scores.json"] + assert _manager(tmp_path).get("scores", max_age=300) == {"events": []} + + +def test_changed_ttl_writes(cm, clock, writes, tmp_path): + cm.set("scores", DATA, ttl=60) + clock.now += 10 + cm.set("scores", DATA, ttl=600) + clock.now += 10 + cm.set("scores", DATA) + assert writes == ["scores.json"] * 3 + record = _manager(tmp_path).get_cached_data("scores", max_age=300) + assert "ttl" not in record + + +def test_a_record_written_old_stays_old(tmp_path, clock): + # Only a skip may move freshness forward: a record saved with an old + # timestamp on purpose is not made fresh by the write's own mtime. + disk = DiskCache(str(tmp_path)) + disk.set("k", {"timestamp": clock.now - 600, "data": DATA}) + assert disk.get("k", max_age=300) is None + + +def test_a_file_replaced_by_another_process_is_rewritten(tmp_path, clock): + ours, theirs = DiskCache(str(tmp_path)), DiskCache(str(tmp_path)) + ours.set("k", {"timestamp": clock.now, "data": "ours"}) + clock.now += 10 + theirs.set("k", {"timestamp": clock.now, "data": "theirs"}) + clock.now += 10 + ours.set("k", {"timestamp": clock.now, "data": "ours"}) + assert DiskCache(str(tmp_path)).get("k", max_age=300)["data"] == "ours" + + +def test_a_cleared_file_is_rewritten(cm, clock, tmp_path): + cm.set("scores", DATA) + os.remove(cm._get_cache_path("scores")) # e.g. the web UI's delete + clock.now += 10 + cm.set("scores", DATA) + assert _manager(tmp_path).get("scores", max_age=300) == DATA