From c0d97e4867d74c9aec934c7ea66742cef8333a99 Mon Sep 17 00:00:00 2001 From: Chuck <33324927+ChuckBuilds@users.noreply.github.com> Date: Wed, 30 Sep 2026 09:10:41 -0400 Subject: [PATCH] fix(plugins): one hung plugin no longer stops every plugin from updating (#677) The update worker no longer blocks forever on a plugin whose display() never returns. It waits at most PLUGIN_LOCK_TIMEOUT (5s) for a plugin's lock, then skips that plugin's update (a report-only "busy skip" in health) and keeps updating every other plugin. display() frames are timed (slow calls logged and counted; calls past the executor timeout recorded as hangs), a hung update() is recorded, and on_config_change() now runs under the plugin lock or is deferred to the worker. The plugin-facing API is unchanged. Co-Authored-By: Claude Opus 5.5 --- CHANGELOG.md | 26 ++ docs/ARCHITECTURE.md | 3 +- docs/PLUGIN_API_REFERENCE.md | 6 + src/display_controller.py | 53 ++- src/plugin_system/plugin_executor.py | 29 +- src/plugin_system/plugin_health.py | 106 ++++- src/plugin_system/plugin_manager.py | 332 +++++++++++++++- test/test_on_demand_disabled_plugin.py | 5 +- test/test_plugin_config_preparation.py | 21 +- test/test_plugin_hang_containment.py | 524 +++++++++++++++++++++++++ 10 files changed, 1077 insertions(+), 28 deletions(-) create mode 100644 test/test_plugin_hang_containment.py diff --git a/CHANGELOG.md b/CHANGELOG.md index b073bb11..c9e9ac37 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -101,6 +101,32 @@ accepts both, but the store flags the old spelling as deprecated - A stop request now clears an on-demand error. After a failed request, `/display/on-demand/status` kept reporting `status: error` for up to two minutes even after a stop. +- 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) and skips that update; + the other plugins keep updating. The skip is logged (at most once a minute + per plugin) and counted in plugin health as a busy skip (`busy_skip_count`, + `last_busy_skip`), but it is not a failure and never opens the circuit + breaker: Vegas mode holds a plugin's lock for its whole content render, + which on a slow Pi can outlast 5s, and a healthy plugin must not be pulled + from rotation for that. +- 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. Only these real hangs + count toward the breaker. +- 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 diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index bf94e6c8..78e990d1 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -105,7 +105,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 46ff8fd0..9d61a333 100644 --- a/src/display_controller.py +++ b/src/display_controller.py @@ -1029,17 +1029,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.""" @@ -2498,6 +2509,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. @@ -2517,7 +2529,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: @@ -2528,6 +2540,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) @@ -2535,11 +2548,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), @@ -2562,6 +2577,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 @@ -2575,7 +2602,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 @@ -3274,8 +3301,18 @@ class DisplayController: # says disabled, and on_config_change would switch # the instance off mid-session. new_config = {**new_config, 'enabled': True} - _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..2fa9f16f 100644 --- a/src/plugin_system/plugin_executor.py +++ b/src/plugin_system/plugin_executor.py @@ -19,8 +19,25 @@ 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(), on_config_change() or a Vegas content render -- slow, or hung + -- so the caller skipped the plugin rather than wait on it. Report-only: + it is kept as the plugin's state error info and counted as a busy skip in + health, never as a failure, so it cannot open the circuit breaker. + """ + + 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 +134,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 +192,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 +223,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..e30ef644 100644 --- a/src/plugin_system/plugin_health.py +++ b/src/plugin_system/plugin_health.py @@ -254,6 +254,104 @@ 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 display() or update() call that ran past its limit. + + 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". + + Not for an update skipped because the plugin's lock stayed held: the + holder may be a healthy but long render (Vegas prefetch). That is + :meth:`record_busy_skip`, which never touches the breaker. + + Args: + plugin_id: Plugin identifier + operation: What hung: ``"display"`` or ``"update"``. + seconds: How long it had been running 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 or busy-skip + #: counters. The in-memory record is updated every time; 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, + } + self._save_reporting_throttled('slow', plugin_id, state, now) + + def record_busy_skip(self, plugin_id: str, operation: str, seconds: float) -> None: + """Note a call skipped because the plugin's lock stayed held. Reporting only. + + The update worker gives up on a plugin's lock after + ``PluginManager.PLUGIN_LOCK_TIMEOUT``. Whatever held it may be healthy + -- Vegas prefetch holds the lock for a plugin's whole content render, + which on a slow Pi can take longer than that -- so like + :meth:`record_slow_call` this never touches the circuit breaker, the + failure streak or ``last_error``. A real hang is recorded by + :meth:`record_hang` where it is measured. + + Args: + plugin_id: Plugin identifier + operation: What was skipped, e.g. ``"update lock wait"``. + seconds: How long the lock was waited on. + """ + state = self.get_health_state(plugin_id) + count = state.get('busy_skip_count') + state['busy_skip_count'] = (count if isinstance(count, int) and not isinstance(count, bool) + else 0) + 1 + now = time.time() + state['last_busy_skip'] = { + 'operation': operation, + 'seconds': round(float(seconds), 3), + 'time': now, + } + self._save_reporting_throttled('busy', plugin_id, state, now) + + def _save_reporting_throttled(self, kind: str, plugin_id: str, + state: Dict[str, Any], now: float) -> None: + """Persist a reporting-only change at most once per + SLOW_CALL_PERSIST_INTERVAL per plugin and ``kind``. The first one is + saved at once, so the web process (which reads the persisted record) + sees it; repeats in between stay in memory until the next save.""" + saved_at = self.__dict__.setdefault('_reporting_saved_at', {}) + last = saved_at.get((kind, plugin_id)) + if last is None or now - last >= self.SLOW_CALL_PERSIST_INTERVAL: + saved_at[(kind, 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 +443,13 @@ 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'), + 'busy_skip_count': state.get('busy_skip_count', 0), + 'last_busy_skip': state.get('last_busy_skip'), } 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 92175416..87ae4642 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: @@ -686,6 +743,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) @@ -1022,6 +1081,8 @@ class PluginManager: self, plugin_id: str, exc: Optional[Exception] = None, + log: bool = True, + count_failure: bool = True, ) -> None: """Apply the standard failure-recovery path for a plugin update. @@ -1035,6 +1096,11 @@ 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. + log: Log the generic failure line. Callers that already logged + something more specific (rate-limited) pass False. + count_failure: Record the failure in plugin health, where it + counts toward the circuit breaker. A busy skip passes False: + it records itself as a busy skip, reporting only. """ failure_time = time.time() if exc is not None: @@ -1050,13 +1116,91 @@ 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 count_failure and 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. + + PluginHealthTracker.record_hang also counts the hang separately. Never + raises: this runs on the update worker and the render thread. + """ + tracker = self.health_tracker + if tracker is None: + return + try: + tracker.record_hang(plugin_id, operation, seconds, 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. @@ -1142,11 +1286,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) @@ -1212,14 +1358,28 @@ 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, a lingering + executor thread, or a long but healthy Vegas content render -- costs + this worker that long once per attempt, and the plugin's update is + skipped and reported as a busy skip (_skip_busy_update), which never + counts toward the circuit breaker; 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 — @@ -1228,6 +1388,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) @@ -1238,6 +1401,142 @@ 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 with PluginBusyError error info, + last-update stamped so the retry waits a full interval -- but + report-only in health: counted as a busy skip (``busy_skip_count`` / + ``last_busy_skip``), never as a failure or a hang. The lock holder + may be perfectly healthy: Vegas prefetch holds a plugin's lock for its + whole content render, which on a slow Pi can outlast + PLUGIN_LOCK_TIMEOUT, and counting that would pull a healthy plugin + from rotation. Real hangs -- display() or update() past the executor + timeout -- are recorded where they are measured and still open the + breaker. + """ + 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 anything 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(), Vegas render or update() of it is still running); " + "retrying next interval, not counted as a failure", 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 slow or hung display()/update(); update skipped"), + log=False, + count_failure=False) + tracker = self.health_tracker + record_busy = getattr(tracker, 'record_busy_skip', None) if tracker is not None else None + if callable(record_busy): + try: + record_busy(plugin_id, 'update lock wait', waited) + except Exception as e: # pylint: disable=broad-except + self.logger.debug("Could not record busy skip for %s: %s", plugin_id, e) + + 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(): @@ -1348,14 +1647,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_on_demand_disabled_plugin.py b/test/test_on_demand_disabled_plugin.py index 072d8805..a7c52f80 100644 --- a/test/test_on_demand_disabled_plugin.py +++ b/test/test_on_demand_disabled_plugin.py @@ -260,8 +260,11 @@ class TestReleasingThePlugin: callback({}, {'enabled': False, 'color': 'red'}) + # The change goes through the manager's locked apply_config_change + # (which calls on_config_change under the plugin's lock). plugin = controller.plugin_modes['preview_a'] - plugin.on_config_change.assert_called_once_with({'enabled': True, 'color': 'red'}) + controller.plugin_manager.apply_config_change.assert_called_once_with( + 'preview-me', {'enabled': True, 'color': 'red'}, plugin_instance=plugin) class TestRestoredSession: 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..92d4692a --- /dev/null +++ b/test/test_plugin_hang_containment.py @@ -0,0 +1,524 @@ +"""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 and skips the busy plugin. The + skip is report-only (a busy skip in health): the holder may be a healthy + Vegas prefetch render, so it never counts toward the circuit breaker; +- display() calls are timed, slow ones recorded, overlong ones counted as + hangs; update() past the executor timeout is a hang too. Only hangs open + the breaker; +- 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 TestLockTimeoutIsReportOnly: + def test_skip_lands_in_health_and_state_as_a_busy_skip(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')['busy_skip_count'] == 1) + summary = tracker.get_health_summary('hung') + assert summary['last_busy_skip']['operation'] == 'update lock wait' + assert summary['last_busy_skip']['seconds'] >= pm.PLUGIN_LOCK_TIMEOUT - 0.01 + # Reporting only: no failure, no hang, no error, breaker closed. + assert summary['consecutive_failures'] == 0 + assert summary['total_failures'] == 0 + assert summary['hang_count'] == 0 + assert summary['last_error'] is None + assert summary['circuit_state'] == CircuitState.CLOSED.value + # Still visible in the plugin's state for the web UI. + 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_busy_skips_never_open_the_circuit_breaker(self, pm, tracker): + """A Vegas prefetch render can hold a healthy plugin's lock past the + bound on every update; that must not pull it from rotation.""" + plugin = _Plugin('busy') + _install(pm, plugin) + gate, render_thread = _hang_display(pm, plugin) + skips = tracker.failure_threshold * 3 + try: + for attempt in range(1, skips + 1): + pm.run_scheduled_updates() + assert _wait_for( + lambda: pm.state_manager.can_execute('busy') + and tracker.get_health_summary('busy')['busy_skip_count'] == attempt) + time.sleep(0.02) # past the 0.01s interval + summary = tracker.get_health_summary('busy') + assert summary['busy_skip_count'] == skips + assert summary['consecutive_failures'] == 0 + assert summary['hang_count'] == 0 + assert summary['circuit_state'] == CircuitState.CLOSED.value + assert tracker.should_skip_plugin('busy') is False + finally: + gate.set() + render_thread.join(timeout=2) + # Once the lock frees the plugin updates normally. + pm.plugin_last_update.pop('busy', None) + pm.run_scheduled_updates() + assert _wait_for(lambda: plugin.update_calls == 1) + + def test_busy_skips_do_not_add_to_a_real_hang_streak(self, pm, tracker): + plugin = _Plugin('p') + _install(pm, plugin) + tracker.record_hang('p', 'display', 31.0) + gate, render_thread = _hang_display(pm, plugin) + try: + for attempt in range(1, tracker.failure_threshold + 2): + pm.run_scheduled_updates() + assert _wait_for( + lambda: pm.state_manager.can_execute('p') + and tracker.get_health_summary('p')['busy_skip_count'] == attempt) + time.sleep(0.02) + finally: + gate.set() + render_thread.join(timeout=2) + summary = tracker.get_health_summary('p') + assert summary['consecutive_failures'] == 1 + assert summary['hang_count'] == 1 + assert summary['circuit_state'] == CircuitState.CLOSED.value + + def test_repeated_real_hangs_still_open_the_circuit_breaker(self, pm, tracker): + """display() past the executor timeout is a real hang: the breaker + opens after failure_threshold of them, busy skips in between or not.""" + plugin = _Plugin('hung') + _install(pm, plugin) + gate, render_thread = _hang_display(pm, plugin) + too_long = pm.plugin_executor.default_timeout + 1 + 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')['busy_skip_count'] == attempt) + assert tracker.get_health_summary('hung')['circuit_state'] == CircuitState.CLOSED.value + pm.note_display_duration('hung', too_long) + time.sleep(0.02) # past the 0.01s interval + summary = tracker.get_health_summary('hung') + assert summary['hang_count'] == tracker.failure_threshold + assert summary['busy_skip_count'] == tracker.failure_threshold + 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 + 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 + summary = tracker.get_health_summary('hung') + assert summary['hang_count'] == 0 + assert summary['busy_skip_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_busy_skip_persistence_is_rate_limited_and_breaker_free(self): + cache = _Cache() + t = PluginHealthTracker(cache_manager=cache) + t.record_busy_skip('p', 'update lock wait', 5.0) + writes = cache.writes + # The first one is saved at once, so the web process sees it. + assert PluginHealthTracker(cache_manager=cache).get_health_summary( + 'p')['busy_skip_count'] == 1 + for _ in range(50): + t.record_busy_skip('p', 'update lock wait', 5.0) + assert cache.writes == writes + summary = t.get_health_summary('p') + assert summary['busy_skip_count'] == 51 + assert summary['consecutive_failures'] == 0 + assert summary['circuit_state'] == CircuitState.CLOSED.value + assert t.should_skip_plugin('p') is False + + 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