mirror of
https://github.com/ChuckBuilds/LEDMatrix.git
synced 2026-10-04 22:35: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,
|
- 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
|
||||||
|
|
||||||
|
|||||||
Vendored
+113
-20
@@ -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.
|
||||||
|
|||||||
@@ -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