Compare commits

...
Author SHA1 Message Date
Chuck ffbf7b7067 Merge remote-tracking branch 'origin/main' into claude/fix-cache-write-dedup
# Conflicts:
#	CHANGELOG.md
2026-09-29 19:28:21 -04:00
Chuck 8363983f1c Merge remote-tracking branch 'origin/main' into claude/fix-cache-write-dedup
# Conflicts:
#	CHANGELOG.md
2026-09-29 16:56:35 -04:00
ChuckandClaude Opus 5.5 0f39e9a2f3 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>
2026-09-29 13:41:52 -04:00
3 changed files with 280 additions and 20 deletions
+9
View File
@@ -42,6 +42,15 @@ accepts both, but the store flags the old spelling as deprecated
- A stop request now clears an on-demand error. After a failed request, - A stop request now clears an on-demand error. After a failed request,
`/display/on-demand/status` kept reporting `status: error` for up to two `/display/on-demand/status` kept reporting `status: error` for up to two
minutes even after a stop. minutes even after a stop.
- 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 ## 3.7.0
+113 -20
View File
@@ -14,7 +14,7 @@ import tempfile
import logging import logging
import threading import threading
import zlib import zlib
from typing import Dict, Any, Optional, Protocol from typing import Dict, Any, Optional, Protocol, Tuple
from datetime import datetime from datetime import datetime
from src.common.path_safety import safe_path_component 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. """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 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 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. 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) match = _HEAD_RE.match(head)
if not match: if not match:
return False return False
try: try:
timestamp = float(match.group(1)) timestamp = float(match.group(1))
if refreshed is not None:
timestamp = max(timestamp, refreshed)
limit = max_age limit = max_age
if match.group(2) is not None: if match.group(2) is not None:
ttl = float(match.group(2)) ttl = float(match.group(2))
@@ -248,11 +296,13 @@ class DiskCache:
self.cache_dir = cache_dir self.cache_dir = cache_dir
self.logger = logger or logging.getLogger(__name__) self.logger = logger or logging.getLogger(__name__)
self._lock = threading.Lock() self._lock = threading.Lock()
# key -> adler32 of the last payload successfully written to the # key -> ((length, adler32) of the last content written to the
# primary cache path; lets set() skip rewriting identical data # primary cache path, (st_ino, st_size) of the file it left); lets
# (per-process only — worst case another process rewrites, never # set() skip rewriting identical data. The file identity catches
# a missed write). Guarded by _lock. # another process -- the web interface writes and clears keys too --
self._write_digests: Dict[str, int] = {} # 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]: 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 # records (a season schedule is re-fetched when its cache
# expires), and parsing 53MB to throw it away held the GIL # expires), and parsing 53MB to throw it away held the GIL
# for ~1.8s -- a visible freeze on the panel. # 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 return None
f.seek(0) f.seek(0)
record = _loads(f.read()) record = _loads(f.read())
@@ -315,6 +370,11 @@ class DiskCache:
record_ts = None record_ts = None
if isinstance(record, dict): if isinstance(record, dict):
record_ts = record.get('timestamp') 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: if record_ts is None:
try: try:
record_ts = os.path.getmtime(cache_path) 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) self.logger.warning("Cache data for key '%s' not serializable: %s", key, e)
return 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: try:
# Atomic write to avoid partial/corrupt files # Atomic write to avoid partial/corrupt files
with self._lock: 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 # written for this key (plugins re-save unchanged API data
# every update cycle — each write is real SD-card wear). # every update cycle — each write is real SD-card wear).
# Refresh the file mtime so records that rely on it for TTL # Move the file mtime instead, so the record stays as fresh as
# (no embedded 'timestamp') don't expire early; a metadata # the rewrite would have left it; a metadata touch is
# touch is journal-cheap compared to rewriting the data. # journal-cheap compared to rewriting the data.
if self._write_digests.get(key) == digest: known = self._write_digests.get(key)
if known is not None and known[0] == digest:
try: try:
os.utime(cache_path, None) st = os.stat(cache_path)
return 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: except OSError:
# File vanished or perms changed — fall through and write pass
self._write_digests.pop(key, None) # 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) tmp_dir = os.path.dirname(cache_path)
# Try to create temp file in cache directory first # Try to create temp file in cache directory first
@@ -458,7 +531,7 @@ class DiskCache:
# opened it in between was refused. # opened it in between was refused.
_share_open_file(tmp_file.fileno(), _shared_group(tmp_dir)) _share_open_file(tmp_file.fileno(), _shared_group(tmp_dir))
os.replace(tmp_path, cache_path) os.replace(tmp_path, cache_path)
self._write_digests[key] = digest self._remember_write(key, cache_path, digest, stamped_at)
finally: finally:
if os.path.exists(tmp_path): if os.path.exists(tmp_path):
try: try:
@@ -471,7 +544,7 @@ class DiskCache:
with open(cache_path, 'wb') as cache_file: with open(cache_path, 'wb') as cache_file:
cache_file.write(payload) cache_file.write(payload)
_share_open_file(cache_file.fileno(), _shared_group(tmp_dir)) _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) self.logger.debug("Wrote cache for %s directly (non-atomic)", key)
except (IOError, OSError, PermissionError) as write_error: except (IOError, OSError, PermissionError) as write_error:
# If direct write also fails, try fallback location # If direct write also fails, try fallback location
@@ -520,6 +593,26 @@ class DiskCache:
) )
return # Exit gracefully without raising exception 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: def clear(self, key: Optional[str] = None) -> None:
""" """
Clear cache entry or all entries. Clear cache entry or all entries.
+158
View File
@@ -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