diff --git a/CHANGELOG.md b/CHANGELOG.md index 8faf579c..c81c38fd 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -36,17 +36,21 @@ accepts both, but the store flags the old spelling as deprecated 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. + 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. + 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 diff --git a/src/plugin_system/plugin_executor.py b/src/plugin_system/plugin_executor.py index e58eff6b..2fa9f16f 100644 --- a/src/plugin_system/plugin_executor.py +++ b/src/plugin_system/plugin_executor.py @@ -23,8 +23,10 @@ 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. + 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. """ diff --git a/src/plugin_system/plugin_health.py b/src/plugin_system/plugin_health.py index 9c2de935..e30ef644 100644 --- a/src/plugin_system/plugin_health.py +++ b/src/plugin_system/plugin_health.py @@ -256,7 +256,7 @@ class PluginHealthTracker: 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. + """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 @@ -264,12 +264,14 @@ class PluginHealthTracker: 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"``, ``"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. + 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. """ @@ -287,9 +289,9 @@ class PluginHealthTracker: # 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. + #: 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: @@ -308,10 +310,46 @@ class PluginHealthTracker: 'seconds': round(float(seconds), 3), 'time': now, } - saved_at = self.__dict__.setdefault('_slow_call_saved_at', {}) - last = saved_at.get(plugin_id) + 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[plugin_id] = now + saved_at[(kind, plugin_id)] = now self._save_health_state(plugin_id, state) def set_degraded(self, plugin_id: str, reason: Optional[str]) -> None: @@ -410,6 +448,8 @@ class PluginHealthTracker: '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 5b623634..e9089afc 100644 --- a/src/plugin_system/plugin_manager.py +++ b/src/plugin_system/plugin_manager.py @@ -1071,8 +1071,8 @@ class PluginManager: self, plugin_id: str, exc: Optional[Exception] = None, - hang: Optional[Tuple[str, float]] = None, log: bool = True, + count_failure: bool = True, ) -> None: """Apply the standard failure-recovery path for a plugin update. @@ -1086,10 +1086,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. - 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. + 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: @@ -1110,9 +1111,7 @@ class PluginManager: 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 hang is not None: - self._record_hang(plugin_id, hang[0], hang[1], err) - elif 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: @@ -1351,10 +1350,12 @@ class PluginManager: 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. + 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() @@ -1394,29 +1395,42 @@ class PluginManager: """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. + 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 a failure would resurrect it as ENABLED. + # 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() or update() of it is hung or very slow); other " - "plugins keep updating", plugin_id, waited) + "(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 hung or slow display()/update(); update skipped"), - hang=('update lock wait', waited), - log=False) + "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: diff --git a/test/test_plugin_hang_containment.py b/test/test_plugin_hang_containment.py index e886527f..92d4692a 100644 --- a/test/test_plugin_hang_containment.py +++ b/test/test_plugin_hang_containment.py @@ -9,9 +9,12 @@ 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; +- 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(). @@ -167,18 +170,24 @@ class TestHungDisplayDoesNotStallTheWorker: assert _wait_for(lambda: plugin.update_calls == 1) -class TestLockTimeoutIsRecordedAsAHang: - def test_skip_lands_in_health_and_state(self, pm, tracker): +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')['hang_count'] == 1) + assert _wait_for(lambda: tracker.get_health_summary('hung')['busy_skip_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'] + 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. @@ -187,24 +196,78 @@ class TestLockTimeoutIsRecordedAsAHang: gate.set() render_thread.join(timeout=2) - def test_repeated_hangs_open_the_circuit_breaker(self, pm, tracker): + 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')['hang_count'] == attempt) + 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 - assert pm.state_manager.can_execute('hung') finally: gate.set() render_thread.join(timeout=2) @@ -239,7 +302,9 @@ class TestLockTimeoutIsRecordedAsAHang: 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 + 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) @@ -433,6 +498,23 @@ class TestHealthRecords: 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)