mirror of
https://github.com/ChuckBuilds/LEDMatrix.git
synced 2026-10-04 14:25:08 +00:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ffbf7b7067 | ||
|
|
8363983f1c | ||
|
|
0f39e9a2f3 |
@@ -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,
|
||||
`/display/on-demand/status` kept reporting `status: error` for up to two
|
||||
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
|
||||
|
||||
|
||||
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.
|
||||
|
||||
@@ -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
|
||||
Reference in New Issue
Block a user