mirror of
https://github.com/ChuckBuilds/LEDMatrix.git
synced 2026-08-01 16:58:06 +00:00
* feat(plugin-system): activate dormant plugin health & metrics subsystem
PluginManager shipped a fully-built health tracker, resource monitor and
circuit breaker that were never instantiated (health_tracker/resource_monitor
were left as None), so the circuit breaker never engaged and the existing
health/metrics API routes always returned "not available".
- DisplayController now wires a PluginHealthTracker and PluginResourceMonitor
onto the plugin manager, enabling the circuit breaker (a repeatedly-failing
plugin's update() is skipped after consecutive failures, then retried after
a cooldown) and per-plugin execution-time metrics. Both persist to the
shared cache.
- load_plugin() now validates each plugin's config against its JSON schema in
a strictly warn/degrade-only way: a violation logs a warning and flags the
plugin degraded in the health tracker, but never changes whether the plugin
loads or its pass/fail behaviour. Adds PluginHealthTracker.set_degraded(),
which never touches the circuit breaker.
- ResourceMonitor CPU/memory sampling now reuses a cached psutil.Process and
reads cpu_percent(interval=None), so monitoring no longer blocks ~100ms per
call on the display loop's update path.
- Fix DiskCache.get() raising TypeError for max_age=None ("never expires"),
which silently discarded persisted plugin health/metrics on read and thus
broke cross-process and post-restart surfacing.
- Fix two dead PluginManager helpers that called non-existent tracker methods.
Tests: new test_resource_monitor, test_plugin_health,
test_plugin_manager_schema_soft; extended test_cache_manager and
test_display_controller.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01UvTav268UXv44ub9K11LYq
* feat(web-ui): surface plugin health, metrics and load state
With the health/metrics subsystem now active in the display service, expose it
in the web UI (which runs as a separate process from the display loop):
- Wire a health tracker / resource monitor backed by the shared on-disk cache
into the web process so /api/v3/plugins/health and /plugins/metrics read the
data the display service persists.
- Build those route responses per installed plugin id (the tracker's in-memory
view is empty in a fresh web process) so cross-process data is included.
- Add state + error_info to /plugins/installed entries so the UI can show why a
plugin isn't running instead of just loaded:false.
- Add a "Plugin Health" panel to the Tools page (circuit status, avg/max update
time, update count, last error) plus PluginAPI.getPluginMetrics().
Tests: route-level tests for the health/metrics endpoints in test_web_api.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01UvTav268UXv44ub9K11LYq
* fix(plugin-metrics): refresh cross-process health/metrics reads; type hints
Addresses CodeRabbit review on #388:
- Major: the web process's health/resource trackers cached the first persisted
read in an in-memory dict (and the CacheManager memory tier held max_age=None
entries indefinitely), so a long-lived web process showed the first snapshot
and never reflected the display service's later updates. Add an opt-in
force_reload path (get_health_summary/get_health_state/_load_health_state and
get_metrics_summary/get_metrics) that bypasses the in-memory copy and, via a
new memory_ttl passthrough on CacheManager.get, the cache manager's memory
tier — so each /plugins/health and /plugins/metrics poll reads fresh persisted
state. Default behaviour (force_reload=False) is unchanged for the display
process and existing callers.
- Minor: DiskCache.get type hint is now Optional[int] with the None ("never
expires") semantics documented, matching MemoryCache.get.
Tests: new force_reload staleness cases in test_plugin_health and
test_resource_monitor.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01UvTav268UXv44ub9K11LYq
---------
Co-authored-by: Claude <noreply@anthropic.com>
374 lines
15 KiB
Python
374 lines
15 KiB
Python
"""
|
|
Plugin Resource Monitor
|
|
|
|
Tracks resource usage (memory, CPU, execution time) for plugins.
|
|
Provides resource limits and performance monitoring.
|
|
"""
|
|
|
|
import time
|
|
import logging
|
|
import threading
|
|
from typing import Dict, Optional, Any, Callable
|
|
from dataclasses import dataclass, field
|
|
|
|
try:
|
|
import psutil
|
|
PSUTIL_AVAILABLE = True
|
|
except ImportError:
|
|
PSUTIL_AVAILABLE = False
|
|
|
|
|
|
class ResourceLimitExceeded(Exception):
|
|
"""Raised when a plugin exceeds its resource limits."""
|
|
|
|
|
|
@dataclass
|
|
class ResourceLimits:
|
|
"""Resource limits for a plugin."""
|
|
max_memory_mb: Optional[float] = None # Maximum memory in MB
|
|
max_cpu_percent: Optional[float] = None # Maximum CPU percentage
|
|
max_execution_time: Optional[float] = None # Maximum execution time in seconds
|
|
warning_threshold: float = 0.8 # Warning at 80% of limit
|
|
|
|
|
|
@dataclass
|
|
class ResourceMetrics:
|
|
"""Resource usage metrics for a plugin."""
|
|
memory_mb: float = 0.0
|
|
cpu_percent: float = 0.0
|
|
execution_time: float = 0.0
|
|
call_count: int = 0
|
|
total_execution_time: float = 0.0
|
|
max_execution_time: float = 0.0
|
|
min_execution_time: float = float('inf')
|
|
last_update_time: float = field(default_factory=time.time)
|
|
|
|
def update_average_execution_time(self):
|
|
"""Update average execution time."""
|
|
if self.call_count > 0:
|
|
self.total_execution_time = self.total_execution_time / self.call_count
|
|
|
|
|
|
class PluginResourceMonitor:
|
|
"""
|
|
Monitors resource usage for plugins.
|
|
|
|
Tracks:
|
|
- Memory usage (if psutil available)
|
|
- CPU usage (if psutil available)
|
|
- Execution time for update() and display() calls
|
|
- Call counts and statistics
|
|
"""
|
|
|
|
def __init__(self, cache_manager, enable_monitoring: bool = True):
|
|
"""
|
|
Initialize resource monitor.
|
|
|
|
Args:
|
|
cache_manager: Cache manager for persisting metrics
|
|
enable_monitoring: Enable resource monitoring (requires psutil)
|
|
"""
|
|
self.cache_manager = cache_manager
|
|
self.enable_monitoring = enable_monitoring and PSUTIL_AVAILABLE
|
|
self.logger = logging.getLogger(__name__)
|
|
|
|
# Resource metrics per plugin
|
|
self._metrics: Dict[str, ResourceMetrics] = {}
|
|
self._limits: Dict[str, ResourceLimits] = {}
|
|
|
|
# Thread-local storage for execution tracking
|
|
self._local = threading.local()
|
|
|
|
# Lock for thread-safe access
|
|
self._lock = threading.Lock()
|
|
|
|
# Cache a single psutil.Process handle. Reusing the same handle is what
|
|
# lets cpu_percent() be read non-blocking (interval=None): psutil returns
|
|
# the utilisation since the *previous* call on that same object. Creating
|
|
# a fresh Process() per call would force interval-based sampling that
|
|
# blocks the caller — unacceptable on the display loop's update path.
|
|
self._process = None
|
|
if self.enable_monitoring:
|
|
try:
|
|
self._process = psutil.Process()
|
|
# Prime cpu_percent so the first real measurement returns a
|
|
# meaningful delta instead of 0.0.
|
|
self._process.cpu_percent(interval=None)
|
|
except Exception: # pragma: no cover - psutil edge cases
|
|
self._process = None
|
|
|
|
if not PSUTIL_AVAILABLE and enable_monitoring:
|
|
self.logger.warning(
|
|
"psutil not available - resource monitoring will be limited to execution time only"
|
|
)
|
|
|
|
def _get_metrics_key(self, plugin_id: str) -> str:
|
|
"""Get cache key for plugin metrics."""
|
|
return f"plugin_metrics:{plugin_id}"
|
|
|
|
def _get_limits_key(self, plugin_id: str) -> str:
|
|
"""Get cache key for plugin limits."""
|
|
return f"plugin_limits:{plugin_id}"
|
|
|
|
def get_metrics(self, plugin_id: str, force_reload: bool = False) -> ResourceMetrics:
|
|
"""Get current metrics for a plugin.
|
|
|
|
``force_reload=True`` bypasses both the in-memory copy and the cache
|
|
manager's memory tier so a read-only consumer (e.g. the web process)
|
|
sees the writer process's latest persisted metrics rather than a stale
|
|
first snapshot.
|
|
"""
|
|
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
|
|
)
|
|
if cached:
|
|
metrics = ResourceMetrics(**cached)
|
|
else:
|
|
metrics = ResourceMetrics()
|
|
self._metrics[plugin_id] = metrics
|
|
return self._metrics[plugin_id]
|
|
|
|
def set_limits(self, plugin_id: str, limits: ResourceLimits) -> None:
|
|
"""Set resource limits for a plugin."""
|
|
with self._lock:
|
|
self._limits[plugin_id] = limits
|
|
# Persist to cache
|
|
cache_key = self._get_limits_key(plugin_id)
|
|
self.cache_manager.set(cache_key, {
|
|
'max_memory_mb': limits.max_memory_mb,
|
|
'max_cpu_percent': limits.max_cpu_percent,
|
|
'max_execution_time': limits.max_execution_time,
|
|
'warning_threshold': limits.warning_threshold
|
|
})
|
|
|
|
def get_limits(self, plugin_id: str) -> Optional[ResourceLimits]:
|
|
"""Get resource limits for a plugin."""
|
|
with self._lock:
|
|
if plugin_id not in self._limits:
|
|
# Try to load from cache
|
|
cache_key = self._get_limits_key(plugin_id)
|
|
cached = self.cache_manager.get(cache_key, max_age=None)
|
|
if cached:
|
|
self._limits[plugin_id] = ResourceLimits(**cached)
|
|
else:
|
|
return None
|
|
return self._limits[plugin_id]
|
|
|
|
def _get_process_memory_mb(self) -> float:
|
|
"""Get current process memory usage in MB."""
|
|
if not self.enable_monitoring or self._process is None:
|
|
return 0.0
|
|
try:
|
|
return self._process.memory_info().rss / 1024 / 1024
|
|
except Exception:
|
|
return 0.0
|
|
|
|
def _get_process_cpu_percent(self) -> float:
|
|
"""Get current process CPU usage percentage (non-blocking).
|
|
|
|
Reads cpu_percent(interval=None) against the cached process handle, so
|
|
it returns immediately with the utilisation observed since the previous
|
|
call rather than blocking to sample a fresh interval.
|
|
"""
|
|
if not self.enable_monitoring or self._process is None:
|
|
return 0.0
|
|
try:
|
|
return self._process.cpu_percent(interval=None)
|
|
except Exception:
|
|
return 0.0
|
|
|
|
def monitor_call(self, plugin_id: str, func: Callable, *args, **kwargs) -> Any:
|
|
"""
|
|
Monitor a plugin method call.
|
|
|
|
Tracks execution time and resource usage, enforces limits.
|
|
|
|
Args:
|
|
plugin_id: Plugin identifier
|
|
func: Function to call
|
|
*args: Function arguments
|
|
**kwargs: Function keyword arguments
|
|
|
|
Returns:
|
|
Function return value
|
|
|
|
Raises:
|
|
ResourceLimitExceeded: If resource limits are exceeded
|
|
"""
|
|
metrics = self.get_metrics(plugin_id)
|
|
limits = self.get_limits(plugin_id)
|
|
|
|
# Record start time and memory
|
|
start_time = time.time()
|
|
start_memory = self._get_process_memory_mb()
|
|
|
|
try:
|
|
# Execute the function
|
|
result = func(*args, **kwargs)
|
|
|
|
# Calculate execution time
|
|
execution_time = time.time() - start_time
|
|
|
|
# Update metrics
|
|
with self._lock:
|
|
metrics.execution_time = execution_time
|
|
metrics.call_count += 1
|
|
metrics.total_execution_time += execution_time
|
|
metrics.max_execution_time = max(metrics.max_execution_time, execution_time)
|
|
if metrics.min_execution_time == float('inf'):
|
|
metrics.min_execution_time = execution_time
|
|
else:
|
|
metrics.min_execution_time = min(metrics.min_execution_time, execution_time)
|
|
metrics.last_update_time = time.time()
|
|
|
|
# Update memory and CPU if monitoring enabled
|
|
if self.enable_monitoring:
|
|
end_memory = self._get_process_memory_mb()
|
|
metrics.memory_mb = max(metrics.memory_mb, end_memory - start_memory)
|
|
# CPU is harder to measure per-call, so we track it separately
|
|
metrics.cpu_percent = self._get_process_cpu_percent()
|
|
|
|
# Persist metrics
|
|
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
|
|
})
|
|
|
|
# Check limits
|
|
if limits:
|
|
self._check_limits(plugin_id, metrics, limits, execution_time)
|
|
|
|
return result
|
|
|
|
except ResourceLimitExceeded:
|
|
raise
|
|
except Exception:
|
|
# Still record execution time even on error
|
|
execution_time = time.time() - start_time
|
|
with self._lock:
|
|
metrics.execution_time = execution_time
|
|
metrics.last_update_time = time.time()
|
|
raise
|
|
|
|
def _check_limits(self, plugin_id: str, metrics: ResourceMetrics,
|
|
limits: ResourceLimits, execution_time: float) -> None:
|
|
"""Check if plugin has exceeded resource limits."""
|
|
warnings = []
|
|
errors = []
|
|
|
|
# Check execution time
|
|
if limits.max_execution_time and execution_time > limits.max_execution_time:
|
|
errors.append(
|
|
f"Execution time {execution_time:.2f}s exceeds limit {limits.max_execution_time:.2f}s"
|
|
)
|
|
elif limits.max_execution_time and execution_time > limits.max_execution_time * limits.warning_threshold:
|
|
warnings.append(
|
|
f"Execution time {execution_time:.2f}s approaching limit {limits.max_execution_time:.2f}s"
|
|
)
|
|
|
|
# Check memory
|
|
if limits.max_memory_mb and metrics.memory_mb > limits.max_memory_mb:
|
|
errors.append(
|
|
f"Memory usage {metrics.memory_mb:.2f}MB exceeds limit {limits.max_memory_mb:.2f}MB"
|
|
)
|
|
elif limits.max_memory_mb and metrics.memory_mb > limits.max_memory_mb * limits.warning_threshold:
|
|
warnings.append(
|
|
f"Memory usage {metrics.memory_mb:.2f}MB approaching limit {limits.max_memory_mb:.2f}MB"
|
|
)
|
|
|
|
# Check CPU
|
|
if limits.max_cpu_percent and metrics.cpu_percent > limits.max_cpu_percent:
|
|
errors.append(
|
|
f"CPU usage {metrics.cpu_percent:.2f}% exceeds limit {limits.max_cpu_percent:.2f}%"
|
|
)
|
|
elif limits.max_cpu_percent and metrics.cpu_percent > limits.max_cpu_percent * limits.warning_threshold:
|
|
warnings.append(
|
|
f"CPU usage {metrics.cpu_percent:.2f}% approaching limit {limits.max_cpu_percent:.2f}%"
|
|
)
|
|
|
|
# Log warnings
|
|
for warning in warnings:
|
|
self.logger.warning(f"Plugin {plugin_id}: {warning}")
|
|
|
|
# Raise exception for errors
|
|
if errors:
|
|
error_msg = f"Plugin {plugin_id} exceeded resource limits: {'; '.join(errors)}"
|
|
self.logger.error(error_msg)
|
|
raise ResourceLimitExceeded(error_msg)
|
|
|
|
def get_metrics_summary(self, plugin_id: str, force_reload: bool = False) -> Dict[str, Any]:
|
|
"""Get metrics summary for a plugin.
|
|
|
|
``force_reload=True`` refreshes from the persisted cache first so
|
|
cross-process readers reflect the writer's latest metrics.
|
|
"""
|
|
metrics = self.get_metrics(plugin_id, force_reload=force_reload)
|
|
limits = self.get_limits(plugin_id)
|
|
|
|
avg_execution_time = 0.0
|
|
if metrics.call_count > 0:
|
|
avg_execution_time = metrics.total_execution_time / metrics.call_count
|
|
|
|
summary = {
|
|
'plugin_id': plugin_id,
|
|
'memory_mb': round(metrics.memory_mb, 2),
|
|
'cpu_percent': round(metrics.cpu_percent, 2),
|
|
'execution_time': round(metrics.execution_time, 3),
|
|
'avg_execution_time': round(avg_execution_time, 3),
|
|
'min_execution_time': round(metrics.min_execution_time if metrics.min_execution_time != float('inf') else 0.0, 3),
|
|
'max_execution_time': round(metrics.max_execution_time, 3),
|
|
'call_count': metrics.call_count,
|
|
'last_update_time': metrics.last_update_time
|
|
}
|
|
|
|
if limits:
|
|
summary['limits'] = {
|
|
'max_memory_mb': limits.max_memory_mb,
|
|
'max_cpu_percent': limits.max_cpu_percent,
|
|
'max_execution_time': limits.max_execution_time,
|
|
'warning_threshold': limits.warning_threshold
|
|
}
|
|
|
|
# Calculate usage percentages
|
|
if limits.max_memory_mb:
|
|
summary['memory_usage_percent'] = round(
|
|
(metrics.memory_mb / limits.max_memory_mb) * 100, 2
|
|
)
|
|
if limits.max_cpu_percent:
|
|
summary['cpu_usage_percent'] = round(
|
|
(metrics.cpu_percent / limits.max_cpu_percent) * 100, 2
|
|
)
|
|
if limits.max_execution_time:
|
|
summary['execution_time_usage_percent'] = round(
|
|
(avg_execution_time / limits.max_execution_time) * 100, 2
|
|
)
|
|
|
|
return summary
|
|
|
|
def get_all_metrics_summaries(self) -> Dict[str, Dict[str, Any]]:
|
|
"""Get metrics summaries for all tracked plugins."""
|
|
summaries = {}
|
|
for plugin_id in self._metrics.keys():
|
|
summaries[plugin_id] = self.get_metrics_summary(plugin_id)
|
|
return summaries
|
|
|
|
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)
|
|
|