diff --git a/CHANGELOG.md b/CHANGELOG.md index 5319b1dc..e5894d34 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -82,6 +82,45 @@ policies are unchanged. `GET /api/v3/plugins/fetch-stats`. - `fetch_service` is a core config section (`src/core_config_keys.py`). +### Control socket (stage 2: wake-ups, brightness, plugin reload) + +- **Socket commands land at once.** Stage 1's socket was no faster than the + mailbox: a command waited for the static screen's 1 s frame sleep, the + dwell's 0.25 s tick, or Vegas's interrupt check every 10 frames (about + 0.4 s on a Pi 4). The render thread now waits on the socket's queue + instead of sleeping, and Vegas checks the queue every frame, so an + on-demand start or stop is applied within about a millisecond on a static + screen or in a dwell, and at the next frame in Vegas or on a scrolling + screen. Commands still run only on the render thread. The file mailbox + keeps its old delays. Idle CPU is unchanged in practice: the waits are + timed `Event` waits with the same wake-ups as the sleeps they replace + (about 25 µs more per wait, measured). +- **`brightness.set`.** Saving a brightness (`POST /api/v3/config/main`) + also puts it on the panel at once over the socket, instead of when the + display's config watcher next reads `config.json` (up to about 2 s). The + response says `brightness_transport: "socket"`, or `"config"` with + `brightness_socket_error` when the watcher applies it as before. The + command itself writes nothing; the dim schedule still applies on top. +- **`plugin.reload`.** Updating an enabled plugin from the store no longer + asks for a display restart when the display can reload it: the update + route asks the display over the socket, which reloads the plugin on its + render thread at the start of the next screen (its modes keep their place + in the rotation) and answers once the new code runs. The response then + says `restart_required: false`, `reloaded: true` and `reloaded_version`. + Without the socket, with a display older than this command, or when the + reload fails, the route answers `restart_required: true` as before, with + `reload_error` giving the reason. The route now also uses the manifest's + plugin id (the one the display runs it under) for this decision, so an + update through a registry alias of an enabled plugin no longer reports + that no restart is needed. +- Both commands answer with the render thread's outcome, or `pending` when + it did not get to them in time (2 s and 10 s). Socket protocol version is + still 1: new commands are additive, and an older display answers + `unknown_command`, which the web interface falls back from. The security + model is unchanged: the same `0660` group socket and peer-credential + check. `config.reload` was not added; see + `docs/IPC_CONTROL_SOCKET.md` for why. + ### New modules - `src/common/fetch_service.py` -- the fetch service above. Core-internal in diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index a20b6f74..12db5d41 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -86,7 +86,8 @@ How a web-side change reaches the running plugins: | Plugin enabled or disabled | `ConfigService` → `_controller_config_change` flags a reconcile; `_reconcile_enabled_plugins` loads it (fresh from disk) or unloads it on the render thread | | Plugin uninstalled (config removed) | the removed section flips its `enabled` flag, and the reconcile unloads it | | Plugin installed, not enabled | nothing to do until it is enabled, which loads it | -| Plugin installed while already enabled, updated while enabled, or uninstalled with its config kept | **not picked up**: the display keeps running what it loaded. The route answers `restart_required: true` and the UI shows its restart banner | +| Plugin updated while enabled | the update route asks the display over the control socket (`plugin.reload`) to reload it on the render thread, and answers `restart_required: false` once the new code runs. Without the socket, as the next row | +| Plugin installed while already enabled, updated while enabled and not reloaded, or uninstalled with its config kept | **not picked up**: the display keeps running what it loaded. The route answers `restart_required: true` and the UI shows its restart banner | `display_restart_required()` in `plugin_catalog.py` holds that last rule; routes return it as `restart_required` (with the banner's wording in @@ -110,9 +111,11 @@ other web-UI action runs its script as a subprocess. A later, explicit **plugin web-entry contract** -- a declared entry point for plugin web code -- replaces that function. -Next stages: a **control socket** from the web process to the display -(reload one plugin, ask for its state) in place of `restart_required` and -the cache-key mailboxes, and the plugin web-entry contract above. +The **control socket** from the web process to the display +([IPC_CONTROL_SOCKET.md](IPC_CONTROL_SOCKET.md)) carries on-demand +commands and reloads an updated plugin; its next stages stream the +display's state and retire the cache-key mailboxes. The plugin web-entry +contract above is still to come. ### Plugin state: desired, observed, and who owns it diff --git a/docs/IPC_CONTROL_SOCKET.md b/docs/IPC_CONTROL_SOCKET.md index 0b84c6e1..a7ba8367 100644 --- a/docs/IPC_CONTROL_SOCKET.md +++ b/docs/IPC_CONTROL_SOCKET.md @@ -2,14 +2,16 @@ The display process serves a Unix socket that the web interface uses to send it commands and get an answer back. It replaces the cache-file "mailboxes" on -the SD card one command at a time. Stage 1, described here, carries on-demand -start/stop/status. The file mailbox stays as a fallback for one release. +the SD card one command at a time. Stage 1 carries on-demand start, stop and +status. Stage 2 makes those commands land within a frame on every kind of +screen, and adds `brightness.set` and `plugin.reload`. The file mailbox stays +as a fallback for one release. | | | |---|---| | Socket | `/run/ledmatrix/control.sock` (tmpfs) | | Served by | the display process ([`src/ipc/server.py`](../src/ipc/server.py)), started by `DisplayController.run()` | -| Used by | the web interface ([`src/ipc/client.py`](../src/ipc/client.py)): `POST /api/v3/display/on-demand/start` and `/stop` | +| Used by | the web interface ([`src/ipc/client.py`](../src/ipc/client.py)): `POST /api/v3/display/on-demand/start` and `/stop`, `POST /api/v3/plugins/update` (reload), `POST /api/v3/config/main` (brightness) | | Contract | [`src/ipc/contract.py`](../src/ipc/contract.py): messages, versions, framing and the socket path; both sides import it | | Override | `LEDMATRIX_CONTROL_SOCKET=/some/path.sock` for both processes, or `=off` to disable it | @@ -71,6 +73,8 @@ one. Clients branch on `error.code`, never on the message text. | `on_demand.start` | `{plugin_id?, mode?, duration?, pinned?}` (at least one of `plugin_id` and `mode`) | ack | queued | | `on_demand.stop` | — | ack | queued | | `on_demand.status` | — | `{on_demand: {...}, current_mode, display_active}` | answered directly | +| `brightness.set` | `{brightness: int 0-100}` | `{brightness, panel_brightness, dimmed, display_active}` | queued, awaited (2 s) | +| `plugin.reload` | `{plugin_id}` | `{plugin_id, reloaded: true, version, modes}` | queued, awaited (10 s) | `duration` is a number of seconds, or a numeric string. `0`, `null` or `""` mean "until stopped". `pinned` must be a real boolean: the REST route has @@ -78,28 +82,59 @@ already converted strings like `"false"` before it sends the command. The `on_demand` object in `on_demand.status` is the same dict the display publishes to `display_on_demand_state`. -**Acknowledgements.** A queued command is *accepted*, not *done*. +`brightness.set` sets the panel's normal brightness. It is transient: it +writes nothing to `config.json`, and the next config the display's watcher +loads (or a restart) puts the configured value back. The web interface +sends it after it has saved the setting, so the two agree. The dim schedule +still applies on top, so `panel_brightness` is the dim level while the +schedule dims. While the schedule has the display off, the new level is +kept for when it comes back on. + +`plugin.reload` loads a plugin the display is running again from disk, +manifest included: the steps of disabling it live and enabling it again, +with its modes kept in their place in the rotation. Only a running plugin +can be reloaded (`not_loaded` otherwise), so the id never makes the display +import anything new. A plugin loaded only for an on-demand session gets +`busy`. A new version that fails to load gets `failed` and stays out of the +rotation, as it would after a restart. + +**Acknowledgements.** A queued on-demand command is *accepted*, not *done*. `{"accepted": true, "request_id": …}` means the command is waiting in the -render thread's queue, and the render thread will apply it at its next -on-demand check. That is within one frame on a scrolling screen, 0.25 s -during a dwell, and up to 1 s on a static screen, whose frame loop sleeps a -second between frames. Except on a scrolling screen, where the mailbox waits -up to 0.25 s, these are the mailbox's delays too: stage 1 adds -acknowledgements, not speed. Any outcome is published as before +render thread's queue. Since stage 2 the render thread waits on that queue +instead of sleeping, so it applies the command within one frame on every +kind of screen (see below). Any outcome is published as before (`display_on_demand_state`, and `status`/`error` for a bad plugin or mode), and it can be read with `on_demand.status`. +**Awaited commands.** `brightness.set` and `plugin.reload` are answered only +once the render thread has applied them, with their result or their error. +The connection thread waits for that (2 s and 10 s, `AWAIT_SECONDS` in the +contract); the render thread never waits for a client. When the render thread +has not got to the command in time, the answer is `pending`: the command +stays queued and is still applied, so a client treats `pending` as "not known +to be done", not as a refusal. The client's own timeout is one second longer +than the display's wait, so `pending` arrives before the client gives up. + **Versions.** Every request carries `v`. For any command except `hello`, a `v` the display does not speak gets `unsupported_version`. `hello` is checked by its `versions` list instead, and its result names the highest version both sides share, so a client can find out what a display supports before it -relies on anything newer. Stage 1's client sends `v: 1` and falls back to the +relies on anything newer. The client sends `v: 1` and falls back to the mailbox when the display refuses it. It does not send `hello` first, which saves a round trip. +New commands are added within a version, so stage 2 is still version 1. A +display that does not know a command answers `unknown_command`, which the +web interface treats like any other socket failure and falls back from, and +`hello` lists the commands a display knows. The version changes only when the +envelope or the meaning of an existing command changes. + **Error codes:** `bad_json`, `bad_request`, `message_too_large`, `unsupported_version`, `unknown_command`, `invalid_args`, `busy` (queue full, or too many connections), `forbidden` (peer credentials refused), `internal`. +Stage 2 adds `pending` (accepted, not applied in time, still queued), +`not_loaded` (`plugin.reload` of a plugin the display is not running) and +`failed` (the render thread tried, and it did not work). Try it on a device: @@ -107,6 +142,7 @@ Try it on a device: python3 - <<'EOF' from src.ipc import client # run from the project directory print(client.on_demand_status()) +print(client.brightness_set(60)) EOF ``` @@ -116,19 +152,71 @@ The server's threads never touch rendering. A connection thread parses the request, validates it against the contract, and then does one of two things: - For a command that changes the panel, it puts a `QueuedCommand` on a - bounded queue (16 entries) and answers with the ack. + bounded queue (16 entries) and answers with the ack, or, for an awaited + command, with the outcome the render thread reports back through the + command's `CommandOutcome`. - For a query, it answers from a status snapshot the display provides (`DisplayController._control_status`). The snapshot only reads attributes. The render thread drains the queue in `_poll_on_demand_requests()`, the same -place it reads the mailbox, and hands each command to -`_handle_on_demand_request()`, which is the mailbox's own handler. The two -paths share all of their code: activation, the processed-id guard, error -publishing, and resuming the rotation afterwards. The 0.25 s floor on the -mailbox read does not apply to the queue, because draining it costs no disk -read. A queued command also lets `_service_pending_changes()` skip its own -floor, so a long scrolling screen or a Vegas iteration takes the command at -its next frame. +place it reads the mailbox: + +- An on-demand command goes to `_handle_on_demand_request()`, which is the + mailbox's own handler. The two paths share all of their code: activation, + the processed-id guard, error publishing, and resuming the rotation + afterwards. +- `brightness.set` is applied there and then (`_apply_control_brightness`), + and the current frame is pushed again so the panel shows it. +- `plugin.reload` waits for the top of the next loop pass, the place where + plugins are enabled and disabled live, because there no `display()` and no + Vegas iteration is on the stack (`_apply_pending_plugin_reloads`). Until + then the current screen ends early, as it does for a WiFi notice: the + frame loops, the dwell and Vegas's interrupt check all treat a pending + reload as a reason to stop (`_screen_preempted`). The rotation then + advances, and the next pass reloads before it draws. + +The 0.25 s floor on the mailbox read does not apply to the queue, because +draining it costs no disk read. A queued command also lets +`_service_pending_changes()` skip its own floor. + +### Waking the render thread (stage 2) + +Stage 1 made the socket answer, but not land sooner: a queued command waited +for the same polls the mailbox does. Measured on ledpi (Pi 4, 24 fps Vegas), +a start took 1.02 s on a static screen and about 0.4 s in Vegas either way. +Now the queue wakes the render thread: + +- **The waits.** The server sets a `threading.Event` whenever it queues a + command. The render thread waits on it (`ControlServer.wait_for_command`) + where it used to sleep: the static screen's 1 s frame sleep + (`_wait_frame_interval`) and the dwell's 0.25 s ticks + (`_sleep_with_plugin_updates`, which also covers scheduled-off and the + empty-rotation pause). On a wake it applies the command at once. A command + that does not end the screen, such as a brightness, does not cut the frame + short: the wait carries on to the end of the interval, so the plugin is + still drawn once a second. +- **Vegas.** The coordinator still runs its interrupt check every 10 frames, + and now also at any frame where `urgent()` is true. The display passes + "a control socket command is queued", which is one `Event.is_set()` per + frame. +- **Scrolling screens** already service pending changes every frame. + +So a command lands within a millisecond or so on a static screen and in a +dwell, and within one frame in Vegas and on a scrolling screen. The mailbox +keeps its old delays. Commands still run only on the render thread: the +connection threads only queue them and set the event. + +The waits are timed `Event.wait()` calls: no polling, and no more wake-ups +than the sleeps they replace when nothing arrives. Measured under WSL +(Python 3.12, 20 s runs in the order before, after, after, before, with the +socket's accept thread up), the idle process used 0.015–0.018% of a core +before and 0.019–0.021% after on a static screen, and 0.035–0.037% before and +0.047% after in a dwell: about 25 µs more per wait, from `Event.wait`'s own +bookkeeping. A client's send to the render thread waking took 0.72 ms median +(1.04 ms max), and a whole `brightness.set` round trip 0.64 ms median. + +Without a socket (Windows, `LEDMATRIX_CONTROL_SOCKET=off`) the waits are the +plain sleeps they were. **Exactly once.** A command and a mailbox write for the same request share one `request_id`. If the client times out after the display queued the @@ -154,6 +242,11 @@ block the render loop or crash it: - **Full queue.** When the queue is full, the client gets `busy` and falls back to the mailbox. A full queue means the render thread is stuck, and the systemd watchdog deals with that. +- **Awaited commands.** The wait for an awaited command's outcome happens on + its connection thread and is bounded (`AWAIT_SECONDS`), so a stuck render + thread costs that client `pending` and one connection slot for at most + 10 s. The render thread settles an outcome without blocking; one nobody is + waiting for any more is simply dropped. - **Startup.** The server binds under a temporary name, sets the mode and the group, then renames the socket into place, so it never appears with the umask's permissions. It removes a stale socket (a file that nothing is @@ -191,9 +284,13 @@ anything else in the group they share: and is disconnected. This covers a socket mode that someone loosened by hand. -The commands are deliberately narrow. Stage 1 can start or stop on-demand -display and read its state, which anyone who can reach the web UI can already -do. Nothing on the socket runs a shell, writes a file, or names a path. +The commands are deliberately narrow. They start or stop on-demand display, +read its state, set the brightness, and reload a plugin the display is +already running, all of which anyone who can reach the web UI can already do +(the last by restarting the display). Nothing on the socket runs a shell, +writes a file, or names a path, and `plugin.reload` cannot make the display +import a plugin it was not running. Stage 2 changed none of the access rules +above. **Development.** A display that is not root and cannot write to `/run/ledmatrix`, such as `python3 run.py -e` from a checkout, serves the @@ -205,28 +302,49 @@ device never touches the live display. ## Stage plan -1. **On-demand, with acks (this stage).** Contract, server, client. +1. **On-demand, with acks (done, #706).** Contract, server, client. `on_demand.start`/`stop`/`status`, `hello`, `ping`. The REST routes try the socket first and report `transport: "socket" | "mailbox"` (plus `socket_error` on fallback). The mailbox is unchanged, and the plugins that write it directly (birdnet-go, mqtt-notifications, on-air, pomodoro-timer) keep working. -2. **Commands that are restarts or polls today.** - - `brightness.set`, transient and with no `config.json` write. +2. **Commands that were restarts or polls (done).** + - The render thread waits on the queue instead of sleeping, and Vegas + checks it every frame, so a command lands within a frame on every kind + of screen (see "Waking the render thread"). + - `brightness.set`, transient and with no `config.json` write. `POST + /api/v3/config/main` sends it after saving a brightness and reports + `brightness_transport`; without the socket the config watcher applies + the saved value, as before. - `plugin.reload`, which replaces the `restart_required` answer from #688 - with a live reload of the updated plugin on the render thread. - - `config.reload`, which applies a saved config without waiting for the 2 s - mtime poll and acks which sections changed. - - The dwell sleep and the static screen's 1 s frame sleep wait on the - queue instead of sleeping, so a command lands within milliseconds on - every kind of screen. Under WSL, with a static plugin on screen, a stop - takes 1.0 s by either path today. + for a store update of an enabled plugin. `POST /api/v3/plugins/update` + answers `restart_required: false, reloaded: true` once the new code + runs, and falls back to the restart banner (with `reload_error`) + otherwise. + - `config.reload` was left out. Its only gain over the config watcher + would be skipping the watcher's 2 s mtime poll, and the one setting + where those seconds show, brightness, now has its own command. Plugin + settings already reach the running plugin through the watcher, and the + "which sections changed" ack had no reader: the web interface knows + what it saved. A reload from the socket thread would also run every + config subscriber on a second thread beside the watcher's. 3. **A state stream.** A `subscribe` command that keeps the connection open - and pushes events: mode changes, on-demand state, plugin runtime state and - the heartbeat. It replaces the polled `display_current_state`, + and pushes events: mode changes, on-demand state (including the outcome of + an acked on-demand command, which today is only published), plugin + runtime state, the outcome of a reload that answered `pending`, and the + heartbeat. It replaces the polled `display_current_state`, `plugin_runtime_snapshot` (#690) and `display-heartbeat.json` (#687) for readers that hold a connection. The web interface relays it to its existing SSE stream. The files remain for one release for older readers. + - The server's per-connection threads (8 at most) do not suit long-lived + subscribers. A subscriber needs its own bound and a writer that drops + events for a slow reader rather than blocking the display. + - Events are produced on the render thread, so publishing must be a + non-blocking hand-off, like the queue in the other direction. + - The store's install of an already-enabled plugin, and an uninstall that + keeps its config, still answer `restart_required`. With the stream they + can use a load/unload command and report the result the same way the + update route does now. 4. **Retire the mailboxes.** After a release in which every device has had the socket, the web interface stops writing `display_on_demand_request`, and the display stops polling it, logging the plugins that still write it so @@ -248,3 +366,18 @@ If the response says `"transport": "mailbox"`, `socket_error` gives the reason. `no_socket` means the display is stopped or predates the socket. `refused` usually means the web user is not in the socket's group, which takes effect when the web service restarts after the user is added. + +Brightness and a plugin reload: + +```bash +curl -s -X POST localhost:5000/api/v3/config/main \ + -H 'Content-Type: application/json' -d '{"brightness":40}' +# ... "brightness_transport": "socket" +curl -s -X POST localhost:5000/api/v3/plugins/update \ + -H 'Content-Type: application/json' -d '{"plugin_id":"clock-simple"}' +# after a real update of an enabled plugin: "restart_required": false, "reloaded": true +sudo journalctl -u ledmatrix | grep -E "Brightness set|Reload(ing|ed) plugin" +``` + +`unknown_command` in `brightness_socket_error` or `reload_error` means the +display runs a stage-1 build: restart it once to pick up this one. diff --git a/docs/REST_API_REFERENCE.md b/docs/REST_API_REFERENCE.md index 37479f1f..1e08576d 100644 --- a/docs/REST_API_REFERENCE.md +++ b/docs/REST_API_REFERENCE.md @@ -165,6 +165,14 @@ the web UI shows its restart banner on the flag. (Plugin sections saved through this route reach the running plugin live, like `POST /plugins/config`.) +A saved `brightness` is the exception: it reaches the panel without a +restart. The route also sends it to the running display over the control +socket (`brightness.set`), which puts it on the panel at once, and the +response adds `"brightness_transport": "socket"`. Otherwise it is +`"config"`, with `brightness_socket_error` giving the reason, and the +display's config watcher applies the saved value within a few seconds, as +before. + Invalid values (e.g. an out-of-range `target_fps`, a hardware option the Raspberry Pi 5 driver cannot use) are rejected with `400` and nothing is saved. @@ -456,8 +464,9 @@ Request a specific plugin to display on-demand. `service` is `null` when `start_service` is false. `transport` says how the request reached the display: `"socket"` means the -display's control socket acknowledged it (it is queued for the render thread; -see [IPC_CONTROL_SOCKET.md](IPC_CONTROL_SOCKET.md)), `"mailbox"` means it was +display's control socket acknowledged it (it is queued for the render thread, +which wakes for it and applies it within a frame; see +[IPC_CONTROL_SOCKET.md](IPC_CONTROL_SOCKET.md)), `"mailbox"` means it was written to the cache mailbox the display polls, as before the socket existed. With `"mailbox"`, `socket_error` gives the reason the socket was not used (`no_socket` when the display is stopped or predates the socket, `timeout`, @@ -812,8 +821,29 @@ Update a plugin to the latest version. Runs synchronously. ``` `update_status` is `updated`, `up_to_date` or `local_only`. -`restart_required` is true when the plugin changed and is enabled: the -running display keeps the code it loaded until it restarts. + +When the plugin changed and is enabled, the route asks the running display +to reload it over the control socket (`plugin.reload`, see +[IPC_CONTROL_SOCKET.md](IPC_CONTROL_SOCKET.md)). Once the new code is +running, the answer is: + +```json +{ + "status": "success", + "message": "Plugin football-scoreboard updated to version 2.1.0; the display is running the new version", + "restart_required": false, + "reloaded": true, + "reloaded_version": "2.1.0" +} +``` + +If the display could not reload it, `restart_required` is true (the running +display keeps the code it loaded until it restarts) and `reload_error` says +why: `no_socket` (the display is stopped or predates the socket), +`unknown_command` (a display older than this command), `not_loaded`, +`failed` (the new version did not load; it is out of the rotation), +`pending` (not done within 10 s; it will still be reloaded), or another +transport reason. An update this core cannot run answers `409` with `Plugin update refused:` and the reason; the installed version is left as it was. diff --git a/docs/RUN_LOOP_REDESIGN.md b/docs/RUN_LOOP_REDESIGN.md index dafe195a..89452b13 100644 --- a/docs/RUN_LOOP_REDESIGN.md +++ b/docs/RUN_LOOP_REDESIGN.md @@ -33,7 +33,13 @@ on purpose. Each pass, in order: -1. `loop_pass()` (watchdog). Apply a pending plugin enable/disable. +1. `loop_pass()` (watchdog). Apply a pending plugin enable/disable, then + any plugin reloads the control socket asked for + (`_apply_pending_plugin_reloads`; a pending reload ends the screen + before it, like a WiFi notice, through `_screen_preempted`). The static + screen's frame sleep and the dwell wait on the socket's queue instead of + sleeping (`_wait_frame_interval`, `_sleep_with_plugin_updates`); without + a socket, as in the golden traces, they are the plain sleeps. 2. With no modes: dwell 1 s, next pass. 3. Poll on-demand requests and expiry, release plugins loaded only for on-demand, tick plugin updates, drop an expired WiFi notice, evaluate diff --git a/src/display_controller.py b/src/display_controller.py index 6ae459a3..ea42ce52 100644 --- a/src/display_controller.py +++ b/src/display_controller.py @@ -44,7 +44,15 @@ from src.logging_config import get_logger from src.exceptions import PluginError from src.common.frame_timing import HANDOVER_OP from src.common.sync_manager import DisplaySyncManager, SyncRole -from src.ipc.server import ControlServer, start_control_server +from src.ipc.contract import ( + BrightnessResult, + BrightnessSetArgs, + Command as ControlCommand, + ErrorCode as ControlErrorCode, + PluginReloadArgs, + PluginReloadResult, +) +from src.ipc.server import ControlServer, QueuedCommand, start_control_server from src.vegas_mode.render_pipeline import SYNC_SEND_INTERVAL # Get logger with consistent configuration @@ -579,9 +587,13 @@ class DisplayController: # Set up interrupt checker for on-demand/wifi status and follower mode def _vegas_interrupt(): return self._check_vegas_interrupt() or self.sync_manager.is_follower_active() + # Every 10 frames (~80ms at 125 FPS, ~0.4 s at the 24 a Pi 4 + # often manages), or at the next frame when a control socket + # command is queued: that check is one Event read per frame. self.vegas_coordinator.set_interrupt_checker( _vegas_interrupt, - check_interval=10 # Check every 10 frames (~80ms at 125 FPS) + check_interval=10, + urgent=self._control_command_pending, ) # Run plugin updates inside the Vegas loop so the inter-iteration @@ -714,6 +726,10 @@ class DisplayController: if self._check_wifi_status_message(): return True + # A plugin reload waits for the top of the loop, outside the iteration. + if self._plugin_reload_pending: + return True + return False def _timezone(self): @@ -1154,12 +1170,15 @@ class DisplayController: Also services pending changes (see _service_pending_changes), and returns early when one of them changes what the panel should show -- an on-demand start or stop, the display schedule turning the panel - on or off, a WiFi notice arriving, or a live game taking over - (_check_live_takeover) -- so the caller can act on it instead of - finishing a dwell that could be a minute long (sixty seconds while - scheduled off). + on or off, a WiFi notice arriving, a live game taking over + (_check_live_takeover), or a plugin reload waiting for the top of + the loop -- so the caller can act on it instead of finishing a dwell + that could be a minute long (sixty seconds while scheduled off). + + The waits between checks wake for a control socket command, so one + is applied within milliseconds rather than at the next 0.25 s tick. """ - if duration <= 0: + if duration <= 0 or self._plugin_reload_pending: return end_time = time.time() + duration @@ -1177,7 +1196,9 @@ class DisplayController: break sleep_time = min(tick_interval, remaining) - time.sleep(sleep_time) + # Woken early by a control socket command, which + # _service_pending_changes then applies without its floor. + self._wait_for_control(sleep_time) # A dwell can be a minute long (sixty seconds while scheduled # off); the watchdog must hear from this thread throughout. display_watchdog.watchdog.beat() @@ -1188,7 +1209,8 @@ class DisplayController: or self.is_display_active != display_active or self.on_demand_active != on_demand or (not wifi_pending and self.is_display_active - and self._wifi_notice_pending())): + and self._wifi_notice_pending()) + or self._plugin_reload_pending): break def _note_empty_pass(self) -> None: @@ -1657,22 +1679,192 @@ class DisplayController: 'display_active': self.is_display_active} def _drain_control_commands(self) -> None: - """Apply on-demand commands that arrived over the control socket. + """Apply the commands that arrived over the control socket. - Each goes through _handle_on_demand_request, the mailbox's own - handler, so both ways in behave the same, and a request that came - both ways (a client that timed out and fell back) has one request - id and is processed once. + On-demand commands go through _handle_on_demand_request, the + mailbox's own handler, so both ways in behave the same, and a + request that came both ways (a client that timed out and fell back) + has one request id and is processed once. A brightness is applied + 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). """ server = self._control_server if server is None or not server.has_pending: return for command in server.drain(): try: - self._handle_on_demand_request(command.as_on_demand_request()) + if command.cmd == ControlCommand.BRIGHTNESS_SET: + self._apply_control_brightness(command) + elif command.cmd == ControlCommand.PLUGIN_RELOAD: + self._pending_plugin_reloads = self._pending_plugin_reloads + (command,) + else: + self._handle_on_demand_request(command.as_on_demand_request()) except Exception: # pylint: disable=broad-except logger.exception("Failed to apply control socket command %s", command.request_id) + command.fail(ControlErrorCode.INTERNAL, 'the display failed to apply it') + + 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. + """ + server = self._control_server + wait = getattr(server, 'wait_for_command', None) if server is not None else None + if wait is None: + 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.""" + server = self._control_server + return bool(server is not None and server.has_pending) + + def _screen_preempted(self, active_mode: Optional[str]) -> bool: + """What ends a screen mid-way: the checks the frame loops make.""" + return (self.current_display_mode != active_mode + or not self.is_display_active + or self._wifi_notice_pending() + or self._plugin_reload_pending) + + def _wait_frame_interval(self, interval: float, active_mode: Optional[str]) -> bool: + """The static screen's sleep between frames, woken by socket commands. + + Each command that arrives is applied at once (_service_pending_changes + skips its floor while one is queued). True when that ended the + screen; otherwise the wait carries on to the end of the interval, so + the frame cadence is unchanged by a command that does not change the + screen (a brightness, say). + """ + if self._control_server is None: + time.sleep(interval) + return False + deadline = time.monotonic() + interval + while True: + remaining = deadline - time.monotonic() + if remaining <= 0: + return False + if not self._wait_for_control(remaining): + return False + self._service_pending_changes() + if self._screen_preempted(active_mode): + return True + + def _apply_control_brightness(self, command: QueuedCommand) -> None: + """``brightness.set``: the new normal brightness, on the panel now. + + It replaces the configured value in memory only, the way a saved + config would once the config watcher saw it, so the dim schedule and + the scheduled-off rules treat it exactly like the setting; the next + config the watcher loads replaces it again. + """ + args = command.args + if not isinstance(args, BrightnessSetArgs): + command.fail(ControlErrorCode.INTERNAL, 'not a brightness command') + return + self._normal_brightness = args.brightness + # The dim schedule's per-minute cache holds the old normal level. + self._dim_checked_minute = None + # An explicit request retries a level the panel refused before. + self._failed_brightness_target = None + self._apply_brightness_target(repaint=True) + if self.is_display_active: + target = self._check_dim_schedule() # cached for the minute: no new work + if self.current_brightness != target: + command.fail(ControlErrorCode.FAILED, + f'the panel did not take brightness {target}') + return + logger.info("Brightness set to %d%% over the control socket (panel %s%%)", + args.brightness, self.current_brightness) + result: BrightnessResult = { + 'brightness': args.brightness, + 'panel_brightness': int(self.current_brightness), + 'dimmed': bool(self.is_dimmed), + 'display_active': bool(self.is_display_active), + } + command.succeed(dict(result)) + + #: Plugin reloads from the control socket, waiting for the top of the + #: next loop pass. A tuple, replaced rather than mutated. + _pending_plugin_reloads: Tuple[QueuedCommand, ...] = () + + @property + def _plugin_reload_pending(self) -> bool: + return bool(self._pending_plugin_reloads) + + def _apply_pending_plugin_reloads(self) -> None: + """Reload the plugins the control socket asked for. Render thread, + top of the loop pass: no display() and no Vegas iteration on the stack. + + The same steps as disabling and re-enabling the plugin live + (_unregister_plugin, then load and _register_loaded_plugin), with the + manifest re-read from disk, so the running set ends up as a restart + would build it. A plugin that fails to load stays out of the + rotation, as it would after a restart. + """ + commands, self._pending_plugin_reloads = self._pending_plugin_reloads, () + for command in commands: + try: + self._reload_plugin_for_command(command) + except Exception: # pylint: disable=broad-except + logger.exception("Plugin reload over the control socket failed") + command.fail(ControlErrorCode.INTERNAL, 'the display failed to reload it') + + def _reload_plugin_for_command(self, command: QueuedCommand) -> None: + args = command.args + if not isinstance(args, PluginReloadArgs): + command.fail(ControlErrorCode.INTERNAL, 'not a reload command') + return + plugin_id = args.plugin_id + if self.plugin_manager is None or plugin_id not in self.plugin_display_modes: + command.fail(ControlErrorCode.NOT_LOADED, f'{plugin_id} is not running') + return + if plugin_id in self._on_demand_loaded_plugins: + # Loaded only for an on-demand session (config says disabled); + # reloading would have to repeat that special load. + command.fail(ControlErrorCode.BUSY, + f'{plugin_id} is loaded only for on-demand; restart to reload it') + return + + previous_mode = self.current_display_mode + previous_order = {mode: i for i, mode in enumerate(self.available_modes)} + logger.info("Reloading plugin %s over the control socket", plugin_id) + self._unregister_plugin(plugin_id, action='Unloaded') + loaded = bool(self.plugin_manager.reload_plugin(plugin_id)) + modes: List[str] = [] + if loaded: + modes = list(self._register_loaded_plugin(plugin_id)) + # Registering appends; put its modes back where they were in the + # rotation (a mode the new version added goes last). + self.available_modes.sort( + key=lambda mode: previous_order.get(mode, len(previous_order))) + self._apply_plugin_rotation_order() + self._resync_mode_index_after_change(previous_mode) + if not loaded: + logger.error("Plugin %s did not load after its update; it is out of the " + "rotation until it loads", plugin_id) + command.fail(ControlErrorCode.FAILED, + f'{plugin_id} did not load; see the display log') + return + vegas = getattr(self, 'vegas_coordinator', None) + if vegas is not None: + try: + # Fetch its content again rather than scroll the old copy. + vegas.mark_plugin_updated(plugin_id) + except Exception: # pylint: disable=broad-except + logger.debug("Vegas did not take the reload of %s", plugin_id, exc_info=True) + manifest = (getattr(self.plugin_manager, 'plugin_manifests', None) or {}).get(plugin_id) + version = manifest.get('version') if isinstance(manifest, dict) else None + logger.info("Reloaded plugin %s (version %s, modes %s)", plugin_id, version, modes) + result: PluginReloadResult = { + 'plugin_id': plugin_id, 'reloaded': True, + 'version': version if isinstance(version, str) else None, + 'modes': modes, + } + command.succeed(dict(result)) def _poll_on_demand_requests(self) -> None: """Poll cache for new on-demand requests from external controllers.""" @@ -3076,6 +3268,12 @@ class DisplayController: if self._pending_plugin_reconcile and not self.on_demand_active: self._service_pending_reconcile() + # Plugin reloads from the control socket (a store update), + # here for the same reason: nothing of the plugin's is on the + # stack. The screen that was showing ended early for them. + if self._pending_plugin_reloads: + self._apply_pending_plugin_reloads() + if not self.available_modes: # Nothing to render yet. Re-check _pending_plugin_reconcile # every ~1s (rather than a long sleep) so enabling a plugin @@ -3177,6 +3375,10 @@ class DisplayController: # Scheduled off mid-iteration: blank the # panel now rather than render a screen. continue + if self._plugin_reload_pending: + # Reload first (top of the loop), then + # the ticker carries on. + continue if self._wifi_notice_pending(): # It yielded for a WiFi notice: the next # pass shows it, not a rotation screen @@ -3367,9 +3569,7 @@ class DisplayController: # update threads and the web UI are not starved of the GIL. time.sleep(_remaining if _remaining > 0 else 0.001) - if (self.current_display_mode != active_mode - or not self.is_display_active - or self._wifi_notice_pending()): + if self._screen_preempted(active_mode): logger.debug("Mode changed during high-FPS loop, breaking early") break @@ -3400,7 +3600,13 @@ class DisplayController: ) while True: - time.sleep(display_interval) + # Wakes for a control socket command and applies + # it at once, instead of up to a second later. + if self._wait_frame_interval(display_interval, active_mode): + logger.info("Mode changed during display loop from %s to %s, " + "breaking early", active_mode, + self.current_display_mode) + break self._tick_plugin_updates() elapsed = time.time() - start_time @@ -3432,9 +3638,7 @@ class DisplayController: self._service_pending_changes() self._check_live_takeover() - if (self.current_display_mode != active_mode - or not self.is_display_active - or self._wifi_notice_pending()): + if self._screen_preempted(active_mode): logger.info("Mode changed during display loop from %s to %s, breaking early", active_mode, self.current_display_mode) break @@ -3464,6 +3668,9 @@ class DisplayController: or not self.is_display_active or (not loop_completed and self._wifi_notice_pending())): continue + # A screen cut short for a plugin reload is over: the + # make-up dwell below returns at once, the rotation + # advances, and the next pass reloads before it draws. # Ensure we honour minimum duration when not dynamic and loop ended early if ( @@ -3792,9 +3999,10 @@ class DisplayController: self._plugin_accepts_display_mode.pop(plugin_id, None) return display_modes - def _unregister_plugin(self, plugin_id: str) -> None: + def _unregister_plugin(self, plugin_id: str, action: str = 'Disabled') -> None: """Remove a plugin's modes, config subscription and instance, then - unload it. Used by live disable hot-reload.""" + unload it. Used by live disable hot-reload, and by a reload + (``action`` names which in the log line).""" with self._plugin_modes_lock: modes = self.plugin_display_modes.pop(plugin_id, []) for mode in modes: @@ -3826,7 +4034,7 @@ class DisplayController: except Exception as e: logger.error("Error unloading plugin %s: %s", plugin_id, e, exc_info=True) - logger.info("Disabled plugin %s live (removed modes: %s)", plugin_id, modes) + logger.info("%s plugin %s live (removed modes: %s)", action, plugin_id, modes) def _enabled_set_changed(self, old_config: Dict[str, Any], new_config: Dict[str, Any]) -> bool: """True if any top-level section's ``enabled`` flag differs between two diff --git a/src/ipc/client.py b/src/ipc/client.py index a38d48fb..6f178d66 100644 --- a/src/ipc/client.py +++ b/src/ipc/client.py @@ -15,6 +15,7 @@ import uuid from typing import Any, Dict, List, Mapping, Optional, Sequence from src.ipc.contract import ( + AWAIT_SECONDS, MAX_MESSAGE_BYTES, PROTOCOL_VERSION, SUPPORTED_VERSIONS, @@ -188,6 +189,40 @@ def on_demand_status(*, timeout: float = DEFAULT_TIMEOUT_SECONDS, return request(Command.ON_DEMAND_STATUS, {}, timeout=timeout, paths=paths) +#: Headroom over the display's own wait for an awaited command, so its +#: ``pending`` answer arrives before the client gives up. +_AWAIT_MARGIN_SECONDS = 1.0 + + +def _awaited_timeout(cmd: str) -> float: + return AWAIT_SECONDS[cmd] + _AWAIT_MARGIN_SECONDS + + +def brightness_set(brightness: int, *, timeout: Optional[float] = None, + paths: Optional[Sequence[str]] = None) -> Dict[str, Any]: + """Set the panel's normal brightness now (transient: config.json is not + written). Returns the applied :class:`~src.ipc.contract.BrightnessResult`; + raises :class:`ControlError`. + """ + return request(Command.BRIGHTNESS_SET, {'brightness': brightness}, + timeout=_awaited_timeout(Command.BRIGHTNESS_SET) if timeout is None + else timeout, paths=paths) + + +def plugin_reload(plugin_id: str, *, timeout: Optional[float] = None, + paths: Optional[Sequence[str]] = None) -> Dict[str, Any]: + """Have the display reload a running plugin from disk. + + Returns :class:`~src.ipc.contract.PluginReloadResult` once the new code is + running. Raises :class:`ControlError`: ``not_loaded`` (not running it), + ``failed`` (the new version did not load), ``pending`` (not done in + time; it will still happen), or a transport reason. + """ + return request(Command.PLUGIN_RELOAD, {'plugin_id': plugin_id}, + timeout=_awaited_timeout(Command.PLUGIN_RELOAD) if timeout is None + else timeout, paths=paths) + + def ping(*, timeout: float = DEFAULT_TIMEOUT_SECONDS, paths: Optional[Sequence[str]] = None) -> Dict[str, Any]: return request(Command.PING, {}, timeout=timeout, paths=paths) diff --git a/src/ipc/contract.py b/src/ipc/contract.py index 3b1d4ad8..b563d17d 100644 --- a/src/ipc/contract.py +++ b/src/ipc/contract.py @@ -26,7 +26,14 @@ order. Commands that change what the panel shows are *acknowledged*, not completed: ``{"accepted": true, "request_id": ...}`` means the render thread has the command queued and will apply it at its next on-demand check. Its outcome is published the way it always was (``display_on_demand_state``, -later the state stream). +later the state stream). A few commands (:data:`AWAITED_COMMANDS`) are +answered only once the render thread has applied them, or with ``pending`` +when it has not within :data:`AWAIT_SECONDS`. + +New commands are added within a protocol version: a display that does not +know one answers ``unknown_command``, the client falls back, and ``hello`` +lists the commands a display knows. The version changes only when the +envelope or the meaning of an existing command changes. See docs/IPC_CONTROL_SOCKET.md for the full description. """ @@ -133,19 +140,42 @@ class Command: ON_DEMAND_START = 'on_demand.start' ON_DEMAND_STOP = 'on_demand.stop' ON_DEMAND_STATUS = 'on_demand.status' + BRIGHTNESS_SET = 'brightness.set' + PLUGIN_RELOAD = 'plugin.reload' #: Every command version 1 defines, in the order ``hello`` reports them. +#: ``brightness.set`` and ``plugin.reload`` came in stage 2, within version 1 +#: (see the module docstring on adding commands). COMMANDS: Tuple[str, ...] = ( Command.HELLO, Command.PING, Command.ON_DEMAND_START, Command.ON_DEMAND_STOP, Command.ON_DEMAND_STATUS, + Command.BRIGHTNESS_SET, + Command.PLUGIN_RELOAD, ) -#: Commands that are queued for the render thread and answered with an ack. -QUEUED_COMMANDS = frozenset({Command.ON_DEMAND_START, Command.ON_DEMAND_STOP}) +#: Commands that are queued for the render thread. +QUEUED_COMMANDS = frozenset({Command.ON_DEMAND_START, Command.ON_DEMAND_STOP, + Command.BRIGHTNESS_SET, Command.PLUGIN_RELOAD}) + +#: Queued commands whose answer waits for the render thread's outcome +#: instead of being an ack. The value is how long the display waits before +#: answering ``pending``; the command stays queued and is still applied. +#: A plugin reload first lets the current screen end (within a frame on a +#: scrolling screen, at once on a static one) and then imports the plugin, +#: which can take a few seconds on a slow board. +AWAIT_SECONDS: Dict[str, float] = { + Command.BRIGHTNESS_SET: 2.0, + Command.PLUGIN_RELOAD: 10.0, +} +AWAITED_COMMANDS = frozenset(AWAIT_SECONDS) + +#: Brightness, in percent, as the display's hardware setting takes it. +MIN_BRIGHTNESS = 0 +MAX_BRIGHTNESS = 100 class ErrorCode: @@ -159,6 +189,10 @@ class ErrorCode: BUSY = 'busy' # queue full / too many clients FORBIDDEN = 'forbidden' # peer credentials refused INTERNAL = 'internal' # a bug on the display side + # From the awaited commands (stage 2): + PENDING = 'pending' # accepted, not applied within AWAIT_SECONDS; still queued + NOT_LOADED = 'not_loaded' # plugin.reload: the display is not running that plugin + FAILED = 'failed' # the render thread tried, and it did not work class ProtocolError(Exception): @@ -398,7 +432,56 @@ class NoArgs: return cls() -CommandArgs = Union[HelloArgs, OnDemandStartArgs, OnDemandStopArgs, NoArgs] +@dataclass(frozen=True) +class BrightnessSetArgs: + """``brightness.set``: the panel's normal brightness, in percent, now. + + Transient: nothing is written to config.json, and the next config change + the display picks up (or a restart) goes back to the configured value. + The web interface sends it after saving the setting, so the two agree. + The dim schedule still applies on top, as it does to the saved value. + """ + brightness: int + + def to_dict(self) -> Dict[str, Any]: + return {'brightness': self.brightness} + + @classmethod + def from_dict(cls, args: Mapping[str, Any]) -> 'BrightnessSetArgs': + value = args.get('brightness') + if not _is_int(value) or not MIN_BRIGHTNESS <= value <= MAX_BRIGHTNESS: + raise ProtocolError(ErrorCode.INVALID_ARGS, + f'brightness must be an integer from {MIN_BRIGHTNESS} ' + f'to {MAX_BRIGHTNESS}') + return cls(brightness=value) + + +@dataclass(frozen=True) +class PluginReloadArgs: + """``plugin.reload``: load a running plugin again from disk. + + For a plugin the store has just updated. Only a plugin the display is + running can be reloaded (``not_loaded`` otherwise), so the id never + makes the display import anything it was not already running. + """ + plugin_id: str + + def to_dict(self) -> Dict[str, Any]: + return {'plugin_id': self.plugin_id} + + @classmethod + def from_dict(cls, args: Mapping[str, Any]) -> 'PluginReloadArgs': + plugin_id = _optional_name(args, 'plugin_id') + if plugin_id is None: + raise ProtocolError(ErrorCode.INVALID_ARGS, 'plugin_id is required') + return cls(plugin_id=plugin_id) + + +CommandArgs = Union[HelloArgs, OnDemandStartArgs, OnDemandStopArgs, NoArgs, + BrightnessSetArgs, PluginReloadArgs] + +#: The arguments of a command that goes on the render thread's queue. +QueuedArgs = Union[OnDemandStartArgs, OnDemandStopArgs, BrightnessSetArgs, PluginReloadArgs] _ARG_TYPES: Dict[str, Any] = { Command.HELLO: HelloArgs, @@ -406,6 +489,8 @@ _ARG_TYPES: Dict[str, Any] = { Command.ON_DEMAND_START: OnDemandStartArgs, Command.ON_DEMAND_STOP: OnDemandStopArgs, Command.ON_DEMAND_STATUS: NoArgs, + Command.BRIGHTNESS_SET: BrightnessSetArgs, + Command.PLUGIN_RELOAD: PluginReloadArgs, } @@ -457,6 +542,27 @@ class AckResult(TypedDict): queued: int +class BrightnessResult(TypedDict): + """``brightness.set``, once applied. + + ``panel_brightness`` is what the panel shows now: the dim schedule's + level while it dims, and unchanged while the schedule has the display + off (the new level applies when it comes back on). + """ + brightness: int + panel_brightness: int + dimmed: bool + display_active: bool + + +class PluginReloadResult(TypedDict): + """``plugin.reload``, once the plugin is running again.""" + plugin_id: str + reloaded: bool + version: Optional[str] + modes: List[str] + + def negotiate_version(client_versions: Tuple[int, ...]) -> Optional[int]: """The highest version both sides speak, or None.""" common = set(client_versions) & set(SUPPORTED_VERSIONS) diff --git a/src/ipc/server.py b/src/ipc/server.py index 7fc08e50..0240dc86 100644 --- a/src/ipc/server.py +++ b/src/ipc/server.py @@ -8,6 +8,12 @@ where it reads the file mailbox (``DisplayController._poll_on_demand_requests``) handing each command to the same code. Queries (``on_demand.status``) are answered from a snapshot callable the display provides. +The queue also wakes the render thread: :meth:`ControlServer.wait_for_command` +is what it waits on in place of a sleep, so a command lands within a frame on +every kind of screen. An awaited command (``brightness.set``, +``plugin.reload``) carries a :class:`CommandOutcome` that the render thread +fills in; its connection thread waits for that, bounded, before answering. + Robustness rules, because this runs inside the display process: * every connection has its own daemon thread, at most :data:`MAX_CLIENTS` at @@ -39,10 +45,12 @@ import stat import struct import threading import time -from dataclasses import dataclass -from typing import Any, Callable, Dict, FrozenSet, List, Mapping, Optional, Union +from dataclasses import dataclass, field +from typing import Any, Callable, Dict, FrozenSet, List, Mapping, Optional from src.ipc.contract import ( + AWAIT_SECONDS, + AWAITED_COMMANDS, COMMANDS, DEFAULT_SOCKET_DIR, DEFAULT_SOCKET_PATH, @@ -51,6 +59,7 @@ from src.ipc.contract import ( QUEUED_COMMANDS, SUPPORTED_VERSIONS, AckResult, + BrightnessSetArgs, Command, ErrorCode, FrameReader, @@ -58,7 +67,9 @@ from src.ipc.contract import ( HelloResult, OnDemandStartArgs, OnDemandStopArgs, + PluginReloadArgs, ProtocolError, + QueuedArgs, Request, Response, configured_socket_path, @@ -100,19 +111,75 @@ _LISTEN_BACKLOG = 64 # -- queued work --------------------------------------------------------------------- +class CommandOutcome: + """How an awaited command turned out, handed from the render thread back + to the connection thread that is waiting to answer. + + The render thread calls :meth:`succeed` or :meth:`fail` once; the first + call wins. The connection thread may have stopped waiting already (it + answered ``pending``), and then nobody reads it. + """ + + def __init__(self) -> None: + self._done = threading.Event() + self._lock = threading.Lock() + self.result: Optional[Dict[str, Any]] = None + self.error_code: Optional[str] = None + self.error_message = '' + + @property + def done(self) -> bool: + return self._done.is_set() + + def succeed(self, result: Mapping[str, Any]) -> None: + with self._lock: + if self._done.is_set(): + return + self.result = dict(result) + self._done.set() + + def fail(self, code: str, message: str) -> None: + with self._lock: + if self._done.is_set(): + return + self.error_code = code + self.error_message = message + self._done.set() + + def wait(self, timeout: float) -> bool: + return self._done.wait(timeout) + + @dataclass(frozen=True) class QueuedCommand: - """A command waiting for the render thread.""" + """A command waiting for the render thread. + + ``outcome`` is set for an awaited command (``AWAITED_COMMANDS``): the + render thread reports through it, and the client's answer waits for it. + """ request_id: str cmd: str - args: Union[OnDemandStartArgs, OnDemandStopArgs] + args: QueuedArgs received_at: float # time.time() when it was accepted peer_uid: Optional[int] = None + outcome: Optional[CommandOutcome] = field(default=None, compare=False, repr=False) def as_on_demand_request(self) -> Dict[str, Any]: """The mailbox-shaped payload the display's on-demand handler takes.""" + if not isinstance(self.args, (OnDemandStartArgs, OnDemandStopArgs)): + raise TypeError(f'{self.cmd} is not an on-demand command') return on_demand_request(self.request_id, self.args, self.received_at) + def succeed(self, result: Mapping[str, Any]) -> None: + """Report success to a waiting client (a no-op for an acked command).""" + if self.outcome is not None: + self.outcome.succeed(result) + + def fail(self, code: str, message: str) -> None: + """Report failure to a waiting client (a no-op for an acked command).""" + if self.outcome is not None: + self.outcome.fail(code, message) + # -- peer credentials ------------------------------------------------------------------ @@ -247,8 +314,12 @@ class ControlServer: max_clients: int = MAX_CLIENTS, io_timeout: float = IO_TIMEOUT_SECONDS, message_timeout: float = MESSAGE_TIMEOUT_SECONDS, idle_timeout: float = IDLE_TIMEOUT_SECONDS, - check_peer: bool = True): + check_peer: bool = True, + await_seconds: Optional[Mapping[str, float]] = None): self.path = path + self._await_seconds: Dict[str, float] = dict(AWAIT_SECONDS) + if await_seconds: + self._await_seconds.update(await_seconds) self._status_provider = status_provider self._group = group self._queue: 'queue.Queue[QueuedCommand]' = queue.Queue(maxsize=queue_size) @@ -418,6 +489,17 @@ class ControlServer: """Cheap check for queued commands, for the render thread's fast path.""" return self._pending.is_set() + def wait_for_command(self, timeout: float) -> bool: + """Block up to ``timeout`` seconds for a queued command; True if one is. + + The render thread waits here instead of sleeping, in the dwell and + on a static screen, so a command wakes it at once. It is a timed + wait on an Event: no polling, and nothing more than the sleep it + replaces when no command comes. The flag stays set until drain(), + so a caller that does not drain would return at once every time. + """ + return self._pending.wait(timeout) + def drain(self) -> List[QueuedCommand]: """Every queued command, oldest first. Called from the render thread.""" commands: List[QueuedCommand] = [] @@ -596,11 +678,13 @@ class ControlServer: v=request.v) return Response.success(request.id, self._status_provider(), v=request.v) - if request.cmd in QUEUED_COMMANDS and isinstance(args, (OnDemandStartArgs, - OnDemandStopArgs)): + if request.cmd in QUEUED_COMMANDS and isinstance(args, ( + OnDemandStartArgs, OnDemandStopArgs, BrightnessSetArgs, PluginReloadArgs)): + awaited = request.cmd in AWAITED_COMMANDS command = QueuedCommand(request_id=request.id, cmd=request.cmd, args=args, received_at=time.time(), - peer_uid=peer.uid if peer is not None else None) + peer_uid=peer.uid if peer is not None else None, + outcome=CommandOutcome() if awaited else None) try: self._queue.put_nowait(command) except queue.Full: @@ -610,15 +694,37 @@ class ControlServer: 'the display is not taking commands right now', v=request.v) self._pending.set() + logger.info("Control socket accepted %s %s", request.cmd, request.id) + if command.outcome is not None: + return self._await_outcome(request, command.outcome) ack: AckResult = {'accepted': True, 'request_id': request.id, 'queued': self._queue.qsize()} - logger.info("Control socket accepted %s %s", request.cmd, request.id) return Response.success(request.id, dict(ack), v=request.v) # A command in COMMANDS with no handler here is a bug in this module. return Response.failure(request.id, ErrorCode.INTERNAL, f'{request.cmd} is not implemented', v=request.v) + def _await_outcome(self, request: Request, outcome: CommandOutcome) -> Response: + """Answer an awaited command once the render thread has applied it. + + Waits on this connection's thread, never the render thread's. If the + render thread does not get to it in time, the answer is ``pending``: + the command stays queued and is still applied, so a client treats + that as "not known to be done" rather than as a refusal. + """ + timeout = self._await_seconds.get(request.cmd, 0.0) + if not outcome.wait(timeout): + logger.warning("Control socket: %s %s not applied within %.1fs; answering pending", + request.cmd, request.id, timeout) + return Response.failure(request.id, ErrorCode.PENDING, + f'accepted, but not applied within {timeout:g}s; ' + 'the display will still apply it', v=request.v) + if outcome.error_code is not None: + return Response.failure(request.id, outcome.error_code, outcome.error_message, + v=request.v) + return Response.success(request.id, outcome.result or {}, v=request.v) + def start_control_server(status_provider: Optional[StatusProvider] = None, cache_dir: Optional[str] = None, @@ -637,7 +743,7 @@ def start_control_server(status_provider: Optional[StatusProvider] = None, __all__ = [ - 'ControlServer', 'PeerCredentials', 'QueuedCommand', 'StatusProvider', + 'CommandOutcome', 'ControlServer', 'PeerCredentials', 'QueuedCommand', 'StatusProvider', 'peer_allowed', 'peer_credentials', 'process_groups', 'resolve_socket_group', 'server_socket_path', 'start_control_server', 'PROTOCOL_VERSION', ] diff --git a/src/plugin_system/plugin_catalog.py b/src/plugin_system/plugin_catalog.py index ddd6f231..8e0a25f0 100644 --- a/src/plugin_system/plugin_catalog.py +++ b/src/plugin_system/plugin_catalog.py @@ -238,7 +238,9 @@ def display_restart_required(action: str, plugin_enabled: bool, *, config carried over) is not picked up until a restart. - ``update``: the display keeps running the code it loaded until it restarts, if it runs the plugin at all -- only when it is enabled. - ``changed=False`` (already up to date) needs nothing. + ``changed=False`` (already up to date) needs nothing. The update route + first asks the display to reload it over the control socket + (``_reload_after_store_update``); this answer stands when it cannot. - ``uninstall``: removing the plugin's config section flips its enabled flag, and the reconcile unloads it. With ``preserve_config`` the flag stays, and an enabled plugin keeps running until a restart. diff --git a/src/vegas_mode/coordinator.py b/src/vegas_mode/coordinator.py index e8071be0..8eb0f0ab 100644 --- a/src/vegas_mode/coordinator.py +++ b/src/vegas_mode/coordinator.py @@ -152,6 +152,8 @@ class VegasModeCoordinator: # Interrupt checker for yielding control back to display controller self._interrupt_check: Optional[Callable[[], bool]] = None self._interrupt_check_interval: int = 10 # Check every N frames + # Checked every frame; True runs the interrupt check at once. + self._interrupt_urgent: Optional[Callable[[], bool]] = None # Plugin update callback — fired from a background thread inside the loop # so the main loop's _tick_plugin_updates() finds nothing due when Vegas @@ -226,7 +228,8 @@ class VegasModeCoordinator: def set_interrupt_checker( self, checker: Callable[[], bool], - check_interval: int = 10 + check_interval: int = 10, + urgent: Optional[Callable[[], bool]] = None, ) -> None: """ Set the callback for checking if Vegas should yield control. @@ -237,9 +240,25 @@ class VegasModeCoordinator: Args: checker: Callable that returns True if Vegas should yield check_interval: Check every N frames (default 10) + urgent: A cheap per-frame test; when it is True the checker + runs at this frame instead of waiting for the interval (the + display controller passes "a control socket command is + queued", so a command waits one frame, not ten) """ self._interrupt_check = checker self._interrupt_check_interval = max(1, check_interval) + self._interrupt_urgent = urgent + + def _interrupt_is_urgent(self) -> bool: + """The per-frame test set with ``urgent``; never raises.""" + urgent = getattr(self, '_interrupt_urgent', None) + if urgent is None: + return False + try: + return bool(urgent()) # pylint: disable=not-callable + except Exception: # pylint: disable=broad-except + logger.debug("Urgent interrupt test failed", exc_info=True) + return False def set_update_callback(self, callback: Callable[[], None]) -> None: """ @@ -709,7 +728,8 @@ class VegasModeCoordinator: frame_times.clear() if (self._interrupt_check and - frame_count % self._interrupt_check_interval == 0): + (frame_count % self._interrupt_check_interval == 0 + or self._interrupt_is_urgent())): try: if self._interrupt_check(): logger.debug( diff --git a/test/_run_loop_harness.py b/test/_run_loop_harness.py index 79f89eef..2e7f63d5 100644 --- a/test/_run_loop_harness.py +++ b/test/_run_loop_harness.py @@ -391,8 +391,9 @@ class FakeVegas: run_iteration() renders frames at 125 Hz on the fake clock for ``cycle`` seconds and returns True, or returns False as soon as the - interrupt checker (every 10 frames) or the live-priority checker (every - 0.25 s) asks it to yield -- the same cadence the real coordinator uses. + interrupt checker (every 10 frames, or at the next frame when the + ``urgent`` test says so) or the live-priority checker (every 0.25 s) + asks it to yield -- the same cadence the real coordinator uses. A live-priority pause is lifted by the next call, as in the real one. """ @@ -408,14 +409,16 @@ class FakeVegas: self.vegas_config = SimpleNamespace(live_in_ticker=live_in_ticker) self.render_pipeline = None self._interrupt: Optional[Callable[[], bool]] = None + self._urgent: Optional[Callable[[], bool]] = None self._live: Optional[Callable[[], Any]] = None self._paused_for_live = False def set_live_priority_checker(self, fn): self._live = fn - def set_interrupt_checker(self, fn, check_interval=10): + def set_interrupt_checker(self, fn, check_interval=10, urgent=None): self._interrupt = fn + self._urgent = urgent def apply_pending_config_if_idle(self): pass @@ -444,13 +447,68 @@ class FakeVegas: h.log("vegas-frame", None, quiet=True) clock.sleep(self.FRAME) frames += 1 - if self._interrupt and frames % self.INTERRUPT_EVERY == 0 and self._interrupt(): + due = frames % self.INTERRUPT_EVERY == 0 or (self._urgent and self._urgent()) + if self._interrupt and due and self._interrupt(): h.log("vegas-interrupt") return False if clock.now - start >= self.cycle: return True +class FakeControlServer: + """The ControlServer surface run() uses (src/ipc/server.py), on the fake clock. + + ``post()`` queues a real QueuedCommand when the clock reaches ``t``, as a + connection thread would. ``wait_for_command`` advances the clock to the + first of: a command being queued, or the timeout -- the timed Event wait + the real server does, without threads. + """ + + def __init__(self, harness: "RunLoopHarness"): + self._h = harness + self.queue: List[Any] = [] + self.posted: List[Any] = [] + self.closed = False + + @property + def has_pending(self) -> bool: + return bool(self.queue) + + def drain(self) -> List[Any]: + out, self.queue = self.queue, [] + return out + + def wait_for_command(self, timeout: float) -> bool: + clock = self._h.clock + target = clock.now + max(0.0, timeout) + while not self.queue: + nxt = clock._alarms[0][0] if clock._alarms else None + if nxt is None or nxt > target: + clock.sleep(target - clock.now) + return bool(self.queue) + clock.sleep(max(0.0, nxt - clock.now)) + return True + + def post(self, t: float, cmd: str, args: Dict[str, Any], request_id: str = "sock"): + from src.ipc.contract import AWAITED_COMMANDS, parse_args + from src.ipc.server import CommandOutcome, QueuedCommand + + command = QueuedCommand( + request_id=request_id, cmd=cmd, args=parse_args(cmd, args), # type: ignore[arg-type] + received_at=0.0, + outcome=CommandOutcome() if cmd in AWAITED_COMMANDS else None) + self.posted.append(command) + + def enqueue(): + self._h.log("socket", cmd) + self.queue.append(command) + self._h.clock.at(t, enqueue) + return command + + def close(self): + self.closed = True + + # --------------------------------------------------------------------------- # Harness # --------------------------------------------------------------------------- @@ -669,6 +727,12 @@ class RunLoopHarness: encoding="utf-8") self.clock.at(t, write) + def control_socket(self) -> FakeControlServer: + """Serve the control socket (a FakeControlServer) for this run.""" + server = FakeControlServer(self) + self.controller._control_server = server + return server + def enable_vegas(self, cycle: float = 30.0, live_in_ticker: bool = False) -> FakeVegas: """Install FakeVegas, wired up as _initialize_vegas_mode wires the real one.""" dc = self.controller @@ -676,7 +740,7 @@ class RunLoopHarness: vegas.set_live_priority_checker(dc._check_live_priority) vegas.set_interrupt_checker( lambda: dc._check_vegas_interrupt() or dc.sync_manager.is_follower_active(), - check_interval=10) + check_interval=10, urgent=dc._control_command_pending) dc.vegas_coordinator = vegas return vegas diff --git a/test/test_api_v3_brightness_socket.py b/test/test_api_v3_brightness_socket.py new file mode 100644 index 00000000..5c6047df --- /dev/null +++ b/test/test_api_v3_brightness_socket.py @@ -0,0 +1,108 @@ +"""POST /api/v3/config/main with a brightness: on the panel at once. + +The saved brightness used to reach the panel when the display's config +watcher next noticed config.json (it polls every 2 s) and the render thread +then applied it (every 0.25 s). After a successful save the route now also +sends ``brightness.set`` over the control socket (src/ipc), which the +display applies on its render thread at once. When the socket cannot carry +it, nothing changes: the watcher applies the saved value as before. The +response says which (``brightness_transport``), and why +(``brightness_socket_error``). +""" + +import json +import sys +from pathlib import Path +from unittest.mock import patch + +import pytest + +sys.path.insert(0, str(Path(__file__).parent.parent)) + +from test._api_v3_test_helpers import api_v3_client, api_v3_module # noqa: F401,E402 + +from src.ipc import client as control_client # noqa: E402 + +SET = "web_interface.blueprints.api_v3.control_client.brightness_set" + + +@pytest.fixture +def saved(api_v3_module, monkeypatch): + captured = {'ok': True} + api_v3_module.api_v3.config_manager.load_config.return_value = {} + + def fake_save(_manager, config, **_kwargs): + captured['config'] = config + return (True, '') if captured['ok'] else (False, 'disk full') + + monkeypatch.setattr(api_v3_module, '_save_config_atomic', fake_save) + return captured + + +def _post(client, body): + return client.post('/api/v3/config/main', data=json.dumps(body), + content_type='application/json') + + +@pytest.mark.parametrize('value', [55, '55'], ids=['number', 'form-string']) +def test_a_saved_brightness_goes_over_the_socket(api_v3_client, saved, value): + with patch(SET, return_value={'brightness': 55, 'panel_brightness': 55, + 'dimmed': False, 'display_active': True}) as send: + resp = _post(api_v3_client, {'brightness': value}) + assert resp.status_code == 200, resp.get_json() + send.assert_called_once_with(55) + body = resp.get_json() + assert body['brightness_transport'] == 'socket' + assert 'brightness_socket_error' not in body + # Saved as well: the socket's value is transient on the display. + assert saved['config']['display']['hardware']['brightness'] == 55 + + +@pytest.mark.parametrize('reason', ['no_socket', 'unknown_command', 'pending', 'failed', + 'timeout', 'busy']) +def test_without_the_socket_the_config_watcher_applies_it(api_v3_client, saved, reason): + with patch(SET, side_effect=control_client.ControlError(reason, 'x')): + resp = _post(api_v3_client, {'brightness': 40}) + assert resp.status_code == 200 + body = resp.get_json() + assert body['brightness_transport'] == 'config' + assert body['brightness_socket_error'] == reason + assert saved['config']['display']['hardware']['brightness'] == 40 + + +def test_a_client_bug_never_fails_the_save(api_v3_client, saved): + with patch(SET, side_effect=RuntimeError('bug')): + resp = _post(api_v3_client, {'brightness': 40}) + assert resp.status_code == 200 + assert resp.get_json()['brightness_socket_error'] == 'internal' + + +def test_the_test_suite_has_no_socket(api_v3_client, saved): + """LEDMATRIX_CONTROL_SOCKET=off (conftest): the real client says so.""" + body = _post(api_v3_client, {'brightness': 40}).get_json() + assert body['brightness_transport'] == 'config' + assert body['brightness_socket_error'] in ('disabled', 'unsupported') + + +def test_a_save_without_a_brightness_sends_nothing(api_v3_client, saved): + with patch(SET) as send: + resp = _post(api_v3_client, {'rows': 32}) + assert resp.status_code == 200 + send.assert_not_called() + assert 'brightness_transport' not in resp.get_json() + + +@pytest.mark.parametrize('body', [{'brightness': 0}, {'brightness': 101}, {'brightness': 'x'}]) +def test_a_refused_brightness_is_not_sent(api_v3_client, saved, body): + with patch(SET) as send: + resp = _post(api_v3_client, body) + assert resp.status_code == 400 + send.assert_not_called() + + +def test_a_failed_save_is_not_sent(api_v3_client, saved): + saved['ok'] = False + with patch(SET) as send: + resp = _post(api_v3_client, {'brightness': 40}) + assert resp.status_code == 500 + send.assert_not_called() diff --git a/test/test_ipc_contract.py b/test/test_ipc_contract.py index 57be30b1..3d7d3cd8 100644 --- a/test/test_ipc_contract.py +++ b/test/test_ipc_contract.py @@ -165,7 +165,9 @@ class TestOnDemandArgs: @pytest.mark.parametrize('cmd', c.COMMANDS) def test_every_command_has_an_argument_type(self, cmd): - args = {'plugin_id': 'p'} if cmd == Command.ON_DEMAND_START else {} + args = {Command.ON_DEMAND_START: {'plugin_id': 'p'}, + Command.PLUGIN_RELOAD: {'plugin_id': 'p'}, + Command.BRIGHTNESS_SET: {'brightness': 50}}.get(cmd, {}) c.parse_args(cmd, args) def test_hello_versions(self): diff --git a/test/test_ipc_display_stage2.py b/test/test_ipc_display_stage2.py new file mode 100644 index 00000000..33c7be84 --- /dev/null +++ b/test/test_ipc_display_stage2.py @@ -0,0 +1,368 @@ +"""DisplayController's side of control socket stage 2. + +* ``brightness.set`` is applied on the render thread at once, repainted, + and answered with what the panel shows -- including when the dim schedule + holds it lower, when the panel refuses it, and when the schedule has the + display off; +* ``plugin.reload`` is answered from the top of the loop pass, for every + outcome: reloaded (new instance registered, Vegas told), not running, + loaded only for on-demand, failed to load, or raised; +* the render thread wakes for a command in real time, not just on the fake + clock (test_run_loop_socket_wake.py): within milliseconds on a static + screen and in a dwell, while a command that does not end the screen + leaves the frame wait to run its course; +* Vegas checks the queue every frame and runs its interrupt check at once. +""" + +import threading +import time +from datetime import datetime, timedelta, timezone +from types import SimpleNamespace +from unittest.mock import MagicMock + +import pytest + +from src.ipc.contract import BrightnessSetArgs, Command, ErrorCode, PluginReloadArgs +from src.ipc.server import CommandOutcome, ControlServer, QueuedCommand + + +def _command(cmd, args): + return QueuedCommand(request_id='t1', cmd=cmd, args=args, received_at=time.time(), + outcome=CommandOutcome()) + + +def _brightness(value): + return _command(Command.BRIGHTNESS_SET, BrightnessSetArgs(value)) + + +def _reload(plugin_id): + return _command(Command.PLUGIN_RELOAD, PluginReloadArgs(plugin_id)) + + +@pytest.fixture +def dc(test_display_controller): + c = test_display_controller + c.display_manager.set_brightness = MagicMock(return_value=True) + c.display_manager.update_display = MagicMock() + c.config = {'timezone': 'UTC', 'display': {'hardware': {'brightness': 90}}} + c._normal_brightness = 90 + c.current_brightness = 90 + c.is_display_active = True + c._tz = None + return c + + +class TestBrightness: + def test_applied_now_and_repainted(self, dc): + cmd = _brightness(40) + dc._apply_control_brightness(cmd) + dc.display_manager.set_brightness.assert_called_once_with(40) + dc.display_manager.update_display.assert_called() + assert dc.current_brightness == 40 + assert cmd.outcome.result == {'brightness': 40, 'panel_brightness': 40, + 'dimmed': False, 'display_active': True} + + def test_the_dim_schedule_still_applies(self, dc): + now = datetime.now(timezone.utc) + dc.config['dim_schedule'] = { + 'enabled': True, 'dim_brightness': 20, + 'start_time': (now - timedelta(hours=1)).strftime('%H:%M'), + 'end_time': (now + timedelta(hours=1)).strftime('%H:%M')} + # A stale per-minute answer from before the change must not stick. + dc._dim_checked_minute = (now.hour, now.minute) + dc._cached_target_brightness = 90 + cmd = _brightness(70) + dc._apply_control_brightness(cmd) + assert cmd.outcome.result == {'brightness': 70, 'panel_brightness': 20, + 'dimmed': True, 'display_active': True} + assert dc._normal_brightness == 70 + + def test_a_refused_brightness_is_reported(self, dc): + dc.display_manager.set_brightness.return_value = False + cmd = _brightness(40) + dc._apply_control_brightness(cmd) + assert cmd.outcome.error_code == ErrorCode.FAILED + assert dc.current_brightness == 90 + + def test_scheduled_off_keeps_it_for_later(self, dc): + dc.is_display_active = False + cmd = _brightness(40) + dc._apply_control_brightness(cmd) + dc.display_manager.set_brightness.assert_not_called() + assert cmd.outcome.result['display_active'] is False + assert cmd.outcome.result['panel_brightness'] == 90 + # When the schedule turns the display back on, the new level is used. + dc.is_display_active = True + dc._apply_brightness_target() + dc.display_manager.set_brightness.assert_called_once_with(40) + + def test_the_next_config_reload_wins(self, dc): + dc._apply_control_brightness(_brightness(40)) + dc._refresh_config_cache({'display': {'hardware': {'brightness': 75}}}) + assert dc._normal_brightness == 75 + + +class FakeServer: + def __init__(self, *commands): + self.commands = list(commands) + + @property + def has_pending(self): + return bool(self.commands) + + def drain(self): + out, self.commands = self.commands, [] + return out + + def close(self): + pass + + +@pytest.fixture +def running(dc): + """dc running 'clock' (two modes) and 'weather', with a reloadable manager.""" + old = SimpleNamespace(modes=['clock_a', 'clock_b']) + new = SimpleNamespace(modes=['clock_a', 'clock_b']) + weather = SimpleNamespace(modes=['weather']) + pm = MagicMock() + pm.plugins = {'clock': old, 'weather': weather} + pm.plugin_manifests = {'clock': {'version': '1.0.0'}, 'weather': {}} + pm.get_plugin = lambda pid: pm.plugins.get(pid) + calls = [] + + def unload(pid): + calls.append(('unload', pid)) + pm.plugins.pop(pid, None) + return True + + def reload(pid): + calls.append(('reload', pid)) + pm.plugins[pid] = new + pm.plugin_manifests[pid] = {'version': '2.0.0'} + return True + + pm.unload_plugin = MagicMock(side_effect=unload) + pm.reload_plugin = MagicMock(side_effect=reload) + dc.plugin_manager = pm + dc.plugin_display_modes = {'clock': ['clock_a', 'clock_b'], 'weather': ['weather']} + dc.available_modes = ['clock_a', 'clock_b', 'weather'] + dc.plugin_modes = {'clock_a': old, 'clock_b': old, 'weather': weather} + dc.mode_to_plugin_id = {'clock_a': 'clock', 'clock_b': 'clock', 'weather': 'weather'} + dc.current_mode_index = 2 + dc.current_display_mode = 'weather' + dc.vegas_coordinator = MagicMock() + return SimpleNamespace(dc=dc, pm=pm, old=old, new=new, calls=calls) + + +class TestReload: + def _run(self, dc, *commands): + dc._control_server = FakeServer(*commands) + dc._poll_on_demand_requests() # the drain: queued, not applied + assert dc._plugin_reload_pending + dc._apply_pending_plugin_reloads() # the top of the next pass + assert not dc._plugin_reload_pending + + def test_reloaded(self, running): + dc = running.dc + cmd = _reload('clock') + self._run(dc, cmd) + assert running.calls == [('unload', 'clock'), ('reload', 'clock')] + assert cmd.outcome.result == {'plugin_id': 'clock', 'reloaded': True, + 'version': '2.0.0', 'modes': ['clock_a', 'clock_b']} + assert dc.plugin_modes['clock_a'] is running.new + # Its modes keep their place in the rotation. + assert dc.available_modes == ['clock_a', 'clock_b', 'weather'] + # The screen that was showing is still the current one. + assert dc.current_display_mode == 'weather' + assert dc.available_modes[dc.current_mode_index] == 'weather' + dc.vegas_coordinator.mark_plugin_updated.assert_called_once_with('clock') + + def test_not_running(self, running): + cmd = _reload('nope') + self._run(running.dc, cmd) + assert cmd.outcome.error_code == ErrorCode.NOT_LOADED + assert running.calls == [] + + def test_loaded_only_for_on_demand(self, running): + running.dc._on_demand_loaded_plugins = {'clock'} + cmd = _reload('clock') + self._run(running.dc, cmd) + assert cmd.outcome.error_code == ErrorCode.BUSY + assert running.calls == [] + + def test_does_not_load_again(self, running): + running.pm.reload_plugin = MagicMock(return_value=False) + cmd = _reload('clock') + self._run(running.dc, cmd) + assert cmd.outcome.error_code == ErrorCode.FAILED + assert running.dc.available_modes == ['weather'] + assert 'clock' not in running.dc.plugin_display_modes + + def test_raises(self, running): + running.pm.reload_plugin = MagicMock(side_effect=RuntimeError('boom')) + cmd = _reload('clock') + self._run(running.dc, cmd) + assert cmd.outcome.error_code == ErrorCode.INTERNAL + + def test_a_pending_reload_ends_vegas_and_the_dwell(self, running): + dc = running.dc + dc._pending_plugin_reloads = (_reload('clock'),) + dc._check_wifi_status_message = MagicMock(return_value=None) + dc._service_pending_changes = MagicMock() + assert dc._check_vegas_interrupt() is True + started = time.monotonic() + dc._sleep_with_plugin_updates(5) + assert time.monotonic() - started < 0.1 + + def test_a_failing_brightness_does_not_stop_the_rest(self, running): + dc = running.dc + bad = _brightness(40) + dc._apply_control_brightness = MagicMock(side_effect=RuntimeError('x')) + good = _reload('clock') + self._run(dc, bad, good) + assert bad.outcome.error_code == ErrorCode.INTERNAL + assert good.outcome.result['reloaded'] is True + + +class TestRealTimeWake: + """With a real ControlServer and real threads: what the fake clock tests + show, measured. A connection thread queues the command with + handle_line(), as it does for a socket client.""" + + @pytest.fixture + def server(self, dc, tmp_path): + server = ControlServer(str(tmp_path / 'unused.sock')) + dc._control_server = server + dc.current_display_mode = 'clock' + dc.on_demand_active = False + dc._activate_on_demand = MagicMock( + side_effect=lambda request: setattr(dc, 'current_display_mode', 'weather')) + dc._wifi_notice_pending = MagicMock(return_value=False) + dc._tick_plugin_updates = MagicMock() + dc._check_live_takeover = MagicMock() + dc.cache_manager.get = MagicMock(return_value=None) + dc.cache_manager.set = MagicMock() + return server + + @staticmethod + def _post_later(server, line, delay, stamps): + def run(): + time.sleep(delay) + stamps.append(time.monotonic()) + server.handle_line(line) + t = threading.Thread(target=run, daemon=True) + t.start() + return t + + @staticmethod + def _start_line(rid): + import json + return json.dumps({'v': 1, 'id': rid, 'cmd': Command.ON_DEMAND_START, + 'args': {'plugin_id': 'weather'}}).encode() + + def test_static_screen(self, dc, server): + latencies = [] + for i in range(10): + dc.current_display_mode = 'clock' + stamps = [] + t = self._post_later(server, self._start_line(f's{i}'), 0.05, stamps) + ended = dc._wait_frame_interval(1.0, 'clock') + woke = time.monotonic() + t.join() + assert ended is True + latencies.append(woke - stamps[0]) + latencies.sort() + print(f"static-screen wake latency: median {latencies[5] * 1000:.2f} ms, " + f"max {latencies[-1] * 1000:.2f} ms") + assert latencies[-1] < 0.05, latencies + + def test_dwell(self, dc, server): + latencies = [] + for i in range(10): + dc.current_display_mode = 'clock' + stamps = [] + t = self._post_later(server, self._start_line(f'd{i}'), 0.05, stamps) + dc._sleep_with_plugin_updates(5.0) + woke = time.monotonic() + t.join() + assert dc.current_display_mode == 'weather' + latencies.append(woke - stamps[0]) + latencies.sort() + print(f"dwell wake latency: median {latencies[5] * 1000:.2f} ms, " + f"max {latencies[-1] * 1000:.2f} ms") + assert latencies[-1] < 0.05, latencies + + def test_a_brightness_does_not_cut_the_frame_wait_short(self, dc, server): + import json + stamps = [] + line = json.dumps({'v': 1, 'id': 'b1', 'cmd': Command.BRIGHTNESS_SET, + 'args': {'brightness': 33}}).encode() + started = time.monotonic() + t = self._post_later(server, line, 0.1, stamps) + assert dc._wait_frame_interval(0.5, 'clock') is False + assert time.monotonic() - started >= 0.49 + t.join() + dc.display_manager.set_brightness.assert_called_once_with(33) + + def test_without_a_socket_it_is_a_plain_sleep(self, dc): + dc._control_server = None + started = time.monotonic() + assert dc._wait_frame_interval(0.2, 'clock') is False + assert time.monotonic() - started >= 0.19 + + +class TestVegasChecksEveryFrame: + def _coordinator(self, urgent_from): + from src.vegas_mode.config import VegasModeConfig + from src.vegas_mode.coordinator import VegasModeCoordinator + + coord = VegasModeCoordinator.__new__(VegasModeCoordinator) + coord.vegas_config = VegasModeConfig.from_config({'display': {'vegas_scroll': { + 'enabled': True, 'max_cycle_duration': 60, 'continuous_scroll': True}}}) + coord.render_pipeline = MagicMock() + coord.render_pipeline.frame_interval = 0.0 + coord.render_pipeline.target_fps = 90 + coord.display_manager = MagicMock() + coord.stats = {'cycles_completed': 0, 'interruptions': 0} + coord._state_lock = threading.Lock() + coord._is_active = True + coord._is_paused = False + coord._should_stop = False + coord._fps_last_health_log = 0.0 + coord._fps_was_degraded = False + coord._live_priority_check = None + coord._live_priority_active = False + coord._update_callback = None + coord._update_tick_running = False + coord._check_static_plugin_trigger = lambda: None + frames = [0] + + def run_frame(): + frames[0] += 1 + return True + coord.run_frame = run_frame + checked_at = [] + + def checker(): + checked_at.append(frames[0]) + return frames[0] >= urgent_from + coord.set_interrupt_checker(checker, check_interval=10, + urgent=lambda: frames[0] >= urgent_from) + return coord, checked_at + + def test_an_urgent_frame_runs_the_check_at_once(self): + coord, checked_at = self._coordinator(urgent_from=3) + assert coord.run_iteration() is False + assert checked_at == [3] + + def test_without_urgency_the_interval_holds(self): + coord, checked_at = self._coordinator(urgent_from=25) + coord._interrupt_urgent = None + assert coord.run_iteration() is False + assert checked_at == [10, 20, 30] + + def test_a_raising_urgency_test_is_ignored(self): + coord, checked_at = self._coordinator(urgent_from=10) + coord._interrupt_urgent = MagicMock(side_effect=RuntimeError('x')) + assert coord.run_iteration() is False + assert checked_at == [10] diff --git a/test/test_ipc_stage2.py b/test/test_ipc_stage2.py new file mode 100644 index 00000000..74d4aad8 --- /dev/null +++ b/test/test_ipc_stage2.py @@ -0,0 +1,299 @@ +"""Control socket stage 2 (src/ipc): waking the render thread, and the +awaited commands ``brightness.set`` and ``plugin.reload``. + +What these pin, on every platform unless marked: + +* the new commands' arguments refuse what the display could not act on; +* an awaited command is answered with the render thread's outcome -- its + result, its error code, or ``pending`` when the render thread did not get + to it in time (the command stays queued and is still applied); +* ``wait_for_command`` returns as soon as a command is queued, and otherwise + sleeps for the whole timeout without spinning; +* over a real socket (Linux/macOS), the client waits for that outcome and + maps each error code to its ``ControlError`` reason. +""" + +import json +import os +import threading +import time + +import pytest + +from src.ipc import client +from src.ipc import contract as c +from src.ipc.contract import ( + BrightnessSetArgs, Command, ErrorCode, PluginReloadArgs, ProtocolError, +) +from src.ipc.server import CommandOutcome, ControlServer + +needs_unix_sockets = pytest.mark.skipif(not c.socket_supported(), + reason='AF_UNIX sockets are Linux/macOS only') + + +def _req(cmd, args=None, rid='r1'): + return json.dumps({'v': 1, 'id': rid, 'cmd': cmd, 'args': args or {}}).encode() + + +def _server(tmp_path, **kw): + kw.setdefault('await_seconds', {Command.BRIGHTNESS_SET: 2.0, Command.PLUGIN_RELOAD: 2.0}) + return ControlServer(str(tmp_path / 'unused.sock'), **kw) + + +class RenderThread: + """Stands in for the display's render thread: waits on the server the + way DisplayController does and settles each command with ``handler``.""" + + def __init__(self, server, handler): + self.server = server + self.handler = handler + self.stop = threading.Event() + self.woke_at = [] + self.thread = threading.Thread(target=self._run, daemon=True) + + def _run(self): + while not self.stop.is_set(): + if self.server.wait_for_command(0.05): + self.woke_at.append(time.monotonic()) + for command in self.server.drain(): + self.handler(command) + + def __enter__(self): + self.thread.start() + return self + + def __exit__(self, *exc): + self.stop.set() + self.thread.join(2) + + +class TestArguments: + @pytest.mark.parametrize('value', [0, 1, 55, 100]) + def test_brightness_in_range(self, value): + assert BrightnessSetArgs.from_dict({'brightness': value}).brightness == value + + @pytest.mark.parametrize('value', [-1, 101, 50.0, '50', True, None, [50]]) + def test_brightness_refused(self, value): + with pytest.raises(ProtocolError) as e: + BrightnessSetArgs.from_dict({'brightness': value}) + assert e.value.code == ErrorCode.INVALID_ARGS + + def test_reload_needs_a_plugin_id(self): + assert PluginReloadArgs.from_dict({'plugin_id': 'clock'}).plugin_id == 'clock' + for bad in ({}, {'plugin_id': ''}, {'plugin_id': 5}, {'plugin_id': 'x' * 200}, + {'plugin_id': 'a\nb'}): + with pytest.raises(ProtocolError) as e: + PluginReloadArgs.from_dict(bad) + assert e.value.code == ErrorCode.INVALID_ARGS + + def test_the_new_commands_are_queued_and_awaited(self): + for cmd in (Command.BRIGHTNESS_SET, Command.PLUGIN_RELOAD): + assert cmd in c.COMMANDS and cmd in c.QUEUED_COMMANDS and cmd in c.AWAITED_COMMANDS + # On-demand stays an ack: its outcome is published, not awaited. + assert Command.ON_DEMAND_START not in c.AWAITED_COMMANDS + # Additive within version 1 (see the contract's docstring). + assert c.SUPPORTED_VERSIONS == (1,) + + def test_the_client_outwaits_the_display(self): + """The client's timeout must cover the display's wait, or a slow + reload would be reported as a timeout instead of ``pending``.""" + for cmd in c.AWAITED_COMMANDS: + assert client._awaited_timeout(cmd) > c.AWAIT_SECONDS[cmd] + + +class TestOutcome: + def test_first_settlement_wins(self): + outcome = CommandOutcome() + outcome.succeed({'a': 1}) + outcome.fail(ErrorCode.FAILED, 'late') + assert outcome.done and outcome.result == {'a': 1} and outcome.error_code is None + + def test_wait_times_out(self): + assert CommandOutcome().wait(0.01) is False + + +class TestAwaitedCommands: + def test_brightness_answers_with_the_render_thread_result(self, tmp_path): + server = _server(tmp_path) + + def apply(command): + assert command.cmd == Command.BRIGHTNESS_SET + command.succeed({'brightness': command.args.brightness, 'panel_brightness': 40, + 'dimmed': False, 'display_active': True}) + + with RenderThread(server, apply): + resp = server.handle_line(_req(Command.BRIGHTNESS_SET, {'brightness': 40})) + assert resp.ok, resp.error + assert resp.result == {'brightness': 40, 'panel_brightness': 40, + 'dimmed': False, 'display_active': True} + + @pytest.mark.parametrize('code', [ErrorCode.NOT_LOADED, ErrorCode.FAILED, ErrorCode.BUSY]) + def test_a_failure_carries_its_code(self, tmp_path, code): + server = _server(tmp_path) + with RenderThread(server, lambda command: command.fail(code, 'no')): + resp = server.handle_line(_req(Command.PLUGIN_RELOAD, {'plugin_id': 'clock'})) + assert not resp.ok and resp.error.code == code and resp.id == 'r1' + + def test_no_render_thread_means_pending_and_the_command_stays_queued(self, tmp_path): + server = _server(tmp_path, await_seconds={Command.PLUGIN_RELOAD: 0.05}) + started = time.monotonic() + resp = server.handle_line(_req(Command.PLUGIN_RELOAD, {'plugin_id': 'clock'})) + assert time.monotonic() - started < 1.0 + assert not resp.ok and resp.error.code == ErrorCode.PENDING + # Still applied when the render thread gets there; nobody reads it. + queued = server.drain() + assert [(q.cmd, q.args.plugin_id) for q in queued] == [(Command.PLUGIN_RELOAD, 'clock')] + queued[0].succeed({'reloaded': True}) + + def test_bad_arguments_are_refused_before_queueing(self, tmp_path): + server = _server(tmp_path) + resp = server.handle_line(_req(Command.BRIGHTNESS_SET, {'brightness': 400})) + assert not resp.ok and resp.error.code == ErrorCode.INVALID_ARGS + assert not server.has_pending and server.drain() == [] + + def test_a_full_queue_is_busy_without_waiting(self, tmp_path): + server = _server(tmp_path, queue_size=1) + server.handle_line(_req(Command.ON_DEMAND_STOP, rid='fill')) + started = time.monotonic() + resp = server.handle_line(_req(Command.BRIGHTNESS_SET, {'brightness': 10})) + assert not resp.ok and resp.error.code == ErrorCode.BUSY + assert time.monotonic() - started < 0.5 + + def test_on_demand_is_still_acked_without_waiting(self, tmp_path): + server = _server(tmp_path) + resp = server.handle_line(_req(Command.ON_DEMAND_STOP)) + assert resp.ok and resp.result['accepted'] is True + (queued,) = server.drain() + assert queued.outcome is None + queued.succeed({}) # a no-op for an acked command + + def test_an_on_demand_payload_is_only_built_for_on_demand(self, tmp_path): + server = _server(tmp_path, await_seconds={Command.BRIGHTNESS_SET: 0.01}) + server.handle_line(_req(Command.BRIGHTNESS_SET, {'brightness': 10})) + (queued,) = server.drain() + with pytest.raises(TypeError): + queued.as_on_demand_request() + + +class TestWake: + """``wait_for_command`` is what the render thread waits on in place of a + sleep: it must return as soon as a command is queued, and not before.""" + + def test_returns_at_once_when_a_command_is_queued(self, tmp_path): + server = _server(tmp_path) + latencies = [] + for _ in range(20): + queued_at = [] + + def post(): + time.sleep(0.02) + queued_at.append(time.monotonic()) + server.handle_line(_req(Command.ON_DEMAND_STOP)) + + t = threading.Thread(target=post) + t.start() + assert server.wait_for_command(2.0) is True + woke = time.monotonic() + t.join() + latencies.append(woke - queued_at[0]) + server.drain() + latencies.sort() + # Typically well under a millisecond; 50 ms is the budget for a + # static screen, and a loaded CI box needs the slack. + assert latencies[len(latencies) // 2] < 0.01, latencies + assert latencies[-1] < 0.05, latencies + + def test_sleeps_the_whole_timeout_when_nothing_comes(self, tmp_path): + server = _server(tmp_path) + started = time.monotonic() + cpu = time.process_time() + assert server.wait_for_command(0.3) is False + assert time.monotonic() - started >= 0.29 + # A timed wait, not a polling loop. + assert time.process_time() - cpu < 0.05 + + def test_drain_rearms_it(self, tmp_path): + server = _server(tmp_path) + server.handle_line(_req(Command.ON_DEMAND_STOP)) + assert server.wait_for_command(0) is True + server.drain() + assert server.wait_for_command(0.01) is False + + +# -- the real socket ------------------------------------------------------------------ + +@pytest.fixture +def sock_path(): + import shutil + import tempfile + d = tempfile.mkdtemp(prefix='lmipc2-') + yield os.path.join(d, 'control.sock') + shutil.rmtree(d, ignore_errors=True) + + +@pytest.fixture +def live(sock_path): + servers = [] + + def make(**kwargs): + s = ControlServer(sock_path, status_provider=lambda: {}, **kwargs) + assert s.start() + servers.append(s) + return s + + yield make + for s in servers: + s.close() + + +@needs_unix_sockets +class TestOverTheSocket: + def test_brightness_round_trip(self, live, sock_path): + server = live() + + def apply(command): + command.succeed({'brightness': command.args.brightness, 'panel_brightness': 30, + 'dimmed': True, 'display_active': True}) + + with RenderThread(server, apply): + started = time.monotonic() + result = client.brightness_set(70, paths=[sock_path]) + elapsed = time.monotonic() - started + assert result == {'brightness': 70, 'panel_brightness': 30, 'dimmed': True, + 'display_active': True} + assert elapsed < 0.5, elapsed + + def test_reload_round_trip(self, live, sock_path): + server = live() + + def apply(command): + command.succeed({'plugin_id': command.args.plugin_id, 'reloaded': True, + 'version': '2.0.0', 'modes': ['clock']}) + + with RenderThread(server, apply): + result = client.plugin_reload('clock', paths=[sock_path]) + assert result['reloaded'] is True and result['version'] == '2.0.0' + + @pytest.mark.parametrize('code', [ErrorCode.NOT_LOADED, ErrorCode.FAILED]) + def test_reload_errors_reach_the_client(self, live, sock_path, code): + server = live() + with RenderThread(server, lambda command: command.fail(code, 'x')): + with pytest.raises(client.ControlError) as e: + client.plugin_reload('clock', paths=[sock_path]) + assert e.value.reason == code + + def test_pending_reaches_the_client_before_its_own_timeout(self, live, sock_path): + live(await_seconds={Command.PLUGIN_RELOAD: 0.1}) + with pytest.raises(client.ControlError) as e: + client.plugin_reload('clock', paths=[sock_path], timeout=2.0) + assert e.value.reason == ErrorCode.PENDING + + def test_invalid_arguments_never_leave_the_client(self, sock_path): + with pytest.raises(client.ControlError) as e: + client.brightness_set(500, paths=[sock_path]) + assert e.value.reason == 'invalid_request' + + def test_hello_lists_the_new_commands(self, live, sock_path): + live() + commands = client.hello(paths=[sock_path])['commands'] + assert Command.BRIGHTNESS_SET in commands and Command.PLUGIN_RELOAD in commands diff --git a/test/test_run_loop_socket_wake.py b/test/test_run_loop_socket_wake.py new file mode 100644 index 00000000..96bb52ae --- /dev/null +++ b/test/test_run_loop_socket_wake.py @@ -0,0 +1,179 @@ +"""A control socket command wakes the run loop (control socket stage 2). + +Stage 1's socket was no faster than the file mailbox: a queued command +waited for the same polls the mailbox does -- the 1 s frame sleep of a static +screen, the 0.25 s dwell tick, and Vegas's interrupt check every 10 frames +(about 0.4 s at the 24 fps a Pi 4 manages). Now the render thread waits on +the socket's queue instead of sleeping, and Vegas checks the queue every +frame, so a command is applied: + +* at once on a static screen and in a dwell (here: at the instant it is + queued, on the fake clock); +* at the next frame in Vegas (8 ms here, at 125 fps). + +Each test runs the real run() loop on the fake clock of +test/_run_loop_harness.py and compares the socket with the mailbox for the +same request. The latencies are the fake clock's, so they are exact. +""" + +import os + +import pytest + +os.environ.setdefault("EMULATOR", "true") + +from src.ipc.contract import Command, ErrorCode # noqa: E402 +from test._run_loop_harness import FakePlugin, RunLoopHarness # noqa: E402 + +POSTED = 10.3 + + +def _event_time(trace, kind): + return next(e[0] for e in trace["events"] if e[1] == kind) + + +def _on_demand_latency(tmp_path, build, via, horizon=40): + h = RunLoopHarness(tmp_path, horizon=horizon) + build(h) + if via == "socket": + h.control_socket().post(POSTED, Command.ON_DEMAND_START, {"plugin_id": "weather"}) + else: + h.on_demand_request(POSTED, "mb1", plugin_id="weather") + trace = h.run() + return round(_event_time(trace, "on-demand-start") - POSTED, 3), trace + + +def _static(h): + h.add_plugin(FakePlugin("clock", ["clock"], duration=30)) + h.add_plugin(FakePlugin("weather", ["weather"], duration=30)) + + +def _dwell(h): + # Content on the first frame only: the 1 s loop ends at t=1 and the rest + # of the 30 s is made up in the dwell sleep (_sleep_with_plugin_updates). + h.add_plugin(FakePlugin("clock", ["clock"], duration=30, first_frame_only=True)) + h.add_plugin(FakePlugin("weather", ["weather"], duration=30)) + + +def _vegas(h): + _static(h) + h.enable_vegas(cycle=30) + + +@pytest.mark.parametrize("build, mailbox_latency", [ + (_static, 0.7), # the next 1 s frame, at t=11 + (_dwell, 0.2), # the next 0.25 s tick, at t=10.5 + (_vegas, 0.02), # the next 10-frame check: 80 ms at 125 fps, ~0.4 s on a Pi 4 +], ids=["static-screen", "dwell", "vegas"]) +def test_a_socket_command_lands_at_once(tmp_path, build, mailbox_latency): + (tmp_path / "s").mkdir() + (tmp_path / "m").mkdir() + socket_latency, trace = _on_demand_latency(tmp_path / "s", build, "socket") + via_mailbox, _ = _on_demand_latency(tmp_path / "m", build, "mailbox") + + assert via_mailbox == pytest.approx(mailbox_latency, abs=0.002) + if build is _vegas: + # One 8 ms frame: Vegas checks the queue every frame now. + assert socket_latency <= 0.008 + else: + assert socket_latency == 0.0 + # And the requested plugin is what the panel shows next, from then on. + row = next(r for r in trace["screens"] if r[1] == "weather") + assert row[0] == pytest.approx(POSTED + socket_latency, abs=0.001) + + +def test_a_brightness_lands_at_once_without_ending_the_screen(tmp_path): + h = RunLoopHarness(tmp_path, horizon=40) + _static(h) + command = h.control_socket().post(POSTED, Command.BRIGHTNESS_SET, {"brightness": 40}) + trace = h.run() + + assert [e for e in trace["events"] if e[1] == "brightness"] == [[POSTED, "brightness", 40]] + # The clock screen runs its full 30 s and the next starts on time. + assert trace["screens"][0][:4] == [0.0, "clock", 30.0, "duration"] + assert trace["screens"][1][:2] == [30.0, "weather"] + assert command.outcome.done and command.outcome.result["panel_brightness"] == 40 + + +def test_the_static_frame_cadence_is_unchanged_by_a_command(tmp_path): + """The wait resumes after a command that does not end the screen, so the + plugin is still drawn once a second, not once more at the command.""" + h = RunLoopHarness(tmp_path, horizon=31) + _static(h) + quiet = h.run() + + (tmp_path / "b").mkdir() + h2 = RunLoopHarness(tmp_path / "b", horizon=31) + _static(h2) + h2.control_socket().post(POSTED, Command.BRIGHTNESS_SET, {"brightness": 40}) + busy = h2.run() + assert busy["screens"][0] == quiet["screens"][0] + + +def _with_reloadable(h, version="2.0.0", loads=True): + new = FakePlugin("clock", ["clock"], duration=30) + + def reload_plugin(plugin_id): + h.log("reload", plugin_id) + if not loads: + return False + new._h = h + h.pm.plugins[plugin_id] = new + h.pm.plugin_manifests[plugin_id] = {"version": version, "display_modes": ["clock"]} + return True + + h.pm.reload_plugin = reload_plugin + return new + + +def test_a_plugin_reload_ends_the_screen_and_reloads_before_the_next(tmp_path): + h = RunLoopHarness(tmp_path, horizon=70) + _static(h) + new = _with_reloadable(h) + command = h.control_socket().post(POSTED, Command.PLUGIN_RELOAD, {"plugin_id": "clock"}) + trace = h.run() + + # Reloaded the moment it was asked for: the clock screen ends then, and + # the rotation moves on. + assert _event_time(trace, "reload") == POSTED + assert trace["screens"][0][:3] == [0.0, "clock", POSTED] + assert trace["screens"][1][:2] == [POSTED, "weather"] + assert command.outcome.result == {"plugin_id": "clock", "reloaded": True, + "version": "2.0.0", "modes": ["clock"]} + # The next clock screen is drawn by the reloaded instance. + assert h.controller.plugin_modes["clock"] is new + assert trace["screens"][2][1] == "clock" + + +def test_a_plugin_reload_during_vegas_yields_at_the_next_frame(tmp_path): + h = RunLoopHarness(tmp_path, horizon=70) + _vegas(h) + _with_reloadable(h) + command = h.control_socket().post(POSTED, Command.PLUGIN_RELOAD, {"plugin_id": "clock"}) + trace = h.run() + + assert 0 <= _event_time(trace, "reload") - POSTED <= 0.008 + assert command.outcome.result["reloaded"] is True + # The ticker resumes after the reload; no rotation screen in between. + after = [r for r in trace["screens"] if r[0] > POSTED] + assert after and after[0][1] == "" + + +def test_a_plugin_that_does_not_load_again_is_reported(tmp_path): + h = RunLoopHarness(tmp_path, horizon=40) + _static(h) + _with_reloadable(h, loads=False) + command = h.control_socket().post(POSTED, Command.PLUGIN_RELOAD, {"plugin_id": "clock"}) + h.run() + assert command.outcome.error_code == ErrorCode.FAILED + assert "clock" not in h.controller.available_modes + + +def test_reloading_a_plugin_that_is_not_running(tmp_path): + h = RunLoopHarness(tmp_path, horizon=40) + _static(h) + command = h.control_socket().post(POSTED, Command.PLUGIN_RELOAD, {"plugin_id": "nope"}) + trace = h.run() + assert command.outcome.error_code == ErrorCode.NOT_LOADED + # It still cost the clock screen nothing but its end: rotation goes on. + assert [r[1] for r in trace["screens"][:2]] == ["clock", "weather"] diff --git a/test/web_interface/test_web_process_runs_no_plugin_code.py b/test/web_interface/test_web_process_runs_no_plugin_code.py index 9956f02e..53841f25 100644 --- a/test/web_interface/test_web_process_runs_no_plugin_code.py +++ b/test/web_interface/test_web_process_runs_no_plugin_code.py @@ -382,6 +382,95 @@ class TestStoreOperationsReachTheDisplay: display_restart_required("reinstall", True) +RELOAD = "web_interface.blueprints.api_v3.control_client.plugin_reload" + + +class TestAnUpdateIsReloadedByTheDisplay: + """Control socket stage 2: instead of asking for a restart, the update + route asks the running display to reload the plugin (``plugin.reload``), + and only falls back to ``restart_required`` when that does not work. The + web process still runs none of the plugin's code.""" + + def test_reloaded_over_the_socket(self, web): + from unittest.mock import patch + _bump_version(web, "2.0.0") + with patch(RELOAD, return_value={"plugin_id": PLUGIN_ID, "reloaded": True, + "version": "2.0.0", "modes": ["tripwire"]}) as reload: + body = web.post("/api/v3/plugins/update", {"plugin_id": PLUGIN_ID}) + reload.assert_called_once_with(PLUGIN_ID) + assert body["restart_required"] is False + assert body["reloaded"] is True and body["reloaded_version"] == "2.0.0" + assert "restart_message" not in body and "reload_error" not in body + assert "running the new version" in body["message"] + assert web.ran() == [] + + @pytest.mark.parametrize("reason", [ + "no_socket", # display stopped, or predates the socket + "unknown_command", # a stage-1 display + "not_loaded", "failed", "pending", "busy", "timeout", "refused", + ]) + def test_any_failure_falls_back_to_a_restart(self, web, reason): + from unittest.mock import patch + from src.ipc import client as control_client + _bump_version(web, "2.0.0") + with patch(RELOAD, side_effect=control_client.ControlError(reason, "x")): + body = web.post("/api/v3/plugins/update", {"plugin_id": PLUGIN_ID}) + assert body["restart_required"] is True + assert "restart the display" in body["restart_message"] + assert body["reload_error"] == reason + assert "reloaded" not in body + + def test_a_client_bug_falls_back_too(self, web): + from unittest.mock import patch + _bump_version(web, "2.0.0") + with patch(RELOAD, side_effect=RuntimeError("bug")): + body = web.post("/api/v3/plugins/update", {"plugin_id": PLUGIN_ID}) + assert body["restart_required"] is True and body["reload_error"] == "internal" + + def test_without_a_socket_it_is_the_restart_banner_as_before(self, web): + # The test suite runs with LEDMATRIX_CONTROL_SOCKET=off (conftest). + _bump_version(web, "2.0.0") + body = web.post("/api/v3/plugins/update", {"plugin_id": PLUGIN_ID}) + assert body["restart_required"] is True + # 'unsupported' on Windows, which has no Unix sockets at all. + assert body["reload_error"] in ("disabled", "unsupported") + + def test_an_unchanged_plugin_asks_nothing(self, web): + from unittest.mock import patch + with patch(RELOAD) as reload: + body = web.post("/api/v3/plugins/update", {"plugin_id": PLUGIN_ID}) + assert body["data"]["update_status"] == "up_to_date" + assert body["restart_required"] is False + reload.assert_not_called() + + def test_a_disabled_plugin_asks_nothing(self, disabled_web): + from unittest.mock import patch + _bump_version(disabled_web, "2.0.0") + with patch(RELOAD) as reload: + body = disabled_web.post("/api/v3/plugins/update", {"plugin_id": PLUGIN_ID}) + assert body["restart_required"] is False + reload.assert_not_called() + + def test_the_display_is_asked_by_the_manifest_id(self, web): + """A store id can differ from the id the display runs the plugin + under (a registry alias: weather / ledmatrix-weather). The config + section and the reload both use the manifest's.""" + from unittest.mock import patch + manifest_id = "ledmatrix-alias" + _write_plugin(web.plugins_dir, web.log, plugin_id=manifest_id) + web.config_file.write_text(json.dumps({manifest_id: {"enabled": True}}), + encoding="utf-8") + + def update(plugin_id): + _write_plugin(web.plugins_dir, web.log, version="2.0.0", plugin_id=manifest_id) + return True + web.api.plugin_store_manager.update_plugin.side_effect = update + with patch(RELOAD, return_value={"reloaded": True, "version": "2.0.0"}) as reload: + body = web.post("/api/v3/plugins/update", {"plugin_id": "alias"}) + reload.assert_called_once_with(manifest_id) + assert body["restart_required"] is False + + class TestCatalogReadsWhatIsInstalled: """Real plugin directories: the repository's plugin-repos/ and fixtures.""" diff --git a/web_interface/blueprints/api_v3/__init__.py b/web_interface/blueprints/api_v3/__init__.py index 0a24f064..ca9a65b1 100644 --- a/web_interface/blueprints/api_v3/__init__.py +++ b/web_interface/blueprints/api_v3/__init__.py @@ -54,6 +54,7 @@ from src.common.path_safety import resolve_under, safe_path_component from src.core_config_keys import CORE_CONFIG_KEYS, CORE_SECRETS_KEYS from src.backup_manager import BUNDLED_FONTS as _BUNDLED_FONTS from src.device_location import DeviceLocationResolver, apply_device_location +from src.ipc import client as control_client _SUDO = shutil.which('sudo') _JOURNALCTL = shutil.which('journalctl') _GIT = shutil.which('git') @@ -776,6 +777,83 @@ def _store_restart_fields(action: str, plugin_enabled: bool, **kwargs) -> Dict[s return fields +# -- the control socket (src/ipc) --------------------------------------------------- + +#: Socket failures that only mean "this display has no socket": stopped, +#: older than the socket, Windows, or switched off. Not worth a log line. +_QUIET_SOCKET_REASONS = frozenset({'no_socket', 'disabled', 'unsupported'}) + +#: Every reason code a response may echo as a socket error: the client's +#: transport reasons plus the display's ErrorCode values. Anything else is +#: reported as ``other``, so no text taken from an exception reaches a reply. +_REPORTABLE_SOCKET_REASONS = ( + 'disabled', 'unsupported', 'no_socket', 'refused', 'timeout', 'closed', + 'bad_response', 'invalid_request', + 'bad_json', 'bad_request', 'message_too_large', 'unsupported_version', + 'unknown_command', 'invalid_args', 'busy', 'forbidden', 'internal', + 'pending', 'not_loaded', 'failed', +) + + +def _socket_reason_code(reason): + """``reason`` as one of _REPORTABLE_SOCKET_REASONS, else ``'other'``.""" + return next((code for code in _REPORTABLE_SOCKET_REASONS if code == reason), 'other') + + +def _log_socket_failure(what: str, error: Exception, reason: str) -> None: + if reason in _QUIET_SOCKET_REASONS: + logger.debug("%s not sent over the control socket: %s", what, error) + else: + logger.warning("Control socket did not take %s (%s)", what, error) + + +def _reload_after_store_update(plugin_id: str, fields: Dict[str, Any]) -> Dict[str, Any]: + """Reload an updated, enabled plugin on the running display. + + ``fields`` are ``_store_restart_fields('update', ...)``. When they ask + for a restart, the display is asked over the control socket to reload + the plugin instead (``plugin.reload``), and the answer becomes + ``restart_required: false`` with ``reloaded: true`` once the new code is + running. Any failure -- no socket, a display that predates the command, + a plugin it is not running, a load that failed, no answer in time -- + keeps ``fields`` as they were, adding ``reload_error`` with the reason. + """ + if not fields.get('restart_required'): + return fields + try: + result = control_client.plugin_reload(plugin_id) + except control_client.ControlError as e: + reason = _socket_reason_code(e.reason) + _log_socket_failure(f'plugin.reload {plugin_id}', e, reason) + return {**fields, 'reload_error': reason} + except Exception: # never let the socket path break the route + logger.exception("Control socket client failed reloading %s", plugin_id) + return {**fields, 'reload_error': 'internal'} + version = result.get('version') + return {'restart_required': False, 'reloaded': True, + 'reloaded_version': version if isinstance(version, str) else None} + + +def _apply_brightness_on_display(brightness: int) -> Dict[str, Any]: + """Put a just-saved brightness on the panel now, over the control socket. + + Without the socket the display's config watcher applies the saved value + within a few seconds, as it always has; the answer says which happened: + ``brightness_transport`` is ``socket`` or ``config``, and in the second + case ``brightness_socket_error`` gives the reason. + """ + try: + control_client.brightness_set(int(brightness)) + return {'brightness_transport': 'socket'} + except control_client.ControlError as e: + reason = _socket_reason_code(e.reason) + _log_socket_failure('brightness.set', e, reason) + except Exception: # never let the socket path break the route + logger.exception("Control socket client failed setting brightness") + reason = 'internal' + return {'brightness_transport': 'config', 'brightness_socket_error': reason} + + def deep_merge(base_dict, update_dict): """ Deep merge update_dict into base_dict. diff --git a/web_interface/blueprints/api_v3/config.py b/web_interface/blueprints/api_v3/config.py index 1c278a3f..41720c96 100644 --- a/web_interface/blueprints/api_v3/config.py +++ b/web_interface/blueprints/api_v3/config.py @@ -1193,7 +1193,15 @@ def save_main_config(): # Display hardware, rotation/durations and general settings take # effect after a display restart; the UI shows its restart banner on # this flag. - return success_response(message=message, extra={'restart_required': True}) + extra = {'restart_required': True} + # Brightness is the exception: the display applies a saved one + # without a restart. Over the control socket it lands at once, + # instead of when the config watcher next looks (up to ~2 s). + if 'brightness' in data: + saved = (current_config.get('display', {}).get('hardware', {}) or {}).get('brightness') + if isinstance(saved, int) and not isinstance(saved, bool): + extra.update(_pkg._apply_brightness_on_display(saved)) + return success_response(message=message, extra=extra) except Exception as e: logger.error("Error saving config", exc_info=True) return error_response( diff --git a/web_interface/blueprints/api_v3/display.py b/web_interface/blueprints/api_v3/display.py index a2fc74ba..73433c3a 100644 --- a/web_interface/blueprints/api_v3/display.py +++ b/web_interface/blueprints/api_v3/display.py @@ -4,8 +4,9 @@ Routes decorate the shared `api_v3` Blueprint from the package `__init__`, so their endpoint names are unchanged by living here. """ from web_interface.blueprints.api_v3 import ( + _QUIET_SOCKET_REASONS, _REPORTABLE_SOCKET_REASONS, # noqa: F401 - tests read them here _coerce_to_bool, _ensure_display_service_running, - _get_display_service_status, _stop_display_service, api_v3, + _get_display_service_status, _socket_reason_code, _stop_display_service, api_v3, jsonify, logger, request, uuid, ) from web_interface import display_preview @@ -30,24 +31,6 @@ def _cache_manager(): return cache -#: Socket failures that only mean "this display has no socket": stopped, -#: older than the socket, Windows, or switched off. Not worth a log line. -_QUIET_SOCKET_REASONS = frozenset({'no_socket', 'disabled', 'unsupported'}) - -#: Every reason code a response may echo as ``socket_error``: the client's -#: transport reasons plus the display's ErrorCode values. Anything else is -#: reported as ``other``, so no text taken from an exception reaches a reply. -_REPORTABLE_SOCKET_REASONS = ( - 'disabled', 'unsupported', 'no_socket', 'refused', 'timeout', 'closed', - 'bad_response', 'invalid_request', - 'bad_json', 'bad_request', 'message_too_large', 'unsupported_version', - 'unknown_command', 'invalid_args', 'busy', 'forbidden', 'internal', -) - - -def _socket_reason_code(reason): - """``reason`` as one of _REPORTABLE_SOCKET_REASONS, else ``'other'``.""" - return next((code for code in _REPORTABLE_SOCKET_REASONS if code == reason), 'other') def _deliver_on_demand(payload): diff --git a/web_interface/blueprints/api_v3/plugin_store.py b/web_interface/blueprints/api_v3/plugin_store.py index af0cb362..23aca5a9 100644 --- a/web_interface/blueprints/api_v3/plugin_store.py +++ b/web_interface/blueprints/api_v3/plugin_store.py @@ -7,7 +7,7 @@ so their endpoint names do not depend on which module they live in. from web_interface.blueprints.api_v3 import ( ErrorCode, OperationType, Path, _do_transactional_uninstall, _non_plugin_id_error, _get_plugin_version, _plugin_directory, - _plugin_enabled_in_config, _store_restart_fields, api_v3, + _plugin_enabled_in_config, _reload_after_store_update, _store_restart_fields, api_v3, error_response, exception_error_response, json, jsonify, logger, request, success_response, validate_request_json, ) @@ -181,12 +181,20 @@ def update_plugin(): if success: updated_last_updated = current_last_updated updated_version = current_version + # The id the display knows the plugin by (and keys its config + # section on) is the manifest's, which can differ from the + # store id this route was given (an alias such as weather / + # ledmatrix-weather). + display_plugin_id = plugin_id try: if manifest_path is not None and manifest_path.exists(): with open(manifest_path, 'r', encoding='utf-8') as f: manifest = json.load(f) updated_last_updated = manifest.get('last_updated', current_last_updated) updated_version = manifest.get('version', current_version) + manifest_id = manifest.get('id') + if isinstance(manifest_id, str) and safe_path_component(manifest_id): + display_plugin_id = manifest_id except Exception as e: logger.debug("Could not read updated manifest after update: %s", e) @@ -230,12 +238,19 @@ def update_plugin(): api_v3.schema_manager.invalidate_cache(plugin_id) # Rediscover plugins. The web process runs no plugin code, so - # there is nothing here to reload: the display keeps running the - # version it loaded until it restarts, which restart_required - # below asks for. + # there is nothing to reload here: the display reloads it, asked + # over the control socket below, or keeps running the version it + # loaded until a restart, which restart_required then asks for. if api_v3.plugin_catalog: api_v3.plugin_catalog.discover_plugins() + restart_fields = _reload_after_store_update( + display_plugin_id, + _store_restart_fields('update', _plugin_enabled_in_config(display_plugin_id), + changed=update_status == 'updated')) + if restart_fields.get('reloaded'): + message += '; the display is running the new version' + # Record in history (the only record of when it was updated; # the version is the manifest on disk). if api_v3.operation_history: @@ -260,9 +275,7 @@ def update_plugin(): 'update_status': update_status }, message=message, - extra=_store_restart_fields( - 'update', _plugin_enabled_in_config(plugin_id), - changed=update_status == 'updated'), + extra=restart_fields, ) else: refusal = _compatibility_refusal(plugin_id)