From 26d697b5ef775f9f4346b20d00c499a2e7265aa0 Mon Sep 17 00:00:00 2001 From: Chuck <33324927+ChuckBuilds@users.noreply.github.com> Date: Tue, 29 Sep 2026 13:58:16 -0400 Subject: [PATCH] fix(plugins): one hung plugin no longer stops every plugin from updating The single update worker took each plugin's lock with a blocking acquire(), and the render thread holds that lock while it runs the plugin's display(). A display() that never returned -- or a first frame still running on the executor's lingering thread after its 30s timeout -- parked the worker for good, and no plugin updated again. - The worker waits at most PLUGIN_LOCK_TIMEOUT (5s, the bound unload_plugin already uses), skips the busy plugin and records it through the normal update-failure path as a hang (PluginBusyError), so repeats open its circuit breaker. Log lines about it are rate-limited per plugin. - display() is timed on every frame (two monotonic reads). Calls of 2s or more are logged once a minute and counted in plugin health (slow_call_count, last_slow_call); calls past the executor timeout, and a first frame still running at it, are recorded as hangs (hang_count, last_hang) and no longer as successes. An update() still running after its timeout is recorded as a hang too. - on_config_change() runs under the plugin's lock via PluginManager.apply_config_change(); if the lock stays busy the latest change is deferred to the update worker, applied as soon as the lock frees and before the plugin's next update() at the latest. The plugin API is unchanged. Which thread runs each hook is documented in PluginManager.__init__. Co-Authored-By: Claude Opus 5.5 --- CHANGELOG.md | 25 ++ docs/ARCHITECTURE.md | 3 +- docs/PLUGIN_API_REFERENCE.md | 6 + src/display_controller.py | 53 ++- src/plugin_system/plugin_executor.py | 27 +- src/plugin_system/plugin_health.py | 66 +++- src/plugin_system/plugin_manager.py | 323 +++++++++++++++++- test/test_plugin_config_preparation.py | 21 +- test/test_plugin_hang_containment.py | 442 +++++++++++++++++++++++++ 9 files changed, 939 insertions(+), 27 deletions(-) create mode 100644 test/test_plugin_hang_containment.py diff --git a/CHANGELOG.md b/CHANGELOG.md index b9b0001a..9d11936e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -19,6 +19,31 @@ accepts both, but the store flags the old spelling as deprecated ## Unreleased +### Fixes + +- One hung plugin no longer stops every plugin from updating. The single + update worker waited on each plugin's lock with no time limit, and the + render thread holds that lock while it runs the plugin's display(); a + display() that never returned (or a first frame still running after the + executor's 30s timeout) parked the worker for good, so scores, weather and + clocks all froze while the panel kept scrolling. The worker now waits at + most 5s (the bound `unload_plugin()` already uses), skips the busy plugin + and records the skip as a hang, so a plugin that keeps hanging opens its + circuit breaker and drops out of updates and rotation until the cooldown. + The other plugins keep updating. +- display() calls are timed on every frame. One taking 2s or more is logged + (at most once a minute per plugin) and counted in plugin health + (`slow_call_count`, `last_slow_call`); one that runs past the executor's + timeout counts as a hang (`hang_count`, `last_hang`) and as a failure to + the circuit breaker. A first frame that times out is no longer recorded as + a success, and an update() still running after its timeout is recorded as + a hang instead of leaving the plugin silently stuck. +- A plugin's `on_config_change()` no longer runs while its update() is + running on the worker thread. It now runs under the plugin's lock; if the + lock stays busy past the same 5s bound the change is handed to the update + worker, which applies the latest one as soon as the lock frees, and before + the plugin's next update() at the latest. The plugin API is unchanged. + ## 3.7.0 Sports consolidation stage 3 (#672). No behaviour change: nothing in core diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 52c82201..51cf4bd3 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -99,7 +99,8 @@ then normal rotation. changes. The controller refreshes its cached settings; enabling or disabling a plugin queues `_reconcile_enabled_plugins()`, which loads or unloads it on the display thread; each plugin gets `on_config_change()` - for its own section. Set `LEDMATRIX_HOT_RELOAD=false` to turn this off. + for its own section, under its plugin lock + (`PluginManager.apply_config_change()`). Set `LEDMATRIX_HOT_RELOAD=false` to turn this off. Matrix hardware settings are only read at start-up. - **Vegas mode.** [`src/vegas_mode/`](../src/vegas_mode/): the display loop calls `VegasModeCoordinator.run_iteration()` diff --git a/docs/PLUGIN_API_REFERENCE.md b/docs/PLUGIN_API_REFERENCE.md index 36405db2..ec53f223 100644 --- a/docs/PLUGIN_API_REFERENCE.md +++ b/docs/PLUGIN_API_REFERENCE.md @@ -151,6 +151,12 @@ Clean up resources when plugin is unloaded. Override to close connections, stop Called after plugin configuration is updated via web API. +In the display service it runs on the config watcher thread while holding +the plugin's lock, so it never overlaps your `update()` or `display()`. If +the plugin stays busy for more than 5 seconds, the change is applied later +from the update thread: as soon as the plugin is free, and before its next +`update()` at the latest. + #### `on_enable() -> None` Called when plugin is enabled. diff --git a/src/display_controller.py b/src/display_controller.py index e6ed837c..e0a61083 100644 --- a/src/display_controller.py +++ b/src/display_controller.py @@ -1021,17 +1021,28 @@ class DisplayController: accepts_display_mode: Whether display() takes ``display_mode``. force_clear: Passed through to display(). + Each call is timed (two monotonic reads) and handed to + PluginManager.note_display_duration, which logs and records slow + calls and counts one that ran past the executor's timeout as a hang. + Returns: display()'s result, or True when the frame was skipped because the plugin's update() holds its lock (the panel keeps the last frame; that is not a failure). """ - with self._display_lock_or_skip(getattr(plugin, 'plugin_id', None)) as can_display: + plugin_id = getattr(plugin, 'plugin_id', None) + with self._display_lock_or_skip(plugin_id) as can_display: if not can_display: return True - if accepts_display_mode: - return plugin.display(display_mode=mode, force_clear=force_clear) - return plugin.display(force_clear=force_clear) + started = time.monotonic() + try: + if accepts_display_mode: + return plugin.display(display_mode=mode, force_clear=force_clear) + return plugin.display(force_clear=force_clear) + finally: + note = getattr(self.plugin_manager, 'note_display_duration', None) + if note is not None and plugin_id: + note(plugin_id, time.monotonic() - started) def _health_tracker(self): """The plugin circuit breaker, or None when it is not enabled.""" @@ -2328,6 +2339,7 @@ class DisplayController: pm = self.plugin_manager display_lock = pm.get_plugin_lock(plugin_id) if pm else None can_display = display_lock is None or display_lock.acquire(blocking=False) + display_hung = False if display_lock is None: # Only when plugin loading failed part-way. @@ -2347,7 +2359,7 @@ class DisplayController: # thread actually finishes it, rather than # here when this dispatch merely returns. release_guard = threading.Lock() - released = {'done': False} + released = {'done': False, 'started': False} def _release_display_lock(): with release_guard: @@ -2358,6 +2370,7 @@ class DisplayController: if _accepts_display_mode: def _display_target(display_mode=None, force_clear=False): + released['started'] = True try: return manager_to_display.display( display_mode=display_mode, force_clear=force_clear) @@ -2365,11 +2378,13 @@ class DisplayController: _release_display_lock() else: def _display_target(force_clear=False): + released['started'] = True try: return manager_to_display.display(force_clear=force_clear) finally: _release_display_lock() + dispatch_start = time.monotonic() try: result = pm.plugin_executor.execute_display( types.SimpleNamespace(display=_display_target), @@ -2392,6 +2407,18 @@ class DisplayController: _release_display_lock() raise + dispatch_seconds = time.monotonic() - dispatch_start + if released['started'] and not released['done']: + # The executor gave up waiting and display() + # is still running on its thread, holding + # the lock. A hang, not a success: recorded + # so repeats open the circuit breaker, and + # the update worker's bounded wait skips it. + display_hung = True + pm.record_display_hang(plugin_id, dispatch_seconds) + else: + pm.note_display_duration(plugin_id, dispatch_seconds) + logger.debug(f"display() returned: {result} (type: {type(result)})") if isinstance(result, bool): display_result = result @@ -2405,7 +2432,7 @@ class DisplayController: # be lost when display() finally does run. if can_display: health_tracker = self._health_tracker() - if health_tracker is not None: + if health_tracker is not None and not display_hung: health_tracker.record_success(plugin_id) self.force_change = False except Exception as exc: # pylint: disable=broad-except @@ -3099,8 +3126,18 @@ class DisplayController: prepared = prepare(_pid, new_config) if callable(prepare) else None if isinstance(prepared, dict): new_config = prepared - _plugin.on_config_change(new_config) - logger.debug("Plugin %s notified of config change", _pid) + # Runs on ConfigService's watcher thread. Under the + # plugin's lock, so it cannot interleave with update() + # on the worker or display() on the render thread; a + # lock held past the bound defers it to the worker. + apply = getattr(self.plugin_manager, 'apply_config_change', None) + if callable(apply): + applied = apply(_pid, new_config, plugin_instance=_plugin) + else: + _plugin.on_config_change(new_config) + applied = True + logger.debug("Plugin %s notified of config change%s", _pid, + "" if applied else " (deferred: plugin busy)") except Exception as e: logger.error("Error in plugin %s config change handler: %s", _pid, e, exc_info=True) diff --git a/src/plugin_system/plugin_executor.py b/src/plugin_system/plugin_executor.py index 7fea423c..e58eff6b 100644 --- a/src/plugin_system/plugin_executor.py +++ b/src/plugin_system/plugin_executor.py @@ -19,8 +19,23 @@ class PluginTimeoutError(Exception): """Raised when a plugin operation times out.""" +class PluginBusyError(PluginTimeoutError): + """A plugin's lock stayed held past its bound. + + Not raised; recorded. The lock is held by the plugin's own display(), + update() or on_config_change() -- one that is hung or far slower than it + should be -- so the caller skipped the plugin rather than wait on it. + """ + + class PluginExecutor: """Handles plugin execution with timeout and error isolation.""" + + #: A display() call at least this long is logged and counted as slow. + #: A frame is milliseconds; two seconds is a plugin doing I/O in display(). + SLOW_DISPLAY_SECONDS = 2.0 + #: An update() call at least this long is logged as slow. + SLOW_UPDATE_SECONDS = 5.0 def __init__( self, @@ -117,15 +132,15 @@ class PluginExecutor: True if update succeeded, False otherwise """ try: - start_time = time.time() + start_time = time.monotonic() self.execute_with_timeout( lambda: plugin.update(), timeout=timeout, plugin_id=plugin_id ) - duration = time.time() - start_time + duration = time.monotonic() - start_time - if duration > 5.0: # Warn if update takes more than 5 seconds + if duration > self.SLOW_UPDATE_SECONDS: self.logger.warning( "Plugin %s update() took %.2fs (consider optimizing)", plugin_id, @@ -175,7 +190,7 @@ class PluginExecutor: True if display succeeded, False otherwise """ try: - start_time = time.time() + start_time = time.monotonic() # Does display() take a display_mode keyword? The caller usually # knows and caches the answer, so prefer what it passed. @@ -206,9 +221,9 @@ class PluginExecutor: plugin_id=plugin_id ) - duration = time.time() - start_time + duration = time.monotonic() - start_time - if duration > 2.0: # Warn if display takes more than 2 seconds + if duration > self.SLOW_DISPLAY_SECONDS: self.logger.warning( "Plugin %s display() took %.2fs (consider optimizing)", plugin_id, diff --git a/src/plugin_system/plugin_health.py b/src/plugin_system/plugin_health.py index ec4b3c0f..9c2de935 100644 --- a/src/plugin_system/plugin_health.py +++ b/src/plugin_system/plugin_health.py @@ -254,6 +254,66 @@ class PluginHealthTracker: self._save_health_state(plugin_id, state) + def record_hang(self, plugin_id: str, operation: str, seconds: float, + error: Optional[Exception] = None) -> None: + """Record a call that ran past its limit, or held the plugin's lock past it. + + Counts as a failure, so the ordinary circuit breaker handles a plugin + that keeps hanging: after ``failure_threshold`` in a row it is skipped + by both the update scheduler and the display rotation until the + cooldown ends. The hang itself is kept alongside (``hang_count``, + ``last_hang``) so the health API can tell "hung" from "raised". + + Args: + plugin_id: Plugin identifier + operation: What hung: ``"display"``, ``"update"``, or the lock + wait that found one of them still running. + seconds: How long it had been running, or how long the lock + was waited on, when this was recorded. + error: The error to store as ``last_error``; one is built from + the other arguments when omitted. + """ + state = self.get_health_state(plugin_id) + count = state.get('hang_count') + state['hang_count'] = (count if isinstance(count, int) and not isinstance(count, bool) + else 0) + 1 + state['last_hang'] = { + 'operation': operation, + 'seconds': round(float(seconds), 3), + 'time': time.time(), + } + if error is None: + error = TimeoutError(f"{operation} still running after {seconds:.1f}s") + # record_failure saves the record, hang fields included. + self.record_failure(plugin_id, error) + + #: Minimum seconds between persisting a plugin's slow-call counters. The + #: in-memory record is updated on every slow call; a plugin that is slow + #: on every frame must not become an SD-card write per frame. + SLOW_CALL_PERSIST_INTERVAL = 60.0 + + def record_slow_call(self, plugin_id: str, operation: str, seconds: float) -> None: + """Note a call that finished, but slowly. Reporting only. + + Unlike :meth:`record_hang` this never touches the circuit breaker: a + slow display() still drew its frame. + """ + state = self.get_health_state(plugin_id) + count = state.get('slow_call_count') + state['slow_call_count'] = (count if isinstance(count, int) and not isinstance(count, bool) + else 0) + 1 + now = time.time() + state['last_slow_call'] = { + 'operation': operation, + 'seconds': round(float(seconds), 3), + 'time': now, + } + saved_at = self.__dict__.setdefault('_slow_call_saved_at', {}) + last = saved_at.get(plugin_id) + if last is None or now - last >= self.SLOW_CALL_PERSIST_INTERVAL: + saved_at[plugin_id] = now + self._save_health_state(plugin_id, state) + def set_degraded(self, plugin_id: str, reason: Optional[str]) -> None: """Flag (or clear) a plugin as degraded without touching the circuit breaker. @@ -345,7 +405,11 @@ class PluginHealthTracker: 'degraded': state.get('degraded', False), 'degraded_reason': state.get('degraded_reason'), 'circuit_opened_time': state.get('circuit_opened_time'), - 'half_open_start_time': state.get('half_open_start_time') + 'half_open_start_time': state.get('half_open_start_time'), + 'hang_count': state.get('hang_count', 0), + 'last_hang': state.get('last_hang'), + 'slow_call_count': state.get('slow_call_count', 0), + 'last_slow_call': state.get('last_slow_call'), } def get_all_health_summaries(self) -> Dict[str, Dict[str, Any]]: diff --git a/src/plugin_system/plugin_manager.py b/src/plugin_system/plugin_manager.py index 6bd798d8..7c44fc66 100644 --- a/src/plugin_system/plugin_manager.py +++ b/src/plugin_system/plugin_manager.py @@ -16,12 +16,14 @@ import time import threading import types from pathlib import Path -from typing import Dict, List, Optional, Any, Tuple +from typing import Dict, List, NamedTuple, Optional, Any, Tuple, Union import logging from src.exceptions import PluginError, ConfigError from src.logging_config import get_logger from src.plugin_system.plugin_loader import PluginLoader -from src.plugin_system.plugin_executor import PluginExecutor +from src.plugin_system.plugin_executor import ( + PluginBusyError, PluginExecutor, PluginTimeoutError, +) from src.plugin_system.plugin_state import PluginStateManager, PluginState from src.plugin_system.schema_manager import ( CORE_VEGAS_TUNING_KEYS, SchemaManager, normalize_legacy_booleans, @@ -36,6 +38,16 @@ from src.common.permission_utils import ( ) +class _DeferredConfigChange(NamedTuple): + """Update-queue item: apply the config change parked for ``plugin_id``. + + Queued by apply_config_change() when the plugin's lock was busy; the + change itself waits in ``PluginManager._deferred_config_changes`` so only + the latest one is ever applied. + """ + plugin_id: str + + class PluginManager: """ Manages plugin discovery, loading, and lifecycle. @@ -56,6 +68,19 @@ class PluginManager: # How long unload_plugin() waits for an in-flight update() to finish # before tearing the instance down anyway. UNLOAD_LOCK_TIMEOUT = 5.0 + + # How long the update worker and apply_config_change() wait for a + # plugin's lock -- the same bound unload already uses for the same lock. + # A display() frame holds it for milliseconds, so this only runs out when + # the holder is hung or pathologically slow. The worker then skips that + # plugin (recorded as a hang, so repeats open its circuit breaker) + # instead of stalling every other plugin's update behind it. + PLUGIN_LOCK_TIMEOUT = UNLOAD_LOCK_TIMEOUT + + # Minimum seconds between repeats of the same hang/slow-call warning for + # one plugin. A hung plugin is re-detected every interval; a slow + # display() can be re-detected every frame. + HANG_LOG_INTERVAL = 60.0 def __init__(self, plugins_dir: str = "plugins", config_manager: Optional[Any] = None, @@ -122,7 +147,33 @@ class PluginManager: # post-timeout window. # Kill switch: plugin_system.synchronous_updates: true restores the # inline path. - self._update_queue: "queue.Queue[Optional[Tuple[str, float]]]" = queue.Queue() + # + # Which thread runs each plugin hook, and what it holds: + # __init__, on_enable the loading thread (main thread at startup, + # the render thread on a live enable). + # update() plugin-update-worker, under the plugin lock, + # via PluginExecutor (whose daemon thread runs + # the call; if it outlives the executor's + # timeout it keeps the lock until it returns). + # Exceptions: the startup pass + # (DisplayController._run_initial_updates, main + # thread, before the display loop starts) and + # the synchronous_updates kill switch (render + # thread) run it without the lock. + # display() the render thread, under a try-lock: a busy + # lock skips the frame. The first frame of a + # screen goes through PluginExecutor. Vegas + # mode's adapter and coordinator take the lock + # with a bounded wait. + # on_config_change() ConfigService's watcher thread, under the + # plugin lock via apply_config_change(); if the + # lock stays busy it is deferred to the update + # worker, which applies it under the lock. + # cleanup(), on_disable() whoever calls unload_plugin(), under the + # lock with UNLOAD_LOCK_TIMEOUT. + # No wait on a plugin lock is unbounded, so one hung plugin can only + # cost the worker PLUGIN_LOCK_TIMEOUT per attempt. + self._update_queue: "queue.Queue[Union[None, Tuple[str, float], _DeferredConfigChange]]" = queue.Queue() self._pending_updates: set = set() self._pending_lock = threading.Lock() # Serializes the "is this plugin eligible?" -> "claim it (RUNNING)" @@ -145,6 +196,12 @@ class PluginManager: # run_scheduled_updates_with_changes(). self._completed_updates: set = set() self._completed_updates_lock = threading.Lock() + # Config changes that found the plugin's lock busy, latest per plugin, + # with the instance they were meant for. See apply_config_change(). + self._deferred_config_changes: Dict[str, Tuple[Any, Dict[str, Any]]] = {} + self._deferred_config_lock = threading.Lock() + # key -> (monotonic time last logged, repeats suppressed since) + self._rate_limited_warnings: Dict[str, Tuple[float, int]] = {} self._synchronous_updates = False if self.config_manager is not None: try: @@ -676,6 +733,8 @@ class PluginManager: # Remove from active plugins del self.plugins[plugin_id] + with self._deferred_config_lock: + self._deferred_config_changes.pop(plugin_id, None) with self._plugin_last_update_lock: self.plugin_last_update.pop(plugin_id, None) self._update_interval_cache.pop(plugin_id, None) @@ -1012,6 +1071,8 @@ class PluginManager: self, plugin_id: str, exc: Optional[Exception] = None, + hang: Optional[Tuple[str, float]] = None, + log: bool = True, ) -> None: """Apply the standard failure-recovery path for a plugin update. @@ -1025,6 +1086,10 @@ class PluginManager: exc: The exception that caused the failure, if any. When None a synthetic ExecutionFailure exception is constructed from the timeout/executor-error path. + hang: ``(operation, seconds)`` when the failure is a hang rather + than an error, so health records it as one. + log: Log the generic failure line. Callers that already logged + something more specific (rate-limited) pass False. """ failure_time = time.time() if exc is not None: @@ -1040,13 +1105,98 @@ class PluginManager: 'timestamp': failure_time, 'recoverable': True, } - self.logger.warning("Plugin %s update() failed; will retry after interval", plugin_id) + if log: + self.logger.warning("Plugin %s update() failed; will retry after interval", plugin_id) with self._plugin_last_update_lock: self.plugin_last_update[plugin_id] = failure_time self.state_manager.set_state_with_error(plugin_id, PluginState.ENABLED, error_info) - if self.health_tracker: + if hang is not None: + self._record_hang(plugin_id, hang[0], hang[1], err) + elif self.health_tracker: self.health_tracker.record_failure(plugin_id, err) + def _warn_rate_limited(self, key: str, message: str, *args: Any) -> None: + """Log a warning at most once per HANG_LOG_INTERVAL for ``key``. + + Repeats in between are counted and the count is appended to the next + one that is logged, so the journal shows the problem continuing + without a line per frame or per scheduler tick. + """ + # setdefault: tests build bare managers with PluginManager.__new__. + seen = self.__dict__.setdefault('_rate_limited_warnings', {}) + now = time.monotonic() + last, suppressed = seen.get(key, (None, 0)) + if last is not None and now - last < self.HANG_LOG_INTERVAL: + seen[key] = (last, suppressed + 1) + return + seen[key] = (now, 0) + if suppressed: + message += " (%d more since the last warning)" + args = args + (suppressed,) + self.logger.warning(message, *args) + + def _record_hang(self, plugin_id: str, operation: str, seconds: float, + err: Exception) -> None: + """Record a hang in plugin health: a failure to the circuit breaker. + + Goes through PluginHealthTracker.record_hang when the tracker has it + (it also counts the hang separately), else plain record_failure. Never + raises: this runs on the update worker and the render thread. + """ + tracker = self.health_tracker + if tracker is None: + return + try: + record_hang = getattr(tracker, 'record_hang', None) + if callable(record_hang): + record_hang(plugin_id, operation, seconds, err) + else: + tracker.record_failure(plugin_id, err) + except Exception as e: # pylint: disable=broad-except + self.logger.debug("Could not record hang for %s: %s", plugin_id, e) + + def note_display_duration(self, plugin_id: str, seconds: float) -> None: + """Account for one display() call that took ``seconds``. + + Called by the render loop for every frame, so the common case is one + comparison. At or above PluginExecutor.SLOW_DISPLAY_SECONDS the call + is logged (rate-limited) and counted as slow in plugin health; at or + above the executor's timeout -- the limit the first frame of a screen + is already held to -- it is recorded as a hang, which the circuit + breaker counts as a failure. + """ + if seconds < PluginExecutor.SLOW_DISPLAY_SECONDS: + return + if seconds >= self.plugin_executor.default_timeout: + self.record_display_hang(plugin_id, seconds) + return + self._warn_rate_limited( + "slow-display:" + plugin_id, + "Plugin %s display() took %.2fs; a frame should take milliseconds " + "(is it fetching or loading files in display()?)", plugin_id, seconds) + tracker = self.health_tracker + record_slow = getattr(tracker, 'record_slow_call', None) if tracker is not None else None + if callable(record_slow): + try: + record_slow(plugin_id, 'display', seconds) + except Exception as e: # pylint: disable=broad-except + self.logger.debug("Could not record slow display for %s: %s", plugin_id, e) + + def record_display_hang(self, plugin_id: str, seconds: float) -> None: + """Record a display() call that ran ``seconds``, past its limit. + + Either it has since returned (note_display_duration) or it is still + running on the executor's lingering thread, holding the plugin's lock + (the render loop's first-frame dispatch). + """ + self._warn_rate_limited( + "hung-display:" + plugin_id, + "Plugin %s display() ran for at least %.1fs (limit %.0fs); recorded " + "as a hang -- repeated hangs open its circuit breaker", + plugin_id, seconds, self.plugin_executor.default_timeout) + self._record_hang(plugin_id, 'display', seconds, PluginTimeoutError( + f"Plugin {plugin_id} display() ran for at least {seconds:.1f}s")) + def run_scheduled_updates(self, current_time: Optional[float] = None) -> None: """ Trigger plugin updates based on their defined update intervals. @@ -1132,11 +1282,13 @@ class PluginManager: self.state_manager.set_state(plugin_id, PluginState.ENABLED) def get_plugin_lock(self, plugin_id: str) -> threading.Lock: - """Per-plugin lock keeping update() and display() mutually exclusive. + """Per-plugin lock keeping update(), display() and on_config_change() + mutually exclusive. The update worker holds it for the duration of a plugin's update(); the display side acquires it non-blocking and skips that frame's - display() call when the plugin is mid-update. + display() call when the plugin is mid-update. Every other waiter uses + a bounded acquire (see the thread notes in __init__). """ with self._plugin_locks_guard: lock = self._plugin_locks.get(plugin_id) @@ -1202,14 +1354,26 @@ class PluginManager: real update() call genuinely finishes (see _execute_update_now), which can be after this dispatch returns if PluginExecutor's own timeout elapses first. + + The lock wait is bounded by PLUGIN_LOCK_TIMEOUT. Whatever holds it + past that -- a hung display() on the render thread or a lingering + executor thread -- costs this worker that long once per attempt, and + the plugin is skipped and recorded as hung (_skip_busy_update); the + other plugins' queued updates carry on. """ while True: item = self._update_queue.get() if item is None: # shutdown sentinel return + if isinstance(item, _DeferredConfigChange): + self._apply_deferred_config_change(item.plugin_id) + continue plugin_id, scheduled_time = item lock = self.get_plugin_lock(plugin_id) - lock.acquire() + wait_start = time.monotonic() + if not lock.acquire(timeout=self.PLUGIN_LOCK_TIMEOUT): + self._skip_busy_update(plugin_id, time.monotonic() - wait_start) + continue plugin_instance = self.plugins.get(plugin_id) if plugin_instance is None: # unloaded while queued; its # lifecycle state was already cleared by unload_plugin — @@ -1218,6 +1382,9 @@ class PluginManager: with self._pending_lock: self._pending_updates.discard(plugin_id) continue + # A config change that found the lock busy goes in first, so + # this update() runs against the settings the user saved. + self._apply_deferred_config_locked(plugin_id, plugin_instance) try: self._execute_update_now(plugin_id, plugin_instance, scheduled_time, lock=lock) @@ -1228,6 +1395,129 @@ class PluginManager: self.logger.exception("update worker: unexpected error for %s", plugin_id) + def _skip_busy_update(self, plugin_id: str, waited: float) -> None: + """Give up on a queued update whose plugin lock stayed held. + + Same bookkeeping as a failed update() -- pending slot dropped before + the state returns to ENABLED, last-update stamped so the retry waits + a full interval -- recorded as a hang, so a plugin stuck this way + opens its circuit breaker after the usual number of attempts and + stops being scheduled or displayed until the cooldown. + """ + with self._pending_lock: + self._pending_updates.discard(plugin_id) + if plugin_id not in self.plugins: + # Unloaded while we waited: its lifecycle state is already + # cleared; recording a failure would resurrect it as ENABLED. + return + self._warn_rate_limited( + "busy-update:" + plugin_id, + "Plugin %s update skipped: its lock was still held after %.1fs " + "(a display() or update() of it is hung or very slow); other " + "plugins keep updating", plugin_id, waited) + self._record_update_failure( + plugin_id, + exc=PluginBusyError( + f"Plugin {plugin_id} busy: its lock was held for over {waited:.1f}s " + "by a hung or slow display()/update(); update skipped"), + hang=('update lock wait', waited), + log=False) + + def apply_config_change(self, plugin_id: str, new_config: Dict[str, Any], + plugin_instance: Optional[Any] = None) -> bool: + """Call ``on_config_change(new_config)`` without racing update()/display(). + + Runs on the calling thread -- ConfigService's watcher, for the display + service -- holding the plugin's lock, waited on for at most + PLUGIN_LOCK_TIMEOUT. If the lock is still busy (an update() mid-fetch + can outlast that) the change is parked and handed to the update + worker, which applies it under the same lock once it is free, and at + the latest just before the plugin's next update(). A later change for + the same plugin replaces a parked one. + + Exceptions from on_config_change propagate on the immediate path, + as they did when the caller invoked it directly. + + Args: + plugin_id: Plugin identifier. + new_config: The prepared config to hand the plugin. + plugin_instance: The instance to notify; defaults to the loaded one. + + Returns: + True if on_config_change ran now, False if it was deferred or there + is no loaded plugin to notify. + """ + if plugin_instance is None: + plugin_instance = self.plugins.get(plugin_id) + if plugin_instance is None or not hasattr(plugin_instance, 'on_config_change'): + return False + lock = self.get_plugin_lock(plugin_id) + if lock.acquire(timeout=self.PLUGIN_LOCK_TIMEOUT): + try: + with self._deferred_config_lock: + # This change supersedes any older one still parked. + self._deferred_config_changes.pop(plugin_id, None) + plugin_instance.on_config_change(new_config) + finally: + lock.release() + return True + + with self._deferred_config_lock: + self._deferred_config_changes[plugin_id] = (plugin_instance, new_config) + self._warn_rate_limited( + "busy-config:" + plugin_id, + "Plugin %s is busy (lock held for over %.1fs); its config change " + "will be applied by the update worker once it is free", + plugin_id, self.PLUGIN_LOCK_TIMEOUT) + try: + self._ensure_update_worker() + self._update_queue.put(_DeferredConfigChange(plugin_id)) + except Exception as exc: # pylint: disable=broad-except + # No worker (thread start refused): still parked, so the next + # update() of this plugin applies it. + self.logger.error( + "Could not queue the config change for plugin %s (%s: %s); it " + "will be applied before its next update()", + plugin_id, type(exc).__name__, exc) + return False + + def _apply_deferred_config_change(self, plugin_id: str) -> None: + """Worker side of a parked config change: take the lock, apply it.""" + with self._deferred_config_lock: + if plugin_id not in self._deferred_config_changes: + return # applied or superseded meanwhile + lock = self.get_plugin_lock(plugin_id) + wait_start = time.monotonic() + if not lock.acquire(timeout=self.PLUGIN_LOCK_TIMEOUT): + self._warn_rate_limited( + "busy-config:" + plugin_id, + "Plugin %s still busy after %.1fs; its config change stays " + "parked until its next update()", + plugin_id, time.monotonic() - wait_start) + return + try: + self._apply_deferred_config_locked(plugin_id, self.plugins.get(plugin_id)) + finally: + lock.release() + + def _apply_deferred_config_locked(self, plugin_id: str, + current_instance: Optional[Any]) -> None: + """Apply the parked config change for plugin_id; caller holds its lock.""" + with self._deferred_config_lock: + entry = self._deferred_config_changes.pop(plugin_id, None) + if entry is None: + return + instance, new_config = entry + if current_instance is None or instance is not current_instance: + # Unloaded, or reloaded as a new instance built from the current + # config: nothing left to tell. + return + try: + instance.on_config_change(new_config) + self.logger.info("Applied deferred config change for plugin %s", plugin_id) + except Exception: # pylint: disable=broad-except + self.logger.exception("Error in plugin %s config change handler", plugin_id) + def stop_update_worker(self, timeout: float = 5.0) -> None: """Signal the worker to exit (used by cleanup; thread is a daemon).""" if self._update_worker is not None and self._update_worker.is_alive(): @@ -1338,14 +1628,29 @@ class PluginManager: else: _finish(True) + started = time.monotonic() try: - self.plugin_executor.execute_update( + success = self.plugin_executor.execute_update( types.SimpleNamespace(update=_target_update), plugin_id) except Exception as exc: # pragma: no cover - defensive; execute_update # catches everything internally, but guarantee _finish still # runs (releasing the lock) if something unexpected slips through. self.logger.exception("Unexpected error dispatching update for %s: %s", plugin_id, exc) _finish(False, exc=exc) + return + if not success and not finished['done']: + # The executor stopped waiting but update() is still running: it + # keeps the lock and the RUNNING state until it returns (then + # _finish records the outcome). Say so now, rather than leave the + # plugin silently stuck; record_success on a late return clears it. + elapsed = time.monotonic() - started + self._warn_rate_limited( + "hung-update:" + plugin_id, + "Plugin %s update() still running after %.1fs; it keeps its " + "lock until it returns, and is not rescheduled until then", + plugin_id, elapsed) + self._record_hang(plugin_id, 'update', elapsed, PluginTimeoutError( + f"Plugin {plugin_id} update() still running after {elapsed:.1f}s")) def run_scheduled_updates_with_changes(self, current_time: Optional[float] = None) -> List[str]: """ diff --git a/test/test_plugin_config_preparation.py b/test/test_plugin_config_preparation.py index e1ffb6bb..c510d65c 100644 --- a/test/test_plugin_config_preparation.py +++ b/test/test_plugin_config_preparation.py @@ -84,6 +84,19 @@ def plugin_manager(plugin_dir, tmp_path): return manager +@pytest.fixture +def delivering_manager(tmp_path): + """A real PluginManager whose apply_config_change() the mocked one uses. + + The hot-reload callback hands on_config_change to the plugin manager, so + it runs under the plugin's lock; delegate to the real method so these + tests still see the plugin called. + """ + manager = PluginManager(plugins_dir=str(tmp_path / "plugin-repos")) + yield manager + manager.stop_update_worker() + + class TestPluginManagerPreparation: def test_legacy_boolean_and_defaults(self, plugin_manager): prepared = plugin_manager.prepare_plugin_config( @@ -107,11 +120,13 @@ class TestPluginManagerPreparation: class TestHotReload: - def test_on_config_change_gets_the_prepared_config(self, test_display_controller, plugin_manager): + def test_on_config_change_gets_the_prepared_config(self, test_display_controller, plugin_manager, + delivering_manager): controller = test_display_controller plugin = MagicMock() plugin.modes = ["demo"] pm = controller.plugin_manager + pm.apply_config_change.side_effect = delivering_manager.apply_config_change pm.discover_plugins.return_value = ["demo"] pm.load_plugin.return_value = True pm.plugin_manifests = {} @@ -128,11 +143,13 @@ class TestHotReload: "enabled": True, "max_duration_seconds": 300} assert new_config["nhl"]["show_records"] is True - def test_raw_section_still_delivered_without_a_preparer(self, test_display_controller): + def test_raw_section_still_delivered_without_a_preparer(self, test_display_controller, + delivering_manager): controller = test_display_controller plugin = MagicMock() plugin.modes = ["demo"] pm = controller.plugin_manager + pm.apply_config_change.side_effect = delivering_manager.apply_config_change pm.discover_plugins.return_value = ["demo"] pm.load_plugin.return_value = True pm.plugin_manifests = {} diff --git a/test/test_plugin_hang_containment.py b/test/test_plugin_hang_containment.py new file mode 100644 index 00000000..e886527f --- /dev/null +++ b/test/test_plugin_hang_containment.py @@ -0,0 +1,442 @@ +"""One hung plugin must not stop every other plugin from updating. + +The update worker is a single thread. It took each plugin's lock with a +blocking ``acquire()``, and the render thread holds that same lock while it +runs the plugin's display(). A display() that never returned -- or a first +frame that outlived PluginExecutor's timeout and kept running on its lingering +thread -- parked the worker on that acquire for good, and from then on no +plugin updated at all: scores, weather and clocks all froze while the panel +kept scrolling stale data. + +Now: +- the worker waits PLUGIN_LOCK_TIMEOUT at most, skips the busy plugin and + records the skip as a hang, so repeats open that plugin's circuit breaker; +- display() calls are timed, slow ones recorded, overlong ones counted as hangs; +- on_config_change runs under the plugin's lock, deferred to the worker if the + lock stays busy, so it never interleaves with update(). + +Timeouts here are fractions of a second so the suite stays fast. +""" + +import copy +import threading +import time +import types +from unittest.mock import MagicMock + +import pytest + +from src.plugin_system.plugin_health import CircuitState, PluginHealthTracker +from src.plugin_system.plugin_manager import PluginManager +from src.plugin_system.plugin_state import PluginState + + +class _Cache: + """Serialising stand-in for CacheManager; counts writes.""" + + def __init__(self): + self.store = {} + self.writes = 0 + + def set(self, key, data, ttl=None, **kwargs): + self.writes += 1 + self.store[key] = copy.deepcopy(data) + + def get(self, key, max_age=None, memory_ttl=None, **kwargs): + return copy.deepcopy(self.store.get(key)) + + +class _Plugin: + """Counts update() calls; display() can be made to block on an Event.""" + + def __init__(self, plugin_id, update_seconds=0.0): + self.plugin_id = plugin_id + self.enabled = True + self.update_seconds = update_seconds + self.update_calls = 0 + self.in_update = False + self.display_gate = None # threading.Event: display() blocks on it + self.config_changes = [] + self.overlap = False + self.events = [] + + def update(self): + self.in_update = True + self.update_calls += 1 + self.events.append('update') + time.sleep(self.update_seconds) + self.in_update = False + + def display(self, force_clear=False): + if self.display_gate is not None: + self.display_gate.wait(timeout=10) + return True + + def on_config_change(self, new_config): + if self.in_update: + self.overlap = True + self.events.append('config') + self.config_changes.append(new_config) + + +@pytest.fixture +def pm(tmp_path): + manager = PluginManager(plugins_dir=str(tmp_path), config_manager=None, + display_manager=None, cache_manager=None) + manager.PLUGIN_LOCK_TIMEOUT = 0.15 + yield manager + manager.stop_update_worker() + + +@pytest.fixture +def tracker(pm): + t = PluginHealthTracker(cache_manager=_Cache()) + pm.health_tracker = t + return t + + +def _install(pm, plugin, interval=0.01): + pm.plugins[plugin.plugin_id] = plugin + pm._update_interval_cache[plugin.plugin_id] = interval + pm.state_manager.set_state(plugin.plugin_id, PluginState.ENABLED) + return plugin.plugin_id + + +def _hang_display(pm, plugin): + """Run the plugin's display() the way the render loop does, on its own + thread, holding the plugin's lock while display() blocks on its gate.""" + plugin.display_gate = threading.Event() + entered = threading.Event() + + def render(): + lock = pm.get_plugin_lock(plugin.plugin_id) + with lock: + entered.set() + plugin.display() + + thread = threading.Thread(target=render, daemon=True) + thread.start() + assert entered.wait(timeout=2) + return plugin.display_gate, thread + + +def _wait_for(predicate, timeout=3.0): + deadline = time.monotonic() + timeout + while time.monotonic() < deadline: + if predicate(): + return True + time.sleep(0.01) + return predicate() + + +class TestHungDisplayDoesNotStallTheWorker: + def test_other_plugins_keep_updating_on_schedule(self, pm): + hung = _Plugin('hung') + healthy = _Plugin('healthy') + _install(pm, hung) + _install(pm, healthy) + gate, render_thread = _hang_display(pm, hung) + try: + # 'hung' is dict-ordered first, so its item is dequeued ahead of + # 'healthy' every round. With a blocking acquire the worker parks + # on it forever and 'healthy' never updates. + for _ in range(8): + pm.run_scheduled_updates() + time.sleep(0.05) + assert _wait_for(lambda: healthy.update_calls >= 2), ( + f"healthy plugin updated {healthy.update_calls} time(s) while " + "another plugin's display() was hung") + assert hung.update_calls == 0 + assert pm._update_worker.is_alive() + finally: + gate.set() + render_thread.join(timeout=2) + + def test_hung_plugin_is_retried_once_its_display_returns(self, pm): + plugin = _Plugin('hung') + _install(pm, plugin) + gate, render_thread = _hang_display(pm, plugin) + pm.run_scheduled_updates() + # Skipped, and handed back rather than left RUNNING. + assert _wait_for(lambda: pm.state_manager.can_execute('hung')) + assert 'hung' not in pm._pending_updates + gate.set() + render_thread.join(timeout=2) + pm.plugin_last_update.pop('hung', None) # due again now + pm.run_scheduled_updates() + assert _wait_for(lambda: plugin.update_calls == 1) + + +class TestLockTimeoutIsRecordedAsAHang: + def test_skip_lands_in_health_and_state(self, pm, tracker): + plugin = _Plugin('hung') + _install(pm, plugin) + gate, render_thread = _hang_display(pm, plugin) + try: + pm.run_scheduled_updates() + assert _wait_for(lambda: tracker.get_health_summary('hung')['hang_count'] == 1) + summary = tracker.get_health_summary('hung') + assert summary['consecutive_failures'] == 1 + assert summary['last_hang']['operation'] == 'update lock wait' + assert 'busy' in summary['last_error'] + error_info = pm.state_manager.get_error_info('hung') + assert error_info['error_type'] == 'PluginBusyError' + # Stamped like any failed update, so the retry waits an interval. + assert pm.plugin_last_update.get('hung', 0) > 0 + finally: + gate.set() + render_thread.join(timeout=2) + + def test_repeated_hangs_open_the_circuit_breaker(self, pm, tracker): + plugin = _Plugin('hung') + _install(pm, plugin) + gate, render_thread = _hang_display(pm, plugin) + try: + for attempt in range(1, tracker.failure_threshold + 1): + pm.run_scheduled_updates() + assert _wait_for( + lambda: pm.state_manager.can_execute('hung') + and tracker.get_health_summary('hung')['hang_count'] == attempt) + time.sleep(0.02) # past the 0.01s interval + summary = tracker.get_health_summary('hung') + assert summary['circuit_state'] == CircuitState.OPEN.value + assert tracker.should_skip_plugin('hung') is True + # Circuit open: the scheduler no longer queues it at all. + pm.run_scheduled_updates() + assert 'hung' not in pm._pending_updates + assert pm.state_manager.can_execute('hung') + finally: + gate.set() + render_thread.join(timeout=2) + + def test_skip_is_logged_once_per_interval(self, pm): + pm.logger = MagicMock() + plugin = _Plugin('hung') + _install(pm, plugin) + gate, render_thread = _hang_display(pm, plugin) + try: + for _ in range(3): + pm.run_scheduled_updates() + assert _wait_for(lambda: pm.state_manager.can_execute('hung')) + time.sleep(0.02) + finally: + gate.set() + render_thread.join(timeout=2) + busy_warnings = [c for c in pm.logger.warning.call_args_list + if 'update skipped' in c.args[0]] + assert len(busy_warnings) == 1 + + def test_unloaded_while_waiting_is_not_resurrected(self, pm, tracker): + plugin = _Plugin('hung') + _install(pm, plugin) + gate, render_thread = _hang_display(pm, plugin) + try: + pm.PLUGIN_LOCK_TIMEOUT = 0.3 + pm.UNLOAD_LOCK_TIMEOUT = 0.05 + pm.run_scheduled_updates() + time.sleep(0.05) # worker is now waiting on the lock + assert pm.unload_plugin('hung') is True + assert _wait_for(lambda: 'hung' not in pm._pending_updates) + time.sleep(0.35) + assert pm.state_manager.get_state('hung') == PluginState.UNLOADED + assert tracker.get_health_summary('hung')['hang_count'] == 0 + finally: + gate.set() + render_thread.join(timeout=2) + + +class TestUpdateOutlivingTheExecutorTimeout: + def test_recorded_as_a_hang_then_cleared_when_it_returns(self, pm, tracker): + plugin = _Plugin('slow', update_seconds=0.4) + _install(pm, plugin) + pm.plugin_executor.default_timeout = 0.05 + pm.run_scheduled_updates() + assert _wait_for(lambda: tracker.get_health_summary('slow')['hang_count'] == 1) + summary = tracker.get_health_summary('slow') + assert summary['last_hang']['operation'] == 'update' + assert summary['consecutive_failures'] == 1 + # The real update() then returns: success resets the streak. + assert _wait_for(lambda: pm.state_manager.can_execute('slow')) + assert _wait_for(lambda: tracker.get_health_summary('slow')['consecutive_failures'] == 0) + + +class TestDisplayTiming: + def test_fast_frames_record_nothing(self, pm, tracker): + pm.logger = MagicMock() + for _ in range(100): + pm.note_display_duration('p', 0.004) + assert tracker.get_health_summary('p')['slow_call_count'] == 0 + pm.logger.warning.assert_not_called() + + def test_slow_display_is_recorded_and_warned_once(self, pm, tracker): + pm.logger = MagicMock() + for _ in range(5): + pm.note_display_duration('p', 2.5) + summary = tracker.get_health_summary('p') + assert summary['slow_call_count'] == 5 + assert summary['last_slow_call']['operation'] == 'display' + assert summary['last_slow_call']['seconds'] == 2.5 + # Slow is not failing: the circuit breaker is untouched. + assert summary['consecutive_failures'] == 0 + assert summary['circuit_state'] == CircuitState.CLOSED.value + assert pm.logger.warning.call_count == 1 + + def test_suppressed_repeats_are_counted_in_the_next_warning(self, pm): + pm.logger = MagicMock() + for _ in range(4): + pm.note_display_duration('p', 2.5) + pm.HANG_LOG_INTERVAL = 0.0 + pm.note_display_duration('p', 2.5) + assert pm.logger.warning.call_count == 2 + last = pm.logger.warning.call_args + assert '3 more since the last warning' in (last.args[0] % last.args[1:]) + + def test_display_past_the_executor_timeout_is_a_hang(self, pm, tracker): + pm.plugin_executor.default_timeout = 3.0 + for _ in range(tracker.failure_threshold): + pm.note_display_duration('p', 3.5) + summary = tracker.get_health_summary('p') + assert summary['hang_count'] == tracker.failure_threshold + assert summary['last_hang']['operation'] == 'display' + assert summary['circuit_state'] == CircuitState.OPEN.value + + def test_render_loop_times_each_frame(self, monkeypatch, emulator_mode): + """DisplayController._display_once hands every frame's duration on.""" + from src import display_controller as dc_module + ticks = iter([100.0, 103.25]) + monkeypatch.setattr(dc_module, 'time', types.SimpleNamespace( + monotonic=lambda: next(ticks))) + controller = dc_module.DisplayController.__new__(dc_module.DisplayController) + controller.plugin_manager = MagicMock() + lock = threading.Lock() + controller.plugin_manager.get_plugin_lock.return_value = lock + plugin = _Plugin('ticker') + + assert controller._display_once(plugin, 'ticker', False) is True + + controller.plugin_manager.note_display_duration.assert_called_once_with('ticker', 3.25) + assert not lock.locked() + + def test_skipped_frame_is_not_timed(self, emulator_mode): + from src import display_controller as dc_module + controller = dc_module.DisplayController.__new__(dc_module.DisplayController) + controller.plugin_manager = MagicMock() + lock = threading.Lock() + lock.acquire() # update() in flight + controller.plugin_manager.get_plugin_lock.return_value = lock + try: + assert controller._display_once(_Plugin('ticker'), 'ticker', False) is True + finally: + lock.release() + controller.plugin_manager.note_display_duration.assert_not_called() + + +class TestConfigChangeIsSerialised: + def test_does_not_run_concurrently_with_update(self, pm): + pm.PLUGIN_LOCK_TIMEOUT = 2.0 + plugin = _Plugin('p', update_seconds=0.3) + _install(pm, plugin) + pm.run_scheduled_updates() + assert _wait_for(lambda: plugin.in_update) + + applied = pm.apply_config_change('p', {'enabled': True, 'n': 1}) + + assert applied is True + assert plugin.overlap is False + assert plugin.events == ['update', 'config'] + assert plugin.config_changes == [{'enabled': True, 'n': 1}] + + def test_runs_holding_the_plugin_lock(self, pm): + plugin = _Plugin('p') + _install(pm, plugin) + held = {} + plugin.on_config_change = lambda cfg: held.setdefault( + 'locked', pm.get_plugin_lock('p').locked()) + assert pm.apply_config_change('p', {}) is True + assert held == {'locked': True} + assert not pm.get_plugin_lock('p').locked() + + def test_busy_lock_defers_the_latest_change_to_before_the_next_update(self, pm): + plugin = _Plugin('p') + _install(pm, plugin) + gate, render_thread = _hang_display(pm, plugin) + try: + start = time.monotonic() + assert pm.apply_config_change('p', {'n': 1}) is False + assert time.monotonic() - start < 1.0, "the wait must be bounded" + assert pm.apply_config_change('p', {'n': 2}) is False + time.sleep(0.4) # the worker's own bounded retry also finds it busy + assert plugin.config_changes == [] + finally: + gate.set() + render_thread.join(timeout=2) + + pm.run_scheduled_updates() + assert _wait_for(lambda: plugin.update_calls == 1) + # Only the latest parked change, applied before the update ran. + assert plugin.config_changes == [{'n': 2}] + assert plugin.events == ['config', 'update'] + assert plugin.overlap is False + + def test_deferred_change_applied_by_worker_once_lock_frees(self, pm): + """No update due to piggyback on: the worker's own queued retry + applies it. Timeline: the watcher gives up at 0.2s, the worker waits + from 0.2s to 0.4s, the holder lets go at 0.3s.""" + pm.PLUGIN_LOCK_TIMEOUT = 0.2 + plugin = _Plugin('p') + _install(pm, plugin, interval=3600) + pm.plugin_last_update['p'] = time.time() # not due + lock = pm.get_plugin_lock('p') + lock.acquire() + releaser = threading.Timer(0.3, lock.release) + releaser.start() + try: + assert pm.apply_config_change('p', {'n': 1}) is False + assert _wait_for(lambda: plugin.config_changes == [{'n': 1}]) + finally: + releaser.join(timeout=2) + assert plugin.update_calls == 0 + assert 'p' not in pm._deferred_config_changes + + def test_stale_deferred_change_is_dropped_after_reload(self, pm): + plugin = _Plugin('p') + _install(pm, plugin) + pm._deferred_config_changes['p'] = (plugin, {'n': 1}) + replacement = _Plugin('p') + pm.plugins['p'] = replacement + pm._apply_deferred_config_change('p') + assert plugin.config_changes == [] + assert replacement.config_changes == [] + assert 'p' not in pm._deferred_config_changes + + def test_exceptions_still_reach_the_caller(self, pm): + plugin = _Plugin('p') + _install(pm, plugin) + + def boom(cfg): + raise ValueError("bad config") + + plugin.on_config_change = boom + with pytest.raises(ValueError): + pm.apply_config_change('p', {}) + assert not pm.get_plugin_lock('p').locked() + + +class TestHealthRecords: + def test_slow_call_persistence_is_rate_limited(self): + cache = _Cache() + t = PluginHealthTracker(cache_manager=cache) + t.record_slow_call('p', 'display', 2.5) + writes = cache.writes + for _ in range(50): + t.record_slow_call('p', 'display', 2.5) + assert cache.writes == writes + assert t.get_health_summary('p')['slow_call_count'] == 51 + + def test_hang_survives_a_restart(self): + cache = _Cache() + PluginHealthTracker(cache_manager=cache).record_hang('p', 'display', 31.0) + summary = PluginHealthTracker(cache_manager=cache).get_health_summary('p') + assert summary['hang_count'] == 1 + assert summary['last_hang']['operation'] == 'display' + assert summary['consecutive_failures'] == 1