Skip to content
40 changes: 40 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
218 changes: 181 additions & 37 deletions src/cache/disk_cache.py

Large diffs are not rendered by default.

59 changes: 49 additions & 10 deletions src/cache_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,9 @@
# 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."""

Expand Down Expand Up @@ -73,21 +76,19 @@ def __init__(self) -> None:
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
Expand Down Expand Up @@ -115,6 +116,44 @@ def __init__(self) -> None:
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)
Expand Down
2 changes: 1 addition & 1 deletion src/error_aggregator.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.
#
Expand Down
141 changes: 107 additions & 34 deletions src/plugin_system/resource_monitor.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@
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
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -140,10 +155,15 @@ def __init__(self, cache_manager, enable_monitoring: bool = True):
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()
Expand Down Expand Up @@ -247,10 +267,14 @@ def get_metrics(self, plugin_id: str, force_reload: bool = False) -> ResourceMet
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:
Expand Down Expand Up @@ -498,12 +522,70 @@ def get_all_metrics_summaries(self) -> Dict[str, Dict[str, Any]]:
summaries[plugin_id] = self.get_metrics_summary(plugin_id)
return summaries

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:
"""Write a plugin's metrics to the cache, at most once per interval.
"""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
Expand All @@ -515,35 +597,26 @@ def _persist_metrics(self, plugin_id: str, metrics: ResourceMetrics,
# 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

Loading
Loading