From 515248b34e49e76b9c9feefb81d68241af651abb Mon Sep 17 00:00:00 2001 From: Chuck <33324927+ChuckBuilds@users.noreply.github.com> Date: Sat, 3 Oct 2026 12:56:47 -0400 Subject: [PATCH] feat(ipc): control socket stage 3 - a state stream replaces polled cache keys (#735) Adds state.get / state.subscribe to the display's control socket (StateHub in src/ipc/server.py). The web interface holds one subscription per process (web_interface/display_state.py) and reads current-status, on-demand status, plugin runtime and /health's display_loop from it, falling back to the cache keys and heartbeat file. While the socket serves readers, display_current_state and plugin_runtime_snapshot are written less often (about 1.5 instead of 5 cache writes a minute for 15 s screens). Co-Authored-By: Claude Opus 5.5 --- CHANGELOG.md | 30 ++ docs/IPC_CONTROL_SOCKET.md | 266 +++++++++- docs/REST_API_REFERENCE.md | 31 +- src/display_controller.py | 114 ++++- src/display_watchdog.py | 17 + src/ipc/client.py | 230 ++++++++- src/ipc/contract.py | 188 +++++++- src/ipc/server.py | 358 +++++++++++++- src/plugin_system/plugin_runtime.py | 128 ++++- test/test_ipc_state_stream.py | 491 +++++++++++++++++++ test/test_state_stream_readers.py | 506 ++++++++++++++++++++ web_interface/app.py | 6 +- web_interface/blueprints/api_v3/__init__.py | 10 +- web_interface/blueprints/api_v3/display.py | 48 +- web_interface/blueprints/api_v3/misc.py | 16 +- web_interface/display_state.py | 145 ++++++ 16 files changed, 2487 insertions(+), 97 deletions(-) create mode 100644 test/test_ipc_state_stream.py create mode 100644 test/test_state_stream_readers.py create mode 100644 web_interface/display_state.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 7b00a972..c481f0f1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -19,6 +19,36 @@ accepts both, but the store flags the old spelling as deprecated ## Unreleased +### Control socket stage 3: the display's state over the socket + +- Two new commands, still protocol version 1. `state.get` returns a + versioned snapshot of what the display is doing: the current mode and + plugin, the on-demand session, the brightness, the plugin runtime + snapshot and the render loop's heartbeat age. With `since`/`epoch` it + returns a short "unchanged" answer. `state.subscribe` returns the same + snapshot, then pushes a `state` event on every change (always the latest + version) and a `tick` at least every 5 s. The display serves all of it + from memory (`StateHub` in `src/ipc/server.py`), and publishing never + waits for a reader. Subscribers have their own bound (4), separate from + the 8 request slots, and one that stops reading is dropped after the 2 s + IO timeout. See `docs/IPC_CONTROL_SOCKET.md`, "The state stream". +- The web interface holds one subscription per process + (`web_interface/display_state.py`). `/display/current-status`, + `/display/on-demand/status`, the plugin runtime fields of + `/plugins/installed` and `/plugins/state`, the reconciliations and + `/health`'s `display_loop` read it first. When the socket is missing (a + stopped or older display, Windows), they fall back to the cache keys and + the heartbeat file. Each answer has a `source` (`socket`, `cache` or + `heartbeat_file`). The stale and stalled rules from #726 apply the same + way to both. +- Fewer SD-card writes while the socket serves those readers. + `display_current_state` is written once a minute and on a flag change, + not on every mode change. The `plugin_runtime_snapshot` refresh goes from + 60 s to 120 s. For a rotation of 15 s screens, that is 1.5 cache writes a + minute instead of 5. Both keys keep being written for one release. +- `RenderWatchdog.liveness()` reports the heartbeat age from memory. + `PluginRuntimeView` has a `source`, and `describe()` includes it. + ### Web UI: four more tabs are ES-module pages (stage 2) - Rotation, Operation History, Config Editor and Backup & Restore follow the diff --git a/docs/IPC_CONTROL_SOCKET.md b/docs/IPC_CONTROL_SOCKET.md index 96919bc6..56f98e71 100644 --- a/docs/IPC_CONTROL_SOCKET.md +++ b/docs/IPC_CONTROL_SOCKET.md @@ -4,14 +4,17 @@ 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 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. +screen, and adds `brightness.set` and `plugin.reload`. Stage 3 adds a state +stream (`state.get`, `state.subscribe`), so the web interface reads what the +display is doing from the socket instead of from cache files the display +wrote to the SD card. The file mailbox and the cache keys stay 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`, `POST /api/v3/plugins/update` (reload), `POST /api/v3/config/main` (brightness) | +| 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); and through [`web_interface/display_state.py`](../web_interface/display_state.py) (the state stream), `GET /api/v3/display/current-status`, `/display/on-demand/status`, `/plugins/installed` (`runtime`), `/plugins/state` and the reconciliations, `/health` (`display_loop`) | | 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 | @@ -36,6 +39,11 @@ running, the socket does not exist, and the web interface knows right away. ## Protocol (version 1) +Stage 3 is still version 1: `state.get` and `state.subscribe` are new +commands, and a stage-2 display answers them `unknown_command`, which the +web interface treats as "no socket" and falls back from. + + **Framing.** One JSON object per line (newline-delimited JSON), UTF-8, at most 64 KiB per line (`MAX_MESSAGE_BYTES`). Senders encode with `ensure_ascii`, so a newline never appears inside a message. A connection @@ -75,6 +83,8 @@ one. Clients branch on `error.code`, never on the message text. | `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) | +| `state.get` | `{since?, epoch?}` | a state snapshot (see "The state stream") | answered directly | +| `state.subscribe` | — | a state snapshot, then pushed `state` / `tick` events | answered directly, then a stream | `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 @@ -130,6 +140,19 @@ 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. +**Events.** `state.subscribe` is the one command with more than one message +in reply. After its response, the display pushes events on the same +connection until either side hangs up: + +```json +{"v": 1, "id": "", "event": "state", "result": {...a state snapshot...}} +{"v": 1, "id": "", "event": "tick", "result": {"version": 7, "epoch": "…", "pid": 812, "served_at": 1790000000.1, "changed": false, "loop": {...}}} +``` + +An event has `event` where a response has `ok`, which is how a reader tells +them apart. The client sends nothing after the subscribe; anything it does +send is ignored. + **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`. @@ -147,6 +170,169 @@ print(client.brightness_set(60)) EOF ``` +## The state stream (stage 3) + +Before stage 3 the web interface learned what the display was doing by +reading files the display kept writing: + +| What | Written by the display | How often | Medium | +|---|---|---|---| +| current mode, plugin, `is_display_active`, `on_demand_active` | `display_current_state` | every mode change, every flag change, and every 30 s | cache (SD card) | +| on-demand session | `display_on_demand_state` | on each on-demand event | cache (SD card) | +| plugin runtime snapshot (#690) | `plugin_runtime_snapshot` | on a change (at most every 10 s), else every 60 s | cache (SD card) | +| render-loop liveness (#687) | `display-heartbeat.json` | every 5 s | tmpfs | + +Now the display also keeps the same state in memory and serves it on the +socket. + +**The snapshot.** `state.get` and `state.subscribe` answer with one object: + +```json +{"schema": 1, "version": 42, "epoch": "3f9c0d1e2a4b5c6d", "pid": 812, + "served_at": 1790000000.1, "changed": true, + "loop": {"heartbeat_age_seconds": 1.8, "armed": true, "stale_after": 60.0}, + "state": { + "display": {"mode": "nfl_live", "plugin_id": "football-scoreboard", "mode_index": 3, + "total_modes": 9, "on_demand_active": false, "is_display_active": true, + "last_updated": 1790000000.0}, + "on_demand": {"active": false, "status": "idle", "...": "as display_on_demand_state"}, + "brightness": {"brightness": 80, "panel_brightness": 40, "dimmed": true}, + "plugins": {"schema": 1, "running": true, "published_at": 1789999998.5, "...": "as plugin_runtime_snapshot"}, + "loop": {"heartbeat_age_seconds": 1.8, "armed": true, "stale_after": 60.0} + }} +``` + +- `display` and `on_demand` are the dicts the cache keys hold, `plugins` is + the runtime snapshot (`build_runtime_snapshot`), and `brightness` is the + configured level, what the panel shows now, and whether the dim schedule + has it dimmed. A section not published yet is `null`. +- `loop` is not published: the display measures it when it answers, from + the render thread's last beat in memory (`RenderWatchdog.liveness()`), + the same beat that writes the heartbeat file. So it keeps ageing while the + render thread is stuck, and the socket's connection threads still answer. + `heartbeat_age_seconds` is `null` until the loop has drawn its first frame. +- `version` goes up whenever a section changes, ignoring the timestamps that + move on every publish (`last_updated`, `remaining`, `published_at`). It + counts within an `epoch`, one run of the display process, so a reader that + sees a new `epoch` has a restarted display. +- `state.get` with `since` and `epoch` from an earlier answer gets just + `{changed: false, version, epoch, pid, served_at, loop}` while nothing has + changed. +- A snapshot that would not fit in a message (hundreds of plugins) is sent + without `plugins`, and `truncated: ["plugins"]` says so. Readers then use + the cache for that section only. + +**The stream.** `state.subscribe` answers with the snapshot, then: + +- a `state` event (a full snapshot) whenever the version changes, and +- a `tick` at least every 5 s (`SUBSCRIBE_KEEPALIVE_SECONDS`) when nothing + changed. It carries `loop`, so a stalled render loop shows up within one + tick, and it tells the reader the connection is alive. + +A slow reader is never sent a backlog: each event is the latest version, so +one that falls behind skips the versions in between. A reader that has heard +nothing for 15 s (three keepalives) stops trusting its copy. + +**Who publishes, and when.** All of it is in memory, with no disk writes: + +- the render thread, at the places it already published the cache keys: + `display` and `brightness` on every pass of + `_publish_current_mode_state_if_changed()` (every loop pass, and every + `_service_pending_changes()` in a dwell, a scrolling screen or Vegas), and + `on_demand` in `_publish_on_demand_state()`. Every pass refreshes + `display.last_updated`, so a reader can tell when the render thread has + stopped publishing, just as the cache key's 120 s `max_age` does. +- the plugin runtime publisher's thread, on every 5 s tick: the snapshot is + rebuilt when the state machine changed, otherwise only its `published_at` + moves. A change reaches subscribers within a tick, without the cache's + 10 s throttle. + +Publishing is a hand-off, as the command queue is in the other direction. +The hub (`StateHub` in [`src/ipc/server.py`](../src/ipc/server.py)) holds a +lock only to swap a dict reference, compare it with the last one and bump the +version. Every socket write happens on the subscriber's own connection +thread. The render thread never waits for a reader. + +### Readers in the web interface + +[`web_interface/display_state.py`](../web_interface/display_state.py) holds +one `state.subscribe` connection per web process +(`src.ipc.client.StateSubscription`, a daemon thread, started on the first +read and reconnecting with a backoff of 1 s up to 30 s). A route answers +from the latest pushed snapshot in memory. Before the subscription has one, +the route asks once with `state.get` (0.5 s timeout). When neither works, it +reads the cache keys and the heartbeat file as before: + +| Route | From the socket | Fallback | +|---|---|---| +| `GET /api/v3/display/current-status` | `state.display` | `display_current_state` | +| `GET /api/v3/display/on-demand/status` | `state.on_demand`, with `remaining` worked out from `expires_at` now | `display_on_demand_state` | +| `GET /api/v3/plugins/installed` (`runtime`), `/plugins/state`, `POST /plugins/state/reconcile` and the startup reconciliation | `state.plugins` + `state.loop` | `plugin_runtime_snapshot` + `display-heartbeat.json` | +| `GET /api/v3/health` (`checks.display_loop`) | `state.loop` | `display-heartbeat.json` | + +Each answer says where it came from: `source: "socket" | "cache"` (or +`"heartbeat_file"` for the health check). + +The SSE display stream (`/api/v3/stream/display`) reads the preview frame +file, not a cache key, so it does not change. + +**The same verdicts either way.** The socket's answers are judged by the +rules the cache readers apply (#726): + +- the runtime view is `stalled` when the render loop's heartbeat age is at + least `HEARTBEAT_STALE_SECONDS` (60 s, the health check's threshold), and + then reports no per-plugin facts; +- it is `stale` when the snapshot is older than its `stale_after` (the + publisher thread stopped); +- with no beat yet, the snapshot is judged on its own; +- there is no pid check, because the display that answered is alive; +- a `display` section the render thread has not refreshed for 120 s reads + as unknown, as the cache key does once it ages out. + +The age a reader uses is the age the display measured, plus the time since +the snapshot arrived. + +### Fewer SD writes + +The cache keys are still written, for one release, as the fallback. While +the socket serves the readers, the display writes two of them less often. +"Serves the readers" means a subscriber is connected, or a `state.get` came +within the last 60 s (`StateHub.readers_active()`): + +- `display_current_state` is no longer written on every mode change: once + every 60 s (`CURRENT_STATE_RELAXED_REFRESH_SECONDS`, inside the readers' + 120 s `max_age`), and at once when `is_display_active` or + `on_demand_active` changes. +- `plugin_runtime_snapshot`'s refresh goes from 60 s to 120 s + (`RELAXED_REFRESH_INTERVAL`), and the snapshot says so in its own + `refresh_interval` and `stale_after` (360 s). Changes are still written at + once, at most every 10 s. + +`display_on_demand_state` is written only on events, so it is unchanged. +The heartbeat file is on tmpfs, so it costs no SD writes, and it stays: the +automatic update's health check reads it. + +This is safe because the relaxed rate only applies while readers are using +the socket. If they stop (the web interface loses the socket, or is stopped), +the next publish after the reader window writes a changed mode at once, and +the runtime refresh goes back to 60 s. A fallback reader in that window sees +a mode up to 60 s old, never one older than its `max_age`. + +Measured with fake clocks (`test_cache_writes_per_minute_with_and_without_socket_readers` +in `test/test_state_stream_readers.py`), for a rotation of 15 s screens: + +| Key | Writes/min, no socket readers | Writes/min, socket readers | +|---|---|---| +| `display_current_state` | 4.0 | 1.0 | +| `plugin_runtime_snapshot` | 1.0 | 0.5 | +| Total | 5.0 | 1.5 | + +That is 70% fewer writes for these keys: about 2,200 a day instead of 7,200. +Shorter screens save more, because the old rate followed the mode changes. +A display that rarely changes mode (one plugin, a long live game) saves less. Plugin +data caches, the error snapshot and font usage are written by other code +and are not affected. + ## How the display applies a command The server's threads never touch rendering. A connection thread parses the @@ -279,7 +465,16 @@ block the render loop or crash it: process created. - **Never fatal.** If the server cannot start (Windows, no `AF_UNIX`, a bind failure, `LEDMATRIX_CONTROL_SOCKET=off`), it logs that and the display runs - as before. The web interface then uses the mailbox. + as before. The web interface then uses the mailbox, and reads the cache + keys and the heartbeat file. +- **Subscribers (stage 3).** A `state.subscribe` connection gives its request + slot back and takes one of 4 subscriber slots (`MAX_SUBSCRIBERS`). A fifth + gets `busy`. So a few browsers' web processes holding streams can never + use up the 8 slots that commands need. Each subscriber has its own thread. + A send that cannot finish within the 2 s IO timeout (a reader that stopped + reading) drops that subscriber. Nothing else waits for it, and the render + thread only publishes to the hub. `close()` wakes every subscriber, so + they end at once. ## Security model @@ -313,8 +508,10 @@ 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. +import a plugin it was not running. Stages 2 and 3 changed none of the +access rules above. The state stream carries what the cache keys already +held, and those are readable by the same group. A subscriber goes through +the same connect-time and peer-credential checks as any other connection. **Development.** A display that is not root and cannot write to `/run/ledmatrix`, such as `python3 run.py -e` from a checkout, serves the @@ -352,29 +549,34 @@ device never touches the live display. "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 (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. +3. **A state stream (done).** `state.get` (a versioned snapshot) and + `state.subscribe` (the snapshot, then pushed changes and keepalive ticks) + carry the current mode, the on-demand state (including the outcome of an + acked on-demand command), the brightness, the plugin runtime snapshot and + the render loop's liveness, all served from memory (see "The state + stream"). The web interface's readers use it and fall back to the cache + keys and the heartbeat file. `display_current_state` and + `plugin_runtime_snapshot` are written less often while it serves them. + The keys remain for one release. + - Left for later: the outcome of a `plugin.reload` that answered + `pending` is visible only as the plugin's new `loaded_version` in + `state.plugins`, not as an event of its own. + - Left for later: the SSE display stream reads the preview frame, not + state, so nothing relays the stream to the browser yet. A browser still + polls the REST routes, which now answer from memory. + - Left for later: the store's install of an already-enabled plugin, and + an uninstall that keeps its config, still answer `restart_required`. + They can now use a load/unload command and report the result the same + way the update route does. 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 they can move to an in-process `request_display()`. The other cache keys used as messages (`plugin_error_clear_request` and the remaining - `display_*` keys) move to the socket or to tmpfs. + `display_*` keys) move to the socket or to tmpfs. The display also stops + writing `display_current_state`, `display_on_demand_state` and + `plugin_runtime_snapshot` once the web interface no longer falls back to + them. ## Checking it on a device @@ -405,3 +607,19 @@ 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. + +The state stream: + +```bash +curl -s localhost:5000/api/v3/display/current-status # ... "source": "socket" +curl -s localhost:5000/api/v3/health | python3 -m json.tool | grep -A3 display_loop +python3 - <<'EOF' +from src.ipc import client # run from the project directory +snap = client.state_get() +print(snap['version'], snap['epoch'], snap['loop'], snap['state']['display']) +EOF +``` + +`"source": "cache"` means the web interface could not use the socket: the +display is stopped, predates stage 3, or the web user is not in the +socket's group. diff --git a/docs/REST_API_REFERENCE.md b/docs/REST_API_REFERENCE.md index b6e08422..0dce2c66 100644 --- a/docs/REST_API_REFERENCE.md +++ b/docs/REST_API_REFERENCE.md @@ -331,12 +331,19 @@ by the display process (stale after 120 seconds). "data": { "mode": "nfl_live", "plugin_id": "football-scoreboard", - "last_updated": 1234567890.123 + "last_updated": 1234567890.123, + "source": "socket" } } ``` -When nothing has been published, every field is `null`. +When nothing has been published, every field is `null`. `source` is +`socket` when the answer came from the display's state stream over the +control socket ([IPC_CONTROL_SOCKET.md](IPC_CONTROL_SOCKET.md)), and `cache` +when it came from the `display_current_state` cache key (no socket: the +display is stopped or older, or this is Windows). A display whose render +loop has not refreshed its state for 120 seconds is reported with every +field `null`, either way. ### List Display Modes @@ -413,11 +420,16 @@ Get the current on-demand display state. "returncode": 0, "stdout": "active", "stderr": "" - } + }, + "source": "socket" } } ``` +`source` is `socket` (the display's state stream, with `remaining` worked +out at the time of the request) or `cache` (the `display_on_demand_state` +cache key). + With no on-demand request, `state` is `{"active": false, "status": "idle", "last_updated": null}`. @@ -552,7 +564,8 @@ List all installed plugins with their status and metadata. "published_at": 1790000030.0, "age_seconds": 12.4, "stale_after": 180.0, - "heartbeat_age_seconds": 2.1 + "heartbeat_age_seconds": 2.1, + "source": "socket" } } } @@ -581,7 +594,10 @@ the display is hung or died), `stopped` (the display shut down) or `unknown` (nothing published yet). Unless it is `live`, every one of those fields is `null`. `heartbeat_age_seconds` is the heartbeat's age when it was taken into account, `null` otherwise (no heartbeat, as on the dev server, or one from -another process). Health and metrics are at [`/plugins/health`](#get-plugin-health) +another process). `runtime.source` is `socket` when the snapshot and the +heartbeat age came from the display's state stream over the control socket, +and `cache` when they came from the `plugin_runtime_snapshot` cache key and +the heartbeat file; the rules above are the same for both. Health and metrics are at [`/plugins/health`](#get-plugin-health) and `/plugins/metrics`. `vegas_participation` is what Vegas mode does with the plugin: `"scroll"`, @@ -2265,7 +2281,10 @@ display snapshot. `data.status` is `healthy` or `degraded`, with (with `heartbeat_age_seconds`), `stalled` (no heartbeat for 60s: the panel is frozen even if the service is active; the status turns `degraded`), or `not_reported` when the display writes none (not started yet, the dev server, -Windows), which does not affect the status. +Windows), which does not affect the status. Its `source` is `socket` when the +age came from the display's state stream over the control socket (measured +in memory by the display) and `heartbeat_file` when it came from +`/run/ledmatrix/display-heartbeat.json`. Open even when the web login is on, for uptime monitors; a caller that is not logged in (and has no token) then gets only `{"status": "success", "data": diff --git a/src/display_controller.py b/src/display_controller.py index f076c05d..d081fbaf 100644 --- a/src/display_controller.py +++ b/src/display_controller.py @@ -55,7 +55,7 @@ from src.ipc.contract import ( PluginReloadArgs, PluginReloadResult, ) -from src.ipc.server import ControlServer, QueuedCommand, start_control_server +from src.ipc.server import ControlServer, QueuedCommand, StateHub, start_control_server from src.vegas_mode.render_pipeline import SYNC_SEND_INTERVAL # Get logger with consistent configuration @@ -65,6 +65,13 @@ logger = get_logger(__name__) # treats display_current_state older than 120 s as unknown. CURRENT_STATE_REFRESH_SECONDS = 30 +# While the control socket serves the web interface's state readers +# (StateHub.readers_active), display_current_state is only their fallback: +# it is then rewritten at this interval and on a change of the flags, not on +# every mode change. Below the readers' 120 s max_age, so the fallback copy +# never reads as unknown. +CURRENT_STATE_RELAXED_REFRESH_SECONDS = 60 + # How long startup will wait for plugins to fetch their first data before # showing anything. Each plugin's update blocks for up to the executor's 30s # timeout and they run one after another, so the uncapped total is the sum of @@ -1430,18 +1437,60 @@ class DisplayController: return None return max(0.0, expires_at - time.time()) + #: The control socket's state stream (src/ipc/server.StateHub), while the + #: socket is served. Class-level default for controllers built without + #: __init__ (tests) and for a display with no socket. + _state_hub: Optional[StateHub] = None + + def _current_mode_state(self) -> Dict[str, Any]: + """What display_current_state and the socket's ``display`` section hold.""" + return { + 'mode': self.current_display_mode, + 'plugin_id': self.mode_to_plugin_id.get(self.current_display_mode), + 'mode_index': self.current_mode_index, + 'total_modes': len(self.available_modes), + 'on_demand_active': self.on_demand_active, + 'is_display_active': self.is_display_active, + 'last_updated': time.time(), + } + + def _push_live_state(self, display_state: Optional[Dict[str, Any]] = None) -> None: + """Hand the current mode and the brightness to the control socket's + state stream. In memory, no disk: the hub only bumps its version (and + wakes subscribers) when something other than ``last_updated`` changed. + + Called on every pass of the publish points below, so + ``display.last_updated`` doubles as the render thread's proof of life + for the socket's readers, as the cache key's max_age does today. + """ + hub = self._state_hub + if hub is None: + return + try: + hub.publish('display', display_state or self._current_mode_state(), + volatile=('last_updated',)) + hub.publish('brightness', { + 'brightness': getattr(self, '_normal_brightness', None), + 'panel_brightness': getattr(self, 'current_brightness', None), + 'dimmed': bool(getattr(self, 'is_dimmed', False)), + }) + except Exception as err: # pylint: disable=broad-except + logger.debug("Could not publish the display state to the control socket: %s", + err, exc_info=True) + + def _state_readers_on_socket(self) -> bool: + """Is the control socket serving the web interface's state readers?""" + hub = self._state_hub + try: + return bool(hub is not None and hub.readers_active()) + except Exception: # pylint: disable=broad-except + return False + def _publish_current_mode_state(self) -> None: """Publish the currently active display mode/plugin to cache for the web UI.""" try: - state = { - 'mode': self.current_display_mode, - 'plugin_id': self.mode_to_plugin_id.get(self.current_display_mode), - 'mode_index': self.current_mode_index, - 'total_modes': len(self.available_modes), - 'on_demand_active': self.on_demand_active, - 'is_display_active': self.is_display_active, - 'last_updated': time.time(), - } + state = self._current_mode_state() + self._push_live_state(state) self.cache_manager.set('display_current_state', state) self._last_published_mode = self.current_display_mode self._last_published_flags = self._current_state_flags() @@ -1467,11 +1516,23 @@ class DisplayController: priority, a single enabled plugin -- has to be republished or the UI reports it as unknown. Otherwise this writes only on a change, not on every render tick. + + While the control socket serves the web interface's state readers, + the socket's in-memory copy is updated on every call and the cache + key is only their fallback: a mode change alone is then written at + the relaxed refresh (CURRENT_STATE_RELAXED_REFRESH_SECONDS), still + inside the readers' max_age. The flags are still written at once. + When the socket stops serving them, the next call writes a changed + mode again. """ - if (self.current_display_mode != self._last_published_mode + relaxed = self._state_readers_on_socket() + refresh = CURRENT_STATE_RELAXED_REFRESH_SECONDS if relaxed else CURRENT_STATE_REFRESH_SECONDS + if ((not relaxed and self.current_display_mode != self._last_published_mode) or self._current_state_flags() != getattr(self, '_last_published_flags', None) - or time.monotonic() - self._last_published_at >= CURRENT_STATE_REFRESH_SECONDS): - self._publish_current_mode_state() + or time.monotonic() - self._last_published_at >= refresh): + self._publish_current_mode_state() # pushes to the socket as well + else: + self._push_live_state() def _on_demand_state(self) -> Dict[str, Any]: """The on-demand state as published to the cache and the control socket.""" @@ -1494,6 +1555,11 @@ class DisplayController: """Publish current on-demand state to cache for external consumers.""" try: state = self._on_demand_state() + hub = self._state_hub + if hub is not None: + # In memory, first: a subscriber hears the outcome of an + # on-demand command even if the cache write below fails. + hub.publish('on_demand', state, volatile=('last_updated', 'remaining')) self.cache_manager.set('display_on_demand_state', state) except (OSError, RuntimeError, ValueError, TypeError) as err: logger.error("Failed to publish on-demand state: %s", err, exc_info=True) @@ -1738,12 +1804,31 @@ class DisplayController: """ if self._control_server is not None: return + hub = StateHub(loop_probe=display_watchdog.watchdog.liveness) try: self._control_server = start_control_server( status_provider=self._control_status, - cache_dir=getattr(self.cache_manager, 'cache_dir', None)) + cache_dir=getattr(self.cache_manager, 'cache_dir', None), + state_hub=hub) except Exception: # pylint: disable=broad-except logger.exception("Control socket not started; using the file mailbox only") + if self._control_server is not None: + self._start_state_stream(hub) + + def _start_state_stream(self, hub: StateHub) -> None: + """Start publishing to the socket's state stream (``state.get`` and + ``state.subscribe``): everything a reader would see, now, then on + every publish. Never raises; without it readers use the cache keys.""" + try: + self._state_hub = hub + self._push_live_state() + hub.publish('on_demand', self._on_demand_state(), + volatile=('last_updated', 'remaining')) + publisher = getattr(self, '_plugin_runtime_publisher', None) + if publisher is not None: + publisher.attach_hub(hub) + except Exception: # pylint: disable=broad-except + logger.exception("Control socket state stream not started; readers use the cache") def _control_status(self) -> Dict[str, Any]: """The socket's on_demand.status answer. Runs on the socket's thread: reads only.""" @@ -4532,6 +4617,7 @@ class DisplayController: except Exception as e: logger.warning("Error closing the control socket: %s", e) self._control_server = None + self._state_hub = None # Stop the async update worker first so no in-flight update() call # is still touching display/cache-backed resources while they're # torn down below. diff --git a/src/display_watchdog.py b/src/display_watchdog.py index c32f82df..b5c4db06 100644 --- a/src/display_watchdog.py +++ b/src/display_watchdog.py @@ -216,6 +216,23 @@ class RenderWatchdog: def armed(self) -> bool: return self._armed + def liveness(self) -> Dict[str, Any]: + """The heartbeat, in memory: what the control socket's state stream + reports as ``loop``. + + ``heartbeat_age_seconds`` is the age of the render thread's last beat, + the beat that writes the heartbeat file, so it ages at the same rate + and is judged by the same ``HEARTBEAT_STALE_SECONDS``. None until the + loop has drawn its first frame. Any thread may call this: it only + reads two attributes. + """ + last = self._last_beat + age = None + if self._armed and last is not None: + age = max(self._clock() - last, 0.0) + return {'heartbeat_age_seconds': age, 'armed': self._armed, + 'stale_after': HEARTBEAT_STALE_SECONDS} + def _on_render_thread(self) -> bool: return self._render_thread is not None and threading.get_ident() == self._render_thread diff --git a/src/ipc/client.py b/src/ipc/client.py index 6f178d66..430859cd 100644 --- a/src/ipc/client.py +++ b/src/ipc/client.py @@ -10,20 +10,24 @@ blocks for longer than ``timeout`` in total. from __future__ import annotations import socket +import threading import time import uuid -from typing import Any, Dict, List, Mapping, Optional, Sequence +from typing import Any, Callable, Dict, List, Mapping, Optional, Sequence from src.ipc.contract import ( AWAIT_SECONDS, MAX_MESSAGE_BYTES, PROTOCOL_VERSION, + SUBSCRIBE_KEEPALIVE_SECONDS, SUPPORTED_VERSIONS, Command, FrameReader, ProtocolError, Request, Response, + StateEvent, + StateEventKind, client_socket_paths, decode_message, encode_message, @@ -233,3 +237,227 @@ def hello(client: str = 'web', *, timeout: float = DEFAULT_TIMEOUT_SECONDS, """Version negotiation: the result's ``version`` is the one both sides speak.""" return request(Command.HELLO, {'versions': list(SUPPORTED_VERSIONS), 'client': client}, timeout=timeout, paths=paths) + + +# -- the state stream (stage 3) --------------------------------------------------------- + +def state_get(since: Optional[int] = None, epoch: Optional[str] = None, *, + timeout: float = DEFAULT_TIMEOUT_SECONDS, + paths: Optional[Sequence[str]] = None) -> Dict[str, Any]: + """The display's state now, as a :class:`~src.ipc.contract.StateSnapshot`. + + With ``since``/``epoch`` from an earlier answer, an unchanged state comes + back in the short ``changed: false`` form. Raises :class:`ControlError` + (``unknown_command`` from a display older than stage 3). + """ + args: Dict[str, Any] = {} + if since is not None: + args['since'] = since + if epoch is not None: + args['epoch'] = epoch + return request(Command.STATE_GET, args, timeout=timeout, paths=paths) + + +def snapshot_age(snapshot: Mapping[str, Any], now_mono: Optional[float] = None) -> float: + """Seconds since ``snapshot`` arrived: ``received_mono`` (set by + :meth:`StateSubscription.latest`) to now; 0 for a one-shot answer.""" + received = snapshot.get('received_mono') + if isinstance(received, (int, float)) and not isinstance(received, bool): + now_mono = time.monotonic() if now_mono is None else now_mono + return max(now_mono - float(received), 0.0) + return 0.0 + + +def snapshot_loop_age(snapshot: Mapping[str, Any], + now_mono: Optional[float] = None) -> Optional[float]: + """The render loop's heartbeat age now, from a state snapshot: the age the + display measured when it answered, plus the time since the answer + arrived. None when the display has no beat to report yet.""" + loop = snapshot.get('loop') + if not isinstance(loop, dict): + state = snapshot.get('state') + loop = state.get('loop') if isinstance(state, dict) else None + age = loop.get('heartbeat_age_seconds') if isinstance(loop, dict) else None + if not isinstance(age, (int, float)) or isinstance(age, bool): + return None + return max(float(age), 0.0) + snapshot_age(snapshot, now_mono) + + +#: A subscription that has heard nothing for this long is not trusted: the +#: display sends a tick at least every SUBSCRIBE_KEEPALIVE_SECONDS. +SUBSCRIPTION_SILENCE_SECONDS = 3 * SUBSCRIBE_KEEPALIVE_SECONDS + +#: Reconnect backoff: the first retry, and the cap. A display that does not +#: know state.subscribe (stage 2 or older) is retried at the cap. +_RECONNECT_MIN_SECONDS = 1.0 +_RECONNECT_MAX_SECONDS = 30.0 + +#: Failures that another try soon will not fix. +_SLOW_RETRY_REASONS = frozenset({'unknown_command', 'unsupported_version', 'disabled', + 'unsupported'}) + + +class StateSubscription: + """One ``state.subscribe`` connection, held on a daemon thread. + + Keeps the latest snapshot the display pushed, so a reader answers from + memory (:meth:`latest`). Reconnects with a backoff when the display goes + away. Never raises into the caller: :meth:`latest` is None whenever the + copy cannot be vouched for (not connected, or silent for longer than + ``silence``), and the caller falls back. + """ + + def __init__(self, paths: Optional[Sequence[str]] = None, *, + silence: float = SUBSCRIPTION_SILENCE_SECONDS, + connect_timeout: float = DEFAULT_TIMEOUT_SECONDS, + clock: Callable[[], float] = time.monotonic): + self._paths = list(paths) if paths is not None else None + self._silence = silence + self._connect_timeout = connect_timeout + self._clock = clock + self._lock = threading.Lock() + self._snapshot: Optional[Dict[str, Any]] = None + self._received: Optional[float] = None + self._connected = False + self._stop = threading.Event() + self._sock: Optional[socket.socket] = None + self._thread: Optional[threading.Thread] = None + #: The reason the last connection ended (a ControlError reason). + self.last_error: Optional[str] = None + #: Full snapshots received: the subscribe answer and each state event. + self.snapshots = 0 + + # -- the reader's side --------------------------------------------------- + + @property + def connected(self) -> bool: + return self._connected + + def latest(self) -> Optional[Dict[str, Any]]: + """A copy of the latest snapshot, with ``received_mono`` (this + process's monotonic clock when it arrived); None when not trusted.""" + with self._lock: + if not self._connected or self._snapshot is None or self._received is None: + return None + if self._clock() - self._received > self._silence: + return None + snap = dict(self._snapshot) + snap['received_mono'] = self._received + return snap + + # -- lifecycle ----------------------------------------------------------- + + def start(self) -> 'StateSubscription': + if self._thread is None or not self._thread.is_alive(): + self._stop.clear() + self._thread = threading.Thread(target=self._run, name='ledmatrix-state-feed', + daemon=True) + self._thread.start() + return self + + def stop(self, timeout: float = 2.0) -> None: + self._stop.set() + sock = self._sock + if sock is not None: + try: + sock.shutdown(socket.SHUT_RDWR) + except OSError: + pass + thread = self._thread + if thread is not None and thread is not threading.current_thread(): + thread.join(timeout) + self._thread = None + + # -- the feed thread ----------------------------------------------------- + + def _run(self) -> None: + backoff = _RECONNECT_MIN_SECONDS + while not self._stop.is_set(): + try: + self._follow() + backoff = _RECONNECT_MIN_SECONDS + except ControlError as e: + self.last_error = e.reason + if e.reason in _SLOW_RETRY_REASONS: + backoff = _RECONNECT_MAX_SECONDS + except Exception as e: # pylint: disable=broad-except + self.last_error = type(e).__name__ + finally: + with self._lock: + self._connected = False + sock, self._sock = self._sock, None + if sock is not None: + try: + sock.close() + except OSError: + pass + if self._stop.wait(backoff): + return + backoff = min(backoff * 2, _RECONNECT_MAX_SECONDS) + + def _follow(self) -> None: + """Subscribe, then read events until the connection ends. Raises ControlError.""" + if not socket_supported(): + raise ControlError('unsupported', 'no Unix sockets on this platform') + candidates = list(self._paths) if self._paths is not None else client_socket_paths() + if not candidates: + raise ControlError('disabled', 'the control socket is turned off') + request_id = str(uuid.uuid4()) + payload = encode_message(Request(id=request_id, cmd=Command.STATE_SUBSCRIBE, + args={}).to_dict()) + sock = _connect(candidates, time.monotonic() + self._connect_timeout) + self._sock = sock + try: + sock.settimeout(self._connect_timeout) + sock.sendall(payload) + # A read waits for the next event; the display sends one at least + # every keepalive, so this much silence means it is gone. + sock.settimeout(self._silence) + reader = FrameReader(MAX_MESSAGE_BYTES) + first = True + while not self._stop.is_set(): + data = sock.recv(65536) + if not data: + raise ControlError('closed', 'the display closed the connection') + for line in reader.feed(data): + obj = decode_message(line) + if first: + response = Response.from_dict(obj) + if not response.ok: + error = response.error + raise ControlError(error.code if error else 'bad_response', + error.message if error else '') + self._store(dict(response.result or {}), full=True) + first = False + continue + event = StateEvent.from_dict(obj) + self._store(event.result, full=event.event == StateEventKind.STATE) + except socket.timeout: + raise ControlError('timeout', 'the display went quiet') from None + except ProtocolError as e: + raise ControlError('bad_response', e.message) from None + except OSError as e: + if self._stop.is_set(): + return + raise ControlError('closed', str(e)) from None + + def _store(self, result: Dict[str, Any], full: bool) -> None: + now = self._clock() + with self._lock: + if full and isinstance(result.get('state'), dict): + self._snapshot = result + self.snapshots += 1 + elif (self._snapshot is not None + and result.get('epoch') == self._snapshot.get('epoch')): + # A tick: nothing changed but the render loop's liveness. + snap = dict(self._snapshot) + loop = result.get('loop') + if isinstance(loop, dict): + snap['state'] = dict(snap.get('state') or {}, loop=loop) + snap['loop'] = loop + snap['served_at'] = result.get('served_at', snap.get('served_at')) + self._snapshot = snap + else: + return # a tick before any state, or from another epoch + self._received = now + self._connected = True diff --git a/src/ipc/contract.py b/src/ipc/contract.py index b563d17d..ee35c0ae 100644 --- a/src/ipc/contract.py +++ b/src/ipc/contract.py @@ -30,6 +30,11 @@ 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`. +``state.subscribe`` is the one exception to "one response per request": its +response is followed, on the same connection, by :class:`StateEvent` lines +the display pushes until either side hangs up. Events carry ``event`` +instead of ``ok``. + 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 @@ -142,11 +147,14 @@ class Command: ON_DEMAND_STATUS = 'on_demand.status' BRIGHTNESS_SET = 'brightness.set' PLUGIN_RELOAD = 'plugin.reload' + STATE_GET = 'state.get' + STATE_SUBSCRIBE = 'state.subscribe' #: 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). +#: ``brightness.set`` and ``plugin.reload`` came in stage 2, and ``state.get`` +#: and ``state.subscribe`` in stage 3, all within version 1 (see the module +#: docstring on adding commands). COMMANDS: Tuple[str, ...] = ( Command.HELLO, Command.PING, @@ -155,6 +163,8 @@ COMMANDS: Tuple[str, ...] = ( Command.ON_DEMAND_STATUS, Command.BRIGHTNESS_SET, Command.PLUGIN_RELOAD, + Command.STATE_GET, + Command.STATE_SUBSCRIBE, ) #: Commands that are queued for the render thread. @@ -173,6 +183,25 @@ AWAIT_SECONDS: Dict[str, float] = { } AWAITED_COMMANDS = frozenset(AWAIT_SECONDS) +#: The state stream (stage 3). ``state.subscribe`` turns its connection into +#: a one-way stream of :class:`StateEvent` lines. Subscribers have their own +#: bound, separate from the short request connections, so they can never +#: take the slots a command needs. +MAX_SUBSCRIBERS = 4 + +#: A subscriber hears from the display at least this often: a ``state`` +#: event when something changed, else a ``tick`` carrying the render loop's +#: liveness. A client that has heard nothing for a few of these treats its +#: copy as unknown. +SUBSCRIBE_KEEPALIVE_SECONDS = 5.0 + +#: The shape of the ``state`` object in a state snapshot. Bumped only when a +#: field changes meaning; new fields are added within a schema. +STATE_SCHEMA = 1 + +#: The sections of a state snapshot, in the order they are documented. +STATE_SECTIONS: Tuple[str, ...] = ('display', 'on_demand', 'brightness', 'plugins', 'loop') + #: Brightness, in percent, as the display's hardware setting takes it. MIN_BRIGHTNESS = 0 MAX_BRIGHTNESS = 100 @@ -477,8 +506,65 @@ class PluginReloadArgs: return cls(plugin_id=plugin_id) +def _optional_version(args: Mapping[str, Any], key: str) -> Optional[int]: + value = args.get(key) + if value is None: + return None + if not _is_int(value) or value < 0: + raise ProtocolError(ErrorCode.INVALID_ARGS, f'{key} must be a non-negative integer') + return value + + +def _optional_epoch(args: Mapping[str, Any]) -> Optional[str]: + value = args.get('epoch') + if value is None or value == '': + return None + if not _valid_id(value): + raise ProtocolError(ErrorCode.INVALID_ARGS, + f'epoch must be a printable string of 1-{MAX_ID_LENGTH} characters') + return str(value) + + +@dataclass(frozen=True) +class StateGetArgs: + """``state.get``: the display's state, as a versioned snapshot. + + With ``since`` and the ``epoch`` it came from, the answer is only + ``{changed: false, version, epoch, served_at, loop}`` while the state is + still at that version, so a poller that already has it is sent no state. + """ + since: Optional[int] = None + epoch: Optional[str] = None + + def to_dict(self) -> Dict[str, Any]: + return {'since': self.since, 'epoch': self.epoch} + + @classmethod + def from_dict(cls, args: Mapping[str, Any]) -> 'StateGetArgs': + return cls(since=_optional_version(args, 'since'), epoch=_optional_epoch(args)) + + +@dataclass(frozen=True) +class StateSubscribeArgs: + """``state.subscribe``: the snapshot now, then a push stream of changes. + + The response is the snapshot ``state.get`` returns. After it the + connection carries only :class:`StateEvent` lines from the display: a + ``state`` event whenever the state changes (always the latest version, + so a reader that falls behind skips versions instead of queueing them), + and a ``tick`` at least every :data:`SUBSCRIBE_KEEPALIVE_SECONDS`. + """ + + def to_dict(self) -> Dict[str, Any]: + return {} + + @classmethod + def from_dict(cls, args: Mapping[str, Any]) -> 'StateSubscribeArgs': + return cls() + + CommandArgs = Union[HelloArgs, OnDemandStartArgs, OnDemandStopArgs, NoArgs, - BrightnessSetArgs, PluginReloadArgs] + BrightnessSetArgs, PluginReloadArgs, StateGetArgs, StateSubscribeArgs] #: The arguments of a command that goes on the render thread's queue. QueuedArgs = Union[OnDemandStartArgs, OnDemandStopArgs, BrightnessSetArgs, PluginReloadArgs] @@ -491,6 +577,8 @@ _ARG_TYPES: Dict[str, Any] = { Command.ON_DEMAND_STATUS: NoArgs, Command.BRIGHTNESS_SET: BrightnessSetArgs, Command.PLUGIN_RELOAD: PluginReloadArgs, + Command.STATE_GET: StateGetArgs, + Command.STATE_SUBSCRIBE: StateSubscribeArgs, } @@ -563,6 +651,100 @@ class PluginReloadResult(TypedDict): modes: List[str] +class LoopState(TypedDict): + """``loop``: is the render loop still going round? + + ``heartbeat_age_seconds`` is the age of the render thread's last beat, + measured in memory by the display when it answered -- the same beat that + writes ``display-heartbeat.json``. None until the loop has drawn its first + frame. At ``stale_after`` or more the loop is stalled: the threshold + ``/api/v3/health`` uses. + """ + heartbeat_age_seconds: Optional[float] + armed: bool + stale_after: float + + +class StateSnapshot(TypedDict, total=False): + """The answer to ``state.get`` and ``state.subscribe``, and the + ``result`` of a ``state`` event. + + ``version`` counts changes to the state within one ``epoch`` (one run of + the display process): a reader that sees a new epoch starts over. + ``changed`` is False only for a ``state.get`` whose ``since`` is still + current, and then ``state`` is absent. ``served_at`` is the display's + wall clock when it answered. ``loop`` is measured at that moment, so it + is also inside ``state``. + + ``state`` holds the sections in :data:`STATE_SECTIONS`: + + * ``display``: what ``display_current_state`` holds (mode, plugin_id, + mode_index, total_modes, on_demand_active, is_display_active, + last_updated); + * ``on_demand``: what ``display_on_demand_state`` holds; + * ``brightness``: ``{brightness, panel_brightness, dimmed}``; + * ``plugins``: the plugin runtime snapshot (``plugin_runtime_snapshot``), + or None when there is none (or it was too large to send); + * ``loop``: :class:`LoopState`. + + A section the display has not published yet is None. + """ + schema: int + version: int + epoch: str + pid: int + served_at: float + changed: bool + state: Dict[str, Any] + loop: LoopState + + +class StateEventKind: + STATE = 'state' # result: a full StateSnapshot, the latest version + TICK = 'tick' # result: {version, epoch, pid, served_at, loop}; nothing changed + + +@dataclass(frozen=True) +class StateEvent: + """One message the display pushes to a subscriber. + + ``{"v": 1, "id": "", "event": "state" | "tick", + "result": {...}}``. It has no ``ok``, which is how a reader tells it from + a response. + """ + id: str + event: str + result: Dict[str, Any] + v: int = PROTOCOL_VERSION + + def to_dict(self) -> Dict[str, Any]: + return {'v': self.v, 'id': self.id, 'event': self.event, 'result': dict(self.result)} + + @classmethod + def from_dict(cls, obj: Any) -> 'StateEvent': + """Validate an event. Raises :class:`ProtocolError` (BAD_REQUEST).""" + if not isinstance(obj, dict): + raise ProtocolError(ErrorCode.BAD_REQUEST, 'an event must be a JSON object') + version = obj.get('v') + if not _is_int(version): + raise ProtocolError(ErrorCode.BAD_REQUEST, 'v must be an integer') + raw_id = obj.get('id') + if not isinstance(raw_id, str): + raise ProtocolError(ErrorCode.BAD_REQUEST, 'id must be a string') + event = obj.get('event') + if event not in (StateEventKind.STATE, StateEventKind.TICK): + raise ProtocolError(ErrorCode.BAD_REQUEST, 'event must be "state" or "tick"') + result = obj.get('result') + if not isinstance(result, dict): + raise ProtocolError(ErrorCode.BAD_REQUEST, 'result must be a JSON object') + return cls(id=raw_id, event=event, result=result, v=version) + + +def is_event(obj: Any) -> bool: + """Whether a decoded message is a pushed event rather than a response.""" + return isinstance(obj, dict) and 'event' in obj and 'ok' not in obj + + 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 0240dc86..7075b76c 100644 --- a/src/ipc/server.py +++ b/src/ipc/server.py @@ -14,6 +14,13 @@ 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. +The state stream (stage 3): the display publishes what it is doing into a +:class:`StateHub`, in memory, and ``state.get`` / ``state.subscribe`` read +it. A subscriber's connection gives back its request slot, takes one of +:data:`~src.ipc.contract.MAX_SUBSCRIBERS`, and is pushed the latest version +on every change plus a keepalive tick, from its own thread: publishing never +waits for a reader, and a reader that stops reading is dropped. + Robustness rules, because this runs inside the display process: * every connection has its own daemon thread, at most :data:`MAX_CLIENTS` at @@ -37,6 +44,7 @@ server checks them again: root, its own user, or a member of that group. from __future__ import annotations +import json import logging import os import queue @@ -45,8 +53,9 @@ import stat import struct import threading import time +import uuid from dataclasses import dataclass, field -from typing import Any, Callable, Dict, FrozenSet, List, Mapping, Optional +from typing import Any, Callable, Dict, FrozenSet, Iterable, List, Mapping, Optional, Tuple from src.ipc.contract import ( AWAIT_SECONDS, @@ -55,8 +64,12 @@ from src.ipc.contract import ( DEFAULT_SOCKET_DIR, DEFAULT_SOCKET_PATH, MAX_MESSAGE_BYTES, + MAX_SUBSCRIBERS, PROTOCOL_VERSION, QUEUED_COMMANDS, + STATE_SCHEMA, + STATE_SECTIONS, + SUBSCRIBE_KEEPALIVE_SECONDS, SUPPORTED_VERSIONS, AckResult, BrightnessSetArgs, @@ -72,6 +85,9 @@ from src.ipc.contract import ( QueuedArgs, Request, Response, + StateEvent, + StateEventKind, + StateGetArgs, configured_socket_path, decode_message, dev_socket_path, @@ -181,6 +197,214 @@ class QueuedCommand: self.outcome.fail(code, message) +# -- the state stream (stage 3) ---------------------------------------------------------- + +#: How long a ``state.get`` keeps the display counting its readers as served +#: over the socket (:meth:`StateHub.readers_active`). A subscriber counts for +#: as long as it is connected. +READER_WINDOW_SECONDS = 60.0 + +#: Room kept for the envelope (``v``, ``id``, ``event``) around a snapshot, +#: within MAX_MESSAGE_BYTES. +_ENVELOPE_ROOM = 512 + +_MISSING = object() + +LoopProbe = Callable[[], Mapping[str, Any]] + + +def _fingerprint(value: Optional[Mapping[str, Any]], volatile: Iterable[str]) -> Any: + """What a section's version is judged on: the value minus its volatile keys + (timestamps that move on every publish without anything changing).""" + if value is None: + return None + skip = frozenset(volatile) + return {k: v for k, v in value.items() if k not in skip} if skip else dict(value) + + +def _unknown_loop() -> Dict[str, Any]: + return {'heartbeat_age_seconds': None, 'armed': False, 'stale_after': None} + + +def fit_snapshot(snapshot: Dict[str, Any]) -> Dict[str, Any]: + """``snapshot``, or a copy without the plugin runtime section when the + message would be over MAX_MESSAGE_BYTES (hundreds of plugins). The + reader then falls back to the cache for that section only; ``truncated`` + says which was left out.""" + state = snapshot.get('state') + if not isinstance(state, dict) or state.get('plugins') is None: + return snapshot + try: + size = len(json.dumps(snapshot, separators=(',', ':'), ensure_ascii=True, + allow_nan=False)) + except (TypeError, ValueError): + size = MAX_MESSAGE_BYTES + if size <= MAX_MESSAGE_BYTES - _ENVELOPE_ROOM: + return snapshot + logger.warning("State snapshot is %d bytes; sending it without the plugin runtime " + "section", size) + trimmed = dict(snapshot) + trimmed['state'] = dict(state, plugins=None) + trimmed['truncated'] = ['plugins'] + return trimmed + + +class StateHub: + """The display's live state, in memory, for ``state.get`` and ``state.subscribe``. + + Writers publish whole sections (:meth:`publish`): the render thread + publishes ``display``, ``on_demand`` and ``brightness``, and the plugin + runtime publisher's thread publishes ``plugins``. Each section has one + writer. ``loop`` is not published: it is measured when a reader asks + (``loop_probe``), so it keeps ageing while the render thread is stuck. + + The version goes up when a section's value changes, ignoring the keys + the publisher names as volatile (timestamps). Publishing never blocks on + a reader: the lock is held only to swap a dict reference and compare it, + and every socket write happens on the reader's own thread, outside it. + A reader that is slow gets the latest version when it next asks, not + every version in between. + """ + + def __init__(self, loop_probe: Optional[LoopProbe] = None, *, + clock: Callable[[], float] = time.monotonic, + wall_clock: Callable[[], float] = time.time, + epoch: Optional[str] = None, pid: Optional[int] = None, + reader_window: float = READER_WINDOW_SECONDS): + self._cond = threading.Condition(threading.Lock()) + self._sections: Dict[str, Optional[Dict[str, Any]]] = {} + self._fingerprints: Dict[str, Any] = {} + self._version = 0 + self.epoch = epoch or uuid.uuid4().hex[:16] + self.pid = os.getpid() if pid is None else pid + self._loop_probe = loop_probe + self._clock = clock + self._wall_clock = wall_clock + self._reader_window = reader_window + self._last_read: Optional[float] = None + self._subscribers = 0 + + @property + def version(self) -> int: + return self._version + + @property + def subscribers(self) -> int: + return self._subscribers + + # -- writers ------------------------------------------------------------- + + def publish(self, section: str, value: Optional[Mapping[str, Any]], + volatile: Iterable[str] = ()) -> bool: + """Store a section's latest value; True when that is a new version. + + The value is copied (one level), so the caller may reuse its dict. + """ + stored = None if value is None else dict(value) + fingerprint = _fingerprint(stored, volatile) + with self._cond: + self._sections[section] = stored + if self._fingerprints.get(section, _MISSING) == fingerprint: + return False + self._fingerprints[section] = fingerprint + self._version += 1 + self._cond.notify_all() + return True + + def wake(self) -> None: + """Wake every waiting reader (the server is closing).""" + with self._cond: + self._cond.notify_all() + + # -- readers ------------------------------------------------------------- + + def loop(self) -> Dict[str, Any]: + """The render loop's liveness now. Never raises.""" + if self._loop_probe is None: + return _unknown_loop() + try: + return dict(self._loop_probe()) + except Exception: # pylint: disable=broad-except + logger.debug("Render loop liveness probe failed", exc_info=True) + return _unknown_loop() + + def snapshot(self, since: Optional[int] = None, + epoch: Optional[str] = None) -> Dict[str, Any]: + """The :class:`~src.ipc.contract.StateSnapshot` now. + + ``since`` with this hub's ``epoch``, still the current version, gives + the short ``changed: false`` form. + """ + with self._cond: + version = self._version + sections = dict(self._sections) + loop = self.loop() + result: Dict[str, Any] = { + 'schema': STATE_SCHEMA, + 'version': version, + 'epoch': self.epoch, + 'pid': self.pid, + 'served_at': self._wall_clock(), + 'loop': loop, + } + if since is not None and epoch == self.epoch and since == version: + result['changed'] = False + return result + state: Dict[str, Any] = {name: sections.get(name) for name in STATE_SECTIONS + if name != 'loop'} + state['loop'] = loop + result['changed'] = True + result['state'] = state + return result + + def wait_for_change(self, version: int, timeout: float, + stop: Optional[threading.Event] = None) -> bool: + """Block up to ``timeout`` for a version other than ``version``.""" + with self._cond: + self._cond.wait_for( + lambda: self._version != version or (stop is not None and stop.is_set()), + timeout) + return self._version != version + + # -- who is reading ------------------------------------------------------ + + def note_read(self) -> None: + self._last_read = self._clock() + + def subscriber_joined(self) -> None: + with self._cond: + self._subscribers += 1 + + def subscriber_left(self) -> None: + with self._cond: + self._subscribers = max(0, self._subscribers - 1) + self._last_read = self._clock() + + def readers_active(self) -> bool: + """Is the socket serving state readers? A subscriber is connected, or a + ``state.get`` came within the reader window. The display uses this to + write the cache copies of the same state less often.""" + if self._subscribers > 0: + return True + last = self._last_read + return last is not None and self._clock() - last < self._reader_window + + +class _Slot: + """A connection slot, released once (a subscriber gives its back early).""" + + def __init__(self, semaphore: threading.BoundedSemaphore): + self._semaphore = semaphore + self._held = True + self._lock = threading.Lock() + + def release(self) -> None: + with self._lock: + if self._held: + self._held = False + self._semaphore.release() + + # -- peer credentials ------------------------------------------------------------------ @dataclass(frozen=True) @@ -315,8 +539,14 @@ class ControlServer: message_timeout: float = MESSAGE_TIMEOUT_SECONDS, idle_timeout: float = IDLE_TIMEOUT_SECONDS, check_peer: bool = True, - await_seconds: Optional[Mapping[str, float]] = None): + await_seconds: Optional[Mapping[str, float]] = None, + state_hub: Optional[StateHub] = None, + max_subscribers: int = MAX_SUBSCRIBERS, + keepalive: float = SUBSCRIBE_KEEPALIVE_SECONDS): self.path = path + self.state_hub = state_hub + self._subscriber_slots = threading.BoundedSemaphore(max_subscribers) + self._keepalive = keepalive self._await_seconds: Dict[str, float] = dict(AWAIT_SECONDS) if await_seconds: self._await_seconds.update(await_seconds) @@ -378,6 +608,8 @@ class ControlServer: """Stop accepting and remove the socket file (only if it is still ours).""" self._stopping.set() self._close_socket() + if self.state_hub is not None: + self.state_hub.wake() # subscribers see _stopping and hang up thread = self._thread if thread is not None and thread is not threading.current_thread(): thread.join(timeout=2.0) @@ -549,6 +781,7 @@ class ControlServer: def _serve(self, conn: socket.socket) -> None: """One connection: authenticate, then answer requests until it ends.""" + slot = _Slot(self._slots) try: conn.settimeout(self._io_timeout) peer = peer_credentials(conn) @@ -557,7 +790,7 @@ class ControlServer: "this user or group %s", peer.pid, peer.uid, peer.gid, self._group) self._send(conn, Response.failure(None, ErrorCode.FORBIDDEN, 'not permitted')) return - self._read_requests(conn, peer) + self._read_requests(conn, peer, slot) except Exception: # pylint: disable=broad-except logger.exception("Control socket connection failed") finally: @@ -565,7 +798,7 @@ class ControlServer: conn.close() except OSError: pass - self._slots.release() + slot.release() def _peer_ok(self, peer: PeerCredentials) -> bool: groups = None @@ -573,7 +806,8 @@ class ControlServer: groups = process_groups(peer.pid) return peer_allowed(peer, self._own_uid, self._group, groups) - def _read_requests(self, conn: socket.socket, peer: Optional[PeerCredentials]) -> None: + def _read_requests(self, conn: socket.socket, peer: Optional[PeerCredentials], + slot: Optional[_Slot] = None) -> None: reader = FrameReader(MAX_MESSAGE_BYTES) idle_since = time.monotonic() message_started: Optional[float] = None @@ -598,7 +832,13 @@ class ControlServer: self._send(conn, Response.failure(None, e.code, e.message)) return # can't find the next message boundary: hang up for line in lines: - if not self._send(conn, self.handle_line(line, peer)): + response, cmd = self._handle(line, peer) + if cmd == Command.STATE_SUBSCRIBE and response.ok: + # The connection becomes a one-way stream; anything the + # client sent after the subscribe is ignored. + self._subscribe(conn, response, slot) + return + if not self._send(conn, response): return if reader.pending: if message_started is None or lines: @@ -624,20 +864,91 @@ class ControlServer: # -- requests -------------------------------------------------------------------- def handle_line(self, line: bytes, peer: Optional[PeerCredentials] = None) -> Response: - """Answer one request line. Never raises.""" + """Answer one request line. Never raises. + + A ``state.subscribe`` answered here gets its snapshot only; the + stream that follows needs a connection (``_read_requests``). + """ + return self._handle(line, peer)[0] + + def _handle(self, line: bytes, + peer: Optional[PeerCredentials]) -> Tuple[Response, Optional[str]]: + """The response to one line, and the command it answered (when known).""" request_id: Optional[str] = None + cmd: Optional[str] = None try: obj = decode_message(line) raw_id = obj.get('id') request_id = raw_id if isinstance(raw_id, str) and len(raw_id) <= 128 else None request = Request.from_dict(obj) request_id = request.id - return self._dispatch(request, peer) + cmd = request.cmd + return self._dispatch(request, peer), cmd except ProtocolError as e: - return Response.failure(e.request_id or request_id, e.code, e.message) + return Response.failure(e.request_id or request_id, e.code, e.message), cmd except Exception: # pylint: disable=broad-except logger.exception("Control socket handler failed") - return Response.failure(request_id, ErrorCode.INTERNAL, 'internal error') + return Response.failure(request_id, ErrorCode.INTERNAL, 'internal error'), cmd + + # -- the state stream ---------------------------------------------------------- + + def _subscribe(self, conn: socket.socket, response: Response, + slot: Optional[_Slot]) -> None: + """Answer a ``state.subscribe`` and push state events until it ends. + + Subscribers have their own bound (MAX_SUBSCRIBERS) and give their + request slot back, so a few browsers watching never use up the slots + commands need. Everything here runs on this connection's thread: a + reader that does not keep up only stalls its own sends, and one that + stops reading for a whole IO timeout is dropped. The render thread + only ever publishes into the hub. + """ + hub = self.state_hub + if hub is None or not self._subscriber_slots.acquire(blocking=False): + self._send(conn, Response.failure(response.id, ErrorCode.BUSY, + 'too many state subscribers', v=response.v)) + return + if slot is not None: + slot.release() + hub.subscriber_joined() + try: + if not self._send(conn, response): + return + result = response.result or {} + version = result.get('version', -1) + sub_id = response.id or '' + logger.debug("Control socket: state subscriber joined at version %s", version) + while not self._stopping.is_set(): + hub.wait_for_change(version, self._keepalive, self._stopping) + if self._stopping.is_set(): + return + snap = hub.snapshot(since=version, epoch=hub.epoch) + if snap.get('changed'): + version = snap['version'] + event = StateEvent(sub_id, StateEventKind.STATE, fit_snapshot(snap), + v=response.v) + else: + event = StateEvent(sub_id, StateEventKind.TICK, snap, v=response.v) + if not self._send_event(conn, event): + return + finally: + hub.subscriber_left() + self._subscriber_slots.release() + + def _send_event(self, conn: socket.socket, event: StateEvent) -> bool: + try: + data = encode_message(event.to_dict()) + except ProtocolError as e: + logger.error("Control socket state event not sent: %s", e.message) + return False + try: + conn.sendall(data) + return True + except socket.timeout: + logger.info("Control socket: dropping a state subscriber that stopped reading") + return False + except OSError: + return False def _dispatch(self, request: Request, peer: Optional[PeerCredentials]) -> Response: if request.cmd == Command.HELLO: @@ -678,6 +989,18 @@ class ControlServer: v=request.v) return Response.success(request.id, self._status_provider(), v=request.v) + if request.cmd in (Command.STATE_GET, Command.STATE_SUBSCRIBE): + hub = self.state_hub + if hub is None: + return Response.failure(request.id, ErrorCode.INTERNAL, 'no state available', + v=request.v) + if isinstance(args, StateGetArgs): + hub.note_read() + snap = hub.snapshot(since=args.since, epoch=args.epoch) + else: + snap = hub.snapshot() + return Response.success(request.id, fit_snapshot(snap), v=request.v) + if request.cmd in QUEUED_COMMANDS and isinstance(args, ( OnDemandStartArgs, OnDemandStopArgs, BrightnessSetArgs, PluginReloadArgs)): awaited = request.cmd in AWAITED_COMMANDS @@ -728,22 +1051,25 @@ class ControlServer: def start_control_server(status_provider: Optional[StatusProvider] = None, cache_dir: Optional[str] = None, - environ: Optional[Mapping[str, str]] = None) -> Optional[ControlServer]: + environ: Optional[Mapping[str, str]] = None, + state_hub: Optional[StateHub] = None) -> Optional[ControlServer]: """Start the display's control socket, or return None when it can't run. None covers Windows, ``LEDMATRIX_CONTROL_SOCKET=off`` and any failure to - bind; in every case the web interface falls back to the file mailbox. + bind; in every case the web interface falls back to the file mailbox + and to the cache keys the display still writes. """ path = server_socket_path(environ) if path is None: logger.debug("Control socket disabled or unsupported here; using the file mailbox only") return None - server = ControlServer(path, status_provider, resolve_socket_group(cache_dir)) + server = ControlServer(path, status_provider, resolve_socket_group(cache_dir), + state_hub=state_hub) return server if server.start() else None __all__ = [ - 'CommandOutcome', 'ControlServer', 'PeerCredentials', 'QueuedCommand', 'StatusProvider', - 'peer_allowed', 'peer_credentials', 'process_groups', 'resolve_socket_group', - 'server_socket_path', 'start_control_server', 'PROTOCOL_VERSION', + 'CommandOutcome', 'ControlServer', 'PeerCredentials', 'QueuedCommand', 'StateHub', + 'StatusProvider', 'fit_snapshot', 'peer_allowed', 'peer_credentials', 'process_groups', + 'resolve_socket_group', 'server_socket_path', 'start_control_server', 'PROTOCOL_VERSION', ] diff --git a/src/plugin_system/plugin_runtime.py b/src/plugin_system/plugin_runtime.py index 50796734..2a541396 100644 --- a/src/plugin_system/plugin_runtime.py +++ b/src/plugin_system/plugin_runtime.py @@ -36,6 +36,13 @@ tmpfs. A missing heartbeat (dev server, emulator, Windows, a display still starting up) or one from another process (a display restarted after a watchdog kill) says nothing, and the snapshot is judged on its own. +The control socket. Where the display serves its state stream (stage 3, +docs/IPC_CONTROL_SOCKET.md), every tick also hands the snapshot to it, in +memory, and the web interface reads it there first +(``view_from_socket_state``, judged by the same rules). While the socket +serves those readers, the cache copy is their fallback and an unchanged +snapshot is rewritten every ``RELAXED_REFRESH_INTERVAL`` instead. + A dead publisher. systemd removes the heartbeat's directory when the service stops, so after a watchdog kill there is no heartbeat to go stale. The reader then asks whether the snapshot's ``pid`` still exists (POSIX @@ -48,7 +55,7 @@ import math import os import threading import time -from dataclasses import dataclass, field +from dataclasses import dataclass, field, replace from typing import Any, Callable, Dict, Optional from src import display_watchdog @@ -75,6 +82,15 @@ TICK_INTERVAL = 5.0 #: A snapshot older than this is stale: three missed refreshes. STALE_AFTER = 3 * REFRESH_INTERVAL +#: The refresh while the control socket serves the web interface's readers +#: (``StateHub.readers_active``). The cache copy is then only their fallback, +#: so an unchanged snapshot is rewritten half as often; the snapshot says so +#: in its own ``refresh_interval`` and ``stale_after``. +RELAXED_REFRESH_INTERVAL = 2 * REFRESH_INTERVAL + +#: The control socket's section for this snapshot (``state.plugins``). +STATE_SECTION = "plugins" + #: Bounds on a published ``stale_after``, so a corrupt value can make a #: reader neither trust a dead display for hours nor distrust a live one. _STALE_AFTER_MIN = 30.0 @@ -140,7 +156,8 @@ def summarize_error(error_info: Optional[Dict[str, Any]]) -> Optional[Dict[str, def build_runtime_snapshot(state_manager: Any, *, started_at: float, now: Optional[float] = None, - running: bool = True) -> Dict[str, Any]: + running: bool = True, + refresh_interval: float = REFRESH_INTERVAL) -> Dict[str, Any]: """The snapshot for ``state_manager`` (a plugin_state.PluginStateManager). A stopped snapshot (``running=False``) lists no plugins: nothing is @@ -162,8 +179,8 @@ def build_runtime_snapshot(state_manager: Any, *, started_at: float, "running": running, "published_at": time.time() if now is None else now, "started_at": started_at, - "refresh_interval": REFRESH_INTERVAL, - "stale_after": STALE_AFTER, + "refresh_interval": refresh_interval, + "stale_after": 3 * refresh_interval, "pid": os.getpid(), "plugins": plugins, } @@ -197,30 +214,87 @@ class PluginRuntimePublisher: self._tick_lock = threading.Lock() self._stop = threading.Event() self._thread: Optional[threading.Thread] = None + # The control socket's state stream (src/ipc/server.StateHub), when + # the display serves one: every tick also hands it the snapshot, in + # memory, and the cache refresh relaxes while it has readers. + self._hub: Any = None + self._hub_change: Optional[int] = None + self._hub_snapshot: Optional[Dict[str, Any]] = None + self.relaxed_refresh_interval = RELAXED_REFRESH_INTERVAL - def _write(self, running: bool) -> None: - snapshot = build_runtime_snapshot(self.state_manager, started_at=self.started_at, - now=self._wall_clock(), running=running) + def attach_hub(self, hub: Any) -> None: + """Also publish to the control socket's state hub, starting now.""" + with self._tick_lock: + self._hub = hub + self._hub_change = None + self._hub_snapshot = None + try: + self._push_to_hub(self.state_manager.change_count) + except Exception as err: # never let reporting break the display + logger.debug("Could not publish the plugin runtime state: %s", err, + exc_info=True) + + def _push_to_hub(self, change: int) -> None: + """The snapshot to the state hub: rebuilt when the state machine + changed, otherwise the last one with a new ``published_at``, which + the hub does not count as a new version. In memory, every tick, so + the socket's copy is never more than a tick old.""" + hub = self._hub + if hub is None: + return + now = self._wall_clock() + if self._hub_snapshot is None or change != self._hub_change: + snapshot = build_runtime_snapshot(self.state_manager, started_at=self.started_at, + now=now) + else: + snapshot = dict(self._hub_snapshot, published_at=now) + hub.publish(STATE_SECTION, snapshot, volatile=("published_at",)) + self._hub_snapshot = snapshot + self._hub_change = change + + def _cache_refresh_interval(self) -> float: + """The cache refresh: relaxed while the socket serves the readers.""" + hub = self._hub + try: + if hub is not None and hub.readers_active(): + return self.relaxed_refresh_interval + except Exception: # pylint: disable=broad-except + # The normal interval is the safe answer: it only writes more. + logger.debug("State hub readers_active() failed; using the normal refresh", exc_info=True) + return self.refresh_interval + + def _write(self, running: bool, refresh_interval: Optional[float] = None) -> None: + snapshot = build_runtime_snapshot( + self.state_manager, started_at=self.started_at, now=self._wall_clock(), + running=running, + refresh_interval=self.refresh_interval if refresh_interval is None + else refresh_interval) self.cache_manager.set(PLUGIN_RUNTIME_KEY, snapshot) def tick(self) -> bool: """Publish if something changed (throttled) or the refresh is due. - True if a snapshot was written.""" + True if a snapshot was written to the cache.""" with self._tick_lock: try: change = self.state_manager.change_count + try: + self._push_to_hub(change) + except Exception as err: # the cache copy still goes out below + logger.debug("Could not publish the plugin runtime state: %s", err, + exc_info=True) now = self._clock() + refresh = self._cache_refresh_interval() since = None if self._last_attempt is None else now - self._last_attempt if since is not None: if change == self._published_change: - if since < self.refresh_interval: + if since < refresh: return False elif since < self.min_interval: return False # Stamp the attempt before writing: a cache that keeps failing # is retried at the throttled rate, not on every tick. self._last_attempt = now - self._write(running=True) + self._write(running=True, refresh_interval=refresh) self._published_change = change return True except Exception as err: # never let reporting break the display @@ -317,6 +391,9 @@ class PluginRuntimeView: plugins: Dict[str, Dict[str, Any]] = field(default_factory=dict) #: Age of the render loop's heartbeat, when it was taken into account. heartbeat_age_seconds: Optional[float] = None + #: Where the snapshot came from: ``cache`` (the shared cache file and the + #: heartbeat file) or ``socket`` (the control socket's state stream). + source: str = "cache" @property def live(self) -> bool: @@ -348,6 +425,7 @@ class PluginRuntimeView: "stale_after": self.stale_after, "heartbeat_age_seconds": (None if self.heartbeat_age_seconds is None else round(self.heartbeat_age_seconds, 1)), + "source": self.source, } @@ -445,6 +523,36 @@ def view_from_snapshot(snapshot: Any, now: Optional[float] = None, ) +def view_from_socket_state(snapshot: Any, now: Optional[float] = None, + now_mono: Optional[float] = None) -> Optional[PluginRuntimeView]: + """Judge the ``plugins`` section of a control-socket state snapshot by + the same rules as the cache copy; None when it has none (an older + display, or a snapshot too large to carry it), so the caller reads the + cache instead. + + The display measured its render loop's heartbeat age when it answered + (``state.loop``); that is the heartbeat here, aged by the time since the + answer arrived. A live snapshot with a stalled loop is ``stalled``, and a + snapshot older than its ``stale_after`` (the publisher thread stopped) + is ``stale``, exactly as for the cache. The display answered, so its + process is alive: there is no pid check. + """ + from src.ipc.client import snapshot_loop_age # stdlib-only module + if not isinstance(snapshot, dict): + return None + state = snapshot.get("state") + plugins = state.get(STATE_SECTION) if isinstance(state, dict) else None + if not isinstance(plugins, dict): + return None + now_mono = time.monotonic() if now_mono is None else now_mono + beat_age = snapshot_loop_age(snapshot, now_mono=now_mono) + heartbeat = None + if beat_age is not None: + heartbeat = {"pid": plugins.get("pid"), "mono": now_mono - beat_age} + view = view_from_snapshot(plugins, now=now, heartbeat=heartbeat, now_mono=now_mono) + return replace(view, source="socket") + + def read_plugin_runtime(cache_manager: Any, now: Optional[float] = None, heartbeat_path: Optional[str] = None) -> PluginRuntimeView: """The display's latest snapshot, judged for staleness and against the diff --git a/test/test_ipc_state_stream.py b/test/test_ipc_state_stream.py new file mode 100644 index 00000000..f72a68f9 --- /dev/null +++ b/test/test_ipc_state_stream.py @@ -0,0 +1,491 @@ +"""Control socket stage 3: the display's state over the socket (state.get, +state.subscribe) instead of polled cache keys. + +* ``StateHub`` and the request handling are plain Python, tested on every + platform: versions move only when something a reader sees changes, the + ``since``/``epoch`` short answer, an oversized snapshot drops only its + plugin section, and publishing never waits for a reader. +* ``TestLiveStream`` needs AF_UNIX (Linux, WSL, a Pi; skipped on Windows): a + real server pushing to real subscribers -- several at once, one that never + reads, one whose display goes away and comes back. +""" + +import json +import socket +import threading +import time + +import pytest + +from src.ipc import client +from src.ipc import contract as c +from src.ipc import server as srv +from src.ipc.contract import Command, ErrorCode +from src.ipc.server import ControlServer, StateHub, fit_snapshot + +needs_unix_sockets = pytest.mark.skipif(not c.socket_supported(), + reason='AF_UNIX sockets are Linux/macOS only') + + +class FakeClock: + def __init__(self, now=1000.0): + self.now = now + + def __call__(self): + return self.now + + +def _req(cmd, args=None, rid='r1', v=1): + return json.dumps({'v': v, 'id': rid, 'cmd': cmd, 'args': args or {}}).encode() + + +def _loop(age=1.0): + return lambda: {'heartbeat_age_seconds': age, 'armed': age is not None, + 'stale_after': 60.0} + + +def _display(mode='clock', active=True, updated=None): + return {'mode': mode, 'plugin_id': mode, 'mode_index': 0, 'total_modes': 2, + 'on_demand_active': False, 'is_display_active': active, + 'last_updated': time.time() if updated is None else updated} + + +@pytest.fixture +def hub(): + h = StateHub(loop_probe=_loop(), epoch='e1', pid=4242) + h.publish('display', _display(), volatile=('last_updated',)) + return h + + +# --- the hub --------------------------------------------------------------------- + +class TestStateHub: + def test_snapshot_has_every_section_and_the_envelope(self, hub): + snap = hub.snapshot() + assert snap['schema'] == c.STATE_SCHEMA + assert (snap['version'], snap['epoch'], snap['pid']) == (1, 'e1', 4242) + assert snap['changed'] is True + assert list(snap['state']) == list(c.STATE_SECTIONS) + assert snap['state']['display']['mode'] == 'clock' + assert snap['state']['on_demand'] is None # not published yet + assert snap['state']['loop'] == snap['loop'] == _loop()() + + def test_only_a_real_change_is_a_new_version(self, hub): + assert not hub.publish('display', _display(updated=1.0), volatile=('last_updated',)) + assert hub.version == 1 + # ...but the latest timestamp is what a reader gets. + assert hub.snapshot()['state']['display']['last_updated'] == 1.0 + assert hub.publish('display', _display(mode='weather'), volatile=('last_updated',)) + assert hub.version == 2 + assert hub.publish('brightness', {'brightness': 50}) + assert not hub.publish('brightness', {'brightness': 50}) + assert hub.version == 3 + + def test_the_published_dict_is_copied(self, hub): + value = {'brightness': 50} + hub.publish('brightness', value) + value['brightness'] = 10 + assert hub.snapshot()['state']['brightness'] == {'brightness': 50} + + def test_since_the_current_version_is_the_short_answer(self, hub): + short = hub.snapshot(since=1, epoch='e1') + assert short['changed'] is False and 'state' not in short + assert short['version'] == 1 and short['loop']['heartbeat_age_seconds'] == 1.0 + # An older version, or another epoch (a restarted display): the full state. + assert hub.snapshot(since=0, epoch='e1')['changed'] is True + assert hub.snapshot(since=1, epoch='other')['changed'] is True + assert hub.snapshot(since=1)['changed'] is True + + def test_each_display_run_has_its_own_epoch(self): + assert StateHub().epoch != StateHub().epoch + + def test_the_loop_is_measured_when_asked(self): + ages = iter([1.0, 61.0]) + hub = StateHub(loop_probe=lambda: {'heartbeat_age_seconds': next(ages), + 'armed': True, 'stale_after': 60.0}) + assert hub.snapshot()['loop']['heartbeat_age_seconds'] == 1.0 + # Nothing was published -- a stuck render thread publishes nothing -- + # and the age still moves. + assert hub.snapshot()['loop']['heartbeat_age_seconds'] == 61.0 + + def test_a_failing_probe_is_unknown_not_raised(self): + def boom(): + raise RuntimeError('x') + assert StateHub(loop_probe=boom).loop()['heartbeat_age_seconds'] is None + assert StateHub().loop()['armed'] is False + + def test_wait_for_change(self, hub): + assert hub.wait_for_change(1, 0.01) is False + threading.Timer(0.05, lambda: hub.publish('brightness', {'b': 1})).start() + start = time.monotonic() + assert hub.wait_for_change(1, 5.0) is True + assert time.monotonic() - start < 2.0 + + def test_wait_ends_on_stop(self, hub): + stop = threading.Event() + + def end(): + stop.set() + hub.wake() + threading.Timer(0.05, end).start() + start = time.monotonic() + assert hub.wait_for_change(1, 5.0, stop) is False + assert time.monotonic() - start < 2.0 + + def test_readers_active(self): + clock = FakeClock() + hub = StateHub(clock=clock, reader_window=60.0) + assert not hub.readers_active() + hub.note_read() + clock.now += 59 + assert hub.readers_active() + clock.now += 2 + assert not hub.readers_active() + hub.subscriber_joined() + clock.now += 1000 + assert hub.readers_active() # for as long as it is connected + hub.subscriber_left() + assert hub.readers_active() # and one window after + clock.now += 61 + assert not hub.readers_active() + + def test_publishing_never_waits_for_readers(self, hub): + """Many readers snapshotting and waiting at once: every publish from + the 'render thread' still returns in well under a frame.""" + stop = threading.Event() + + def reader(): + version = 0 + while not stop.is_set(): + hub.wait_for_change(version, 0.01) + version = hub.snapshot()['version'] + threads = [threading.Thread(target=reader, daemon=True) for _ in range(8)] + for t in threads: + t.start() + worst = 0.0 + try: + for i in range(500): + start = time.perf_counter() + hub.publish('display', _display(mode=f'm{i}'), volatile=('last_updated',)) + worst = max(worst, time.perf_counter() - start) + finally: + stop.set() + for t in threads: + t.join(2) + assert worst < 0.05 + + +class TestFitSnapshot: + def test_small_is_unchanged(self, hub): + snap = hub.snapshot() + assert fit_snapshot(snap) is snap + + def test_too_large_drops_only_the_plugins(self, hub): + plugins = {f'p{i}': {'error': {'message': 'x' * 200}} for i in range(400)} + hub.publish('plugins', {'schema': 1, 'plugins': plugins}) + fitted = fit_snapshot(hub.snapshot()) + assert fitted['state']['plugins'] is None + assert fitted['truncated'] == ['plugins'] + assert fitted['state']['display']['mode'] == 'clock' + c.encode_message({'v': 1, 'id': 'x' * 128, 'ok': True, 'result': fitted}) + + +# --- the commands, without a socket ------------------------------------------------ + +class TestHandleLine: + def test_state_get(self, hub): + server = ControlServer('/unused.sock', state_hub=hub) + resp = server.handle_line(_req(Command.STATE_GET)) + assert resp.ok and resp.result['version'] == 1 + assert resp.result['state']['display']['mode'] == 'clock' + assert hub.readers_active() + + def test_state_get_since(self, hub): + server = ControlServer('/unused.sock', state_hub=hub) + resp = server.handle_line(_req(Command.STATE_GET, {'since': 1, 'epoch': 'e1'})) + assert resp.ok and resp.result['changed'] is False and 'state' not in resp.result + + @pytest.mark.parametrize('args', [{'since': -1}, {'since': 'x'}, {'since': True}, + {'epoch': 5}, {'epoch': 'x' * 200}]) + def test_bad_args(self, hub, args): + server = ControlServer('/unused.sock', state_hub=hub) + resp = server.handle_line(_req(Command.STATE_GET, args)) + assert not resp.ok and resp.error.code == ErrorCode.INVALID_ARGS + + def test_without_a_hub_it_is_an_error_not_a_crash(self): + server = ControlServer('/unused.sock') + for cmd in (Command.STATE_GET, Command.STATE_SUBSCRIBE): + resp = server.handle_line(_req(cmd)) + assert not resp.ok and resp.error.code == ErrorCode.INTERNAL + + def test_hello_lists_the_new_commands(self, hub): + resp = ControlServer('/unused.sock', state_hub=hub).handle_line( + _req(Command.HELLO, {'versions': [1]})) + assert {Command.STATE_GET, Command.STATE_SUBSCRIBE} <= set(resp.result['commands']) + + def test_state_commands_are_not_queued(self, hub): + server = ControlServer('/unused.sock', state_hub=hub) + server.handle_line(_req(Command.STATE_GET)) + assert not server.has_pending and server.drain() == [] + + +class TestEvents: + def test_round_trip(self): + event = c.StateEvent('s1', c.StateEventKind.TICK, {'version': 3}) + obj = json.loads(c.encode_message(event.to_dict())) + assert c.is_event(obj) + assert c.StateEvent.from_dict(obj) == event + + def test_a_response_is_not_an_event(self): + assert not c.is_event(c.Response.success('x', {}).to_dict()) + + @pytest.mark.parametrize('obj', [ + [], {'v': 1, 'id': 's', 'event': 'other', 'result': {}}, + {'v': 1, 'id': 's', 'event': 'state', 'result': []}, + {'v': '1', 'id': 's', 'event': 'state', 'result': {}}, + {'v': 1, 'id': None, 'event': 'state', 'result': {}}, + ]) + def test_bad_events_are_refused(self, obj): + with pytest.raises(c.ProtocolError): + c.StateEvent.from_dict(obj) + + +class TestSubscriptionStore: + """StateSubscription's bookkeeping, without a socket.""" + + def test_latest_is_none_until_a_snapshot_and_after_silence(self, hub): + clock = FakeClock() + sub = client.StateSubscription(paths=['/nowhere'], silence=15.0, clock=clock) + assert sub.latest() is None + sub._store(hub.snapshot(), full=True) + latest = sub.latest() + assert latest['version'] == 1 and latest['received_mono'] == clock.now + clock.now += 16 + assert sub.latest() is None + + def test_a_tick_refreshes_the_loop_and_keeps_the_state(self, hub): + clock = FakeClock() + sub = client.StateSubscription(paths=['/nowhere'], clock=clock) + sub._store(hub.snapshot(), full=True) + clock.now += 10 + tick = {'version': 1, 'epoch': 'e1', 'served_at': 5.0, + 'loop': {'heartbeat_age_seconds': 70.0, 'armed': True, 'stale_after': 60.0}} + sub._store(tick, full=False) + latest = sub.latest() + assert latest['state']['display']['mode'] == 'clock' + assert latest['state']['loop']['heartbeat_age_seconds'] == 70.0 + assert latest['received_mono'] == clock.now + + def test_a_tick_from_another_epoch_is_ignored(self, hub): + clock = FakeClock() + sub = client.StateSubscription(paths=['/nowhere'], clock=clock) + sub._store(hub.snapshot(), full=True) + clock.now += 10 + sub._store({'version': 9, 'epoch': 'other', 'loop': {}}, full=False) + assert sub.latest()['received_mono'] == clock.now - 10 + + def test_loop_age_counts_the_time_since_it_arrived(self, hub): + snap = dict(hub.snapshot(), received_mono=100.0) + assert client.snapshot_loop_age(snap, now_mono=104.0) == 5.0 + snap['loop'] = {'heartbeat_age_seconds': None} + assert client.snapshot_loop_age(snap, now_mono=104.0) is None + + +# --- a real socket ------------------------------------------------------------------ + +def _wait_until(predicate, timeout=5.0): + deadline = time.monotonic() + timeout + while time.monotonic() < deadline: + if predicate(): + return True + time.sleep(0.01) + return predicate() + + +@needs_unix_sockets +class TestLiveStream: + @pytest.fixture + def path(self, tmp_path): + return str(tmp_path / 'control.sock') + + @pytest.fixture + def live(self, path, hub): + server = ControlServer(path, state_hub=hub, keepalive=0.2, io_timeout=0.5, + max_clients=2, max_subscribers=3) + assert server.start() + yield server, hub + server.close() + + def _subscribe(self, path): + sock = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM) + sock.settimeout(5) + sock.connect(path) + sock.sendall(c.encode_message({'v': 1, 'id': 's1', 'cmd': Command.STATE_SUBSCRIBE, + 'args': {}})) + return sock, c.FrameReader() + + def _next(self, sock, reader, pending): + while not pending: + data = sock.recv(65536) + assert data, 'the display hung up' + pending.extend(reader.feed(data)) + return json.loads(pending.pop(0)) + + def test_state_get_over_the_socket(self, live, path): + snap = client.state_get(paths=[path]) + assert snap['version'] == 1 and snap['state']['display']['mode'] == 'clock' + short = client.state_get(since=snap['version'], epoch=snap['epoch'], paths=[path]) + assert short['changed'] is False + + def test_subscribe_answers_then_pushes_changes_in_order(self, live, path): + _, hub = live + sock, reader = self._subscribe(path) + pending = [] + first = self._next(sock, reader, pending) + assert first['ok'] is True and first['id'] == 's1' and first['result']['version'] == 1 + hub.publish('display', _display(mode='weather'), volatile=('last_updated',)) + versions = [] + while True: + msg = self._next(sock, reader, pending) + assert c.is_event(msg) and msg['id'] == 's1' + if msg['event'] == 'state': + versions.append(msg['result']['version']) + assert msg['result']['state']['display']['mode'] == 'weather' + break + hub.publish('brightness', {'brightness': 30}) + while True: + msg = self._next(sock, reader, pending) + if msg['event'] == 'state': + versions.append(msg['result']['version']) + break + assert versions == [2, 3] + sock.close() + + def test_ticks_keep_a_quiet_subscription_alive(self, live, path): + sock, reader = self._subscribe(path) + pending = [] + self._next(sock, reader, pending) + kinds = [self._next(sock, reader, pending)['event'] for _ in range(3)] + assert kinds == ['tick', 'tick', 'tick'] + sock.close() + + def test_a_burst_is_coalesced_to_the_latest(self, live, path): + _, hub = live + sock, reader = self._subscribe(path) + pending = [] + self._next(sock, reader, pending) + for i in range(200): + hub.publish('display', _display(mode=f'm{i}'), volatile=('last_updated',)) + seen = [] + while not seen or seen[-1] != hub.version: + msg = self._next(sock, reader, pending) + if msg['event'] == 'state': + seen.append(msg['result']['version']) + assert seen == sorted(seen) and len(seen) < 200 + assert msg['result']['state']['display']['mode'] == 'm199' + sock.close() + + def test_several_subscribers_and_commands_still_get_a_slot(self, live, path): + server, hub = live + subs = [client.StateSubscription(paths=[path]).start() for _ in range(3)] + try: + assert _wait_until(lambda: all(s.latest() for s in subs)) + assert hub.subscribers == 3 + # max_clients is 2 and three streams are open: subscribers gave + # their request slots back. + for _ in range(4): + assert client.ping(paths=[path]) == {'pong': True} + hub.publish('display', _display(mode='weather'), volatile=('last_updated',)) + assert _wait_until(lambda: all( + s.latest()['state']['display']['mode'] == 'weather' for s in subs)) + # A fourth is over the subscriber bound. + sock, reader = self._subscribe(path) + refused = self._next(sock, reader, []) + assert refused['ok'] is False and refused['error']['code'] == ErrorCode.BUSY + sock.close() + finally: + for s in subs: + s.stop() + assert _wait_until(lambda: hub.subscribers == 0) + + def test_a_subscriber_that_never_reads_blocks_nobody(self, live, path): + """The render thread keeps publishing at full speed, a reading + subscriber keeps up, and the stuck one is dropped.""" + _, hub = live + stuck = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM) + stuck.setsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF, 4096) + stuck.connect(path) + stuck.sendall(c.encode_message({'v': 1, 'id': 'stuck', + 'cmd': Command.STATE_SUBSCRIBE, 'args': {}})) + good = client.StateSubscription(paths=[path]).start() + try: + assert _wait_until(lambda: good.latest() is not None) + assert _wait_until(lambda: hub.subscribers == 2) + big = {f'p{i}': {'state': 'enabled', 'note': 'x' * 100} for i in range(100)} + worst = 0.0 + deadline = time.monotonic() + 3.0 + i = 0 + while time.monotonic() < deadline: + i += 1 + start = time.perf_counter() + hub.publish('plugins', {'schema': 1, 'n': i, 'plugins': big}) + worst = max(worst, time.perf_counter() - start) + time.sleep(0.005) + assert worst < 0.05, f'a publish took {worst:.3f}s' + # io_timeout is 0.5 s: the stuck one is gone, the good one is not. + assert _wait_until(lambda: hub.subscribers == 1) + assert _wait_until(lambda: (good.latest() or {}).get('state', {}) + .get('plugins', {}).get('n') == i) + finally: + stuck.close() + good.stop() + + def test_close_ends_the_streams(self, path, hub): + server = ControlServer(path, state_hub=hub, keepalive=30.0) + assert server.start() + sub = client.StateSubscription(paths=[path]).start() + try: + assert _wait_until(lambda: sub.latest() is not None) + start = time.monotonic() + server.close() + assert _wait_until(lambda: hub.subscribers == 0, timeout=3.0) + assert time.monotonic() - start < 3.0 + assert _wait_until(lambda: sub.latest() is None, timeout=3.0) + finally: + sub.stop() + + def test_the_subscription_follows_a_restarted_display(self, path, hub): + server = ControlServer(path, state_hub=hub, keepalive=0.2) + assert server.start() + sub = client.StateSubscription(paths=[path]).start() + try: + assert _wait_until(lambda: sub.latest() is not None) + server.close() + assert _wait_until(lambda: sub.latest() is None) + hub2 = StateHub(loop_probe=_loop(), epoch='e2') + hub2.publish('display', _display(mode='restarted'), volatile=('last_updated',)) + server2 = ControlServer(path, state_hub=hub2, keepalive=0.2) + assert server2.start() + try: + assert _wait_until(lambda: (sub.latest() or {}).get('epoch') == 'e2', + timeout=8.0) + assert sub.latest()['state']['display']['mode'] == 'restarted' + finally: + server2.close() + finally: + sub.stop() + + def test_a_server_without_state_answers_an_error(self, path): + server = ControlServer(path) # no hub + assert server.start() + try: + with pytest.raises(client.ControlError) as err: + client.state_get(paths=[path]) + assert err.value.reason == ErrorCode.INTERNAL + finally: + server.close() + + def test_no_display_is_no_socket(self, path): + with pytest.raises(client.ControlError) as err: + client.state_get(paths=[path]) + assert err.value.reason == 'no_socket' diff --git a/test/test_state_stream_readers.py b/test/test_state_stream_readers.py new file mode 100644 index 00000000..6402264e --- /dev/null +++ b/test/test_state_stream_readers.py @@ -0,0 +1,506 @@ +"""Control socket stage 3, both ends: what the display publishes into the +state stream, what it stops writing to the SD card, and the web readers +that use the stream and fall back to the cache keys. + +* The display: DisplayController's publish points feed the StateHub; while + the socket serves readers, display_current_state and the plugin runtime + snapshot are written less often (measured below, in cache writes per + minute of a simulated rotation). +* The readers: /display/current-status, /display/on-demand/status, + /plugins/installed's runtime and /health's display_loop take the socket's + answer when there is one, judged by the same stale / stalled rules as the + cache (#726), and the cache keys and heartbeat file when there is not. +* End to end (AF_UNIX only): a real server, the web routes reading it, and + the fallback once it is gone. +""" + +import json +import os +import sys +import time +from pathlib import Path +from unittest.mock import MagicMock, patch + +import pytest + +sys.path.insert(0, str(Path(__file__).parent.parent)) + +os.environ.setdefault("EMULATOR", "true") + +from src import display_watchdog # noqa: E402 +from src.ipc import contract as c # noqa: E402 +from src.ipc.server import ControlServer, StateHub # noqa: E402 +from src.plugin_system import plugin_runtime as rt # noqa: E402 +from src.plugin_system.plugin_runtime import ( # noqa: E402 + PluginRuntimePublisher, view_from_socket_state, +) +from src.plugin_system.plugin_state import PluginState, PluginStateManager # noqa: E402 +from test._api_v3_test_helpers import api_v3_client, api_v3_module # noqa: F401,E402 +from web_interface import display_state # noqa: E402 + +needs_unix_sockets = pytest.mark.skipif(not c.socket_supported(), + reason='AF_UNIX sockets are Linux/macOS only') + + +class FakeClock: + def __init__(self, now=1000.0): + self.now = now + + def __call__(self): + return self.now + + +def _loop(age): + return lambda: {'heartbeat_age_seconds': age, 'armed': age is not None, + 'stale_after': display_watchdog.HEARTBEAT_STALE_SECONDS} + + +def _states(): + states = PluginStateManager() + states.set_state("clock", PluginState.ENABLED) + states.record_loaded("clock", "1.0.0", loaded_at=10.0) + return states + + +def _hub_with_everything(loop_age=1.0, display_updated=None, plugins_published=None): + hub = StateHub(loop_probe=_loop(loop_age), epoch='e1') + hub.publish('display', { + 'mode': 'weather', 'plugin_id': 'weather', 'mode_index': 1, 'total_modes': 3, + 'on_demand_active': False, 'is_display_active': True, + 'last_updated': time.time() if display_updated is None else display_updated, + }, volatile=('last_updated',)) + hub.publish('on_demand', {'active': True, 'status': 'active', 'plugin_id': 'clock', + 'expires_at': time.time() + 30, 'remaining': 999.0, + 'last_updated': time.time()}, + volatile=('last_updated', 'remaining')) + snapshot = rt.build_runtime_snapshot(_states(), started_at=1.0, + now=time.time() if plugins_published is None + else plugins_published) + hub.publish('plugins', snapshot, volatile=('published_at',)) + hub.publish('brightness', {'brightness': 80, 'panel_brightness': 40, 'dimmed': True}) + return hub + + +@pytest.fixture(autouse=True) +def _no_live_subscription(): + yield + display_state.stop_subscription() + + +# --- the display's side --------------------------------------------------------- + +def _controller(hub=None): + from src.display_controller import DisplayController + dc = object.__new__(DisplayController) + dc.cache_manager = MagicMock() + dc.current_display_mode = "clock" + dc.mode_to_plugin_id = {"clock": "clock", "weather": "weather"} + dc.current_mode_index = 0 + dc.available_modes = ["clock", "weather"] + dc.on_demand_active = False + dc.is_display_active = True + dc._last_published_mode = None + dc._last_published_at = 0.0 + dc._normal_brightness = 80 + dc.current_brightness = 80 + dc.is_dimmed = False + for name in ("on_demand_mode", "on_demand_plugin_id", "on_demand_requested_at", + "on_demand_expires_at", "on_demand_duration", "on_demand_last_error", + "on_demand_last_event"): + setattr(dc, name, None) + dc.on_demand_pinned = False + dc.on_demand_status = 'idle' + dc._state_hub = hub + return dc + + +def _writes(dc, key): + return sum(1 for call in dc.cache_manager.set.call_args_list if call.args[0] == key) + + +class TestDisplayPublishes: + def test_the_publish_points_feed_the_hub(self): + from src import display_controller as dc_module + hub = StateHub(epoch='e1') + dc = _controller(hub) + with patch.object(dc_module.time, "monotonic", lambda: 1000.0): + dc._publish_current_mode_state_if_changed() + state = hub.snapshot()['state'] + assert state['display']['mode'] == 'clock' + assert state['brightness'] == {'brightness': 80, 'panel_brightness': 80, 'dimmed': False} + dc.current_brightness, dc.is_dimmed = 30, True + with patch.object(dc_module.time, "monotonic", lambda: 1001.0): + dc._publish_current_mode_state_if_changed() # no cache write due... + assert hub.snapshot()['state']['brightness']['panel_brightness'] == 30 # ...hub updated + + def test_an_on_demand_outcome_reaches_the_hub(self): + hub = StateHub(epoch='e1') + dc = _controller(hub) + dc.on_demand_status, dc.on_demand_last_error = 'error', 'Plugin nope not found' + dc._publish_on_demand_state() + assert hub.snapshot()['state']['on_demand']['error'] == 'Plugin nope not found' + version = hub.version + dc._publish_on_demand_state() # only last_updated moved + assert hub.version == version + + def test_without_a_socket_nothing_changes(self): + """No hub (Windows, socket off, the golden traces): every write as before.""" + from src import display_controller as dc_module + dc = _controller(None) + now = [1000.0] + with patch.object(dc_module.time, "monotonic", lambda: now[0]): + dc._publish_current_mode_state_if_changed() + now[0] += 1 + dc.current_display_mode = "weather" + dc._publish_current_mode_state_if_changed() + assert _writes(dc, 'display_current_state') == 2 + + def test_with_readers_on_the_socket_a_mode_change_waits_for_the_refresh(self): + from src import display_controller as dc_module + hub = StateHub(epoch='e1') + hub.subscriber_joined() + dc = _controller(hub) + now = [1000.0] + with patch.object(dc_module.time, "monotonic", lambda: now[0]): + dc._publish_current_mode_state_if_changed() # first: written + now[0] += 1 + dc.current_display_mode = "weather" + dc._publish_current_mode_state_if_changed() # mode: hub only + assert _writes(dc, 'display_current_state') == 1 + assert hub.snapshot()['state']['display']['mode'] == 'weather' + dc.is_display_active = False + dc._publish_current_mode_state_if_changed() # a flag: written + assert _writes(dc, 'display_current_state') == 2 + now[0] += dc_module.CURRENT_STATE_RELAXED_REFRESH_SECONDS + dc._publish_current_mode_state_if_changed() # refresh + assert _writes(dc, 'display_current_state') == 3 + + def test_when_the_readers_leave_a_changed_mode_is_written_at_once(self): + from src import display_controller as dc_module + hub = StateHub(epoch='e1', reader_window=60.0) + hub.subscriber_joined() + dc = _controller(hub) + now = [1000.0] + with patch.object(dc_module.time, "monotonic", lambda: now[0]), \ + patch.object(hub, "_clock", lambda: now[0]): + dc._publish_current_mode_state_if_changed() + dc.current_display_mode = "weather" + now[0] += 1 + dc._publish_current_mode_state_if_changed() + assert _writes(dc, 'display_current_state') == 1 + hub.subscriber_left() + now[0] += 61 # past the reader window + dc._publish_current_mode_state_if_changed() + assert _writes(dc, 'display_current_state') == 2 + assert dc.cache_manager.set.call_args.args[1]['mode'] == 'weather' + + def test_relaxed_refresh_is_inside_the_readers_max_age(self): + from src import display_controller as dc_module + assert dc_module.CURRENT_STATE_RELAXED_REFRESH_SECONDS < \ + display_state.CURRENT_STATE_MAX_AGE_SECONDS + + def test_start_state_stream_publishes_everything(self): + hub = StateHub(epoch='e1') + dc = _controller() + publisher = MagicMock() + dc._plugin_runtime_publisher = publisher + dc._start_state_stream(hub) + state = hub.snapshot()['state'] + assert state['display']['mode'] == 'clock' and state['on_demand']['status'] == 'idle' + publisher.attach_hub.assert_called_once_with(hub) + assert dc._state_hub is hub + + +class TestRuntimePublisherHub: + def test_attach_publishes_and_ticks_keep_it_fresh_without_new_versions(self): + wall = FakeClock(1_800_000_000.0) + hub = StateHub(epoch='e1') + publisher = PluginRuntimePublisher(MagicMock(), _states(), clock=FakeClock(), + wall_clock=wall) + publisher.attach_hub(hub) + assert hub.snapshot()['state']['plugins']['plugins']['clock']['version'] == '1.0.0' + version = hub.version + wall.now += 5 + publisher.tick() + assert hub.version == version + assert hub.snapshot()['state']['plugins']['published_at'] == wall.now + + def test_a_change_is_a_new_version_at_once_unthrottled(self): + states = _states() + hub = StateHub(epoch='e1') + mono = FakeClock() + publisher = PluginRuntimePublisher(MagicMock(), states, clock=mono) + publisher.attach_hub(hub) + publisher.tick() + version = hub.version + states.set_state("clock", PluginState.ERROR, error=RuntimeError("boom")) + mono.now += 1 # inside MIN_INTERVAL: no cache write + assert publisher.tick() is False + assert hub.version == version + 1 + assert hub.snapshot()['state']['plugins']['plugins']['clock']['state'] == 'error' + + def test_the_cache_refresh_relaxes_while_the_socket_has_readers(self): + cache = MagicMock() + mono = FakeClock() + hub = StateHub(epoch='e1') + publisher = PluginRuntimePublisher(cache, _states(), clock=mono) + publisher.attach_hub(hub) + publisher.tick() + assert cache.set.call_count == 1 + hub.subscriber_joined() + mono.now += rt.REFRESH_INTERVAL + publisher.tick() + assert cache.set.call_count == 1 # relaxed: not yet + mono.now += rt.RELAXED_REFRESH_INTERVAL - rt.REFRESH_INTERVAL + publisher.tick() + assert cache.set.call_count == 2 + snapshot = cache.set.call_args.args[1] + # The snapshot says how long it may be trusted. + assert snapshot['refresh_interval'] == rt.RELAXED_REFRESH_INTERVAL + assert snapshot['stale_after'] == 3 * rt.RELAXED_REFRESH_INTERVAL + + +def test_watchdog_liveness_is_the_heartbeat_in_memory(): + clock = FakeClock() + wd = display_watchdog.RenderWatchdog(environ={}, clock=clock, heartbeat_dir=None, + send=lambda message: True) + assert wd.liveness()['heartbeat_age_seconds'] is None # no frame yet + wd.bind_render_thread() + wd.note_frame() # arms on the render thread + clock.now += 7 + live = wd.liveness() + assert live['heartbeat_age_seconds'] == 7 and live['armed'] is True + assert live['stale_after'] == display_watchdog.HEARTBEAT_STALE_SECONDS + + +# --- SD writes saved ------------------------------------------------------------ + +def _simulate(readers_on_socket, minutes=10, screen_seconds=15, step=0.25): + """A rotation of two modes, a screen every ``screen_seconds``: the + publish point runs every ``step`` and the runtime publisher ticks every + 5 s, on fake clocks. Returns cache writes per minute for each key.""" + from src import display_controller as dc_module + now = [1000.0] + hub = StateHub(epoch='e1', clock=lambda: now[0]) + if readers_on_socket: + hub.subscriber_joined() + dc = _controller(hub) + cache = dc.cache_manager + publisher = PluginRuntimePublisher(cache, _states(), clock=lambda: now[0]) + publisher.attach_hub(hub) + next_tick = now[0] + end = now[0] + minutes * 60 + with patch.object(dc_module.time, "monotonic", lambda: now[0]): + while now[0] < end: + dc.current_display_mode = ( + "clock" if int((now[0] - 1000.0) // screen_seconds) % 2 == 0 else "weather") + dc._publish_current_mode_state_if_changed() + if now[0] >= next_tick: + publisher.tick() + next_tick += rt.TICK_INTERVAL + now[0] += step + return {key: _writes(dc, key) / minutes + for key in ('display_current_state', rt.PLUGIN_RUNTIME_KEY)} + + +def test_cache_writes_per_minute_with_and_without_socket_readers(): + before = _simulate(readers_on_socket=False) + after = _simulate(readers_on_socket=True) + print(f"\ncache writes/min, rotation of 15 s screens: without socket readers " + f"{before}, with {after}") + # A mode change every 15 s is 4 writes a minute; with readers on the + # socket the key is refreshed once a minute. + assert before['display_current_state'] == pytest.approx(4, abs=0.2) + assert after['display_current_state'] == pytest.approx(1, abs=0.2) + assert before[rt.PLUGIN_RUNTIME_KEY] == pytest.approx(1, abs=0.2) + assert after[rt.PLUGIN_RUNTIME_KEY] == pytest.approx(0.5, abs=0.2) + assert sum(after.values()) < sum(before.values()) / 3 + + +# --- the readers' rules ------------------------------------------------------------ + +class TestSocketRuntimeView: + def test_live(self): + view = view_from_socket_state(_hub_with_everything().snapshot()) + assert view.status == rt.LIVE and view.source == 'socket' + assert view.plugin('clock')['loaded_version'] == '1.0.0' + assert view.describe()['source'] == 'socket' + + def test_a_stalled_loop_is_stalled_like_the_heartbeat_file(self): + """#726: a fresh snapshot with a heartbeat at the health check's + threshold or older is stalled and reports no plugin facts.""" + view = view_from_socket_state(_hub_with_everything( + loop_age=display_watchdog.HEARTBEAT_STALE_SECONDS).snapshot()) + assert view.status == rt.STALLED + assert view.plugin('clock')['loaded'] is None + fresh = view_from_socket_state(_hub_with_everything( + loop_age=display_watchdog.HEARTBEAT_STALE_SECONDS - 0.5).snapshot()) + assert fresh.status == rt.LIVE + + def test_the_time_since_it_arrived_counts(self): + snap = _hub_with_everything(loop_age=50.0).snapshot() + snap['received_mono'] = 100.0 + assert view_from_socket_state(snap, now_mono=105.0).status == rt.LIVE + assert view_from_socket_state(snap, now_mono=111.0).status == rt.STALLED + + def test_no_beat_yet_is_judged_alone(self): + view = view_from_socket_state(_hub_with_everything(loop_age=None).snapshot()) + assert view.status == rt.LIVE and view.heartbeat_age_seconds is None + + def test_a_publisher_that_stopped_ticking_goes_stale(self): + view = view_from_socket_state(_hub_with_everything( + plugins_published=time.time() - rt.STALE_AFTER - 5).snapshot()) + assert view.status == rt.STALE + + @pytest.mark.parametrize('snapshot', [None, {}, {'state': {'plugins': None}}, + {'state': 'x'}]) + def test_no_plugins_section_means_read_the_cache(self, snapshot): + assert view_from_socket_state(snapshot) is None + + +class TestDisplayStateHelpers: + def test_current_status(self): + status = display_state.current_status(_hub_with_everything().snapshot()) + assert status['mode'] == 'weather' and status['is_display_active'] is True + + def test_a_display_section_the_loop_stopped_refreshing_is_unknown(self): + """The cache key reads as unknown after its 120 s max_age; so does this.""" + snap = _hub_with_everything(display_updated=time.time() - 121).snapshot() + assert display_state.current_status(snap) == { + 'mode': None, 'plugin_id': None, 'last_updated': None} + + def test_on_demand_remaining_is_recomputed(self): + state = display_state.on_demand_state(_hub_with_everything().snapshot()) + assert 25 < state['remaining'] <= 30 + + def test_no_snapshot_means_read_the_cache(self): + assert display_state.current_status(None) is None + assert display_state.on_demand_state(None) is None + assert display_state.loop_heartbeat_age(None) is None + + def test_with_the_socket_off_there_is_no_snapshot(self): + # test/conftest.py turns the socket off for the suite. + assert display_state.read_state() is None + + +# --- the routes ----------------------------------------------------------------------- + +@pytest.fixture +def socket_state(monkeypatch): + """Make the routes' read_state return this snapshot (None: socket gone).""" + holder = {'snapshot': None} + monkeypatch.setattr(display_state, 'read_state', lambda: holder['snapshot']) + return holder + + +@pytest.fixture +def web(api_v3_module, api_v3_client, monkeypatch): # noqa: F811 + cache = api_v3_module.api_v3.cache_manager + cached = {} + cache.get.side_effect = lambda key, *a, **kw: cached.get(key) + monkeypatch.setattr("web_interface.blueprints.api_v3.display._get_display_service_status", + lambda: {"active": True}) + monkeypatch.setattr("web_interface.blueprints.api_v3.misc._get_display_service_status", + lambda: {"active": True}) + return api_v3_client, cached + + +def _data(client, url): + response = client.get(url) + assert response.status_code == 200, response.get_data(as_text=True) + return response.get_json()['data'] + + +class TestRoutes: + def test_current_status_from_the_socket(self, web, socket_state): + client, cached = web + cached['display_current_state'] = {'mode': 'stale-cache', 'last_updated': 1} + socket_state['snapshot'] = _hub_with_everything().snapshot() + data = _data(client, '/api/v3/display/current-status') + assert data['mode'] == 'weather' and data['source'] == 'socket' + + def test_current_status_falls_back_to_the_cache(self, web, socket_state): + client, cached = web + cached['display_current_state'] = {'mode': 'clock', 'last_updated': 1} + data = _data(client, '/api/v3/display/current-status') + assert data['mode'] == 'clock' and data['source'] == 'cache' + + def test_on_demand_status_from_the_socket_and_back(self, web, socket_state): + client, cached = web + cached['display_on_demand_state'] = {'active': False, 'status': 'idle'} + socket_state['snapshot'] = _hub_with_everything().snapshot() + data = _data(client, '/api/v3/display/on-demand/status') + assert data['state']['status'] == 'active' and data['source'] == 'socket' + socket_state['snapshot'] = None + data = _data(client, '/api/v3/display/on-demand/status') + assert data['state']['status'] == 'idle' and data['source'] == 'cache' + + @pytest.mark.parametrize('age, status', [ + (3.0, 'running'), + (display_watchdog.HEARTBEAT_STALE_SECONDS + 1, 'stalled'), + (None, 'not_reported'), + ]) + def test_health_display_loop_from_the_socket(self, web, socket_state, age, status): + client, _ = web + socket_state['snapshot'] = _hub_with_everything(loop_age=age).snapshot() + check = _data(client, '/api/v3/health')['checks']['display_loop'] + assert check['status'] == status and check['source'] == 'socket' + + def test_health_falls_back_to_the_heartbeat_file(self, web, socket_state, tmp_path, + monkeypatch): + client, _ = web + path = tmp_path / 'display-heartbeat.json' + path.write_text(json.dumps({'pid': 1, 'mono': time.monotonic() - 300, + 'wall': time.time() - 300})) + monkeypatch.setattr(display_watchdog, 'HEARTBEAT_PATH', str(path)) + check = _data(client, '/api/v3/health')['checks']['display_loop'] + assert check['status'] == 'stalled' and check['source'] == 'heartbeat_file' + + def test_plugin_runtime_view_prefers_the_socket(self, web, socket_state): + import web_interface.blueprints.api_v3 as pkg + socket_state['snapshot'] = _hub_with_everything().snapshot() + view = pkg._plugin_runtime_view() + assert view.source == 'socket' and view.live + socket_state['snapshot'] = _hub_with_everything( + loop_age=display_watchdog.HEARTBEAT_STALE_SECONDS + 1).snapshot() + assert pkg._plugin_runtime_view().status == rt.STALLED + socket_state['snapshot'] = None + assert pkg._plugin_runtime_view().source == 'cache' + + +# --- end to end over a real socket ------------------------------------------------------- + +@needs_unix_sockets +class TestEndToEnd: + def test_routes_read_the_socket_then_fall_back_when_it_goes(self, web, tmp_path, + monkeypatch): + client, cached = web + cached['display_current_state'] = {'mode': 'from-cache', 'last_updated': 1} + path = str(tmp_path / 'control.sock') + monkeypatch.setenv(c.SOCKET_PATH_ENV, path) + hub = _hub_with_everything() + server = ControlServer(path, state_hub=hub, keepalive=0.2) + assert server.start() + try: + data = _data(client, '/api/v3/display/current-status') + assert (data['mode'], data['source']) == ('weather', 'socket') + # The subscription picks up a change the display publishes. + hub.publish('display', dict(hub.snapshot()['state']['display'], mode='stocks'), + volatile=('last_updated',)) + deadline = time.monotonic() + 5 + while time.monotonic() < deadline: + if _data(client, '/api/v3/display/current-status')['mode'] == 'stocks': + break + time.sleep(0.05) + assert _data(client, '/api/v3/display/current-status')['mode'] == 'stocks' + assert hub.subscribers == 1 and hub.readers_active() + finally: + server.close() + # The subscription sees the hang-up, and no socket answers a one-shot. + deadline = time.monotonic() + 5 + while time.monotonic() < deadline: + data = _data(client, '/api/v3/display/current-status') + if data['source'] == 'cache': + break + time.sleep(0.05) + assert (data['mode'], data['source']) == ('from-cache', 'cache') diff --git a/web_interface/app.py b/web_interface/app.py index 66f80594..4f4ca235 100644 --- a/web_interface/app.py +++ b/web_interface/app.py @@ -998,12 +998,14 @@ def _run_startup_reconciliation() -> None: try: from src.plugin_system.state_reconciliation import StateReconciliation - from src.plugin_system.plugin_runtime import read_plugin_runtime + # The display's runtime view: its control socket's state stream when + # it serves one, else the cache snapshot (both judged the same way). + from web_interface.blueprints.api_v3 import _plugin_runtime_view reconciler = StateReconciliation( config_manager=config_manager, plugins_dir=plugins_dir, store_manager=plugin_store_manager, - runtime_source=lambda: read_plugin_runtime(api_v3.cache_manager), + runtime_source=_plugin_runtime_view, ) result = reconciler.reconcile_state() if result.inconsistencies_found: diff --git a/web_interface/blueprints/api_v3/__init__.py b/web_interface/blueprints/api_v3/__init__.py index ca9a65b1..f8a903ab 100644 --- a/web_interface/blueprints/api_v3/__init__.py +++ b/web_interface/blueprints/api_v3/__init__.py @@ -735,8 +735,16 @@ def _plugin_runtime_view(): Only a ``live`` view reports those facts; a stale, stopped or missing snapshot answers None for them (see src/plugin_system/plugin_runtime.py). + + The display's state stream over the control socket comes first + (``source: "socket"``); without it, the cache snapshot and the heartbeat + file (``source: "cache"``). Both are judged by the same rules. """ - from src.plugin_system.plugin_runtime import read_plugin_runtime + from src.plugin_system.plugin_runtime import read_plugin_runtime, view_from_socket_state + from web_interface import display_state + view = view_from_socket_state(display_state.read_state()) + if view is not None: + return view return read_plugin_runtime(getattr(api_v3, 'cache_manager', None)) diff --git a/web_interface/blueprints/api_v3/display.py b/web_interface/blueprints/api_v3/display.py index 73433c3a..cbc83946 100644 --- a/web_interface/blueprints/api_v3/display.py +++ b/web_interface/blueprints/api_v3/display.py @@ -9,7 +9,7 @@ from web_interface.blueprints.api_v3 import ( _get_display_service_status, _socket_reason_code, _stop_display_service, api_v3, jsonify, logger, request, uuid, ) -from web_interface import display_preview +from web_interface import display_preview, display_state import web_interface.blueprints.api_v3 as _pkg from src.ipc import client as control_client # Read through the module rather than bound by value: tests patch these @@ -163,13 +163,22 @@ def get_display_modes(): return jsonify({'status': 'success', 'data': {'modes': modes}}) @api_v3.route('/display/on-demand/status', methods=['GET']) def get_on_demand_status(): - """Return the current on-demand display state.""" - cache = _cache_manager() - # memory_ttl=0: the display service writes this key, so only the file - # is current. This process's memory tier would keep serving the first - # copy it read for the full max_age -- "active" for two minutes after - # the display had already stopped. - state = cache.get('display_on_demand_state', max_age=120, memory_ttl=0) + """Return the current on-demand display state. + + From the display's state stream over the control socket when it is + available (``source: "socket"``), else the cache key it also writes + (``source: "cache"``). + """ + state = display_state.on_demand_state(display_state.read_state()) + source = 'socket' + if state is None: + source = 'cache' + cache = _cache_manager() + # memory_ttl=0: the display service writes this key, so only the file + # is current. This process's memory tier would keep serving the first + # copy it read for the full max_age -- "active" for two minutes after + # the display had already stopped. + state = cache.get('display_on_demand_state', max_age=120, memory_ttl=0) if state is None: state = { 'active': False, @@ -181,7 +190,8 @@ def get_on_demand_status(): 'status': 'success', 'data': { 'state': state, - 'service': service_status + 'service': service_status, + 'source': source, } }) @api_v3.route('/display/on-demand/start', methods=['POST']) @@ -327,18 +337,22 @@ def stop_on_demand_display(): def get_current_display_status(): """Return the display mode/plugin currently intended to be shown. - Published by the display process (display_controller._publish_current_mode_state) - to the shared cache whenever the active mode changes, so the web UI (e.g. the - System Logs page) can show what's on screen without querying the display - process directly. + Read from the display's state stream over the control socket when it is + available (``source: "socket"``). Otherwise from what the display + publishes to the shared cache (display_controller._publish_current_mode_state) + when the active mode changes (``source: "cache"``). """ - cache = _cache_manager() - # memory_ttl=0: written by the display service; see get_on_demand_status. - state = cache.get('display_current_state', max_age=120, memory_ttl=0) + state = display_state.current_status(display_state.read_state()) + source = 'socket' + if state is None: + source = 'cache' + cache = _cache_manager() + # memory_ttl=0: written by the display service; see get_on_demand_status. + state = cache.get('display_current_state', max_age=120, memory_ttl=0) if state is None: state = { 'mode': None, 'plugin_id': None, 'last_updated': None, } - return jsonify({'status': 'success', 'data': state}) + return jsonify({'status': 'success', 'data': dict(state, source=source)}) diff --git a/web_interface/blueprints/api_v3/misc.py b/web_interface/blueprints/api_v3/misc.py index c6961f69..5ef9aea0 100644 --- a/web_interface/blueprints/api_v3/misc.py +++ b/web_interface/blueprints/api_v3/misc.py @@ -17,7 +17,7 @@ from src.common.path_safety import safe_path_component from src import display_watchdog from src.common import sync_manager as _sync from src import error_aggregator as _errors -from web_interface import display_preview +from web_interface import display_preview, display_state from web_interface.auth import request_is_authenticated import web_interface.blueprints.api_v3 as _pkg # Read through the module rather than bound by value: tests patch these @@ -100,20 +100,30 @@ def get_health(): # service is "active". No heartbeat at all (the dev server, an older # display) is not a failure: the preview-frame check below is then # the only signal, as it always was. + # The display reports the same beat's age over the control socket's + # state stream, measured in memory; the file is the fallback. try: - heartbeat = display_watchdog.read_heartbeat(display_watchdog.HEARTBEAT_PATH) - age = display_watchdog.heartbeat_age(heartbeat) if heartbeat else None + snapshot = display_state.read_state() + if snapshot is not None: + source = 'socket' + age = display_state.loop_heartbeat_age(snapshot) + else: + source = 'heartbeat_file' + heartbeat = display_watchdog.read_heartbeat(display_watchdog.HEARTBEAT_PATH) + age = display_watchdog.heartbeat_age(heartbeat) if heartbeat else None if age is None: health_status['checks']['display_loop'] = { 'status': 'not_reported', 'note': 'The display is not writing a heartbeat (not started yet, ' 'or a version or setup without one)', + 'source': source, } else: fresh = age < display_watchdog.HEARTBEAT_STALE_SECONDS health_status['checks']['display_loop'] = { 'status': 'running' if fresh else 'stalled', 'heartbeat_age_seconds': round(age, 1), + 'source': source, } except Exception: logger.warning("Health check could not read the display heartbeat", exc_info=True) diff --git a/web_interface/display_state.py b/web_interface/display_state.py new file mode 100644 index 00000000..766f9e48 --- /dev/null +++ b/web_interface/display_state.py @@ -0,0 +1,145 @@ +"""The display's live state, as the web interface reads it. + +The display process serves its state over the control socket (stage 3, see +docs/IPC_CONTROL_SOCKET.md): what it is showing, the on-demand session, the +brightness, its plugin runtime snapshot and whether its render loop is still +going round. This module holds one ``state.subscribe`` connection for the +web process (:class:`src.ipc.client.StateSubscription`, started on first +use), so a route answers from memory rather than reading a file the display +had to write to the SD card. + +Every reader here returns None when the socket cannot vouch for the answer +-- the socket is off or missing (a stopped display, an older one, Windows, +the test suite), the subscription went quiet, or the display's copy is too +old by the same rules the cache readers apply -- and the route then reads +the cache keys and the heartbeat file exactly as it did before. +""" + +from __future__ import annotations + +import logging +import threading +import time +from typing import Any, Dict, Optional + +from src.ipc import client as control_client +from src.ipc.contract import client_socket_paths, socket_supported + +logger = logging.getLogger(__name__) + +#: The cache readers' max_age for display_current_state: a ``display`` +#: section the render thread has not refreshed for this long is unknown, as +#: the cache key would be. +CURRENT_STATE_MAX_AGE_SECONDS = 120.0 + +#: A one-shot ``state.get``, for a request that arrives before the +#: subscription has its first snapshot. Short: the display answers it from +#: its socket thread, never the render thread. +ONE_SHOT_TIMEOUT_SECONDS = 0.5 + +_feed: Optional[control_client.StateSubscription] = None +_feed_lock = threading.Lock() + + +def _subscription() -> Optional[control_client.StateSubscription]: + """This process's subscription, started on first use; None without a socket.""" + global _feed + if not socket_supported() or not client_socket_paths(): + return None + with _feed_lock: + if _feed is None: + _feed = control_client.StateSubscription().start() + return _feed + + +def stop_subscription() -> None: + """Stop the subscription (tests; a process that is shutting down).""" + global _feed + with _feed_lock: + feed, _feed = _feed, None + if feed is not None: + feed.stop() + + +def read_state() -> Optional[Dict[str, Any]]: + """The display's latest state snapshot, or None (read the cache instead). + + From the subscription when it is live; otherwise one ``state.get``. + """ + feed = _subscription() + if feed is None: + return None + snapshot = feed.latest() + if snapshot is not None: + return snapshot + if feed.last_error == 'unknown_command': + # A display older than stage 3: a one-shot would fail the same way + # on every request. The subscription retries every 30 s. + return None + try: + snapshot = control_client.state_get(timeout=ONE_SHOT_TIMEOUT_SECONDS) + except control_client.ControlError as e: + logger.debug("Display state not available over the control socket: %s", e) + return None + if not isinstance(snapshot.get('state'), dict): + return None + snapshot['received_mono'] = time.monotonic() + return snapshot + + +def _section(snapshot: Optional[Dict[str, Any]], name: str) -> Optional[Dict[str, Any]]: + if not isinstance(snapshot, dict): + return None + state = snapshot.get('state') + value = state.get(name) if isinstance(state, dict) else None + return dict(value) if isinstance(value, dict) else None + + +def current_status(snapshot: Optional[Dict[str, Any]], + now: Optional[float] = None) -> Optional[Dict[str, Any]]: + """``/display/current-status``'s data from a snapshot. + + None when the snapshot has no ``display`` section (fall back to the + cache). A section the render thread last refreshed more than + CURRENT_STATE_MAX_AGE_SECONDS ago -- the loop is stuck -- is reported + as unknown, the same answer the cache key gives once it ages out. + """ + display = _section(snapshot, 'display') + if display is None: + return None + now = time.time() if now is None else now + updated = display.get('last_updated') + if (not isinstance(updated, (int, float)) or isinstance(updated, bool) + or now - updated > CURRENT_STATE_MAX_AGE_SECONDS): + return {'mode': None, 'plugin_id': None, 'last_updated': None} + return display + + +def on_demand_state(snapshot: Optional[Dict[str, Any]], + now: Optional[float] = None) -> Optional[Dict[str, Any]]: + """``/display/on-demand/status``'s state from a snapshot; None to fall back. + + ``remaining`` is worked out again from ``expires_at``: the display set it + when it last published, which may have been minutes ago. + """ + state = _section(snapshot, 'on_demand') + if state is None: + return None + expires_at = state.get('expires_at') + if (state.get('active') and isinstance(expires_at, (int, float)) + and not isinstance(expires_at, bool)): + now = time.time() if now is None else now + state['remaining'] = max(0.0, expires_at - now) + return state + + +def loop_heartbeat_age(snapshot: Optional[Dict[str, Any]]) -> Optional[float]: + """The render loop's heartbeat age now; None when the display has no + beat to report yet (or there is no snapshot).""" + if not isinstance(snapshot, dict): + return None + return control_client.snapshot_loop_age(snapshot) + + +__all__ = ['current_status', 'loop_heartbeat_age', 'on_demand_state', 'read_state', + 'stop_subscription']