mirror of
https://github.com/ChuckBuilds/LEDMatrix.git
synced 2026-10-04 06:15:09 +00:00
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 <noreply@anthropic.com>
This commit is contained in:
@@ -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
|
||||
|
||||
|
||||
@@ -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()`
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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]]:
|
||||
|
||||
@@ -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]:
|
||||
"""
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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 = {}
|
||||
|
||||
@@ -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
|
||||
Reference in New Issue
Block a user