perf(plugins): run scheduled updates off the render thread (#407)

* perf(plugins): run scheduled updates off the render thread

plugin update() executed inline in the render loop — execute_update's
internal thread.join(timeout=30) blocked it, so one slow plugin HTTP
fetch froze scrolling for the whole fetch (up to 30s; DNS-retry storms
made this a regular occurrence on flaky networks).

Scheduling stays on the render thread and keeps every existing gate
(enabled, circuit breaker, can_execute, interval); due updates are now
enqueued to a single background worker (serialized — same one-at-a-time
execution as before, no thundering herd). RUNNING is set at enqueue so
can_execute blocks re-entry alongside the pending-set dedup.

Per-plugin locks make the old implicit update/display no-overlap
guarantee explicit: the worker holds the plugin's lock through its
update; the display side try-locks and, when the plugin is mid-update,
holds the last frame for that iteration — reported as success so a
mid-update skip never advances the rotation. Unlike before, the
guarantee now also holds across the post-timeout window (previously the
lingering update thread overlapped display()). Deadlock-free by
construction: the worker takes one lock; display never blocks.

Timeout semantics unchanged (lingering daemon thread documented).
Kill switch: plugin_system.synchronous_updates: true restores the
inline path.

8 new concurrency tests (non-blocking scheduler, overlap assertion
under a hammering display loop, lock release on failure/timeout paths,
dedup, kill switch); 4-min devpi soak clean (updates completing,
rotation advancing, no stuck RUNNING states).

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01FqzC1nzTWL4kaqgMaQZFam

* fix(plugins): fix skipped-frame health/force_change tracking, config validation, and lock-lifetime gaps in async updates

Addresses PR #407 review findings:
- display_controller: only clear force_change / record health success
  when display() actually ran this frame, not when the frame was skipped
  because the plugin's lock was busy (a skip must preserve a pending
  mode-switch force_clear).
- plugin_manager: replace bool() coercion of synchronous_updates with
  explicit isinstance validation of plugin_system/synchronous_updates,
  failing safe to synchronous mode (with a logged reason) on malformed
  config instead of silently defaulting to async.
- plugin_manager: _update_worker_loop now acquires the plugin lock before
  looking up its instance and re-checks under the lock, so an unloaded
  plugin's lifecycle state is never resurrected to ENABLED.
- plugin_manager + display_controller: move lock ownership (and, for
  updates, RUNNING/pending lifecycle bookkeeping) into the actual update()/
  display() call itself rather than the timeout-wrapped caller, so the
  lock stays held for the real operation's duration even after
  PluginExecutor's own join(timeout) elapses and a lingering daemon thread
  keeps running in the background.
- DisplayController.cleanup() now stops the update worker before tearing
  down display/cache resources; stop_update_worker() logs when the join
  times out instead of failing silently.
- test_async_plugin_updates: rewrite test_unloaded_while_queued_is_harmless
  to exercise the public unload_plugin() lifecycle (via a deterministic
  blocker) instead of deleting pm.plugins directly, and add a regression
  test proving the plugin lock stays held through PluginExecutor's own
  timeout while the real update() call is still running.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01KEZK1P1Q1fu5pcuVrkrCFZ

---------

