mirror of
https://github.com/ChuckBuilds/LEDMatrix.git
synced 2026-10-04 06:15:09 +00:00
perf(cache): skip rewriting unchanged CacheManager.set() records (#730)
The disk cache's unchanged-payload skip now ignores a CacheManager.set() record's timestamp, so unchanged re-saves are skipped; a skip moves the file's mtime to the new timestamp instead, and readers take a record's age from the newer of the two (never more than an hour past the embedded timestamp). Per-plugin plugin_metrics:<id> records become one plugin_metrics_snapshot written at most once a minute, and CacheManager builds its ConfigManager on first use. On hdpi, cache file writes went from ~37 to 8.6 a minute. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
@@ -19,6 +19,46 @@ accepts both, but the store flags the old spelling as deprecated
|
||||
|
||||
## Unreleased
|
||||
|
||||
### Fewer SD-card writes from the cache
|
||||
|
||||
- **An unchanged `CacheManager.set()` no longer rewrites the file.**
|
||||
`DiskCache` already skipped a payload identical to the last one it wrote,
|
||||
but `set()` stamps every record with the current time, so for `set()` the
|
||||
payload never matched and every unchanged re-save was a full rewrite. The
|
||||
comparison now leaves out a header-first record's timestamp (the `ttl` and
|
||||
the data still count), and the newer timestamp is kept in the file's mtime
|
||||
instead: a skipped save touches the file to the record's timestamp, and a
|
||||
real write pins mtime to the record's own timestamp. Every reader ages a
|
||||
record from the newer of the two -- `DiskCache.get`, its header-only
|
||||
staleness check, and the record it returns, whose `timestamp` is the newer
|
||||
value, so `CacheManager.get`, the memory tier and plugins reading
|
||||
`record['timestamp']` all agree; the retention sweep and the web UI's cache
|
||||
list already used mtime. The mtime is trusted at most an hour past the
|
||||
record's own timestamp, and unchanged data is rewritten once an hour, so a
|
||||
file copied without its mtime reads at most an hour fresher than its
|
||||
contents. 100 identical `set()` calls of a 32 KB record: 100 writes before,
|
||||
1 after.
|
||||
- **Plugin metrics are one record, written at most once a minute.** The
|
||||
resource monitor wrote a `plugin_metrics:<id>` record per plugin, each at
|
||||
most every 30 s: two writes a minute per plugin, 28 on a fourteen-plugin
|
||||
rig. Every plugin's metrics now go in one `plugin_metrics_snapshot` record
|
||||
(`{"schema": 1, "plugins": {id: record}}`, each record shaped as before),
|
||||
written at most once a minute. `GET /api/v3/plugins/metrics` and
|
||||
`/plugins/metrics/<id>` return the same fields; the numbers can be up to a
|
||||
minute old instead of 30 s. A plugin the snapshot does not have yet is
|
||||
still read from its old `plugin_metrics:<id>` record, which nothing writes
|
||||
any more and the cache's retention removes. Each write starts from the
|
||||
snapshot on disk, so plugins the display has not run since a restart keep
|
||||
their numbers, and a reset from the web UI sticks for a plugin the display
|
||||
is not running, as it did. A plugin with no call for 30 days is dropped from
|
||||
the snapshot, as its record used to age out.
|
||||
- **`CacheManager` no longer loads the config when it is built.** Every
|
||||
manager built a `ConfigManager` and loaded the whole config for a cache
|
||||
strategy that stopped reading it. `cache_manager.config_manager` is still
|
||||
there -- the sports plugins resolve the global timezone through it -- and
|
||||
is now built and loaded on first access; assigning it still replaces it.
|
||||
`CacheStrategy` is given no config manager (it reads none).
|
||||
|
||||
### Plugin update tick: a few times a second, not every frame
|
||||
|
||||
- The frame loops and the dwell sleep ran
|
||||
|
||||
Vendored
+181
-37
@@ -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,91 @@ _HEAD_RE = re.compile(
|
||||
)
|
||||
|
||||
|
||||
def _stale_from_head(head: bytes, max_age: Optional[int], now: float) -> bool:
|
||||
# UNCHANGED RE-SAVES: THE FILE'S MTIME CARRIES THE NEWER TIMESTAMP
|
||||
# ----------------------------------------------------------------
|
||||
# Plugins re-save unchanged API data every update cycle, and every one of
|
||||
# those saves was a full rewrite on the SD card. DiskCache.set skips the write
|
||||
# when the payload matches the last one it wrote for the key -- but
|
||||
# CacheManager.set stamps each record with time.time(), so for set() the
|
||||
# payload never matched and the skip never fired.
|
||||
#
|
||||
# The digest now leaves out a header-first record's timestamp, so an unchanged
|
||||
# set() is skipped. What the skip must not do is make the record look older
|
||||
# than it is: the timestamp inside the file is from the last real write, and
|
||||
# a reader in another process (the web interface, with memory_ttl=0) or after
|
||||
# a restart would call fresh data stale. So the newer timestamp goes where it
|
||||
# costs no data write -- the file's mtime -- and readers take a record's age
|
||||
# from the newer of the two. The invariant that makes that safe:
|
||||
#
|
||||
# a file's mtime is the timestamp of the newest record saved for its key
|
||||
#
|
||||
# real write mtime is set to the record's own timestamp, so a record saved
|
||||
# with an old timestamp (data as of some earlier time) cannot
|
||||
# borrow freshness from the moment it hit the disk
|
||||
# skip mtime is set to the skipped record's timestamp -- exactly what
|
||||
# a rewrite would have stored, without the rewrite
|
||||
#
|
||||
# Readers of the on-disk timestamp, all of which go through _effective_timestamp:
|
||||
# DiskCache.get (the header check and the full parse; it also returns the
|
||||
# record with 'timestamp' set to the effective value, so CacheManager.get's
|
||||
# max_age path, the memory tier hydrated from disk, and any plugin reading
|
||||
# record['timestamp'] all see it). Readers that use mtime alone already see the
|
||||
# newer value: the retention sweep below, CacheManager.list_cache_files (the
|
||||
# web UI's cache list). Nothing else opens cache files: web_interface and
|
||||
# scripts reach them only through CacheManager.
|
||||
#
|
||||
# Something other than this class can also move an mtime forward -- a copy
|
||||
# without -p, an rsync without -t, a `touch`. (backup_manager.py does not
|
||||
# back up or restore the cache directory, so the in-tree restore cannot.) That
|
||||
# must not make old data fresh, so the lift is bounded: a reader never takes
|
||||
# the mtime as more than _MAX_TIMESTAMP_LIFT past the embedded timestamp, and
|
||||
# set() rewrites the file for real once a skip would need more than that, so
|
||||
# an honest lift never reaches the bound. A file copied a day after it was
|
||||
# written therefore reads at most an hour fresher than its contents say, and a
|
||||
# 30-second live-score record from yesterday stays stale. CacheManager.set
|
||||
# records written before this change have mtime == write time == embedded
|
||||
# timestamp, give or take the write itself, and read exactly as before; a
|
||||
# file an older version wrote or touched later than its embedded timestamp
|
||||
# says reads at most the same hour fresher, once, until it is next saved.
|
||||
|
||||
#: Longest a skipped write may stand in for a real one, and so the furthest a
|
||||
#: file's mtime is ever trusted past the record's own timestamp. Unchanged data
|
||||
#: is rewritten at least this often, at most once an hour per key instead of
|
||||
#: once per update cycle.
|
||||
_MAX_TIMESTAMP_LIFT = 3600.0
|
||||
|
||||
|
||||
def _record_timestamp(value: Any) -> Optional[float]:
|
||||
"""A record's timestamp as a finite float, or None if it has no usable one."""
|
||||
if isinstance(value, bool) or not isinstance(value, (int, float)):
|
||||
return None
|
||||
value = float(value)
|
||||
return value if math.isfinite(value) else None
|
||||
|
||||
|
||||
def _effective_timestamp(embedded: float, mtime: Optional[float]) -> float:
|
||||
"""When a record was last saved: its timestamp, or the file's mtime if a
|
||||
later unchanged save moved that forward -- never by more than
|
||||
_MAX_TIMESTAMP_LIFT. See "UNCHANGED RE-SAVES" above."""
|
||||
if mtime is None:
|
||||
return embedded
|
||||
return max(embedded, min(mtime, embedded + _MAX_TIMESTAMP_LIFT))
|
||||
|
||||
|
||||
def _stale_from_head(head: bytes, max_age: Optional[int], now: float,
|
||||
mtime: 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.
|
||||
header cannot be read, so the full parse decides as it always did. ``mtime``
|
||||
is the file's, which may carry a newer save than the header does.
|
||||
"""
|
||||
match = _HEAD_RE.match(head)
|
||||
if not match:
|
||||
return False
|
||||
try:
|
||||
timestamp = float(match.group(1))
|
||||
timestamp = _effective_timestamp(float(match.group(1)), mtime)
|
||||
limit = max_age
|
||||
if match.group(2) is not None:
|
||||
ttl = float(match.group(2))
|
||||
@@ -179,7 +252,7 @@ else:
|
||||
# --------------------------------------------
|
||||
# The display service runs as root and the web interface as the installing
|
||||
# user, and the web interface reads records only the display writes
|
||||
# (display_current_state, display_on_demand_state, plugin_metrics:*). Files are
|
||||
# (display_current_state, display_on_demand_state, plugin_metrics_snapshot). Files are
|
||||
# written 0660, so the web interface can read one only through its group.
|
||||
#
|
||||
# The installers rely on the directory's setgid bit to set that group. That is
|
||||
@@ -248,11 +321,14 @@ 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 -> what set() last put at the primary cache path: the adler32 of
|
||||
# the payload (less a header-first timestamp), the timestamp the file
|
||||
# holds (None for records without one), and the file's inode and size.
|
||||
# Lets set() skip rewriting identical data. Per-process only, and the
|
||||
# inode/size check means another process's write is never mistaken
|
||||
# for ours -- worst case a redundant write, never a missed one.
|
||||
# Guarded by _lock.
|
||||
self._write_digests: Dict[str, Tuple[int, Optional[float], int, int]] = {}
|
||||
|
||||
def get_cache_path(self, key: str) -> Optional[str]:
|
||||
"""
|
||||
@@ -301,32 +377,41 @@ class DiskCache:
|
||||
try:
|
||||
with self._lock:
|
||||
with open(cache_path, 'rb') as f:
|
||||
# The open file's mtime, not the path's: the file a skipped
|
||||
# write touched is the one being read.
|
||||
mtime = os.fstat(f.fileno()).st_mtime
|
||||
# Decide staleness from the header before paying for the
|
||||
# parse. A stale read is the common case for the biggest
|
||||
# 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()):
|
||||
if _stale_from_head(f.read(_HEAD_BYTES), max_age, time.time(), mtime):
|
||||
return None
|
||||
f.seek(0)
|
||||
record = _loads(f.read())
|
||||
|
||||
# Determine record timestamp (prefer embedded, else file mtime)
|
||||
|
||||
# Determine record timestamp: the embedded one, moved forward by a
|
||||
# later unchanged save if there was one (see "UNCHANGED RE-SAVES"),
|
||||
# else the file mtime.
|
||||
record_ts = None
|
||||
if isinstance(record, dict):
|
||||
record_ts = record.get('timestamp')
|
||||
if record_ts is None:
|
||||
try:
|
||||
record_ts = os.path.getmtime(cache_path)
|
||||
except OSError:
|
||||
record_ts = None
|
||||
|
||||
if record_ts is not None:
|
||||
try:
|
||||
record_ts = float(record_ts)
|
||||
except (TypeError, ValueError):
|
||||
record_ts = None
|
||||
|
||||
record_ts = mtime
|
||||
else:
|
||||
embedded_ts = _record_timestamp(record_ts)
|
||||
if embedded_ts is None:
|
||||
try:
|
||||
record_ts = float(record_ts)
|
||||
except (TypeError, ValueError):
|
||||
record_ts = None
|
||||
else:
|
||||
record_ts = _effective_timestamp(embedded_ts, mtime)
|
||||
if record_ts != embedded_ts:
|
||||
# Hand the record back as a rewrite would have left it,
|
||||
# so callers that age it themselves agree with us.
|
||||
record['timestamp'] = record_ts
|
||||
|
||||
now = time.time()
|
||||
|
||||
# An explicit per-entry ttl wins over the caller's max_age. The
|
||||
@@ -403,7 +488,12 @@ class DiskCache:
|
||||
self.logger.warning("Cache data for key '%s' not serializable: %s", key, e)
|
||||
return
|
||||
|
||||
digest = zlib.adler32(payload)
|
||||
timestamp = _record_timestamp(data.get('timestamp')) if isinstance(data, dict) else None
|
||||
# A header-first record (CacheManager.set's layout) is compared without
|
||||
# its timestamp, which differs on every save; see "UNCHANGED RE-SAVES".
|
||||
# Any other layout is compared whole, as before.
|
||||
head = _HEAD_RE.match(payload) if timestamp is not None else None
|
||||
digest = zlib.adler32(memoryview(payload)[head.end(1):] if head else payload)
|
||||
|
||||
try:
|
||||
# Atomic write to avoid partial/corrupt files
|
||||
@@ -411,16 +501,10 @@ class DiskCache:
|
||||
# Skip the disk entirely when this exact payload 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:
|
||||
try:
|
||||
os.utime(cache_path, None)
|
||||
return
|
||||
except OSError:
|
||||
# File vanished or perms changed — fall through and write
|
||||
self._write_digests.pop(key, None)
|
||||
# A metadata touch is journal-cheap compared to rewriting
|
||||
# the data.
|
||||
if self._skip_unchanged(key, cache_path, digest, timestamp):
|
||||
return
|
||||
|
||||
tmp_dir = os.path.dirname(cache_path)
|
||||
# Try to create temp file in cache directory first
|
||||
@@ -458,7 +542,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, timestamp)
|
||||
finally:
|
||||
if os.path.exists(tmp_path):
|
||||
try:
|
||||
@@ -471,7 +555,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, timestamp)
|
||||
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 +604,66 @@ class DiskCache:
|
||||
)
|
||||
return # Exit gracefully without raising exception
|
||||
|
||||
def _skip_unchanged(self, key: str, cache_path: str, digest: int,
|
||||
timestamp: Optional[float]) -> bool:
|
||||
"""Stand in for a write of an unchanged record by touching the file.
|
||||
|
||||
True when the file already holds this record bar its timestamp and the
|
||||
touch landed; False means write it. The touch sets mtime to the
|
||||
record's timestamp -- what a rewrite would have stored -- or to now
|
||||
for a record without one, whose age readers already take from mtime.
|
||||
Caller holds _lock.
|
||||
"""
|
||||
last = self._write_digests.get(key)
|
||||
if last is None or last[0] != digest:
|
||||
return False
|
||||
_, written_ts, ino, size = last
|
||||
if (timestamp is None) != (written_ts is None):
|
||||
return False
|
||||
if timestamp is not None and written_ts is not None:
|
||||
# Never backwards (a rewrite would make the record older), and
|
||||
# never further than readers will trust the mtime: past that the
|
||||
# record is rewritten, so its own timestamp catches up.
|
||||
if not written_ts <= timestamp <= written_ts + _MAX_TIMESTAMP_LIFT:
|
||||
return False
|
||||
try:
|
||||
st = os.stat(cache_path)
|
||||
if (st.st_ino, st.st_size) != (ino, size):
|
||||
# Replaced since our write (another process, a restore):
|
||||
# its contents are not the ones the digest describes.
|
||||
self._write_digests.pop(key, None)
|
||||
return False
|
||||
# Setting an explicit time needs the file's owner; a file someone
|
||||
# else wrote fails here and is rewritten (as our own file) instead.
|
||||
os.utime(cache_path, None if timestamp is None else (timestamp, timestamp))
|
||||
return True
|
||||
except OSError:
|
||||
# File vanished or perms changed — fall through and write
|
||||
self._write_digests.pop(key, None)
|
||||
return False
|
||||
|
||||
def _remember_write(self, key: str, cache_path: str, digest: int,
|
||||
timestamp: Optional[float]) -> None:
|
||||
"""After a real write: pin mtime to the record's timestamp and note
|
||||
what was written, so the next unchanged save can be skipped.
|
||||
|
||||
Pinning keeps a record saved with an older timestamp from looking as
|
||||
fresh as the moment it was written (see "UNCHANGED RE-SAVES"); for
|
||||
CacheManager.set's records the two differ only by the write itself.
|
||||
A timestamp in the future is left alone, mtime already being older.
|
||||
Never raises: the data is on disk, and anything failing here only
|
||||
costs the next save its skip. Caller holds _lock.
|
||||
"""
|
||||
self._write_digests.pop(key, None)
|
||||
try:
|
||||
if timestamp is not None and timestamp <= time.time():
|
||||
os.utime(cache_path, (timestamp, timestamp))
|
||||
st = os.stat(cache_path)
|
||||
except OSError as e:
|
||||
self.logger.debug("Could not pin mtime of %s: %s", cache_path, e)
|
||||
return
|
||||
self._write_digests[key] = (digest, timestamp, st.st_ino, st.st_size)
|
||||
|
||||
def clear(self, key: Optional[str] = None) -> None:
|
||||
"""
|
||||
Clear cache entry or all entries.
|
||||
|
||||
+49
-10
@@ -43,6 +43,9 @@ from src.logging_config import get_logger
|
||||
# it from either path.
|
||||
from src.cache.disk_cache import DateTimeEncoder # noqa: F401 - deliberate re-export
|
||||
|
||||
# CacheManager.config_manager not built yet (None means "not available").
|
||||
_UNSET: Any = object()
|
||||
|
||||
class CacheManager:
|
||||
"""Manages caching of API responses to reduce API calls."""
|
||||
|
||||
@@ -73,21 +76,19 @@ class CacheManager:
|
||||
self.logger.error("Could not find or create a writable cache directory. Caching will be disabled.")
|
||||
self.cache_dir = None
|
||||
|
||||
# Initialize config manager for sport-specific intervals
|
||||
try:
|
||||
from src.config_manager import ConfigManager
|
||||
self.config_manager: Optional[Any] = ConfigManager()
|
||||
self.config_manager.load_config()
|
||||
except ImportError:
|
||||
self.config_manager: Optional[Any] = None
|
||||
self.logger.warning("ConfigManager not available, using default cache intervals")
|
||||
|
||||
# The config manager is built on first use of self.config_manager; see
|
||||
# the property. Nothing in the cache reads it any more.
|
||||
self._config_manager: Any = _UNSET
|
||||
self._config_manager_lock = threading.Lock()
|
||||
|
||||
# Initialize cache components using composition
|
||||
self._memory_cache_component = MemoryCache(
|
||||
max_size=default_max_size(), cleanup_interval=300.0
|
||||
)
|
||||
self._disk_cache_component = DiskCache(cache_dir=self.cache_dir, logger=self.logger)
|
||||
self._strategy_component = CacheStrategy(config_manager=self.config_manager, logger=self.logger)
|
||||
# No config manager: CacheStrategy keeps the parameter for callers but
|
||||
# reads nothing from it, and passing ours would build it eagerly.
|
||||
self._strategy_component = CacheStrategy(logger=self.logger)
|
||||
self._metrics_component = CacheMetrics(logger=self.logger)
|
||||
|
||||
# Disk cleanup configuration
|
||||
@@ -115,6 +116,44 @@ class CacheManager:
|
||||
if self.cache_dir:
|
||||
self.start_cleanup_thread()
|
||||
|
||||
@property
|
||||
def config_manager(self) -> Optional[Any]:
|
||||
"""A loaded ConfigManager, built the first time it is asked for.
|
||||
|
||||
Every CacheManager used to build one and load the whole config in
|
||||
__init__, for a cache strategy that stopped reading it -- startup paid
|
||||
a config load (and the web interface another) per manager for nothing.
|
||||
It is still public: the sports plugins resolve the global timezone and
|
||||
display settings through ``cache_manager.config_manager``, and they get
|
||||
the same object they always did, on first access instead of at
|
||||
construction. None when ConfigManager cannot be imported, as before.
|
||||
Assigning replaces it, as assigning the attribute always did.
|
||||
"""
|
||||
# getattr: a manager made with __new__ (some tests) has no slot yet.
|
||||
value = getattr(self, '_config_manager', _UNSET)
|
||||
if value is not _UNSET:
|
||||
return value
|
||||
lock = getattr(self, '_config_manager_lock', None) or threading.Lock()
|
||||
with lock:
|
||||
value = getattr(self, '_config_manager', _UNSET)
|
||||
if value is _UNSET:
|
||||
try:
|
||||
from src.config_manager import ConfigManager
|
||||
except ImportError:
|
||||
self.logger.warning("ConfigManager not available, using default cache intervals")
|
||||
value = None
|
||||
else:
|
||||
value = ConfigManager()
|
||||
# Raises as it did from __init__; nothing is kept, so the
|
||||
# next access tries again.
|
||||
value.load_config()
|
||||
self._config_manager = value
|
||||
return value
|
||||
|
||||
@config_manager.setter
|
||||
def config_manager(self, value: Optional[Any]) -> None:
|
||||
self._config_manager = value
|
||||
|
||||
def _get_writable_cache_dir(self) -> Optional[str]:
|
||||
"""Tries to find or create a writable cache directory, preferring a system path when available."""
|
||||
# Attempt 1: System-wide persistent cache directory (preferred for services)
|
||||
|
||||
@@ -485,7 +485,7 @@ def record_error(
|
||||
# and only the display service's ever records anything (plugin_executor runs
|
||||
# the plugins there). The web interface therefore reads a snapshot the display
|
||||
# service publishes to the shared cache directory -- the same channel, and the
|
||||
# same file permissions, as display_current_state and plugin_metrics:*: files
|
||||
# same file permissions, as display_current_state and plugin_metrics_snapshot: files
|
||||
# are 0660 and carry the cache directory's group, so root writes and the web
|
||||
# user reads, and the other way round for the clear request.
|
||||
#
|
||||
|
||||
@@ -8,7 +8,7 @@ Provides resource limits and performance monitoring.
|
||||
import math
|
||||
import time
|
||||
import threading
|
||||
from typing import Dict, Optional, Any, Callable, cast
|
||||
from typing import Dict, Optional, Any, Callable, Set, cast
|
||||
from dataclasses import dataclass, field, fields
|
||||
|
||||
from src.logging_config import get_logger
|
||||
@@ -99,18 +99,33 @@ class ResourceMetrics:
|
||||
last_update_time: float = field(default_factory=time.time)
|
||||
|
||||
|
||||
#: How often a plugin's metrics are written to the cache, in seconds.
|
||||
#: How often the metrics snapshot is written to the cache, in seconds.
|
||||
#:
|
||||
#: Persisting on every call meant a small file rewritten roughly nine times a
|
||||
#: minute per plugin. On a rig with fourteen active plugins that was ~126
|
||||
#: writes a minute for metrics alone, and since each ~350-byte file costs a
|
||||
#: 4KB block plus an ext4 journal entry, it dominated the device's write
|
||||
#: volume -- on an SD card, which wears out.
|
||||
#: volume -- on an SD card, which wears out. Throttling each plugin's own
|
||||
#: record to once per 30 s still left two writes a minute per plugin, so all
|
||||
#: plugins now share one record (METRICS_SNAPSHOT_KEY), written at most once
|
||||
#: a minute: one write a minute however many plugins there are.
|
||||
#:
|
||||
#: The in-memory copy stays authoritative and exact; only the cross-process
|
||||
#: snapshot the web UI reads is delayed, and telemetry up to half a minute old
|
||||
#: is still a fair description of a long-running plugin.
|
||||
_METRICS_PERSIST_INTERVAL = 30.0
|
||||
#: snapshot the web UI reads is delayed, and telemetry up to a minute old is
|
||||
#: still a fair description of a long-running plugin.
|
||||
_METRICS_PERSIST_INTERVAL = 60.0
|
||||
|
||||
#: The one cache record holding every plugin's metrics:
|
||||
#: ``{"schema": 1, "plugins": {plugin_id: <metrics record>}}``, each metrics
|
||||
#: record shaped as the per-plugin ``plugin_metrics:<id>`` records were. Those
|
||||
#: older records are still read for a plugin the snapshot does not have yet
|
||||
#: (an upgrade, or a plugin that has not run since), never written.
|
||||
METRICS_SNAPSHOT_KEY = "plugin_metrics_snapshot"
|
||||
_METRICS_SNAPSHOT_SCHEMA = 1
|
||||
|
||||
#: A plugin with no call for this long is dropped from the snapshot -- what
|
||||
#: the cache's 30-day default retention did to its own record before.
|
||||
_METRICS_SNAPSHOT_ENTRY_MAX_AGE = 30 * 86400
|
||||
|
||||
|
||||
class PluginResourceMonitor:
|
||||
@@ -140,10 +155,15 @@ class PluginResourceMonitor:
|
||||
self._metrics: Dict[str, ResourceMetrics] = {}
|
||||
self._limits: Dict[str, ResourceLimits] = {}
|
||||
self._bad_limits_warned: set = set()
|
||||
# When each plugin's metrics last reached the cache. Metrics change on
|
||||
# every call, so they cannot be de-duplicated the way health state can;
|
||||
# they are rate-limited instead. See _METRICS_PERSIST_INTERVAL.
|
||||
self._metrics_persisted_at: Dict[str, float] = {}
|
||||
# When the metrics snapshot last reached the cache (monotonic), None
|
||||
# until it has. Metrics change on every call, so they cannot be
|
||||
# de-duplicated the way health state can; they are rate-limited
|
||||
# instead. See _METRICS_PERSIST_INTERVAL.
|
||||
self._snapshot_persisted_at: Optional[float] = None
|
||||
# Plugins whose metrics this process recorded since the last snapshot
|
||||
# write: only their entries are overwritten, the rest are kept as
|
||||
# found on disk.
|
||||
self._metrics_dirty: Set[str] = set()
|
||||
|
||||
# Lock for thread-safe access
|
||||
self._lock = threading.Lock()
|
||||
@@ -247,10 +267,14 @@ class PluginResourceMonitor:
|
||||
with self._lock:
|
||||
if force_reload or plugin_id not in self._metrics:
|
||||
# Try to load from cache
|
||||
cache_key = self._get_metrics_key(plugin_id)
|
||||
cached = self.cache_manager.get(
|
||||
cache_key, max_age=None, memory_ttl=0 if force_reload else None
|
||||
)
|
||||
memory_ttl = 0 if force_reload else None
|
||||
cached = self._read_snapshot(memory_ttl).get(plugin_id)
|
||||
if cached is None:
|
||||
# Not in the snapshot: the per-plugin record an older
|
||||
# version wrote, if there is one.
|
||||
cached = self.cache_manager.get(
|
||||
self._get_metrics_key(plugin_id), max_age=None,
|
||||
memory_ttl=memory_ttl)
|
||||
if cached:
|
||||
metrics = self._metrics_from_cache(plugin_id, cached)
|
||||
else:
|
||||
@@ -498,12 +522,70 @@ class PluginResourceMonitor:
|
||||
summaries[plugin_id] = self.get_metrics_summary(plugin_id)
|
||||
return summaries
|
||||
|
||||
def _persist_metrics(self, plugin_id: str, metrics: ResourceMetrics,
|
||||
force: bool = False) -> None:
|
||||
"""Write a plugin's metrics to the cache, at most once per interval.
|
||||
def _read_snapshot(self, memory_ttl: Optional[int] = None) -> Dict[str, Any]:
|
||||
"""The snapshot's per-plugin records, or {} if there is none usable.
|
||||
|
||||
Caller must hold ``self._lock``.
|
||||
"""
|
||||
cached = self.cache_manager.get(
|
||||
METRICS_SNAPSHOT_KEY, max_age=None, memory_ttl=memory_ttl)
|
||||
if not isinstance(cached, dict) or cached.get('schema') != _METRICS_SNAPSHOT_SCHEMA:
|
||||
return {}
|
||||
plugins = cached.get('plugins')
|
||||
return plugins if isinstance(plugins, dict) else {}
|
||||
|
||||
@staticmethod
|
||||
def _metrics_record(metrics: ResourceMetrics) -> Dict[str, Any]:
|
||||
"""One plugin's entry in the snapshot."""
|
||||
return {
|
||||
'memory_mb': metrics.memory_mb,
|
||||
'cpu_percent': metrics.cpu_percent,
|
||||
'execution_time': metrics.execution_time,
|
||||
'call_count': metrics.call_count,
|
||||
'total_execution_time': metrics.total_execution_time,
|
||||
'max_execution_time': metrics.max_execution_time,
|
||||
'min_execution_time': (metrics.min_execution_time
|
||||
if metrics.min_execution_time != float('inf')
|
||||
else 0.0),
|
||||
'last_update_time': metrics.last_update_time,
|
||||
}
|
||||
|
||||
def _write_snapshot(self, drop: Optional[str] = None) -> None:
|
||||
"""Write the snapshot: what is on disk, with this process's recorded
|
||||
plugins updated and ``drop`` removed.
|
||||
|
||||
Starting from the disk copy rather than from memory keeps the entries
|
||||
of plugins this process has not run -- disabled ones, which the web UI
|
||||
still shows -- and a reset made from the other process.
|
||||
|
||||
Caller must hold ``self._lock``.
|
||||
"""
|
||||
plugins = dict(self._read_snapshot(memory_ttl=0))
|
||||
if drop is not None:
|
||||
plugins.pop(drop, None)
|
||||
for plugin_id in self._metrics_dirty:
|
||||
if plugin_id in self._metrics:
|
||||
plugins[plugin_id] = self._metrics_record(self._metrics[plugin_id])
|
||||
cutoff = time.time() - _METRICS_SNAPSHOT_ENTRY_MAX_AGE
|
||||
for plugin_id, record in list(plugins.items()):
|
||||
last = record.get('last_update_time') if isinstance(record, dict) else None
|
||||
if isinstance(last, (int, float)) and last < cutoff:
|
||||
del plugins[plugin_id]
|
||||
self.cache_manager.set(METRICS_SNAPSHOT_KEY, {
|
||||
'schema': _METRICS_SNAPSHOT_SCHEMA,
|
||||
'plugins': plugins,
|
||||
})
|
||||
# Only once the write has landed, so a failed one is retried in full.
|
||||
self._metrics_dirty.clear()
|
||||
|
||||
def _persist_metrics(self, plugin_id: str, metrics: ResourceMetrics,
|
||||
force: bool = False) -> None:
|
||||
"""Record that a plugin's metrics changed, and write the snapshot if
|
||||
the last write is at least an interval old.
|
||||
|
||||
Caller must hold ``self._lock``.
|
||||
"""
|
||||
self._metrics_dirty.add(plugin_id)
|
||||
# Monotonic, not wall clock: these devices have no RTC, so the clock
|
||||
# jumps by however far off boot-time was the moment NTP first syncs.
|
||||
# A forward jump would allow an early write, a backward one would
|
||||
@@ -515,35 +597,26 @@ class PluginResourceMonitor:
|
||||
# single run -- the throttle swallowed the very first snapshot, which
|
||||
# is the one that matters most after a restart.
|
||||
now = time.monotonic()
|
||||
last_written = self._metrics_persisted_at.get(plugin_id)
|
||||
last_written = self._snapshot_persisted_at
|
||||
if (not force and last_written is not None
|
||||
and now - last_written < _METRICS_PERSIST_INTERVAL):
|
||||
return
|
||||
cache_key = self._get_metrics_key(plugin_id)
|
||||
self.cache_manager.set(cache_key, {
|
||||
'memory_mb': metrics.memory_mb,
|
||||
'cpu_percent': metrics.cpu_percent,
|
||||
'execution_time': metrics.execution_time,
|
||||
'call_count': metrics.call_count,
|
||||
'total_execution_time': metrics.total_execution_time,
|
||||
'max_execution_time': metrics.max_execution_time,
|
||||
'min_execution_time': (metrics.min_execution_time
|
||||
if metrics.min_execution_time != float('inf')
|
||||
else 0.0),
|
||||
'last_update_time': metrics.last_update_time,
|
||||
})
|
||||
self._write_snapshot()
|
||||
# Only after the write lands. Marking it first would mean a failed
|
||||
# set() bought the next interval's silence without leaving a snapshot.
|
||||
self._metrics_persisted_at[plugin_id] = now
|
||||
self._snapshot_persisted_at = now
|
||||
|
||||
def reset_metrics(self, plugin_id: str) -> None:
|
||||
"""Reset metrics for a plugin."""
|
||||
with self._lock:
|
||||
if plugin_id in self._metrics:
|
||||
self._metrics[plugin_id] = ResourceMetrics()
|
||||
cache_key = self._get_metrics_key(plugin_id)
|
||||
self.cache_manager.delete(cache_key)
|
||||
self._metrics_dirty.discard(plugin_id)
|
||||
self._write_snapshot(drop=plugin_id)
|
||||
# The record an older version wrote, so the reader's fallback
|
||||
# cannot bring the old numbers back.
|
||||
self.cache_manager.delete(self._get_metrics_key(plugin_id))
|
||||
# Let the next call persist immediately rather than leaving the
|
||||
# deleted key absent for the rest of the interval.
|
||||
self._metrics_persisted_at.pop(plugin_id, None)
|
||||
# plugin absent from the snapshot for the rest of the interval.
|
||||
self._snapshot_persisted_at = None
|
||||
|
||||
|
||||
@@ -0,0 +1,98 @@
|
||||
"""CacheManager builds its ConfigManager on first use, not in __init__.
|
||||
|
||||
Every CacheManager built a ConfigManager and loaded the whole config for a
|
||||
cache strategy that no longer reads it. The attribute stays public -- the
|
||||
sports plugins resolve the global timezone through
|
||||
``cache_manager.config_manager`` -- so it is now built on first access.
|
||||
"""
|
||||
|
||||
from unittest.mock import MagicMock, patch
|
||||
|
||||
import pytest
|
||||
|
||||
import src.config_manager as config_manager_module
|
||||
from src.cache_manager import CacheManager
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def built(monkeypatch):
|
||||
"""Count ConfigManager constructions and load_config calls."""
|
||||
made = []
|
||||
|
||||
class CountingConfigManager:
|
||||
def __init__(self):
|
||||
made.append(self)
|
||||
self.loads = 0
|
||||
|
||||
def load_config(self):
|
||||
self.loads += 1
|
||||
return {}
|
||||
|
||||
monkeypatch.setattr(config_manager_module, "ConfigManager", CountingConfigManager)
|
||||
return made
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def manager(tmp_path):
|
||||
with patch('src.cache_manager.CacheManager._get_writable_cache_dir',
|
||||
return_value=str(tmp_path)):
|
||||
cm = CacheManager()
|
||||
cm.stop_cleanup_thread()
|
||||
return cm
|
||||
|
||||
|
||||
def test_construction_does_not_load_the_config(built, manager):
|
||||
manager.set("k", {"v": 1})
|
||||
assert manager.get("k") == {"v": 1}
|
||||
assert manager.get_cache_strategy("sports_live")["max_age"] > 0
|
||||
assert built == []
|
||||
|
||||
|
||||
def test_first_access_builds_and_loads_it_once(built, manager):
|
||||
first = manager.config_manager
|
||||
assert manager.config_manager is first
|
||||
assert getattr(manager, "config_manager", None) is first
|
||||
assert len(built) == 1 and first.loads == 1
|
||||
|
||||
|
||||
def test_assignment_still_wins(built, manager):
|
||||
replacement = MagicMock()
|
||||
manager.config_manager = replacement
|
||||
assert manager.config_manager is replacement
|
||||
assert built == []
|
||||
|
||||
|
||||
def test_an_unimportable_config_manager_is_none(manager, monkeypatch):
|
||||
import builtins
|
||||
real_import = builtins.__import__
|
||||
|
||||
def refuse(name, *args, **kwargs):
|
||||
if name == "src.config_manager":
|
||||
raise ImportError("no config manager here")
|
||||
return real_import(name, *args, **kwargs)
|
||||
|
||||
monkeypatch.setattr(builtins, "__import__", refuse)
|
||||
assert manager.config_manager is None
|
||||
|
||||
|
||||
def test_a_failed_load_is_retried_on_the_next_access(manager, monkeypatch):
|
||||
attempts = []
|
||||
|
||||
class Flaky:
|
||||
def load_config(self):
|
||||
attempts.append(1)
|
||||
if len(attempts) == 1:
|
||||
raise RuntimeError("config.json unreadable")
|
||||
return {}
|
||||
|
||||
monkeypatch.setattr(config_manager_module, "ConfigManager", Flaky)
|
||||
with pytest.raises(RuntimeError):
|
||||
manager.config_manager
|
||||
assert isinstance(manager.config_manager, Flaky)
|
||||
assert len(attempts) == 2
|
||||
|
||||
|
||||
def test_a_manager_made_without_init_still_answers(built):
|
||||
bare = CacheManager.__new__(CacheManager)
|
||||
bare.logger = MagicMock()
|
||||
assert bare.config_manager is built[0]
|
||||
@@ -0,0 +1,242 @@
|
||||
"""An unchanged CacheManager.set() does not rewrite the file, and the skip
|
||||
never changes how old the record looks.
|
||||
|
||||
DiskCache.set already skipped a payload identical to the last one it wrote,
|
||||
but CacheManager.set stamps every record with time.time(), so for set() the
|
||||
payload was never identical and every unchanged re-save was a full rewrite on
|
||||
the SD card. The digest now leaves the header timestamp out, and the newer
|
||||
timestamp lives in the file's mtime instead ("UNCHANGED RE-SAVES" in
|
||||
src/cache/disk_cache.py). These tests pin both halves: the write is skipped,
|
||||
and every reader still ages the record from its newest save -- across a
|
||||
restart, and with a bound on what a foreign mtime can claim.
|
||||
"""
|
||||
|
||||
import os
|
||||
import shutil
|
||||
import tempfile
|
||||
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,
|
||||
_MAX_TIMESTAMP_LIFT,
|
||||
_effective_timestamp,
|
||||
)
|
||||
from src.cache_manager import CacheManager
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def writes(monkeypatch):
|
||||
"""Count real writes: every atomic write starts with mkstemp."""
|
||||
calls = []
|
||||
real = tempfile.mkstemp
|
||||
|
||||
def counting(*args, **kwargs):
|
||||
calls.append(kwargs.get("prefix"))
|
||||
return real(*args, **kwargs)
|
||||
|
||||
monkeypatch.setattr(disk_cache_module.tempfile, "mkstemp", counting)
|
||||
return calls
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def clock(monkeypatch):
|
||||
"""time.time() for the cache modules, advanced by hand."""
|
||||
now = [time.time()]
|
||||
fake = type("FakeTime", (), {"time": staticmethod(lambda: now[0])})
|
||||
monkeypatch.setattr(disk_cache_module, "time", fake)
|
||||
import src.cache_manager as cache_manager_module
|
||||
monkeypatch.setattr(cache_manager_module, "time", fake)
|
||||
return now
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def cm(tmp_path):
|
||||
with patch('src.cache_manager.CacheManager._get_writable_cache_dir',
|
||||
return_value=str(tmp_path)):
|
||||
manager = CacheManager()
|
||||
manager.stop_cleanup_thread()
|
||||
yield manager
|
||||
|
||||
|
||||
def _record(ts, data=None, ttl=None):
|
||||
"""A record laid out the way CacheManager.set writes it."""
|
||||
rec = {"timestamp": ts}
|
||||
if ttl is not None:
|
||||
rec["ttl"] = ttl
|
||||
rec["data"] = data if data is not None else {"games": [1, 2, 3]}
|
||||
return rec
|
||||
|
||||
|
||||
class TestTheWriteIsSkipped:
|
||||
def test_repeated_identical_set_writes_once(self, cm, writes):
|
||||
path = cm._get_cache_path("scores")
|
||||
for _ in range(20):
|
||||
cm.set("scores", {"games": [1, 2, 3]}, ttl=60)
|
||||
assert len(writes) == 1
|
||||
# The data is the same and the file is the same file.
|
||||
assert cm.get("scores", max_age=60, memory_ttl=0) == {"games": [1, 2, 3]}
|
||||
assert os.stat(path).st_nlink == 1
|
||||
|
||||
def test_the_file_is_not_replaced(self, tmp_path, clock):
|
||||
disk = DiskCache(str(tmp_path))
|
||||
disk.set("k", _record(clock[0]))
|
||||
before = os.stat(disk.get_cache_path("k"))
|
||||
clock[0] += 30
|
||||
disk.set("k", _record(clock[0]))
|
||||
after = os.stat(disk.get_cache_path("k"))
|
||||
assert after.st_ino == before.st_ino
|
||||
assert after.st_mtime == pytest.approx(clock[0], abs=1e-3)
|
||||
|
||||
def test_changed_data_rewrites(self, cm, writes):
|
||||
cm.set("scores", {"games": [1]})
|
||||
cm.set("scores", {"games": [2]})
|
||||
assert len(writes) == 2
|
||||
assert cm.get("scores", memory_ttl=0) == {"games": [2]}
|
||||
|
||||
def test_a_changed_ttl_rewrites(self, cm, writes):
|
||||
cm.set("scores", {"games": [1]}, ttl=60)
|
||||
cm.set("scores", {"games": [1]}, ttl=600)
|
||||
cm.set("scores", {"games": [1]})
|
||||
assert len(writes) == 3
|
||||
|
||||
def test_a_timestamp_going_backwards_rewrites(self, tmp_path, writes, clock):
|
||||
disk = DiskCache(str(tmp_path))
|
||||
disk.set("k", _record(clock[0]))
|
||||
disk.set("k", _record(clock[0] - 100))
|
||||
assert len(writes) == 2
|
||||
assert disk.get("k", max_age=None)["timestamp"] == pytest.approx(clock[0] - 100)
|
||||
|
||||
def test_unchanged_data_is_still_rewritten_once_the_lift_runs_out(
|
||||
self, tmp_path, writes, clock):
|
||||
disk = DiskCache(str(tmp_path))
|
||||
start = clock[0]
|
||||
disk.set("k", _record(start))
|
||||
clock[0] = start + _MAX_TIMESTAMP_LIFT - 1
|
||||
disk.set("k", _record(clock[0]))
|
||||
assert len(writes) == 1
|
||||
clock[0] = start + _MAX_TIMESTAMP_LIFT + 1
|
||||
disk.set("k", _record(clock[0]))
|
||||
assert len(writes) == 2
|
||||
# ...and the embedded timestamp caught up.
|
||||
with open(disk.get_cache_path("k"), "rb") as f:
|
||||
assert disk_cache_module._loads(f.read())["timestamp"] == clock[0]
|
||||
|
||||
def test_another_writer_replacing_the_file_forces_a_rewrite(self, tmp_path, writes):
|
||||
mine, theirs = DiskCache(str(tmp_path)), DiskCache(str(tmp_path))
|
||||
now = time.time()
|
||||
mine.set("k", _record(now, {"v": "mine"}))
|
||||
theirs.set("k", _record(now + 1, {"v": "theirs"}))
|
||||
mine.set("k", _record(now + 2, {"v": "mine"}))
|
||||
assert len(writes) == 3
|
||||
assert mine.get("k", max_age=None)["data"] == {"v": "mine"}
|
||||
|
||||
def test_a_record_without_a_timestamp_still_skips(self, tmp_path, writes):
|
||||
disk = DiskCache(str(tmp_path))
|
||||
disk.set("k", {"plain": True})
|
||||
disk.set("k", {"plain": True})
|
||||
assert len(writes) == 1
|
||||
|
||||
|
||||
class TestAgeAfterSkippedWrites:
|
||||
def test_a_skipped_save_keeps_the_record_fresh(self, tmp_path, clock):
|
||||
disk = DiskCache(str(tmp_path))
|
||||
disk.set("k", _record(clock[0], ttl=60))
|
||||
for _ in range(10): # ten minutes of unchanged 50 s re-saves
|
||||
clock[0] += 50
|
||||
disk.set("k", _record(clock[0], ttl=60))
|
||||
clock[0] += 50
|
||||
record = disk.get("k", max_age=300)
|
||||
assert record is not None
|
||||
# Handed back as a rewrite would have left it.
|
||||
assert record["timestamp"] == pytest.approx(clock[0] - 50, abs=1e-3)
|
||||
|
||||
def test_and_it_expires_on_time_once_the_saves_stop(self, tmp_path, clock, monkeypatch):
|
||||
disk = DiskCache(str(tmp_path))
|
||||
disk.set("k", _record(clock[0]))
|
||||
clock[0] += 200
|
||||
disk.set("k", _record(clock[0]))
|
||||
clock[0] += 59
|
||||
assert disk.get("k", max_age=60) is not None
|
||||
clock[0] += 2
|
||||
parses = []
|
||||
real = disk_cache_module._loads
|
||||
monkeypatch.setattr(disk_cache_module, "_loads",
|
||||
lambda raw: parses.append(1) or real(raw))
|
||||
assert disk.get("k", max_age=60) is None
|
||||
assert parses == [] # still decided from the header
|
||||
|
||||
def test_the_ttl_is_honoured_the_same_way(self, tmp_path, clock):
|
||||
disk = DiskCache(str(tmp_path))
|
||||
disk.set("k", _record(clock[0], ttl=30))
|
||||
clock[0] += 100
|
||||
disk.set("k", _record(clock[0], ttl=30))
|
||||
clock[0] += 20
|
||||
assert disk.get("k", max_age=5) is not None # ttl wins, 20 < 30
|
||||
clock[0] += 20
|
||||
assert disk.get("k", max_age=3600) is None # 40 > 30
|
||||
|
||||
def test_a_restart_sees_the_newest_save(self, tmp_path, clock, writes):
|
||||
disk = DiskCache(str(tmp_path))
|
||||
disk.set("k", _record(clock[0]))
|
||||
clock[0] += 250
|
||||
disk.set("k", _record(clock[0]))
|
||||
restarted = DiskCache(str(tmp_path)) # empty digest map
|
||||
assert restarted.get("k", max_age=60) is not None
|
||||
# It rewrites once (it cannot know what is on disk), then skips.
|
||||
clock[0] += 10
|
||||
restarted.set("k", _record(clock[0]))
|
||||
clock[0] += 10
|
||||
restarted.set("k", _record(clock[0]))
|
||||
assert len(writes) == 2
|
||||
|
||||
def test_the_cache_manager_reads_it_across_processes(self, cm, tmp_path, clock):
|
||||
cm.set("display_state", {"mode": "clock"})
|
||||
clock[0] += 100
|
||||
cm.set("display_state", {"mode": "clock"})
|
||||
# The web interface: its own manager, memory tier bypassed.
|
||||
with patch('src.cache_manager.CacheManager._get_writable_cache_dir',
|
||||
return_value=str(tmp_path)):
|
||||
web = CacheManager()
|
||||
web.stop_cleanup_thread()
|
||||
clock[0] += 60
|
||||
assert web.get("display_state", max_age=120, memory_ttl=0) == {"mode": "clock"}
|
||||
clock[0] += 70
|
||||
assert web.get("display_state", max_age=120, memory_ttl=0) is None
|
||||
|
||||
def test_retention_sees_the_newest_save(self, cm, clock):
|
||||
cm.set("odds_x", {"line": 1})
|
||||
path = cm._get_cache_path("odds_x")
|
||||
clock[0] += 3000
|
||||
cm.set("odds_x", {"line": 1})
|
||||
assert os.path.getmtime(path) == pytest.approx(clock[0], abs=1e-3)
|
||||
|
||||
|
||||
class TestFreshnessCannotBeBorrowed:
|
||||
def test_a_record_written_with_an_old_timestamp_reads_old(self, tmp_path):
|
||||
disk = DiskCache(str(tmp_path))
|
||||
old = time.time() - 600
|
||||
disk.set("k", _record(old))
|
||||
assert os.path.getmtime(disk.get_cache_path("k")) == pytest.approx(old, abs=1e-3)
|
||||
assert disk.get("k", max_age=300) is None
|
||||
|
||||
def test_a_copy_without_mtime_is_bounded(self, tmp_path):
|
||||
disk = DiskCache(str(tmp_path))
|
||||
day_old = time.time() - 86400
|
||||
disk.set("k", _record(day_old))
|
||||
copy_dir = tmp_path / "restored"
|
||||
copy_dir.mkdir()
|
||||
shutil.copyfile(disk.get_cache_path("k"), copy_dir / "k.json") # mtime = now
|
||||
restored = DiskCache(str(copy_dir))
|
||||
assert restored.get("k", max_age=300) is None
|
||||
record = restored.get("k", max_age=None)
|
||||
assert record["timestamp"] == pytest.approx(day_old + _MAX_TIMESTAMP_LIFT)
|
||||
|
||||
def test_effective_timestamp(self):
|
||||
assert _effective_timestamp(1000.0, None) == 1000.0
|
||||
assert _effective_timestamp(1000.0, 900.0) == 1000.0 # mtime older
|
||||
assert _effective_timestamp(1000.0, 1500.0) == 1500.0 # a skipped save
|
||||
assert _effective_timestamp(1000.0, 10 ** 9) == 1000.0 + _MAX_TIMESTAMP_LIFT
|
||||
@@ -13,6 +13,7 @@ No network: sessions are fakes, and the fetch service is a fresh one per test.
|
||||
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import threading
|
||||
import time
|
||||
from datetime import date, datetime
|
||||
@@ -321,6 +322,10 @@ class TestWithARealCacheManager:
|
||||
record["timestamp"] = time.time() - seconds
|
||||
with open(path, "w", encoding="utf-8") as fh:
|
||||
json.dump(record, fh)
|
||||
# Rewriting the file moves its mtime to now, and the disk cache takes a
|
||||
# record's age from the newer of its timestamp and its mtime (an
|
||||
# unchanged re-save only touches the file), so age the mtime as well.
|
||||
os.utime(path, (record["timestamp"], record["timestamp"]))
|
||||
cm._memory_cache_component.clear()
|
||||
|
||||
def test_a_writers_long_ttl_does_not_outlast_the_readers(self, cm, service):
|
||||
|
||||
@@ -18,6 +18,7 @@ from src.plugin_system.resource_monitor import (
|
||||
ResourceLimits,
|
||||
ResourceLimitExceeded,
|
||||
PSUTIL_AVAILABLE,
|
||||
METRICS_SNAPSHOT_KEY,
|
||||
)
|
||||
|
||||
|
||||
@@ -174,7 +175,7 @@ class TestMetricsPersistenceChurn:
|
||||
with patch.object(rm.time, "monotonic", return_value=12.0):
|
||||
mon.monitor_call("p", lambda: None)
|
||||
writes = [c for c in cache.set.call_args_list
|
||||
if "plugin_metrics:" in str(c)]
|
||||
if c.args and c.args[0] == rm.METRICS_SNAPSHOT_KEY]
|
||||
assert writes, \
|
||||
"the first snapshot was dropped because the process was young"
|
||||
|
||||
@@ -184,7 +185,7 @@ class TestMetricsPersistenceChurn:
|
||||
for _ in range(50):
|
||||
mon.monitor_call("p", lambda: None)
|
||||
writes = [c for c in cache.set.call_args_list
|
||||
if c.args and str(c.args[0]).startswith("plugin_metrics:")]
|
||||
if c.args and c.args[0] == METRICS_SNAPSHOT_KEY]
|
||||
assert len(writes) == 1, (
|
||||
f"50 calls produced {len(writes)} metric writes; expected 1")
|
||||
|
||||
@@ -194,10 +195,10 @@ class TestMetricsPersistenceChurn:
|
||||
mon = PluginResourceMonitor(cache, enable_monitoring=False)
|
||||
mon.monitor_call("p", lambda: None)
|
||||
# pretend the interval has passed
|
||||
mon._metrics_persisted_at["p"] -= rm._METRICS_PERSIST_INTERVAL + 1
|
||||
mon._snapshot_persisted_at -= rm._METRICS_PERSIST_INTERVAL + 1
|
||||
mon.monitor_call("p", lambda: None)
|
||||
writes = [c for c in cache.set.call_args_list
|
||||
if c.args and str(c.args[0]).startswith("plugin_metrics:")]
|
||||
if c.args and c.args[0] == METRICS_SNAPSHOT_KEY]
|
||||
assert len(writes) == 2
|
||||
|
||||
def test_in_memory_metrics_stay_exact_while_writes_are_skipped(self):
|
||||
@@ -213,8 +214,11 @@ class TestMetricsPersistenceChurn:
|
||||
mon.reset_metrics("p")
|
||||
mon.monitor_call("p", lambda: None)
|
||||
writes = [c for c in cache.set.call_args_list
|
||||
if c.args and str(c.args[0]).startswith("plugin_metrics:")]
|
||||
assert len(writes) == 2, "reset should clear the throttle timestamp"
|
||||
if c.args and c.args[0] == METRICS_SNAPSHOT_KEY]
|
||||
# The first call, the reset (which drops the plugin), the next call.
|
||||
assert len(writes) == 3, "reset should clear the throttle timestamp"
|
||||
assert "p" not in writes[1].args[1]["plugins"]
|
||||
assert writes[2].args[1]["plugins"]["p"]["call_count"] == 1
|
||||
|
||||
def test_a_failed_write_does_not_buy_the_next_interval_of_silence(self):
|
||||
"""A set() that raises must not count as having persisted.
|
||||
@@ -230,5 +234,107 @@ class TestMetricsPersistenceChurn:
|
||||
# the very next call must try again rather than skip the interval
|
||||
mon.monitor_call("p", lambda: None)
|
||||
writes = [c for c in cache.set.call_args_list
|
||||
if c.args and str(c.args[0]).startswith("plugin_metrics:")]
|
||||
if c.args and c.args[0] == METRICS_SNAPSHOT_KEY]
|
||||
assert len(writes) == 2, "a failed write should be retried, not skipped"
|
||||
|
||||
|
||||
class TestOneSnapshotForAllPlugins:
|
||||
"""Every plugin's metrics share one record, written at most once a minute.
|
||||
|
||||
A record per plugin, each throttled to 30 s, was still two writes a minute
|
||||
per plugin. The web UI's output must not change: it reads the same numbers
|
||||
for the same plugins, from the snapshot or, for a plugin the snapshot does
|
||||
not have yet, from the per-plugin record an older version left.
|
||||
"""
|
||||
|
||||
@pytest.fixture
|
||||
def cache_dir(self, tmp_path):
|
||||
return str(tmp_path)
|
||||
|
||||
def _manager(self, cache_dir):
|
||||
from src.cache_manager import CacheManager
|
||||
with patch('src.cache_manager.CacheManager._get_writable_cache_dir',
|
||||
return_value=cache_dir):
|
||||
manager = CacheManager()
|
||||
manager.stop_cleanup_thread()
|
||||
return manager
|
||||
|
||||
def test_many_plugins_one_write(self):
|
||||
cache = _cache()
|
||||
mon = PluginResourceMonitor(cache, enable_monitoring=False)
|
||||
for _ in range(10):
|
||||
for pid in ("a", "b", "c", "d"):
|
||||
mon.monitor_call(pid, lambda: None)
|
||||
sets = cache.set.call_args_list
|
||||
assert len(sets) == 1
|
||||
assert not any(str(c.args[0]).startswith("plugin_metrics:") for c in sets)
|
||||
|
||||
def test_the_web_reads_what_the_display_has(self, cache_dir):
|
||||
display = PluginResourceMonitor(self._manager(cache_dir), enable_monitoring=False)
|
||||
for pid in ("a", "b"):
|
||||
display.monitor_call(pid, lambda: None)
|
||||
display._snapshot_persisted_at = None # let the next call publish
|
||||
display.monitor_call("a", lambda: None)
|
||||
web = PluginResourceMonitor(self._manager(cache_dir), enable_monitoring=False)
|
||||
for pid in ("a", "b"):
|
||||
assert web.get_metrics_summary(pid, force_reload=True) == \
|
||||
display.get_metrics_summary(pid)
|
||||
assert web.get_metrics_summary("a", force_reload=True)["call_count"] == 2
|
||||
|
||||
def test_a_per_plugin_record_from_an_older_version_is_still_read(self, cache_dir):
|
||||
old = self._manager(cache_dir)
|
||||
old.set("plugin_metrics:legacy", {"call_count": 9, "total_execution_time": 1.8,
|
||||
"last_update_time": time.time()})
|
||||
display = PluginResourceMonitor(self._manager(cache_dir), enable_monitoring=False)
|
||||
display.monitor_call("other", lambda: None)
|
||||
web = PluginResourceMonitor(self._manager(cache_dir), enable_monitoring=False)
|
||||
assert web.get_metrics_summary("legacy", force_reload=True)["call_count"] == 9
|
||||
# The display carries the count on from it, into the snapshot.
|
||||
display.monitor_call("legacy", lambda: None)
|
||||
display._snapshot_persisted_at = None
|
||||
display.monitor_call("other", lambda: None)
|
||||
assert web.get_metrics_summary("legacy", force_reload=True)["call_count"] == 10
|
||||
|
||||
def test_a_restart_keeps_plugins_it_has_not_run(self, cache_dir):
|
||||
first = PluginResourceMonitor(self._manager(cache_dir), enable_monitoring=False)
|
||||
first.monitor_call("disabled_later", lambda: None)
|
||||
restarted = PluginResourceMonitor(self._manager(cache_dir), enable_monitoring=False)
|
||||
restarted.monitor_call("running", lambda: None)
|
||||
web = PluginResourceMonitor(self._manager(cache_dir), enable_monitoring=False)
|
||||
assert web.get_metrics_summary("disabled_later", force_reload=True)["call_count"] == 1
|
||||
assert web.get_metrics_summary("running", force_reload=True)["call_count"] == 1
|
||||
|
||||
def test_a_reset_from_the_web_sticks_for_a_plugin_the_display_is_not_running(
|
||||
self, cache_dir):
|
||||
display = PluginResourceMonitor(self._manager(cache_dir), enable_monitoring=False)
|
||||
display.monitor_call("idle", lambda: None)
|
||||
web = PluginResourceMonitor(self._manager(cache_dir), enable_monitoring=False)
|
||||
assert web.get_metrics_summary("idle", force_reload=True)["call_count"] == 1
|
||||
web.reset_metrics("idle")
|
||||
display._snapshot_persisted_at = None
|
||||
display.monitor_call("busy", lambda: None)
|
||||
assert web.get_metrics_summary("idle", force_reload=True)["call_count"] == 0
|
||||
|
||||
def test_a_long_idle_plugin_is_dropped(self, cache_dir):
|
||||
import src.plugin_system.resource_monitor as rm
|
||||
manager = self._manager(cache_dir)
|
||||
manager.set(METRICS_SNAPSHOT_KEY, {"schema": 1, "plugins": {
|
||||
"gone": {"call_count": 3, "last_update_time":
|
||||
time.time() - rm._METRICS_SNAPSHOT_ENTRY_MAX_AGE - 10},
|
||||
"recent": {"call_count": 4, "last_update_time": time.time() - 60},
|
||||
}})
|
||||
display = PluginResourceMonitor(manager, enable_monitoring=False)
|
||||
display.monitor_call("p", lambda: None)
|
||||
plugins = manager.get(METRICS_SNAPSHOT_KEY, max_age=None, memory_ttl=0)["plugins"]
|
||||
assert set(plugins) == {"recent", "p"}
|
||||
|
||||
@pytest.mark.parametrize("junk", [[1, 2], {"schema": 99, "plugins": {"p": {}}},
|
||||
{"schema": 1, "plugins": "nope"}])
|
||||
def test_an_unusable_snapshot_is_ignored(self, junk):
|
||||
cache = MagicMock()
|
||||
cache.get.side_effect = lambda key, **kw: junk if key == METRICS_SNAPSHOT_KEY else None
|
||||
mon = PluginResourceMonitor(cache, enable_monitoring=False)
|
||||
assert mon.get_metrics_summary("p", force_reload=True)["call_count"] == 0
|
||||
mon.monitor_call("p", lambda: None)
|
||||
written = cache.set.call_args.args[1]
|
||||
assert written["plugins"]["p"]["call_count"] == 1
|
||||
|
||||
Reference in New Issue
Block a user