mirror of
https://github.com/ChuckBuilds/LEDMatrix.git
synced 2026-10-05 14:55:08 +00:00
feat(plugins): request_on_demand() / end_on_demand() -- plugins ask for the screen in-process (#768)
* feat(plugins): request_on_demand() / end_on_demand() -- plugins ask for the screen in-process Four plugins (birdnet-go, mqtt-notifications, on-air, pomodoro-timer) take the screen by writing the display_on_demand_request mailbox, which the display reads once a second while the control socket is up and which stage 5 removes. This is the in-process way in that stage needed. - BasePlugin.request_on_demand(mode=None, duration=None, pinned=False) and end_on_demand(), safe from any thread, go through PluginManager to DisplayController.submit_plugin_on_demand, which only queues (at most 32) and wakes the render thread through ControlServer.wake(). The render thread applies them in _drain_control_commands, after socket commands, through _handle_on_demand_request, so they land within a frame; without a socket, on the next pending-changes pass. - A plugin's stop ends only its own session; a mailbox stop still ends any. - Both answer the request id, or None with no display in the process (web interface, check_plugin.py), a full queue, or a mock manager -- a plugin's cue to write the mailbox, which the display still reads. - docs/PLUGIN_API_REFERENCE.md documents the hasattr pattern for plugins that must keep working on older cores; IPC_CONTROL_SOCKET.md and the CHANGELOG are updated. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(display): wire the on-demand handler only on a manager that has it Tests and the golden traces stand in simpler plugin managers. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
+104
-14
@@ -394,6 +394,12 @@ class DisplayController:
|
||||
plugin_time = time.time()
|
||||
self.plugin_manager = None
|
||||
self._plugin_runtime_publisher = None
|
||||
# On-demand requests plugins make in this process (BasePlugin.
|
||||
# request_on_demand / end_on_demand), from any thread; the render
|
||||
# thread drains them with the socket's commands. Created before the
|
||||
# plugins load, because a plugin may ask from its first thread.
|
||||
self._plugin_on_demand: deque = deque()
|
||||
self._plugin_on_demand_lock = threading.Lock()
|
||||
self.plugin_modes = {} # mode -> plugin_instance mapping for plugin-first dispatch
|
||||
self.mode_to_plugin_id: Dict[str, str] = {}
|
||||
self.plugin_display_modes: Dict[str, List[str]] = {}
|
||||
@@ -508,6 +514,12 @@ class DisplayController:
|
||||
cache_manager=self.cache_manager,
|
||||
font_manager=self.font_manager
|
||||
)
|
||||
# BasePlugin.request_on_demand() / end_on_demand() land here.
|
||||
# Before any plugin loads: a plugin may ask from its first thread.
|
||||
# getattr: tests and golden traces stand in simpler managers.
|
||||
set_handler = getattr(self.plugin_manager, 'set_on_demand_handler', None)
|
||||
if callable(set_handler):
|
||||
set_handler(self.submit_plugin_on_demand)
|
||||
|
||||
# The web UI's loaded / state / error_info for each plugin read
|
||||
# what this publishes. Started before loading, so the loads that
|
||||
@@ -1882,6 +1894,10 @@ class DisplayController:
|
||||
_on_demand_mailbox: Optional[MailboxWatch] = None
|
||||
#: Writers whose mailbox requests have been logged (_note_mailbox_request).
|
||||
_mailbox_writers_logged: FrozenSet[str] = frozenset()
|
||||
#: Most plugin on-demand requests waiting for the render thread at once.
|
||||
#: A plugin that asks faster than the display drains (four times a
|
||||
#: second at worst) is refused, not queued without end.
|
||||
PLUGIN_ON_DEMAND_QUEUE_SIZE = 32
|
||||
|
||||
def _service_pending_changes(self) -> None:
|
||||
"""Apply changes made elsewhere while the display thread is busy.
|
||||
@@ -1906,7 +1922,7 @@ class DisplayController:
|
||||
# A command queued on the control socket skips the floor: it is in
|
||||
# memory, so applying it now costs no disk read.
|
||||
if (last is not None and now - last < self.PENDING_CHANGES_INTERVAL
|
||||
and not (self._control_server and self._control_server.has_pending)):
|
||||
and not self._control_command_pending()):
|
||||
return
|
||||
self._last_pending_service = now
|
||||
|
||||
@@ -2105,11 +2121,15 @@ class DisplayController:
|
||||
here. A plugin reload waits for the top of the next loop pass, where
|
||||
no plugin is on the stack (_apply_pending_plugin_reloads); until
|
||||
then the current screen ends early (_plugin_reload_pending).
|
||||
|
||||
Plugins' own on-demand requests (submit_plugin_on_demand) are
|
||||
applied here too, after the socket's, with or without a socket.
|
||||
"""
|
||||
server = self._control_server
|
||||
if server is None or not server.has_pending:
|
||||
return
|
||||
for command in server.drain():
|
||||
# drain() clears the wake flag before the plugin queue is read below,
|
||||
# so a plugin request queued from here on wakes the next wait.
|
||||
commands = server.drain() if server is not None and server.has_pending else []
|
||||
for command in commands:
|
||||
try:
|
||||
if command.cmd == ControlCommand.BRIGHTNESS_SET:
|
||||
self._apply_control_brightness(command)
|
||||
@@ -2121,24 +2141,78 @@ class DisplayController:
|
||||
logger.exception("Failed to apply control socket command %s",
|
||||
command.request_id)
|
||||
command.fail(ControlErrorCode.INTERNAL, 'the display failed to apply it')
|
||||
self._drain_plugin_on_demand()
|
||||
|
||||
# -- plugins' in-process on-demand requests ---------------------------------
|
||||
|
||||
def submit_plugin_on_demand(self, request: Dict[str, Any]) -> bool:
|
||||
"""Queue a plugin's on-demand request for the render thread. Any thread.
|
||||
|
||||
``PluginManager.request_on_demand`` / ``end_on_demand`` (which
|
||||
BasePlugin's methods of the same names call) build ``request``: the
|
||||
mailbox's shape, with ``source: 'plugin'`` and the asking plugin's
|
||||
id. Nothing here touches the panel or the on-demand state; the render
|
||||
thread applies the request where it applies a socket command
|
||||
(_drain_control_commands), through _handle_on_demand_request, and is
|
||||
woken for it when the control socket is up. True when it was queued;
|
||||
False (logged) when the queue is full.
|
||||
"""
|
||||
pending = self.__dict__.get('_plugin_on_demand')
|
||||
lock = self.__dict__.get('_plugin_on_demand_lock')
|
||||
if pending is None or lock is None:
|
||||
return False # a controller built without __init__ (tests)
|
||||
with lock:
|
||||
if len(pending) >= self.PLUGIN_ON_DEMAND_QUEUE_SIZE:
|
||||
logger.warning("Plugin on-demand queue full; refusing %s %s from %s",
|
||||
request.get('action'), request.get('request_id'),
|
||||
request.get('plugin_id'))
|
||||
return False
|
||||
pending.append(dict(request))
|
||||
server = self._control_server
|
||||
if server is not None:
|
||||
server.wake()
|
||||
return True
|
||||
|
||||
def _plugin_on_demand_pending(self) -> bool:
|
||||
"""A plugin's on-demand request is waiting. A length check: no lock."""
|
||||
return bool(self.__dict__.get('_plugin_on_demand'))
|
||||
|
||||
def _drain_plugin_on_demand(self) -> None:
|
||||
"""Apply the plugins' queued on-demand requests, oldest first. Render thread."""
|
||||
pending = self.__dict__.get('_plugin_on_demand')
|
||||
while pending:
|
||||
try:
|
||||
request = pending.popleft()
|
||||
except IndexError:
|
||||
break
|
||||
try:
|
||||
self._handle_on_demand_request(request)
|
||||
except Exception: # pylint: disable=broad-except
|
||||
logger.exception("Failed to apply on-demand request %s from plugin %s",
|
||||
request.get('request_id'), request.get('plugin_id'))
|
||||
|
||||
def _wait_for_control(self, timeout: float) -> bool:
|
||||
"""Sleep up to ``timeout``, waking early for a control socket command.
|
||||
|
||||
True when a command is waiting. Without a socket (Windows, switched
|
||||
off, tests) this is the plain sleep it replaces.
|
||||
off, tests) this is the plain sleep it replaces, unless a plugin's
|
||||
on-demand request is already waiting.
|
||||
"""
|
||||
server = self._control_server
|
||||
wait = getattr(server, 'wait_for_command', None) if server is not None else None
|
||||
if wait is None:
|
||||
if self._plugin_on_demand_pending():
|
||||
return True
|
||||
time.sleep(timeout)
|
||||
return False
|
||||
return bool(wait(timeout))
|
||||
|
||||
def _control_command_pending(self) -> bool:
|
||||
"""A socket command is queued: Vegas checks this every frame."""
|
||||
"""A socket command or a plugin's on-demand request is queued: Vegas
|
||||
checks this every frame."""
|
||||
server = self._control_server
|
||||
return bool(server is not None and server.has_pending)
|
||||
return bool((server is not None and server.has_pending)
|
||||
or self._plugin_on_demand_pending())
|
||||
|
||||
def _wait_frame_interval(self, interval: float, screen: Screen) -> Optional[ScreenPlan]:
|
||||
"""The static screen's sleep between frames, woken by socket commands.
|
||||
@@ -2450,21 +2524,36 @@ class DisplayController:
|
||||
def _handle_on_demand_request(self, request: Dict[str, Any]) -> None:
|
||||
"""Process one on-demand request, from the mailbox or the control socket.
|
||||
|
||||
A socket command carries ``source: 'socket'``. Only a mailbox request
|
||||
is removed from the mailbox afterwards: a socket command never put
|
||||
anything there, so that would be a disk read and maybe a delete for
|
||||
nothing.
|
||||
A socket command carries ``source: 'socket'``, and a plugin's own
|
||||
request (submit_plugin_on_demand) ``source: 'plugin'``. Only a
|
||||
mailbox request is removed from the mailbox afterwards: the others
|
||||
never put anything there, so that would be a disk read and maybe a
|
||||
delete for nothing.
|
||||
|
||||
A plugin's stop ends only that plugin's own session: a plugin
|
||||
releasing the screen must not end one the user started for
|
||||
another plugin. (A stop through the mailbox ends any session, as it
|
||||
always has.)
|
||||
"""
|
||||
request_id = request.get('request_id')
|
||||
if not request_id:
|
||||
return
|
||||
from_mailbox = request.get('source') != 'socket'
|
||||
source = request.get('source')
|
||||
from_mailbox = source not in ('socket', 'plugin')
|
||||
|
||||
action = request.get('action')
|
||||
|
||||
# For stop requests, always process them (don't check processed_id)
|
||||
# This allows stopping even if the same stop request was sent before
|
||||
if action == 'stop':
|
||||
if source == 'plugin' and not (
|
||||
self.on_demand_active
|
||||
and self.on_demand_plugin_id == request.get('plugin_id')):
|
||||
logger.debug("On-demand stop %s from plugin %s ignored: it does not own "
|
||||
"the screen (on-demand %s, plugin %s)", request_id,
|
||||
request.get('plugin_id'), self.on_demand_status,
|
||||
self.on_demand_plugin_id)
|
||||
return
|
||||
logger.info("Received on-demand stop request %s", request_id)
|
||||
# Always process stop requests, even if same request_id (user might click multiple times)
|
||||
if self.on_demand_active:
|
||||
@@ -2507,8 +2596,9 @@ class DisplayController:
|
||||
self._consume_on_demand_request(request_id)
|
||||
return
|
||||
|
||||
logger.info("Received on-demand request %s: %s (plugin_id=%s, mode=%s)",
|
||||
request_id, action, request.get('plugin_id'), request.get('mode'))
|
||||
logger.info("Received on-demand request %s: %s (plugin_id=%s, mode=%s, via %s)",
|
||||
request_id, action, request.get('plugin_id'), request.get('mode'),
|
||||
'mailbox' if from_mailbox else source)
|
||||
|
||||
# Mark as processed BEFORE processing (to prevent duplicate processing)
|
||||
self.cache_manager.set('display_on_demand_processed_id', request_id, ttl=3600)
|
||||
|
||||
+15
-1
@@ -772,8 +772,22 @@ class ControlServer:
|
||||
"""
|
||||
return self._pending.wait(timeout)
|
||||
|
||||
def wake(self) -> None:
|
||||
"""Wake the render thread as a queued command would, with nothing queued.
|
||||
|
||||
For work that reaches the display another way in the same process (a
|
||||
plugin's on-demand request, ``DisplayController.submit_plugin_on_demand``):
|
||||
the render thread returns from :meth:`wait_for_command` and drains,
|
||||
and reads the caller's own queue there. Safe from any thread.
|
||||
"""
|
||||
self._pending.set()
|
||||
|
||||
def drain(self) -> List[QueuedCommand]:
|
||||
"""Every queued command, oldest first. Called from the render thread."""
|
||||
"""Every queued command, oldest first. Called from the render thread.
|
||||
|
||||
Clears the wake flag first, so anything queued (or woken for) while
|
||||
this runs wakes the next wait again.
|
||||
"""
|
||||
commands: List[QueuedCommand] = []
|
||||
self._pending.clear()
|
||||
while True:
|
||||
|
||||
@@ -1070,6 +1070,65 @@ class BasePlugin(ABC):
|
||||
if callable(notify):
|
||||
notify(self.plugin_id)
|
||||
|
||||
def request_on_demand(self, mode: Optional[str] = None,
|
||||
duration: Optional[float] = None,
|
||||
pinned: bool = False) -> Optional[str]:
|
||||
"""
|
||||
Take the screen now: show this plugin on demand. Safe from any thread.
|
||||
|
||||
For a plugin that reacts to something outside the rotation -- an MQTT
|
||||
message, a timer, a detection -- and wants the panel for it. The
|
||||
request goes straight to the display in this process and is applied
|
||||
on its render thread within a frame or so, exactly like an on-demand
|
||||
start from the web interface.
|
||||
|
||||
Args:
|
||||
mode: One of this plugin's display modes; None for its first.
|
||||
duration: Seconds to show it before the rotation resumes; None
|
||||
(or zero) for no limit, until end_on_demand() or the user
|
||||
stops it.
|
||||
pinned: Stay on ``mode`` instead of cycling through the
|
||||
plugin's other modes.
|
||||
|
||||
Returns:
|
||||
The request id once the display has queued it, or None when
|
||||
there is no display in this process to ask (the web interface,
|
||||
scripts/check_plugin.py) or its queue is full. A plugin that
|
||||
also runs on cores without this method writes the
|
||||
``display_on_demand_request`` mailbox on None, as before; see
|
||||
"On-demand display" in docs/PLUGIN_API_REFERENCE.md.
|
||||
|
||||
Example::
|
||||
|
||||
if not (hasattr(self, 'request_on_demand')
|
||||
and self.request_on_demand(mode='my_alert', duration=15)):
|
||||
self._write_on_demand_mailbox(...) # older cores
|
||||
"""
|
||||
request = getattr(getattr(self, 'plugin_manager', None), 'request_on_demand', None)
|
||||
if not callable(request):
|
||||
return None
|
||||
request_id = request(self.plugin_id, mode=mode, duration=duration, pinned=pinned)
|
||||
# Only a real id counts: a test's MagicMock manager answers a mock,
|
||||
# which must read as "not taken" so the plugin's fallback runs.
|
||||
return request_id if isinstance(request_id, str) else None
|
||||
|
||||
def end_on_demand(self) -> Optional[str]:
|
||||
"""
|
||||
Give the screen back: end this plugin's on-demand session. Any thread.
|
||||
|
||||
Ends only a session this plugin owns. One the user started for
|
||||
another plugin, or a session that already ended, is left alone. The
|
||||
rotation resumes where it left off.
|
||||
|
||||
Returns:
|
||||
The request id once queued, or None as request_on_demand() does.
|
||||
"""
|
||||
end = getattr(getattr(self, 'plugin_manager', None), 'end_on_demand', None)
|
||||
if not callable(end):
|
||||
return None
|
||||
request_id = end(self.plugin_id)
|
||||
return request_id if isinstance(request_id, str) else None
|
||||
|
||||
def get_vegas_participation(self) -> str:
|
||||
"""
|
||||
How this plugin takes part in Vegas mode: ``'scroll'``, ``'pause'`` or
|
||||
|
||||
@@ -15,6 +15,7 @@ import sys
|
||||
import time
|
||||
import threading
|
||||
import types
|
||||
import uuid
|
||||
from pathlib import Path
|
||||
from typing import Callable, Dict, List, NamedTuple, Optional, Any, Tuple, Union
|
||||
import logging
|
||||
@@ -214,6 +215,9 @@ class PluginManager:
|
||||
# add_update_listener(). A tuple, replaced rather than mutated, so the
|
||||
# worker can iterate it without a lock.
|
||||
self._update_listeners: Tuple[Callable[[str], None], ...] = ()
|
||||
# Where plugins' on-demand requests go: the display controller's
|
||||
# submit_plugin_on_demand. See set_on_demand_handler().
|
||||
self._on_demand_handler: Optional[Callable[[Dict[str, Any]], bool]] = None
|
||||
# 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]]] = {}
|
||||
@@ -1844,3 +1848,73 @@ class PluginManager:
|
||||
done = sorted(self._completed_updates)
|
||||
self._completed_updates.clear()
|
||||
return done
|
||||
|
||||
# -- on-demand requests from plugins -------------------------------------
|
||||
|
||||
def set_on_demand_handler(
|
||||
self, handler: Optional[Callable[[Dict[str, Any]], bool]]) -> None:
|
||||
"""Route plugins' on-demand requests to ``handler`` (None: nowhere).
|
||||
|
||||
The display controller sets its ``submit_plugin_on_demand`` here
|
||||
before any plugin loads. The handler takes a mailbox-shaped request
|
||||
from any thread, queues it for the render thread and returns True,
|
||||
or False when it could not. A plugin manager with no handler (the
|
||||
web interface's, a test's, scripts/check_plugin.py's) has no screen
|
||||
to give, so request_on_demand() there answers None.
|
||||
"""
|
||||
self._on_demand_handler = handler
|
||||
|
||||
def request_on_demand(self, plugin_id: str, mode: Optional[str] = None,
|
||||
duration: Optional[float] = None,
|
||||
pinned: bool = False) -> Optional[str]:
|
||||
"""Ask the display to show ``plugin_id`` now. Safe from any thread.
|
||||
|
||||
BasePlugin.request_on_demand() lands here; see it for the arguments.
|
||||
Returns the request id once the display has queued the request (it
|
||||
is applied on the render thread within a frame or so), or None when
|
||||
this process has no display to ask or its queue is full.
|
||||
"""
|
||||
if not isinstance(plugin_id, str) or not plugin_id:
|
||||
raise ValueError('plugin_id is required')
|
||||
if mode is not None and (not isinstance(mode, str) or not mode):
|
||||
raise ValueError('mode must be a non-empty string or None')
|
||||
if duration is not None:
|
||||
if isinstance(duration, bool) or not isinstance(duration, (int, float)):
|
||||
raise ValueError('duration must be a number of seconds or None')
|
||||
if not math.isfinite(duration) or duration <= 0:
|
||||
duration = None # the display reads these as "no limit" too
|
||||
else:
|
||||
duration = float(duration)
|
||||
return self._submit_on_demand({
|
||||
'action': 'start', 'plugin_id': plugin_id, 'mode': mode,
|
||||
'duration': duration, 'pinned': bool(pinned)})
|
||||
|
||||
def end_on_demand(self, plugin_id: str) -> Optional[str]:
|
||||
"""Give the screen back, if ``plugin_id``'s on-demand session has it.
|
||||
|
||||
BasePlugin.end_on_demand() lands here. A session the plugin does not
|
||||
own (the user started another plugin from the web interface, say) is
|
||||
left alone. Returns the request id once queued, or None as
|
||||
request_on_demand() does.
|
||||
"""
|
||||
if not isinstance(plugin_id, str) or not plugin_id:
|
||||
raise ValueError('plugin_id is required')
|
||||
return self._submit_on_demand({'action': 'stop', 'plugin_id': plugin_id})
|
||||
|
||||
def _submit_on_demand(self, request: Dict[str, Any]) -> Optional[str]:
|
||||
# __dict__.get: tests build bare managers with PluginManager.__new__.
|
||||
handler = self.__dict__.get('_on_demand_handler')
|
||||
if handler is None:
|
||||
return None
|
||||
request_id = str(uuid.uuid4())
|
||||
request.update({'request_id': request_id, 'timestamp': time.time(),
|
||||
'source': 'plugin'})
|
||||
try:
|
||||
accepted = handler(request)
|
||||
except Exception as exc: # pylint: disable=broad-except
|
||||
self._warn_rate_limited(
|
||||
"on-demand-handler",
|
||||
"The on-demand request from plugin %s failed: %r",
|
||||
request.get('plugin_id'), exc)
|
||||
return None
|
||||
return request_id if accepted else None
|
||||
|
||||
Reference in New Issue
Block a user