diff --git a/CHANGELOG.md b/CHANGELOG.md index 735bd356..3dc65368 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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:` 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/` 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:` 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 diff --git a/src/cache/disk_cache.py b/src/cache/disk_cache.py index fb776b56..a31782b9 100644 --- a/src/cache/disk_cache.py +++ b/src/cache/disk_cache.py @@ -14,7 +14,7 @@ import tempfile import logging import threading import zlib -from typing import Dict, Any, Optional, Protocol +from typing import Dict, Any, Optional, Protocol, Tuple from datetime import datetime from src.common.path_safety import safe_path_component @@ -111,18 +111,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. diff --git a/src/cache_manager.py b/src/cache_manager.py index 48dfb726..56fbdf9a 100644 --- a/src/cache_manager.py +++ b/src/cache_manager.py @@ -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) diff --git a/src/error_aggregator.py b/src/error_aggregator.py index a683fe9c..cc3bab12 100644 --- a/src/error_aggregator.py +++ b/src/error_aggregator.py @@ -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. # diff --git a/src/plugin_system/resource_monitor.py b/src/plugin_system/resource_monitor.py index 01b3204e..b04e1b7b 100644 --- a/src/plugin_system/resource_monitor.py +++ b/src/plugin_system/resource_monitor.py @@ -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: }}``, each metrics +#: record shaped as the per-plugin ``plugin_metrics:`` 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 diff --git a/test/test_cache_manager_lazy_config.py b/test/test_cache_manager_lazy_config.py new file mode 100644 index 00000000..ee855154 --- /dev/null +++ b/test/test_cache_manager_lazy_config.py @@ -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] diff --git a/test/test_cache_unchanged_resave.py b/test/test_cache_unchanged_resave.py new file mode 100644 index 00000000..b07c2cf8 --- /dev/null +++ b/test/test_cache_unchanged_resave.py @@ -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 diff --git a/test/test_espn_scoreboard_cache.py b/test/test_espn_scoreboard_cache.py index fe7332aa..79d92869 100644 --- a/test/test_espn_scoreboard_cache.py +++ b/test/test_espn_scoreboard_cache.py @@ -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): diff --git a/test/test_resource_monitor.py b/test/test_resource_monitor.py index 5e913210..988c0afe 100644 --- a/test/test_resource_monitor.py +++ b/test/test_resource_monitor.py @@ -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