mirror of
https://github.com/ChuckBuilds/LEDMatrix.git
synced 2026-10-06 07:15:09 +00:00
fix(plugins): a busy-lock update skip is report-only, never a breaker failure
When the update worker gives up waiting for a plugin's lock (PLUGIN_LOCK_TIMEOUT, 5s) the skip was recorded as a hang, so three in a row opened the circuit breaker. Vegas prefetch holds a plugin's lock for its whole content render, which on a slow Pi can outlast 5s, so a healthy plugin could be pulled from rotation. The skip is now report-only: still logged (rate-limited) and still left as PluginBusyError state error info with last-update stamped, but in health it is counted as a busy skip (busy_skip_count / last_busy_skip, via the new PluginHealthTracker.record_busy_skip) and never touches the failure streak, last_error or the breaker. Only real hangs -- display() or update() running past the executor timeout -- still count toward the breaker. Busy-skip persistence shares the slow-call throttle (first one saved at once so the web process sees it, then at most once a minute per plugin). Tests prove repeated busy skips never open the breaker and repeated real hangs still do, with busy skips interleaved. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
+9
-5
@@ -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
|
||||
|
||||
@@ -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.
|
||||
"""
|
||||
|
||||
|
||||
|
||||
@@ -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]]:
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user