From 0577c807ebd9840e994e52223fdd2f256c44b5aa Mon Sep 17 00:00:00 2001 From: Chuck <33324927+ChuckBuilds@users.noreply.github.com> Date: Mon, 5 Oct 2026 01:39:02 -0400 Subject: [PATCH] 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 * 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 --------- Co-authored-by: Claude Opus 5.5 --- CHANGELOG.md | 24 ++ docs/IPC_CONTROL_SOCKET.md | 31 ++- docs/PLUGIN_API_REFERENCE.md | 80 ++++++ src/display_controller.py | 118 +++++++- src/ipc/server.py | 16 +- src/plugin_system/base_plugin.py | 59 ++++ src/plugin_system/plugin_manager.py | 74 +++++ test/test_plugin_on_demand_api.py | 418 ++++++++++++++++++++++++++++ 8 files changed, 799 insertions(+), 21 deletions(-) create mode 100644 test/test_plugin_on_demand_api.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 4c94c537..4607b5e0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -19,6 +19,30 @@ accepts both, but the store flags the old spelling as deprecated ## Unreleased +### Plugins ask for the screen in-process: `request_on_demand()` / `end_on_demand()` + +The in-process way in that stage 5 of the control socket needed +(`docs/IPC_CONTROL_SOCKET.md`, "Plugins in the display process"). + +- **`BasePlugin.request_on_demand(mode=None, duration=None, pinned=False)`** + shows the plugin now, and **`BasePlugin.end_on_demand()`** gives the + screen back. Both are safe from any thread (an MQTT callback, a timer + thread): `PluginManager.request_on_demand()` / `end_on_demand()` hand the + request to `DisplayController.submit_plugin_on_demand()`, which only + queues it (at most 32) and wakes the render thread through the control + socket's flag (`ControlServer.wake()`). The render thread applies it with + the socket's commands, through the same handler as a web on-demand + request, so it lands within a frame rather than on the mailbox's + once-a-second look. Both return the request id, or `None` when no display + runs in the process (the web interface, `scripts/check_plugin.py`) or the + queue is full. +- **A plugin's stop ends only its own session.** A mailbox stop still ends + any session, whoever started it. +- **Older cores.** Plugins detect the methods with `hasattr` and write the + `display_on_demand_request` mailbox when they are missing or answer + `None`; the pattern is in `docs/PLUGIN_API_REFERENCE.md` ("On-demand + display"). The display still reads the mailbox for plugins that write it. + ### Web UI: Schedule and General are ES-module pages (stage 3) - The Schedule and General tabs follow stage 2 (#727): their inline diff --git a/docs/IPC_CONTROL_SOCKET.md b/docs/IPC_CONTROL_SOCKET.md index 38505772..b377a374 100644 --- a/docs/IPC_CONTROL_SOCKET.md +++ b/docs/IPC_CONTROL_SOCKET.md @@ -489,7 +489,7 @@ restart banner, as before. | Mailbox | Written by | Read by the display | While the socket is up | |---|---|---|---| -| `display_on_demand_request` | the web interface, only on fallback; four plugins directly (birdnet-go, mqtt-notifications, on-air, pomodoro-timer) | the render thread, `_poll_on_demand_requests()` | looked at every 1 s (`MAILBOX_POLL_INTERVAL_WITH_SOCKET`), 0.25 s without a socket | +| `display_on_demand_request` | the web interface, only on fallback; plugins that predate `BasePlugin.request_on_demand()`, or run on a core without it | the render thread, `_poll_on_demand_requests()` | looked at every 1 s (`MAILBOX_POLL_INTERVAL_WITH_SOCKET`), 0.25 s without a socket | | `plugin_error_clear_request` | the web interface, only on fallback | the error publisher's thread, every 5 s tick | unchanged rate | A look is one `stat()` of the mailbox file (`CacheManager.file_signature`): @@ -502,8 +502,25 @@ the mailbox instead of being re-read until it expires. A request that comes through the on-demand mailbox while the socket is up is logged once per writer (`came through the file mailbox although the -control socket is up`), which names the plugins that still need an -in-process way in before the mailbox is removed. +control socket is up`), which names the plugins that still write it. + +### Plugins in the display process + +A plugin asks for the screen with `BasePlugin.request_on_demand()` and gives +it back with `end_on_demand()` (see "On-demand display" in +[PLUGIN_API_REFERENCE.md](PLUGIN_API_REFERENCE.md)). Neither goes through +the socket or a file: `PluginManager` hands the mailbox-shaped request, +marked `source: 'plugin'`, to `DisplayController.submit_plugin_on_demand`, +which queues it in memory (at most `PLUGIN_ON_DEMAND_QUEUE_SIZE`, 32) from +whatever thread the plugin called on, and wakes the render thread through +the socket's queue flag (`ControlServer.wake()`). The render thread applies +it in `_drain_control_commands`, after the socket's commands, through the +same `_handle_on_demand_request`, so it lands within a frame like a socket +command. Without a socket it lands on the next pending-changes pass (typically +within 0.25 s). A plugin's stop ends only a session that plugin owns. The four +plugins that wrote the mailbox (birdnet-go, mqtt-notifications, on-air, +pomodoro-timer) use it where the core has it and write the mailbox +otherwise. ## Robustness @@ -659,9 +676,11 @@ device never touches the live display. running display (their routes say so); they are not mailboxes. 5. **Remove the mailboxes (next release).** Once every device has run a display with stage 4, the web interface stops writing both mailboxes and - the display stops reading them. The four plugins that write - `display_on_demand_request` need an in-process way to ask for the screen - first. The display also stops writing `display_current_state`, + the display stops reading them. The four plugins that wrote + `display_on_demand_request` now have an in-process way to ask for the + screen (`BasePlugin.request_on_demand()` / `end_on_demand()`, see + "Plugins in the display process"); they keep the mailbox write only as + their fallback on older cores. The display also stops writing `display_current_state`, `display_on_demand_state` and `plugin_runtime_snapshot` once the web interface no longer falls back to them. diff --git a/docs/PLUGIN_API_REFERENCE.md b/docs/PLUGIN_API_REFERENCE.md index 200dba87..4341a62a 100644 --- a/docs/PLUGIN_API_REFERENCE.md +++ b/docs/PLUGIN_API_REFERENCE.md @@ -488,6 +488,78 @@ working for the plugin itself. `get_vegas_segment_width()` read the `vegas_panel_count` config value, which has never affected Vegas — a card's width comes from `get_vegas_content()` and `vegas_width_pct`. +### On-demand display + +A plugin that reacts to something outside the rotation (an MQTT message, a +timer, a detection) can take the screen for it, and give it back. Both +methods are safe from any thread, including an MQTT callback: they only +queue the request, and the display applies it on its render thread within a +frame or so, exactly like an on-demand start or stop from the web interface. + +#### `request_on_demand(mode=None, duration=None, pinned=False) -> Optional[str]` + +Show this plugin now. + +- `mode`: one of the plugin's display modes; `None` for its first. +- `duration`: seconds before the rotation resumes; `None` (or `0`) 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's plugin +manager, `scripts/check_plugin.py`) or its queue is full. A bad argument +(a `mode` that is not a string, a `duration` that is not a number) raises +`ValueError`. + +#### `end_on_demand() -> Optional[str]` + +Give the screen back. Ends only a session this plugin owns: a session the +user started for another plugin, or one that already ended, is left alone. +Returns the request id once queued, or `None` as above. + +#### Older cores: feature detection + +These methods are new after core 3.8.0 (see `CHANGELOG.md`). Before them, +plugins wrote the `display_on_demand_request` cache key (the "mailbox") +themselves. The display reads it only once a second while the control +socket is up, and it will be removed in a future release (see +[IPC_CONTROL_SOCKET.md](IPC_CONTROL_SOCKET.md), stage 5). A plugin that +must keep working on older cores checks for the method, and writes the +mailbox only when the method is missing or answers `None`: + +```python +import time, uuid + +def _show_alert(self): + if hasattr(self, "request_on_demand") and self.request_on_demand( + mode="my_alert", duration=15): + return + # Older core, or no display in this process: the mailbox, as before. + self.cache_manager.set("display_on_demand_request", { + "request_id": str(uuid.uuid4()), "action": "start", + "plugin_id": self.plugin_id, "mode": "my_alert", + "duration": 15, "pinned": False, "timestamp": time.time(), + }) + +def _release(self): + if hasattr(self, "end_on_demand") and self.end_on_demand(): + return + self.cache_manager.set("display_on_demand_request", { + "request_id": str(uuid.uuid4()), "action": "stop", + "plugin_id": self.plugin_id, "timestamp": time.time(), + }) +``` + +Keep `ledmatrix_min_version` where it is: the fallback is what keeps the +plugin working on older cores. A mailbox stop ends any on-demand session, +whoever started it; `end_on_demand()` ends only the plugin's own. + +Both methods answer a request id only when the plugin manager returned a +string, so a test that gives the plugin a `MagicMock()` plugin manager gets +`None` and exercises the mailbox path. To test the new path, set +`plugin_manager.request_on_demand.return_value = "some-id"`. + > The full source for `BasePlugin` lives in > `src/plugin_system/base_plugin.py`. If a method here disagrees with the > source, the source wins — please open an issue or PR to fix the doc. @@ -966,6 +1038,14 @@ if info: self.logger.info(f"Plugin: {info['name']}, Version: {info.get('version')}") ``` +#### `request_on_demand(plugin_id, mode=None, duration=None, pinned=False)` / `end_on_demand(plugin_id)` + +What `BasePlugin.request_on_demand()` and `end_on_demand()` call, with the +plugin's own id. Call those instead; see +[On-demand display](#on-demand-display). The display controller routes them +to itself with `set_on_demand_handler()`; a plugin manager without a +display behind it answers `None`. + #### `get_all_plugin_info() -> List[Dict[str, Any]]` Get information for all plugins. diff --git a/src/display_controller.py b/src/display_controller.py index dd64c6dd..84bfbe8c 100644 --- a/src/display_controller.py +++ b/src/display_controller.py @@ -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) diff --git a/src/ipc/server.py b/src/ipc/server.py index 1b81ca04..60b3d060 100644 --- a/src/ipc/server.py +++ b/src/ipc/server.py @@ -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: diff --git a/src/plugin_system/base_plugin.py b/src/plugin_system/base_plugin.py index 54314cd9..2b676eb1 100644 --- a/src/plugin_system/base_plugin.py +++ b/src/plugin_system/base_plugin.py @@ -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 diff --git a/src/plugin_system/plugin_manager.py b/src/plugin_system/plugin_manager.py index 69c9548b..1d2961fb 100644 --- a/src/plugin_system/plugin_manager.py +++ b/src/plugin_system/plugin_manager.py @@ -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 diff --git a/test/test_plugin_on_demand_api.py b/test/test_plugin_on_demand_api.py new file mode 100644 index 00000000..667724a6 --- /dev/null +++ b/test/test_plugin_on_demand_api.py @@ -0,0 +1,418 @@ +"""Plugins asking for the screen in-process: BasePlugin.request_on_demand() +and end_on_demand(). + +A plugin running in the display process used to write the +``display_on_demand_request`` mailbox, which the display reads once a second +while the control socket is up. These tests pin the way in that replaces it: + +* BasePlugin -> PluginManager -> DisplayController.submit_plugin_on_demand, + which only queues, from any thread; +* the render thread applies the queue where it applies socket commands, + through the mailbox's own handler, without the mailbox's read floor, and + woken by the control socket when it is up; +* a plugin's stop ends only its own session; +* no display to ask (the web interface's plugin manager, an old core's + plugin manager) answers None, which is a plugin's cue to fall back to the + mailbox; +* the mailbox still works for plugins that write it. +""" + +import logging +import threading +import time +from unittest.mock import MagicMock + +import pytest + +from src.ipc.server import ControlServer +from src.plugin_system.base_plugin import BasePlugin +from src.plugin_system.plugin_manager import PluginManager + + +class _Plugin(BasePlugin): + def update(self): + pass + + def display(self, force_clear=False): + pass + + +def _plugin(plugin_id, manager): + plugin = _Plugin.__new__(_Plugin) + plugin.plugin_id = plugin_id + plugin.plugin_manager = manager + return plugin + + +def _manager(handler=None): + manager = PluginManager.__new__(PluginManager) + manager.logger = logging.getLogger('test.plugin_on_demand') + if handler is not None: + manager.set_on_demand_handler(handler) + return manager + + +class _WakeServer: + """The parts of ControlServer the controller uses, with no socket.""" + + def __init__(self): + self.woken = 0 + self.has_pending = False + + def wake(self): + self.woken += 1 + self.has_pending = True + + def drain(self): + self.has_pending = False + return [] + + +@pytest.fixture +def controller(test_display_controller): + c_ = test_display_controller + c_.on_demand_active = False + c_.on_demand_request_id = None + c_._last_on_demand_poll = None + mailbox = {'value': None} + + def fake_get(key, *a, **kw): + if key == 'display_on_demand_request': + return mailbox['value'] + return None + + c_.cache_manager.get = MagicMock(side_effect=fake_get) + c_.cache_manager.set = MagicMock() + c_.cache_manager.delete = MagicMock() + c_._activate_on_demand = MagicMock() + c_.mailbox = mailbox + return c_ + + +@pytest.fixture +def wired(controller): + """A real PluginManager wired to the controller, as __init__ wires it.""" + manager = _manager(controller.submit_plugin_on_demand) + return controller, manager + + +class TestWiring: + def test_the_controller_wires_its_plugin_manager(self, controller): + controller.plugin_manager.set_on_demand_handler.assert_called_once_with( + controller.submit_plugin_on_demand) + + def test_a_start_reaches_the_mailbox_handler(self, wired): + controller, manager = wired + rid = _plugin('pomodoro-timer', manager).request_on_demand( + mode='pomodoro', duration=30, pinned=True) + assert isinstance(rid, str) and rid + controller._activate_on_demand.assert_not_called() # only queued + controller._poll_on_demand_requests() + controller._activate_on_demand.assert_called_once() + request = controller._activate_on_demand.call_args.args[0] + assert request['request_id'] == rid + assert request['action'] == 'start' + assert request['plugin_id'] == 'pomodoro-timer' + assert request['mode'] == 'pomodoro' + assert request['duration'] == 30.0 and request['pinned'] is True + assert request['source'] == 'plugin' + assert controller.on_demand_request_id == rid + + def test_a_plugin_request_never_touches_the_mailbox(self, wired): + controller, manager = wired + _plugin('on-air', manager).request_on_demand(mode='on_air') + controller._drain_control_commands() + controller._activate_on_demand.assert_called_once() + mailbox_reads = [call for call in controller.cache_manager.get.call_args_list + if call.args[0] == 'display_on_demand_request'] + assert mailbox_reads == [] + controller.cache_manager.delete.assert_not_called() + + def test_requests_apply_in_order(self, wired): + controller, manager = wired + seen = [] + controller._activate_on_demand = MagicMock( + side_effect=lambda r: seen.append(r['mode'])) + plugin = _plugin('p', manager) + for mode in ('a', 'b', 'c'): + plugin.request_on_demand(mode=mode) + controller._poll_on_demand_requests() + assert seen == ['a', 'b', 'c'] + + def test_the_mailbox_still_works_for_older_plugins(self, wired): + controller, manager = wired + controller.mailbox['value'] = {'request_id': 'mb', 'action': 'start', + 'plugin_id': 'birdnet-go'} + _plugin('on-air', manager).request_on_demand() + controller._poll_on_demand_requests() + ids = [call.args[0]['request_id'] for call in controller._activate_on_demand.call_args_list] + assert 'mb' in ids and len(ids) == 2 + + def test_a_failing_request_is_contained(self, wired): + controller, manager = wired + calls = [] + + def activate(request): + calls.append(request['mode']) + if request['mode'] == 'bad': + raise RuntimeError('plugin exploded') + + controller._activate_on_demand = MagicMock(side_effect=activate) + plugin = _plugin('p', manager) + plugin.request_on_demand(mode='bad') + plugin.request_on_demand(mode='good') + controller._poll_on_demand_requests() + assert calls == ['bad', 'good'] + + +class TestPromptness: + def test_a_plugin_request_skips_the_pending_changes_floor(self, wired): + controller, manager = wired + controller._control_server = None + controller._service_pending_changes() + _plugin('p', manager).request_on_demand() + controller._service_pending_changes() # well inside the 0.25 s floor + controller._activate_on_demand.assert_called_once() + + def test_a_plugin_request_skips_the_mailbox_floor(self, wired): + controller, manager = wired + controller._control_server = _WakeServer() + controller._poll_on_demand_requests() # sets the 1 s mailbox floor + _plugin('p', manager).request_on_demand() + controller._poll_on_demand_requests() + controller._activate_on_demand.assert_called_once() + + def test_it_wakes_the_control_socket_wait(self, wired): + controller, manager = wired + server = ControlServer('/nonexistent/control.sock') # never started + controller._control_server = server + assert not server.wait_for_command(0) + _plugin('p', manager).request_on_demand() + assert server.has_pending + assert controller._wait_for_control(5.0) is True # returns at once + assert controller._control_command_pending() + controller._poll_on_demand_requests() + controller._activate_on_demand.assert_called_once() + assert not server.has_pending + assert not controller._control_command_pending() + + def test_without_a_socket_a_waiting_request_cuts_the_sleep(self, wired): + controller, manager = wired + controller._control_server = None + _plugin('p', manager).request_on_demand() + started = time.monotonic() + assert controller._wait_for_control(5.0) is True + assert time.monotonic() - started < 1.0 + assert controller._control_command_pending() + + def test_nothing_waiting_keeps_the_floor(self, controller): + controller._control_server = None + controller._poll_on_demand_requests = MagicMock() + controller._service_pending_changes() + controller._service_pending_changes() + assert controller._poll_on_demand_requests.call_count == 1 + + +class TestThreads: + def test_requests_from_many_threads_all_land_in_order_per_thread(self, wired): + controller, manager = wired + seen = [] + controller._activate_on_demand = MagicMock( + side_effect=lambda r: seen.append(r['mode'])) + controller.PLUGIN_ON_DEMAND_QUEUE_SIZE = 10_000 + threads_n, each = 8, 50 + barrier = threading.Barrier(threads_n) + + def ask(n): + plugin = _plugin(f'p{n}', manager) + barrier.wait() + for i in range(each): + assert plugin.request_on_demand(mode=f'{n}:{i}') + + threads = [threading.Thread(target=ask, args=(n,)) for n in range(threads_n)] + for t in threads: + t.start() + # Drain while they ask, as the render thread would. + while any(t.is_alive() for t in threads): + controller._drain_control_commands() + for t in threads: + t.join() + controller._drain_control_commands() + assert len(seen) == threads_n * each + for n in range(threads_n): + mine = [int(m.split(':')[1]) for m in seen if m.startswith(f'{n}:')] + assert mine == list(range(each)) + + def test_a_full_queue_refuses(self, wired, caplog): + controller, manager = wired + controller.PLUGIN_ON_DEMAND_QUEUE_SIZE = 2 + plugin = _plugin('p', manager) + assert plugin.request_on_demand() + assert plugin.request_on_demand() + assert plugin.request_on_demand() is None + assert 'queue full' in caplog.text + controller._poll_on_demand_requests() + assert controller._activate_on_demand.call_count == 2 + assert plugin.request_on_demand() # room again + + +class TestStop: + def test_a_plugin_ends_its_own_session(self, wired): + controller, manager = wired + controller.on_demand_active = True + controller.on_demand_plugin_id = 'on-air' + controller._clear_on_demand = MagicMock() + assert _plugin('on-air', manager).end_on_demand() + controller._poll_on_demand_requests() + controller._clear_on_demand.assert_called_once_with(reason='requested-stop') + controller.cache_manager.delete.assert_not_called() + + def test_a_plugin_cannot_end_another_plugins_session(self, wired): + controller, manager = wired + controller.on_demand_active = True + controller.on_demand_plugin_id = 'clock' # the user started it + controller.on_demand_request_id = 'user' + controller._clear_on_demand = MagicMock() + _plugin('pomodoro-timer', manager).end_on_demand() + controller._poll_on_demand_requests() + controller._clear_on_demand.assert_not_called() + assert controller.on_demand_request_id == 'user' + + def test_a_stop_with_no_session_does_nothing(self, wired): + controller, manager = wired + controller.on_demand_status = 'error' + controller._clear_on_demand = MagicMock() + _plugin('on-air', manager).end_on_demand() + controller._poll_on_demand_requests() + controller._clear_on_demand.assert_not_called() + + def test_a_mailbox_stop_still_ends_any_session(self, wired): + controller, _ = wired + controller.on_demand_active = True + controller.on_demand_plugin_id = 'clock' + controller._clear_on_demand = MagicMock() + controller.mailbox['value'] = {'request_id': 's', 'action': 'stop', + 'plugin_id': 'on-air'} + controller._poll_on_demand_requests() + controller._clear_on_demand.assert_called_once_with(reason='requested-stop') + + def test_start_then_stop_from_one_thread_ends_the_session(self, wired): + controller, manager = wired + + def activate(request): + controller.on_demand_active = True + controller.on_demand_plugin_id = request['plugin_id'] + + controller._activate_on_demand = MagicMock(side_effect=activate) + controller._clear_on_demand = MagicMock() + plugin = _plugin('pomodoro-timer', manager) + plugin.request_on_demand(mode='pomodoro', pinned=True) + plugin.end_on_demand() + controller._poll_on_demand_requests() + controller._activate_on_demand.assert_called_once() + controller._clear_on_demand.assert_called_once_with(reason='requested-stop') + + +class TestNoDisplay: + """None is a plugin's cue to write the mailbox instead.""" + + def test_a_manager_with_no_handler_answers_none(self): + plugin = _plugin('p', _manager()) + assert plugin.request_on_demand() is None + assert plugin.end_on_demand() is None + + def test_no_plugin_manager_answers_none(self): + plugin = _plugin('p', None) + assert plugin.request_on_demand() is None + assert plugin.end_on_demand() is None + + def test_an_old_cores_plugin_manager_answers_none(self): + class OldManager: + plugin_manifests = {} + + plugin = _plugin('p', OldManager()) + assert plugin.request_on_demand() is None + assert plugin.end_on_demand() is None + + def test_a_handler_that_raises_answers_none(self): + def broken(request): + raise RuntimeError('boom') + + plugin = _plugin('p', _manager(broken)) + assert plugin.request_on_demand() is None + assert plugin.end_on_demand() is None + + def test_a_handler_that_refuses_answers_none(self): + plugin = _plugin('p', _manager(lambda request: False)) + assert plugin.request_on_demand() is None + + def test_a_controller_built_without_init_refuses(self): + from src.display_controller import DisplayController + bare = DisplayController.__new__(DisplayController) + assert bare.submit_plugin_on_demand({'action': 'start'}) is False + assert bare._plugin_on_demand_pending() is False + bare._drain_plugin_on_demand() # nothing to do, no error + + def test_the_feature_detection_pattern(self): + """The hasattr pattern from docs/PLUGIN_API_REFERENCE.md.""" + writes = [] + + class OldCorePlugin: # an older core's BasePlugin has no such method + pass + + for plugin, expect_mailbox in ((OldCorePlugin(), True), + (_plugin('p', _manager()), True), + (_plugin('p', _manager(lambda r: True)), False)): + writes.clear() + if not (hasattr(plugin, 'request_on_demand') + and plugin.request_on_demand(mode='m')): + writes.append('mailbox') + assert (writes == ['mailbox']) is expect_mailbox + + +class TestArguments: + def test_the_manager_shapes_the_request(self): + got = [] + plugin = _plugin('p', _manager(lambda r: got.append(r) or True)) + plugin.request_on_demand() + plugin.end_on_demand() + start, stop = got + assert start['plugin_id'] == 'p' and start['mode'] is None + assert start['duration'] is None and start['pinned'] is False + assert start['source'] == 'plugin' and start['timestamp'] > 0 + assert stop == {'action': 'stop', 'plugin_id': 'p', 'request_id': stop['request_id'], + 'timestamp': stop['timestamp'], 'source': 'plugin'} + assert start['request_id'] != stop['request_id'] + + @pytest.mark.parametrize('duration', [0, -5, float('inf'), float('nan')]) + def test_no_positive_duration_means_no_limit(self, duration): + got = [] + _plugin('p', _manager(lambda r: got.append(r) or True)).request_on_demand( + duration=duration) + assert got[0]['duration'] is None + + @pytest.mark.parametrize('kwargs', [{'mode': 5}, {'mode': ''}, {'duration': '30'}, + {'duration': True}]) + def test_bad_arguments_raise(self, kwargs): + plugin = _plugin('p', _manager(lambda r: True)) + with pytest.raises(ValueError): + plugin.request_on_demand(**kwargs) + + +class TestMockManagers: + def test_a_magicmock_manager_reads_as_not_taken(self): + """A plugin's test with a MagicMock manager keeps its mailbox path.""" + plugin = _plugin('p', MagicMock()) + assert plugin.request_on_demand(mode='m') is None + assert plugin.end_on_demand() is None + plugin.plugin_manager.request_on_demand.assert_called_once_with( + 'p', mode='m', duration=None, pinned=False) + plugin.plugin_manager.end_on_demand.assert_called_once_with('p') + + def test_a_mocked_id_is_passed_through(self): + manager = MagicMock() + manager.request_on_demand.return_value = 'rid' + manager.end_on_demand.return_value = 'rid2' + plugin = _plugin('p', manager) + assert plugin.request_on_demand() == 'rid' + assert plugin.end_on_demand() == 'rid2'