Co-authored-by: Chuck <chuck@example.com>
Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
Chuck
2026-07-14 08:08:51 -04:00
committed by GitHub
co-authored by Claude Sonnet 5 Chuck
parent 9837315308
commit 0aca40cf3a
4 changed files with 637 additions and 68 deletions
+228 -41
View File
@@ -8,12 +8,13 @@ API Version: 1.0.0
"""
import json
import queue
import sys
import time
import threading
import types
from pathlib import Path
from typing import Dict, List, Optional, Any
from typing import Dict, List, Optional, Any, Tuple
import logging
from src.exceptions import PluginError, ConfigError
from src.logging_config import get_logger
@@ -97,6 +98,50 @@ class PluginManager:
# Health tracking (optional, set by display_controller if available)
self.health_tracker = None
self.resource_monitor = None
# --- Asynchronous plugin updates -------------------------------
# update() used to run inline in the render loop (execute_update's
# internal thread.join(timeout=30) blocked it), so one slow plugin
# HTTP fetch froze scrolling for the whole fetch. Scheduling still
# happens on the render thread (run_scheduled_updates), but
# execution moves to this single background worker. Per-plugin
# locks keep a plugin's update() and display() mutually exclusive —
# today's implicit guarantee, now explicit (and, unlike today,
# also held across the 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()
self._pending_updates: set = set()
self._pending_lock = threading.Lock()
self._plugin_locks: Dict[str, threading.Lock] = {}
self._plugin_locks_guard = threading.Lock()
self._update_worker: Optional[threading.Thread] = None
self._synchronous_updates = False
if self.config_manager is not None:
try:
cfg = self.config_manager.get_config() or {}
except (OSError, ValueError) as exc:
self.logger.warning(
"Could not load config to check plugin_system.synchronous_updates "
"(%s: %s); defaulting to synchronous updates", type(exc).__name__, exc)
self._synchronous_updates = True
else:
plugin_system_cfg = cfg.get('plugin_system', {})
if not isinstance(plugin_system_cfg, dict):
self.logger.warning(
"config plugin_system must be a mapping, got %s; "
"defaulting to synchronous updates",
type(plugin_system_cfg).__name__)
self._synchronous_updates = True
else:
sync_value = plugin_system_cfg.get('synchronous_updates', False)
if not isinstance(sync_value, bool):
self.logger.warning(
"config plugin_system.synchronous_updates must be a boolean, "
"got %r; defaulting to synchronous updates", sync_value)
self._synchronous_updates = True
else:
self._synchronous_updates = sync_value
# Ensure plugins directory exists with proper permissions
try:
@@ -744,47 +789,189 @@ class PluginManager:
last_update = self.plugin_last_update.get(plugin_id, 0.0)
if last_update == 0.0 or (current_time - last_update) >= interval:
# Update state to RUNNING
self.state_manager.set_state(plugin_id, PluginState.RUNNING)
try:
# Use PluginExecutor for safe execution
success = False
if self.resource_monitor:
# If resource monitor exists, wrap the call
def monitored_update():
self.resource_monitor.monitor_call(plugin_id, plugin_instance.update)
# SimpleNamespace stores `update` as an *instance*
# attribute, so attribute lookup returns the plain
# function object as-is. A dynamically-built class
# (`type(..., {'update': monitored_update})`) instead
# stores it as a *class* attribute, which the
# descriptor protocol turns into a bound method on
# access -- silently prepending the instance as an
# implicit first argument to a function that takes
# none, raising "monitored_update() takes 0
# positional arguments but 1 was given" on every call.
success = self.plugin_executor.execute_update(
types.SimpleNamespace(update=monitored_update),
plugin_id
)
else:
success = self.plugin_executor.execute_update(plugin_instance, plugin_id)
if success:
with self._plugin_last_update_lock:
self.plugin_last_update[plugin_id] = current_time
self.state_manager.record_update(plugin_id)
# Update state back to ENABLED
self.state_manager.set_state(plugin_id, PluginState.ENABLED)
# Record success
if self.health_tracker:
self.health_tracker.record_success(plugin_id)
else:
self._record_update_failure(plugin_id)
except Exception as exc: # pylint: disable=broad-except
self.logger.exception("Error updating plugin %s: %s", plugin_id, exc)
if self._synchronous_updates:
# Kill-switch path: the original inline execution
# (blocks the caller until update() completes/times out)
self.state_manager.set_state(plugin_id, PluginState.RUNNING)
self._execute_update_now(plugin_id, plugin_instance, current_time)
else:
self._enqueue_update(plugin_id, current_time)
def get_plugin_lock(self, plugin_id: str) -> threading.Lock:
"""Per-plugin lock keeping update() and display() 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.
"""
with self._plugin_locks_guard:
lock = self._plugin_locks.get(plugin_id)
if lock is None:
lock = threading.Lock()
self._plugin_locks[plugin_id] = lock
return lock
def _enqueue_update(self, plugin_id: str, scheduled_time: float) -> None:
"""Queue a due update for the background worker (dedup while pending)."""
with self._pending_lock:
if plugin_id in self._pending_updates:
return
self._pending_updates.add(plugin_id)
# RUNNING is set at enqueue time so can_execute() blocks re-entry and
# the web UI shows the truthful state while the item waits its turn.
self.state_manager.set_state(plugin_id, PluginState.RUNNING)
self._ensure_update_worker()
self._update_queue.put((plugin_id, scheduled_time))
def _ensure_update_worker(self) -> None:
if self._update_worker is not None and self._update_worker.is_alive():
return
self._update_worker = threading.Thread(
target=self._update_worker_loop, name='plugin-update-worker',
daemon=True)
self._update_worker.start()
def _update_worker_loop(self) -> None:
"""Single worker: dispatches queued updates off the render thread
(matching the old inline behavior — no thundering herd of
concurrent fetches).
The plugin's lock is acquired here, before its instance is looked
up, and the instance is re-fetched under the lock — a concurrent
unload_plugin() can't leave this loop about to run update() on an
instance that's already been torn down. The lock — and RUNNING/
pending lifecycle state — is released by the update itself once the
real update() call genuinely finishes (see _execute_update_now),
which can be after this dispatch returns if PluginExecutor's own
timeout elapses first.
"""
while True:
item = self._update_queue.get()
if item is None: # shutdown sentinel
return
plugin_id, scheduled_time = item
lock = self.get_plugin_lock(plugin_id)
lock.acquire()
plugin_instance = self.plugins.get(plugin_id)
if plugin_instance is None: # unloaded while queued; its
# lifecycle state was already cleared by unload_plugin —
# leave it alone rather than resurrecting it to ENABLED
lock.release()
with self._pending_lock:
self._pending_updates.discard(plugin_id)
continue
try:
self._execute_update_now(plugin_id, plugin_instance,
scheduled_time, lock=lock)
except Exception: # pylint: disable=broad-except
# _execute_update_now guarantees the lock/pending bookkeeping
# is released via its own _finish() before returning or
# raising; this is a last-resort log only.
self.logger.exception("update worker: unexpected error for %s",
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():
self._update_queue.put(None)
self._update_worker.join(timeout=timeout)
if self._update_worker.is_alive():
self.logger.warning(
"Update worker did not stop within %.1fs; it is a daemon "
"thread and will be abandoned on shutdown", timeout)
def _execute_update_now(self, plugin_id: str, plugin_instance: Any,
scheduled_time: float,
lock: Optional[threading.Lock] = None) -> None:
"""Execute a plugin's update() via PluginExecutor, then bookkeep.
Caller is responsible for having set RUNNING state.
On the synchronous path (``lock=None``) this is the original,
unchanged inline behavior. On the async worker path, PluginExecutor's
internal thread.join(timeout) blocks only the calling thread -- on
timeout the lingering daemon update-thread keeps running the real
plugin.update() call unkillable in the background. So that the
plugin's lock (and its RUNNING/pending lifecycle state) stays held
for that real duration rather than just this bounded wait, ownership
of both is carried by the wrapped update callable itself, released
from whichever thread actually finishes it -- see _finish() below.
"""
finish_guard = threading.Lock()
finished = {'done': False}
def _finish(success: bool, exc: Optional[Exception] = None) -> None:
with finish_guard:
if finished['done']:
return
finished['done'] = True
try:
if success:
with self._plugin_last_update_lock:
self.plugin_last_update[plugin_id] = scheduled_time
self.state_manager.record_update(plugin_id)
self.state_manager.set_state(plugin_id, PluginState.ENABLED)
if self.health_tracker:
self.health_tracker.record_success(plugin_id)
else:
self._record_update_failure(plugin_id, exc=exc)
finally:
if lock is not None:
lock.release()
with self._pending_lock:
self._pending_updates.discard(plugin_id)
if lock is None:
# Synchronous / no-lock path: unchanged behavior.
try:
if self.resource_monitor:
def monitored_update():
self.resource_monitor.monitor_call(plugin_id, plugin_instance.update)
# SimpleNamespace stores `update` as an *instance*
# attribute, so attribute lookup returns the plain
# function object as-is. A dynamically-built class
# (`type(..., {'update': monitored_update})`) instead
# stores it as a *class* attribute, which the
# descriptor protocol turns into a bound method on
# access -- silently prepending the instance as an
# implicit first argument to a function that takes
# none, raising "monitored_update() takes 0
# positional arguments but 1 was given" on every call.
success = self.plugin_executor.execute_update(
types.SimpleNamespace(update=monitored_update),
plugin_id
)
else:
success = self.plugin_executor.execute_update(plugin_instance, plugin_id)
_finish(success)
except Exception as exc: # pylint: disable=broad-except
self.logger.exception("Error updating plugin %s: %s", plugin_id, exc)
_finish(False, exc=exc)
return
# Async worker path: the real update() call -- through the resource
# monitor, if configured -- owns finishing the lock/lifecycle
# bookkeeping, from whichever thread actually runs it to completion.
def _target_update() -> None:
try:
if self.resource_monitor:
self.resource_monitor.monitor_call(plugin_id, plugin_instance.update)
else:
plugin_instance.update()
except Exception as exc:
_finish(False, exc=exc)
raise
else:
_finish(True)
try:
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)
def run_scheduled_updates_with_changes(self, current_time: Optional[float] = None) -> List[str]:
"""