diff --git a/CHANGELOG.md b/CHANGELOG.md index 223cc50c..c027e85c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -19,6 +19,50 @@ accepts both, but the store flags the old spelling as deprecated ## Unreleased +### The control socket carries every web command; the mailboxes are a fallback + +Stage 4 of the web → display control socket (`docs/IPC_CONTROL_SOCKET.md`). + +- **Mailbox only when the socket cannot carry it.** The on-demand routes + (`POST /api/v3/display/on-demand/start` and `/stop`) write the + `display_on_demand_request` mailbox only when the display never had the + request: no socket (a stopped display, one older than the socket), a + refused or timed-out connect, or a display too old to know the command. + A display that had it and refused or did not answer (a full queue, bad + arguments, silence after the send) is answered `503` (`400` for bad + arguments) with `socket_error`, and no mailbox copy is written: the + display may have applied it, or would refuse the copy too. A stop with + `stop_service` still stops the service. `src.ipc.client.should_fall_back()` + holds the rule; `ControlError.sent` says whether the display had the + request. +- **`errors.clear`.** `POST /api/v3/errors/clear` goes over the socket: the + display clears its error records and republishes its error snapshot + before it answers, so the response says `applied: true` with the + display's own `cleared_count`. The `plugin_error_clear_request` mailbox + is written only on the same fallback rule (a display from before this + release answers `unknown_command`, and gets the mailbox). A display that + had it and failed answers `503`. The error snapshot gains + `applied_clear_cutoff`, so an older mailbox request is not shown as + pending once a wider clear has been applied. +- **The display looks at the mailboxes less, and more cheaply.** While the + control socket is up, the on-demand mailbox is looked at once a second + instead of every 0.25 s (`MAILBOX_POLL_INTERVAL_WITH_SOCKET`), and both + mailboxes are read only when their file changed since the last look: + otherwise a look is one `stat()` (`CacheManager.file_signature`, + `MailboxWatch`). A socket command no longer reads or deletes the mailbox + file. A duplicate already processed is taken out of the mailbox, rather + than re-read for an hour. Without a socket (Windows, + `LEDMATRIX_CONTROL_SOCKET=off`) the mailbox is read every 0.25 s as before. +- **Kept for one release.** The display still reads both mailboxes, so an + older web interface (or a web user not yet in the socket's group) keeps + working during an upgrade, and still writes `display_current_state`, + `display_on_demand_state` and `plugin_runtime_snapshot` for the readers' + fallback. A request that comes through the on-demand mailbox while the + socket is up is logged once per writer: plugins that write + `display_on_demand_request` themselves (birdnet-go, mqtt-notifications, + on-air, pomodoro-timer) now get the screen within a second rather than a + quarter second, and need an in-process way in before the mailbox goes. + ### Display loop stage 3: a ScreenRunner, and the Arbiter decides every screen Internal; no behaviour change. Stage 3 of `docs/RUN_LOOP_REDESIGN.md`. diff --git a/docs/ADVANCED_FEATURES.md b/docs/ADVANCED_FEATURES.md index a90cb8b3..ba3a1cab 100644 --- a/docs/ADVANCED_FEATURES.md +++ b/docs/ADVANCED_FEATURES.md @@ -649,7 +649,9 @@ When nothing is running on demand, `data.state` is > on-demand machinery is internal — drive it through the REST endpoints > above (or the web UI buttons). The API handlers > (`start_on_demand_display()` / `stop_on_demand_display()` in -> `web_interface/blueprints/api_v3/display.py`) write a request into the cache +> `web_interface/blueprints/api_v3/display.py`) send the request over the +> display's control socket ([IPC_CONTROL_SOCKET.md](IPC_CONTROL_SOCKET.md)). +> Only when the socket cannot carry it do they write it into the cache > manager under the `display_on_demand_request` key, which > `DisplayController._poll_on_demand_requests()` > (`src/display_controller.py`) picks up. A separate @@ -747,8 +749,14 @@ keys helps troubleshoot stuck states. "timestamp": 1234567890.123 } ``` -**Purpose:** Communication from web interface to display controller -**When Set:** API endpoint receives request +**Purpose:** Communication from web interface to display controller, as the +fallback when the control socket cannot carry the request (deprecated; it +will be removed in a later release) +**When Set:** API endpoint receives a request and the display's control +socket is unavailable (display stopped, or older than the socket or the +command); some plugins also write it directly +**Read:** once a second while the display serves the control socket (0.25 s +without it), and only when the file changed since the last look **Auto-Cleared:** After processing or 1 hour TTL **2. display_on_demand_config** (No TTL) diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 12db5d41..7a095786 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -42,11 +42,11 @@ each other. They share three things: | State | Where | Written by | Read by | |---|---|---|---| | On-demand command | control socket `/run/ledmatrix/control.sock` ([IPC_CONTROL_SOCKET.md](IPC_CONTROL_SOCKET.md)) | web: `start_on_demand_display()` / `stop_on_demand_display()` in [`api_v3/display.py`](../web_interface/blueprints/api_v3/display.py), via [`src/ipc/client.py`](../src/ipc/client.py) | display: [`src/ipc/server.py`](../src/ipc/server.py) acks; the render thread applies it in `_poll_on_demand_requests()` | -| On-demand request (fallback) | cache `display_on_demand_request` | web, when the socket fails; four plugins write it directly | display: `_poll_on_demand_requests()` | +| On-demand request (fallback) | cache `display_on_demand_request` | web, only when the socket could not carry the request (`should_fall_back`); four plugins write it directly | display: `_poll_on_demand_requests()`, a `stat()` every 1 s while the socket is up (0.25 s without), read only when the file changed | | On-demand state | cache `display_on_demand_state` | display: `_publish_on_demand_state()` | web: `/api/v3/display/on-demand/status` | | Current screen | cache `display_current_state` | display | web: `/api/v3/display/current-status` | | Plugin errors | cache `plugin_error_snapshot` | display: `ErrorSnapshotPublisher` ([`src/error_aggregator.py`](../src/error_aggregator.py)) | web: `read_error_report()` for `/api/v3/errors/*` | -| Error clear | cache `plugin_error_clear_request` | web | display | +| Error clear | control socket `errors.clear`; cache `plugin_error_clear_request` as the fallback | web: `POST /api/v3/errors/clear` | display: applied before the socket answers; the mailbox on the error publisher's 5 s tick, read only when the file changed | | Font usage | cache `font_usage_snapshot` | display: `FontUsagePublisher` ([`src/font_usage.py`](../src/font_usage.py)) | web: Fonts tab | | Fetch statistics (requests per plugin and host) | cache `fetch_stats_snapshot` | display: `FetchStatsPublisher` ([`src/common/fetch_service.py`](../src/common/fetch_service.py)), at most once a minute on change | web: `read_fetch_stats()` for `/api/v3/plugins/fetch-stats` | | Plugin health | cache `plugin_health:` | display (web writes on reset) | web: `/api/v3/plugins/health` | @@ -58,11 +58,16 @@ each other. They share three things: The on-demand start route starts `ledmatrix.service` when it is not running (`start_service`, on by default) but never restarts a running one. The routes -send the command over the display's control socket and get an ack; when that -fails (a stopped display, one older than the socket) they write the mailbox -instead, which the display reads every `ON_DEMAND_POLL_INTERVAL` (0.25s), from -its dwell sleep, its render loops and Vegas's interrupt check as well as the -main loop. Both ways end in the same handler, `_handle_on_demand_request()`. +send the command over the display's control socket and get an ack; only when +the socket could not carry it (a stopped display, one older than the socket +or the command) do they write the mailbox instead. A display that had the +request and refused it is answered with the error, not posted a mailbox +copy. The display looks at the mailbox every +`MAILBOX_POLL_INTERVAL_WITH_SOCKET` (1 s) while it serves the socket, and +every `ON_DEMAND_POLL_INTERVAL` (0.25 s) without one, from its dwell sleep, +its render loops and Vegas's interrupt check as well as the main loop; a +look is one `stat()` unless the file changed. Both ways end in the same +handler, `_handle_on_demand_request()`. The socket's handlers only queue; see [IPC_CONTROL_SOCKET.md](IPC_CONTROL_SOCKET.md) for the protocol, the permission model and the plan to retire the mailboxes. diff --git a/docs/IPC_CONTROL_SOCKET.md b/docs/IPC_CONTROL_SOCKET.md index ffb6988d..38505772 100644 --- a/docs/IPC_CONTROL_SOCKET.md +++ b/docs/IPC_CONTROL_SOCKET.md @@ -7,14 +7,18 @@ status. Stage 2 makes those commands land within a frame on every kind of 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. +wrote to the SD card. Stage 4 makes the socket the only way a command goes +while it works: the web interface writes a mailbox only when the socket +cannot carry the request, `errors.clear` replaces the last command that +always went through a mailbox, and the display looks at the mailboxes once +a second, with a `stat()`. The file mailboxes 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); 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`) | +| 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), `POST /api/v3/errors/clear`; 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 | @@ -85,6 +89,17 @@ one. Clients branch on `error.code`, never on the message text. | `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 | +| `errors.clear` | `{cutoff: number}` (epoch seconds, finite, ≥ 0) | `{request_id, cutoff, cleared}` | answered directly, once applied | + +`errors.clear` (stage 4) forgets the plugin errors the display recorded at +or before `cutoff` and rewrites its error snapshot (`plugin_error_snapshot`) +before it answers, so the web interface's next read already has it. The +request `id` is the clear's id, which the snapshot reports as +`applied_clear_id`. It is answered on the connection thread by a handler the +display registers (`ControlServer(handlers=...)`, the contract's +`DIRECT_COMMANDS`): the error aggregator and its publisher have their own +locks, and nothing the render thread owns is touched. A display that has no +handler answers `unknown_command`, as an older display does. `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 @@ -392,7 +407,7 @@ place it reads the mailbox: unloads it mid-load; a disable saved meanwhile is applied once the reload is done. A second reload of the same plugin runs after the first. -The 0.25 s floor on the mailbox read does not apply to the queue, because +The floor on the mailbox read (0.25 s, 1 s since stage 4) does not apply to the queue, because draining it costs no disk read. A queued command also lets `_service_pending_changes()` skip its own floor. @@ -420,7 +435,7 @@ Now the queue wakes the render thread: So a command lands within a millisecond or so on a static screen and in a dwell, and within one frame in Vegas and on a scrolling screen. The mailbox -keeps its old delays. Commands still run only on the render thread: the +is slower on purpose (see "The mailboxes now"). Commands still run only on the render thread: the connection threads only queue them and set the event. The one exception is the slow half of `plugin.reload` (tearing down and loading the plugin), which runs on its own thread. Every change to the display's state still @@ -438,11 +453,57 @@ bookkeeping. A client's send to the render thread waking took 0.72 ms median Without a socket (Windows, `LEDMATRIX_CONTROL_SOCKET=off`) the waits are the plain sleeps they were. -**Exactly once.** A command and a mailbox write for the same request share -one `request_id`. If the client times out after the display queued the -command and then also writes the mailbox, the display processes the request -once. The existing `on_demand_request_id` and processed-id checks drop the -second copy. +**Exactly once.** Since stage 4 the web interface writes the mailbox only +when the display never had the request (see "When the web interface falls +back"), so a request goes one way or the other, never both. A command and a +mailbox write for the same request still share one `request_id`, and the +`on_demand_request_id` and processed-id checks still drop a second copy: an +older web interface (before stage 4) wrote the mailbox after a reply timed +out, too. The display takes such a copy out of the mailbox when it drops it. + +## When the web interface falls back (stage 4) + +The client tells a request the display never had from one it had and then +failed. `ControlError.sent` is True once the whole request was written to a +connected display; a refusal the display sends before reading anything +(`forbidden`, too many connections) carries no request id, and leaves it +False. `src.ipc.client.should_fall_back()` is the one rule every route uses: + +| What happened | Example reasons | Mailbox? | The route answers | +|---|---|---|---| +| The display never had it | `no_socket`, `refused`, `disabled`, `unsupported`, a connect or send that timed out, `forbidden` / `busy` at the door, `invalid_request` (refused by the client itself) | yes | success, `transport: "mailbox"`, `socket_error` | +| A display too old to know it (the upgrade case) | `unknown_command`, `unsupported_version` | yes | as above | +| The display had it and failed | `busy` (queue full), `invalid_args`, `internal`, a timeout or hang-up after the send, `bad_response` | no | `503` (`400` for `invalid_args`), `socket_error` | + +A display that had the request may have applied it (a reply that timed out), +or would refuse the mailbox copy as well (bad arguments), or is stuck and +would not read the mailbox either (a full queue). Writing the copy anyway +only turned that into a "success". An on-demand stop with `stop_service` +still stops the service, which ends on-demand whatever happened. + +Brightness and plugin reload never had a mailbox: without the socket, the +config watcher applies the saved brightness and a reload becomes the +restart banner, as before. + +### The mailboxes now + +| Mailbox | Written by | Read by the display | While the socket is up | +|---|---|---|---| +| `display_on_demand_request` | the web interface, only on fallback; four plugins directly (birdnet-go, mqtt-notifications, on-air, pomodoro-timer) | the render thread, `_poll_on_demand_requests()` | looked at every 1 s (`MAILBOX_POLL_INTERVAL_WITH_SOCKET`), 0.25 s without a socket | +| `plugin_error_clear_request` | the web interface, only on fallback | the error publisher's thread, every 5 s tick | unchanged rate | + +A look is one `stat()` of the mailbox file (`CacheManager.file_signature`): +`(inode, mtime, size)`, and every write renames a new file into place, so a +new write always looks different. `MailboxWatch` reads the file only when +that changed since the last look, so a mailbox that holds nothing new, or +nothing at all, costs no open and no parse. A socket command never reads or +deletes the on-demand mailbox. A start already processed is taken out of +the mailbox instead of being re-read until it expires. + +A request that comes through the on-demand mailbox while the socket is up +is logged once per writer (`came through the file mailbox although the +control socket is up`), which names the plugins that still need an +in-process way in before the mailbox is removed. ## Robustness @@ -459,9 +520,10 @@ block the render loop or crash it: and the connection is closed, because the next message boundary cannot be found. A client that disconnects mid-message is dropped silently. No exception from a handler leaves the connection thread. -- **Full queue.** When the queue is full, the client gets `busy` and falls - back to the mailbox. A full queue means the render thread is stuck, and the - systemd watchdog deals with that. +- **Full queue.** When the queue is full, the client gets `busy`, and the web + interface answers `503` rather than write the mailbox, which the stuck + render thread would not read either. A full queue means the render thread + is stuck, and the systemd watchdog deals with that. - **Awaited commands.** The wait for an awaited command's outcome happens on its connection thread and is bounded (`AWAIT_SECONDS`), so a stuck render thread costs that client `pending` and one connection slot for at most @@ -578,15 +640,30 @@ device never touches the live display. 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. 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. +4. **The mailboxes become a fallback (done).** + - The web interface writes a mailbox only when the socket could not carry + the request (`should_fall_back`); a display that had it and failed is + answered as that (see "When the web interface falls back"). + - `errors.clear` replaces `plugin_error_clear_request` as the way a clear + reaches the display. + - The display looks at the on-demand mailbox once a second while the + socket is up, reads either mailbox only when its file changed, and logs + who still writes the on-demand one (see "The mailboxes now"). + - Not changed, deliberately: config saves (the schedule, the dim + schedule, plugin settings) still reach the display through + `config.json` and its watcher, which is the setting itself rather than + a message; see `config.reload` under stage 2. The preview viewer marker + (`/tmp/led_matrix_preview_viewer`) is a presence signal the display + already stats at most once a second. Plugin health and metrics resets + write the persisted record the display publishes and do not reach the + running display (their routes say so); they are not mailboxes. +5. **Remove the mailboxes (next release).** Once every device has run a + display with stage 4, the web interface stops writing both mailboxes and + the display stops reading them. The four plugins that write + `display_on_demand_request` need an in-process way to ask for the screen + first. The display also stops writing `display_current_state`, + `display_on_demand_state` and `plugin_runtime_snapshot` once the web + interface no longer falls back to them. ## Checking it on a device @@ -601,7 +678,18 @@ curl -s -X POST localhost:5000/api/v3/display/on-demand/start \ If the response says `"transport": "mailbox"`, `socket_error` gives the reason. `no_socket` means the display is stopped or predates the socket. `refused` usually means the web user is not in the socket's group, which -takes effect when the web service restarts after the user is added. +takes effect when the web service restarts after the user is added. A `503` +with `"transport": "socket"` means the display had the request and did not +take it (`busy`, `timeout`, ...): nothing was written to the mailbox. + +An error clear: + +```bash +curl -s -X POST localhost:5000/api/v3/errors/clear \ + -H 'Content-Type: application/json' -d '{"all":true}' +# ... "applied": true, "transport": "socket" +sudo journalctl -u ledmatrix | grep -E "Cleared .* plugin error|file mailbox" +``` Brightness and a plugin reload: diff --git a/docs/REST_API_REFERENCE.md b/docs/REST_API_REFERENCE.md index a69c331c..46693f7f 100644 --- a/docs/REST_API_REFERENCE.md +++ b/docs/REST_API_REFERENCE.md @@ -464,7 +464,7 @@ Request a specific plugin to display on-demand. - `mode` (string, optional): Display mode name (plugin_id inferred if not provided) - `duration` (number, optional): Duration in seconds (0 = until stopped) - `pinned` (boolean, optional): Pin display (pause rotation) -- `start_service` (boolean, optional): Start the display service if it is not running (default: true). A running service is never restarted: it picks the request up within about a quarter of a second. When false and the service is stopped, the route returns 400. +- `start_service` (boolean, optional): Start the display service if it is not running (default: true). A running service is never restarted: it picks the request up within a frame over its control socket (within about a second through the mailbox fallback). When false and the service is stopped, the route returns 400. **Response**: ```json @@ -489,10 +489,19 @@ display's control socket acknowledged it (it is queued for the render thread, which wakes for it and applies it within a frame; see [IPC_CONTROL_SOCKET.md](IPC_CONTROL_SOCKET.md)), `"mailbox"` means it was written to the cache mailbox the display polls, as before the socket existed. -With `"mailbox"`, `socket_error` gives the reason the socket was not used -(`no_socket` when the display is stopped or predates the socket, `timeout`, -`refused`, `busy`, ...). Either way the request is applied the same way; -`request_id` is the same id in both. +The mailbox is used only when the socket could not carry the request. With +`"mailbox"`, `socket_error` gives the reason (`no_socket` when the display is +stopped or predates the socket, `refused`, a connect `timeout`, +`unknown_command` from a display too old for the command, ...). Either way +the request is applied the same way; `request_id` is the same id in both. + +When the display had the request and did not take it -- a full queue +(`busy`), bad arguments (`invalid_args`), no answer after the request was +sent (`timeout`, `closed`) -- the route answers `503` (`400` for +`invalid_args`) with `status: "error"` and `data: {request_id, transport: +"socket", socket_error}`, and writes nothing to the mailbox. The stop route +does the same, except that with `stop_service: true` it still stops the +service and answers success. ### Stop On-Demand Display @@ -2248,13 +2257,11 @@ error with `"all": true` (`max_age_hours` is then ignored). } ``` -The clear is asynchronous. The web interface records a request -(`plugin_error_clear_request` in the shared cache), and the display service -applies it within about 5 seconds, rebuilding its counts from the errors it -keeps and republishing. Reads hide the cleared errors from the moment the -request is recorded. Until the display service applies an age-based clear, -`recent_errors` and `active_patterns` are already filtered but the counts -are the old ones, and `clear_pending` is `true`. +The clear goes to the display service over its control socket +(`errors.clear`, see [IPC_CONTROL_SOCKET.md](IPC_CONTROL_SOCKET.md)), which +applies it, rebuilding its counts from the errors it keeps, and republishes +before it answers: `applied` is `true`, `transport` is `"socket"`, and +`cleared_count` is the display's own count. ```json { @@ -2262,17 +2269,30 @@ are the old ones, and `clear_pending` is `true`. "data": { "cleared_count": 13, "clear_requested": true, + "applied": true, + "transport": "socket", "request_id": "5f0c1e...", "cutoff": "2026-09-23T09:59:02.310000" }, - "message": "Clear of all errors requested; the display service applies it within about 5 seconds" + "message": "Cleared all errors" } ``` -`cleared_count` is how many of the reported errors the clear hides. It is +When the socket cannot carry it (the display is stopped, or older than +`errors.clear`) the clear is asynchronous, as before: the web interface +records a request (`plugin_error_clear_request` in the shared cache), +`applied` is `false` and `transport` is `"mailbox"`, and the display service +applies it within about 5 seconds. Reads hide the cleared errors from the +moment the request is recorded. Until the display service applies an +age-based clear, `recent_errors` and `active_patterns` are already filtered +but the counts are the old ones, and `clear_pending` is `true`. Then +`cleared_count` is how many of the reported errors the clear hides, and `null` when that cannot be known before the display service applies it (an -age-based clear over more errors than the report lists). A request that -could not be written to the shared cache answers `500`. +age-based clear over more errors than the report lists). + +A request that could not be written to the shared cache answers `500`. A +display that had the request and failed it (`internal`, a timeout after the +request was sent) answers `503`, with `context.socket_error`. --- diff --git a/src/cache_manager.py b/src/cache_manager.py index f18d2380..b922c3e9 100644 --- a/src/cache_manager.py +++ b/src/cache_manager.py @@ -28,7 +28,7 @@ import os import time from datetime import datetime import pytz -from typing import Any, Dict, List, Optional +from typing import Any, Dict, List, Optional, Tuple import logging import threading import tempfile @@ -72,6 +72,43 @@ def _outlived(record: Any, max_age: Optional[float], now: float) -> bool: return False +_NOT_SEEN: Any = object() + + +class MailboxWatch: + """Tells the poller of a mailbox key whether its file changed since the + last look, from one stat() (:meth:`CacheManager.file_signature`). + + The display polls the mailboxes the web interface falls back to. Reading + one is an open and a JSON parse; with this a poll that finds the same file + (or none) costs a stat, and the file is read only after a new write. A + cache without ``file_signature`` (a test double) is read every time. + """ + + def __init__(self, key: str): + self.key = key + self._seen: Any = _NOT_SEEN + + def changed(self, cache_manager: Any) -> bool: + """True when the poller should read the key now.""" + signature = getattr(cache_manager, 'file_signature', None) + sig = signature(self.key) if callable(signature) else _NOT_SEEN + if sig is not None and not isinstance(sig, tuple): + return True # cannot tell: read it + if sig is None: + self._seen = None + return False # no file, nothing to read + if sig == self._seen: + return False + self._seen = sig + return True + + def forget(self) -> None: + """Read the key on the next poll even if its file has not changed + (the last read failed).""" + self._seen = _NOT_SEEN + + class CacheManager: """Manages caching of API responses to reduce API calls.""" @@ -295,7 +332,25 @@ class CacheManager: def _get_cache_path(self, key: str) -> Optional[str]: """Get the path for a cache file.""" return self._disk_cache_component.get_cache_path(key) - + + def file_signature(self, key: str) -> Optional[Tuple[int, int, int]]: + """``(st_ino, st_mtime_ns, st_size)`` of ``key``'s file, or None when + there is no file (the key is absent, or this cache has no disk tier). + + One stat(), no read: a poller of a mailbox another process writes + compares it with the last one it saw and reads the file only when it + changed. Every write replaces the file (a temp file renamed into + place), so a new write always has a new inode, however fast it came. + """ + path = self._get_cache_path(key) + if not path: + return None + try: + st = os.stat(path) + except OSError: + return None + return (st.st_ino, st.st_mtime_ns, st.st_size) + def get_cached_data(self, key: str, max_age: int = 300, memory_ttl: Optional[int] = None) -> Optional[Dict[str, Any]]: """Get data from cache (memory first, then disk) honoring TTLs. diff --git a/src/display_controller.py b/src/display_controller.py index 39b96115..dd64c6dd 100644 --- a/src/display_controller.py +++ b/src/display_controller.py @@ -48,7 +48,7 @@ from src.screen_runner import ( from src.display_manager import DisplayManager from src.config_manager import ConfigManager from src.config_service import ConfigService -from src.cache_manager import CacheManager +from src.cache_manager import CacheManager, MailboxWatch from src.font_manager import FontManager from src.logging_config import get_logger from src.exceptions import PluginError @@ -68,6 +68,10 @@ from src.vegas_mode.render_pipeline import SYNC_SEND_INTERVAL # Get logger with consistent configuration logger = get_logger(__name__) +# The on-demand file mailbox: the fallback for a web interface that cannot +# reach the control socket, and how some plugins still ask for the screen. +ON_DEMAND_MAILBOX_KEY = 'display_on_demand_request' + # How often the unchanged current mode is republished for the web UI, which # treats display_current_state older than 120 s as unknown. CURRENT_STATE_REFRESH_SECONDS = 30 @@ -1858,15 +1862,26 @@ class DisplayController: #: perceptible, and it cuts the read rate by 30x. ON_DEMAND_POLL_INTERVAL = 0.25 - #: Shortest gap between _service_pending_changes passes. The same floor as - #: the mailbox poll, since that read is the only real cost in the pass: - #: the schedule checks are gated to once per clock minute and the rest is - #: attribute compares. Callers run at frame rate, so between passes the - #: whole cost is one monotonic-clock compare. - PENDING_CHANGES_INTERVAL = ON_DEMAND_POLL_INTERVAL + #: The mailbox poll while the control socket is up. The web interface then + #: writes the mailbox only when it could not reach the socket (a display + #: being restarted, a web user not yet in the socket's group), and the + #: plugins that still write it get the screen within this long. A look is + #: one stat() of the mailbox file (MailboxWatch). + MAILBOX_POLL_INTERVAL_WITH_SOCKET = 1.0 - #: Class-level default for controllers built without __init__ (tests). + #: Shortest gap between _service_pending_changes passes. The same floor + #: the mailbox poll had before the socket, since that read was the only + #: real cost in the pass: the schedule checks are gated to once per clock + #: minute and the rest is attribute compares. Callers run at frame rate, + #: so between passes the whole cost is one monotonic-clock compare. + PENDING_CHANGES_INTERVAL = 0.25 + + #: Class-level defaults for controllers built without __init__ (tests). _control_server: Optional[ControlServer] = None + #: Created on the first poll; see _poll_on_demand_requests. + _on_demand_mailbox: Optional[MailboxWatch] = None + #: Writers whose mailbox requests have been logged (_note_mailbox_request). + _mailbox_writers_logged: FrozenSet[str] = frozenset() def _service_pending_changes(self) -> None: """Apply changes made elsewhere while the display thread is busy. @@ -2027,10 +2042,10 @@ class DisplayController: processed_id still guards against reprocessing if the delete fails. """ try: - current = self.cache_manager.get('display_on_demand_request', + current = self.cache_manager.get(ON_DEMAND_MAILBOX_KEY, max_age=3600, memory_ttl=0) if not current or current.get('request_id') == request_id: - self.cache_manager.delete('display_on_demand_request') + self.cache_manager.delete(ON_DEMAND_MAILBOX_KEY) else: logger.debug("Newer on-demand request %s arrived while processing " "%s; leaving it in the mailbox", @@ -2048,10 +2063,12 @@ class DisplayController: return hub = StateHub(loop_probe=display_watchdog.watchdog.liveness) try: + from src.error_aggregator import apply_error_clear self._control_server = start_control_server( status_provider=self._control_status, cache_dir=getattr(self.cache_manager, 'cache_dir', None), - state_hub=hub) + state_hub=hub, + handlers={ControlCommand.ERRORS_CLEAR: apply_error_clear}) except Exception: # pylint: disable=broad-except logger.exception("Control socket not started; using the file mailbox only") if self._control_server is not None: @@ -2354,16 +2371,38 @@ class DisplayController: } command.succeed(dict(result)) + def _mailbox_poll_interval(self) -> float: + """How often the on-demand mailbox is looked at: its old 0.25 s when + it is the only way in, MAILBOX_POLL_INTERVAL_WITH_SOCKET while the + control socket carries the web interface's commands.""" + if self._control_server is not None: + return self.MAILBOX_POLL_INTERVAL_WITH_SOCKET + return self.ON_DEMAND_POLL_INTERVAL + def _poll_on_demand_requests(self) -> None: - """Poll cache for new on-demand requests from external controllers.""" + """Apply on-demand requests: the control socket's, then the mailbox's. + + Socket commands are in memory and are applied at once. The file + mailbox (``display_on_demand_request``) is the fallback for a web + interface that could not reach the socket, and the way four plugins + still ask for the screen. It is looked at once per poll interval + (_mailbox_poll_interval), and read only when its file changed + (MailboxWatch): a look that finds nothing new is one stat(). + """ # Socket commands are already in memory: no disk read, so no floor. self._drain_control_commands() now = time.monotonic() if (self._last_on_demand_poll is not None - and now - self._last_on_demand_poll < self.ON_DEMAND_POLL_INTERVAL): + and now - self._last_on_demand_poll < self._mailbox_poll_interval()): return self._last_on_demand_poll = now + watch = self._on_demand_mailbox + if watch is None: + watch = self._on_demand_mailbox = MailboxWatch(ON_DEMAND_MAILBOX_KEY) + if not watch.changed(self.cache_manager): + return + try: # Use a long max_age (1 hour) to ensure requests aren't expired before processing # The request_id check prevents duplicate processing. @@ -2374,21 +2413,52 @@ class DisplayController: # pinned in memory for the full hour and every later poll returned # that stale copy -- meaning no second on-demand request was honoured # for an hour, while the API still reported success. - request = self.cache_manager.get('display_on_demand_request', + request = self.cache_manager.get(ON_DEMAND_MAILBOX_KEY, max_age=3600, memory_ttl=0) except (OSError, RuntimeError, ValueError, TypeError) as err: + watch.forget() # read it again next time logger.error("Failed to read on-demand request: %s", err, exc_info=True) return - if not request: + if not isinstance(request, dict): return + self._note_mailbox_request(request) self._handle_on_demand_request(request) + def _note_mailbox_request(self, request: Dict[str, Any]) -> None: + """Log, once per writer, an on-demand request that came through the + mailbox while the control socket is up. + + The web interface writes the mailbox only when the socket fails, so + this is mostly a plugin that writes ``display_on_demand_request`` + itself. The mailbox is going away; the log says who still uses it. + """ + if self._control_server is None: + return + writer = request.get('plugin_id') or request.get('mode') or 'unknown' + if not isinstance(writer, str): + writer = 'unknown' + if writer in self._mailbox_writers_logged: + return + self._mailbox_writers_logged = self._mailbox_writers_logged | {writer} + logger.info("On-demand %s request %s (for %s) came through the file mailbox " + "although the control socket is up. The mailbox is deprecated: it " + "is read every %.1fs and will be removed in a future release.", + request.get('action'), request.get('request_id'), writer, + self.MAILBOX_POLL_INTERVAL_WITH_SOCKET) + def _handle_on_demand_request(self, request: Dict[str, Any]) -> None: - """Process one on-demand request, from the mailbox or the control socket.""" + """Process one on-demand request, from the mailbox or the control socket. + + A socket command carries ``source: 'socket'``. Only a mailbox request + is removed from the mailbox afterwards: a socket command never put + anything there, so that would be a disk read and maybe a delete for + nothing. + """ request_id = request.get('request_id') if not request_id: return + from_mailbox = request.get('source') != 'socket' action = request.get('action') @@ -2416,27 +2486,35 @@ class DisplayController: # only thing that ends the request: without it the same stop was # re-read and re-processed on every poll, forever, logging at # ON_DEMAND_POLL_INTERVAL for the life of the process. - self._consume_on_demand_request(request_id) + if from_mailbox: + self._consume_on_demand_request(request_id) return - # For start requests, check if already processed + # For start requests, check if already processed. A duplicate in the + # mailbox (a copy of a socket command, or one read before a restart) + # is taken out of it too, so it is not read again. if request_id == self.on_demand_request_id: logger.debug("On-demand start request %s already processed (instance check)", request_id) + if from_mailbox: + self._consume_on_demand_request(request_id) return - + # Also check persistent processed_id (for restart scenarios) processed_request_id = self.cache_manager.get('display_on_demand_processed_id', max_age=3600) if request_id == processed_request_id: logger.debug("On-demand start request %s already processed (persisted check)", request_id) + if from_mailbox: + self._consume_on_demand_request(request_id) return - - logger.info("Received on-demand request %s: %s (plugin_id=%s, mode=%s)", + + logger.info("Received on-demand request %s: %s (plugin_id=%s, mode=%s)", request_id, action, request.get('plugin_id'), request.get('mode')) - + # Mark as processed BEFORE processing (to prevent duplicate processing) self.cache_manager.set('display_on_demand_processed_id', request_id, ttl=3600) self.on_demand_request_id = request_id - self._consume_on_demand_request(request_id) + if from_mailbox: + self._consume_on_demand_request(request_id) if action == 'start': logger.info("Processing on-demand start request for plugin: %s", request.get('plugin_id')) diff --git a/src/error_aggregator.py b/src/error_aggregator.py index cc3bab12..4d100de9 100644 --- a/src/error_aggregator.py +++ b/src/error_aggregator.py @@ -490,10 +490,13 @@ def record_error( # user reads, and the other way round for the clear request. # # ERROR_SNAPSHOT_KEY written by the display service only -# ERROR_CLEAR_REQUEST_KEY written by the web interface only +# ERROR_CLEAR_REQUEST_KEY written by the web interface only, as a fallback # -# A clear is asynchronous: the web interface records a request, and the -# display service applies it (clear_before) on its next tick and republishes. +# A clear goes over the control socket (``errors.clear``): the display applies +# it (clear_before) and republishes the snapshot before it answers. Only when +# the socket cannot carry it (no socket, or a display older than the command) +# does the web interface record a request in the mailbox, which the display +# applies on its next tick; its tick reads that file only when it changed. # Until it has, the web interface hides whatever the snapshot shows from # before the cutoff, so a clear takes effect for readers immediately and a # snapshot published just before the request cannot bring old errors back. @@ -598,13 +601,43 @@ class ErrorSnapshotPublisher: self._published_version: Optional[int] = None self._last_attempt: Optional[float] = None self._applied_clear_id: Optional[str] = None + # The widest cutoff applied in this process: a clear request at or + # before it has nothing left to clear (see _pending_cutoff). + self._applied_clear_cutoff: Optional[float] = None + from src.cache_manager import MailboxWatch # the display's cache, loaded already + self._mailbox = MailboxWatch(ERROR_CLEAR_REQUEST_KEY) self._tick_lock = threading.Lock() self._stop = threading.Event() self._thread: Optional[threading.Thread] = None + def _clear(self, request_id: str, cutoff: float) -> int: + """Apply one clear and remember it. Caller holds _tick_lock.""" + cleared = 0 + if math.isfinite(cutoff): + cleared = self.aggregator.clear_before(datetime.fromtimestamp(cutoff)) + _snapshot_logger.info("Cleared %d plugin error record(s) as requested (%s)", + cleared, request_id) + if self._applied_clear_cutoff is None or cutoff > self._applied_clear_cutoff: + self._applied_clear_cutoff = cutoff + # A malformed request is acknowledged too, so it is not retried forever. + self._applied_clear_id = request_id + return cleared + def _apply_clear_request(self) -> bool: - """Honour a clear request we have not applied yet. True if one was.""" - request = self.cache_manager.get(ERROR_CLEAR_REQUEST_KEY, max_age=None, memory_ttl=0) + """Honour a mailbox clear request we have not applied yet. True if one was. + + The mailbox is the fallback for a web interface that could not use + the control socket (``errors.clear``, :meth:`clear_now`). It is read + only when its file changed since the last tick; otherwise a tick + costs one stat(). + """ + if not self._mailbox.changed(self.cache_manager): + return False + try: + request = self.cache_manager.get(ERROR_CLEAR_REQUEST_KEY, max_age=None, memory_ttl=0) + except Exception: + self._mailbox.forget() + raise if not isinstance(request, dict): return False request_id = request.get("request_id") @@ -614,14 +647,30 @@ class ErrorSnapshotPublisher: cutoff = float(request.get("cutoff")) except (TypeError, ValueError): cutoff = float("nan") - if math.isfinite(cutoff): - cleared = self.aggregator.clear_before(datetime.fromtimestamp(cutoff)) - _snapshot_logger.info("Cleared %d plugin error record(s) as requested (%s)", - cleared, request_id) - # A malformed request is acknowledged too, so it is not retried forever. - self._applied_clear_id = request_id + self._clear(request_id, cutoff) return True + def clear_now(self, request_id: str, cutoff: float) -> int: + """``errors.clear`` over the control socket: apply a clear at once and + republish the snapshot, so the web interface's next read has it. + Returns how many records were cleared. Raises when the snapshot + could not be written, so the caller is not told it worked.""" + with self._tick_lock: + cleared = self._clear(request_id, float(cutoff)) + self._publish(self.aggregator.version, self._clock()) + return cleared + + def _publish(self, version: int, now: float) -> None: + """Write the snapshot. Caller holds _tick_lock.""" + # Stamp the attempt before writing: a cache that keeps failing + # is retried at the throttled rate, not on every tick. + self._last_attempt = now + snapshot = self.aggregator.build_snapshot() + snapshot["applied_clear_id"] = self._applied_clear_id + snapshot["applied_clear_cutoff"] = self._applied_clear_cutoff + self.cache_manager.set(ERROR_SNAPSHOT_KEY, snapshot) + self._published_version = version + def tick(self) -> bool: """Apply a pending clear and publish if due. True if a snapshot was written.""" with self._tick_lock: @@ -635,13 +684,7 @@ class ErrorSnapshotPublisher: if (self._last_attempt is not None and now - self._last_attempt < 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 - snapshot = self.aggregator.build_snapshot() - snapshot["applied_clear_id"] = self._applied_clear_id - self.cache_manager.set(ERROR_SNAPSHOT_KEY, snapshot) - self._published_version = version + self._publish(version, now) return True except Exception as err: # never let reporting break the display _snapshot_logger.debug("Could not publish the plugin error snapshot: %s", @@ -693,6 +736,22 @@ def start_error_snapshot_publisher(cache_manager: Any) -> Optional[ErrorSnapshot return None +def apply_error_clear(request_id: str, args: Any) -> Dict[str, Any]: + """The display's handler for ``errors.clear`` on the control socket. + + ``args`` is the contract's ErrorsClearArgs (``cutoff``, epoch seconds). + Runs on the socket's connection thread: the aggregator and the publisher + have their own locks, and nothing here touches rendering. Returns + ErrorsClearResult once the clear is applied and the snapshot rewritten. + """ + publisher = _snapshot_publisher + if publisher is None: + raise RuntimeError("the error snapshot publisher is not running") + cutoff = float(args.cutoff) + cleared = publisher.clear_now(request_id, cutoff) + return {"request_id": request_id, "cutoff": cutoff, "cleared": cleared} + + # --- Reading side (web interface) ------------------------------------------- def read_error_report(cache_manager: Any) -> Tuple[Optional[Dict[str, Any]], Optional[Dict[str, Any]]]: @@ -731,7 +790,15 @@ def _pending_cutoff(snapshot: Optional[Dict[str, Any]], cutoff = float(clear_request.get("cutoff")) except (TypeError, ValueError): return None - return cutoff if math.isfinite(cutoff) else None + if not math.isfinite(cutoff): + return None + # A wider clear has been applied since (over the control socket): this + # older request has nothing left to hide. + applied = snapshot.get("applied_clear_cutoff") if snapshot is not None else None + if (isinstance(applied, (int, float)) and not isinstance(applied, bool) + and applied >= cutoff): + return None + return cutoff def _is_after(item: Any, field_name: str, cutoff: float) -> bool: @@ -831,13 +898,31 @@ def _count_cleared(summary: Dict[str, Any], cutoff: float) -> Optional[int]: return None -def request_error_clear(cache_manager: Any, cutoff: float) -> Dict[str, Any]: +#: ``send(request_id, cutoff)`` hands a clear to the display over the control +#: socket and returns its ErrorsClearResult, or None when the socket could +#: not carry it and the mailbox should be written instead. Any exception it +#: raises reaches the caller: the display had the request and failed it. +ClearSender = Callable[[str, float], Optional[Dict[str, Any]]] + + +def request_error_clear(cache_manager: Any, cutoff: float, + send: Optional[ClearSender] = None) -> Dict[str, Any]: """Ask the display service to forget errors recorded at or before ``cutoff``. - Returns ``request_id``, ``cutoff`` (ISO, local time), ``cleared_count`` - (see _count_cleared) and ``clear_requested``. Raises OSError when the - request did not reach the shared cache, since a cache without a usable - directory accepts set() and keeps nothing. + Over the control socket when ``send`` is given and carries it: the + display applies the clear and republishes its snapshot before it + answers, so nothing is written here. Otherwise (no socket, or a display + older than ``errors.clear``) a request is written to the + ``plugin_error_clear_request`` mailbox, which the display applies on + its next tick, and readers hide the cleared errors until then. + + Returns ``request_id``, ``cutoff`` (ISO, local time), ``cleared_count``, + ``clear_requested``, ``applied`` (the display has already cleared them) + and ``transport`` (``socket`` or ``mailbox``). ``cleared_count`` is the + display's own count over the socket, else an estimate from the snapshot + (see _count_cleared). Raises OSError when a mailbox request did not reach + the shared cache, since a cache without a usable directory accepts set() + and keeps nothing. A request the display has not applied yet is only ever widened: a later, narrower one ("older than 24 hours" after "everything") overwriting it @@ -848,8 +933,21 @@ def request_error_clear(cache_manager: Any, cutoff: float) -> Dict[str, Any]: if pending is not None: cutoff = max(cutoff, pending) before = error_summary_from_report(snapshot, clear_request) + request_id = uuid.uuid4().hex + answer = { + "clear_requested": True, + "request_id": request_id, + "cutoff": datetime.fromtimestamp(cutoff).isoformat(), + } + if send is not None: + result = send(request_id, cutoff) + if result is not None: + count = result.get("cleared") + return dict(answer, applied=True, transport="socket", + cleared_count=count if isinstance(count, int) and not isinstance(count, bool) + else _count_cleared(before, cutoff)) request = { - "request_id": uuid.uuid4().hex, + "request_id": request_id, "cutoff": cutoff, "requested_at": time.time(), } @@ -857,9 +955,5 @@ def request_error_clear(cache_manager: Any, cutoff: float) -> Dict[str, Any]: stored = cache_manager.get(ERROR_CLEAR_REQUEST_KEY, max_age=None, memory_ttl=0) if not isinstance(stored, dict) or stored.get("request_id") != request["request_id"]: raise OSError("the clear request was not stored in the shared cache") - return { - "cleared_count": _count_cleared(before, cutoff), - "clear_requested": True, - "request_id": request["request_id"], - "cutoff": datetime.fromtimestamp(cutoff).isoformat(), - } + return dict(answer, applied=False, transport="mailbox", + cleared_count=_count_cleared(before, cutoff)) diff --git a/src/ipc/client.py b/src/ipc/client.py index 82749727..17ab5a19 100644 --- a/src/ipc/client.py +++ b/src/ipc/client.py @@ -3,8 +3,13 @@ Every failure -- no socket (the display is stopped, or predates the socket), a refused or timed-out connection, a reply that breaks the contract, or an error the display returned -- raises :class:`ControlError` with a short -``reason``, and the caller falls back to the file mailbox. Nothing here -blocks for longer than ``timeout`` in total. +``reason``. Nothing here blocks for longer than ``timeout`` in total. + +Whether the caller may then write the file mailbox instead is +:func:`should_fall_back`: only when the display never took the request (it +could not be reached, or it is too old to know the command). A display that +took the request and then failed, refused or went quiet is answered as +that, not posted a second time through the mailbox. """ from __future__ import annotations @@ -22,6 +27,7 @@ from src.ipc.contract import ( SUBSCRIBE_KEEPALIVE_SECONDS, SUPPORTED_VERSIONS, Command, + ErrorCode, FrameReader, ProtocolError, Request, @@ -49,17 +55,49 @@ class ControlError(Exception): ``refused``, ``timeout``, ``closed``, ``bad_response``, ``invalid_request``. When the display answered with an error, ``reason`` is that error's :class:`~src.ipc.contract.ErrorCode` (``busy``, ``unknown_command``, ...). + + ``sent`` is True once the whole request was written to a connected + display, which may then have acted on it. A refusal the display sends + before it reads anything (``forbidden``, too many connections) carries + no request id and leaves ``sent`` False. """ - def __init__(self, reason: str, message: str = ''): + def __init__(self, reason: str, message: str = '', *, sent: bool = False): super().__init__(reason, message) self.reason = reason self.message = message + self.sent = sent def __str__(self) -> str: return f'{self.reason}: {self.message}' if self.message else self.reason +#: Answers from a display that read the request but does not speak it: one +#: older than the command (an upgrade in progress) or the protocol version. +#: It did nothing, so the mailbox is the way to reach it. +UPGRADE_REASONS = frozenset({ErrorCode.UNKNOWN_COMMAND, ErrorCode.UNSUPPORTED_VERSION}) + + +def should_fall_back(error: BaseException) -> bool: + """May the caller write the file mailbox after ``error``? + + Yes when the display never took the request: there is no socket (the + display is stopped, predates the socket, or it is switched off), the + connection was refused or timed out, the display turned the connection + away before reading it, or it is too old to know the command + (:data:`UPGRADE_REASONS`). Also for an error that is not a + :class:`ControlError` (a bug in the client), as before. + + No once the display had the request: a ``busy`` queue, ``invalid_args``, + an ``internal`` error, or a timeout or hang-up after the request was + sent. The display may have applied it, or would refuse it from the + mailbox too, so a second copy there only hides the failure. + """ + if not isinstance(error, ControlError): + return True + return not error.sent or error.reason in UPGRADE_REASONS + + def request(cmd: str, args: Optional[Mapping[str, Any]] = None, *, request_id: Optional[str] = None, timeout: float = DEFAULT_TIMEOUT_SECONDS, @@ -93,11 +131,14 @@ def request(cmd: str, args: Optional[Mapping[str, Any]] = None, *, # A refusal before the request was read (forbidden, too many # connections) carries no id. if response.id != request_id and not (response.id is None and not response.ok): - raise ControlError('bad_response', 'the reply is for a different request') + raise ControlError('bad_response', 'the reply is for a different request', sent=True) if not response.ok: error = response.error + # No id: refused at the door (forbidden, too many connections), + # before the display read the request. raise ControlError(error.code if error else 'bad_response', - error.message if error else '') + error.message if error else '', + sent=response.id is not None) return dict(response.result or {}) @@ -142,26 +183,31 @@ def _connect(paths: Sequence[str], deadline: float) -> socket.socket: def _exchange(sock: socket.socket, payload: bytes, deadline: float) -> Response: + """Send ``payload`` and read the reply. A failure once the whole request + is written raises with ``sent=True``: the display may have it.""" + sent = False try: sock.settimeout(_remaining(deadline)) sock.sendall(payload) + sent = True reader = FrameReader(MAX_MESSAGE_BYTES) while True: sock.settimeout(_remaining(deadline)) data = sock.recv(4096) if not data: - raise ControlError('closed', 'the display closed the connection') + raise ControlError('closed', 'the display closed the connection', sent=sent) lines = reader.feed(data) if lines: return Response.from_dict(decode_message(lines[0])) except socket.timeout: - raise ControlError('timeout', 'no reply in time') from None + raise ControlError('timeout', 'no reply in time', sent=sent) from None except ProtocolError as e: - raise ControlError('bad_response', e.message) from None - except ControlError: + raise ControlError('bad_response', e.message, sent=sent) from None + except ControlError as e: + e.sent = e.sent or sent raise except OSError as e: - raise ControlError('closed', str(e)) from None + raise ControlError('closed', str(e), sent=sent) from None # -- commands --------------------------------------------------------------------------- @@ -227,6 +273,21 @@ def plugin_reload(plugin_id: str, *, timeout: Optional[float] = None, else timeout, paths=paths) +def errors_clear(request_id: str, cutoff: float, *, + timeout: float = DEFAULT_TIMEOUT_SECONDS, + paths: Optional[Sequence[str]] = None) -> Dict[str, Any]: + """Have the display forget the plugin errors recorded at or before + ``cutoff`` (epoch seconds) and publish its error snapshot again. + + Returns :class:`~src.ipc.contract.ErrorsClearResult` once it is done. + Raises :class:`ControlError`: ``unknown_command`` from a display older + than the command, which still reads the ``plugin_error_clear_request`` + mailbox. + """ + return request(Command.ERRORS_CLEAR, {'cutoff': cutoff}, request_id=request_id, + timeout=timeout, paths=paths) + + def ping(*, timeout: float = DEFAULT_TIMEOUT_SECONDS, paths: Optional[Sequence[str]] = None) -> Dict[str, Any]: return request(Command.PING, {}, timeout=timeout, paths=paths) diff --git a/src/ipc/contract.py b/src/ipc/contract.py index ec5ef06b..004dd8bb 100644 --- a/src/ipc/contract.py +++ b/src/ipc/contract.py @@ -149,12 +149,13 @@ class Command: PLUGIN_RELOAD = 'plugin.reload' STATE_GET = 'state.get' STATE_SUBSCRIBE = 'state.subscribe' + ERRORS_CLEAR = 'errors.clear' #: Every command version 1 defines, in the order ``hello`` reports them. -#: ``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). +#: ``brightness.set`` and ``plugin.reload`` came in stage 2, ``state.get`` +#: and ``state.subscribe`` in stage 3, and ``errors.clear`` in stage 4, all +#: within version 1 (see the module docstring on adding commands). COMMANDS: Tuple[str, ...] = ( Command.HELLO, Command.PING, @@ -165,8 +166,15 @@ COMMANDS: Tuple[str, ...] = ( Command.PLUGIN_RELOAD, Command.STATE_GET, Command.STATE_SUBSCRIBE, + Command.ERRORS_CLEAR, ) +#: Commands the connection thread answers itself, through a handler the +#: display registers (``ControlServer(handlers=...)``), because they touch +#: nothing the render thread owns. A display that registered none answers +#: ``unknown_command``, and the client falls back as from an older display. +DIRECT_COMMANDS = frozenset({Command.ERRORS_CLEAR}) + #: Commands that are queued for the render thread. QUEUED_COMMANDS = frozenset({Command.ON_DEMAND_START, Command.ON_DEMAND_STOP, Command.BRIGHTNESS_SET, Command.PLUGIN_RELOAD}) @@ -565,8 +573,32 @@ class StateSubscribeArgs: return cls() +@dataclass(frozen=True) +class ErrorsClearArgs: + """``errors.clear``: forget the plugin errors recorded at or before + ``cutoff`` (seconds since the epoch), as ``POST /api/v3/errors/clear`` + asks. The request id is the clear's id, which the display's error + snapshot then reports as ``applied_clear_id``. + """ + cutoff: float + + def to_dict(self) -> Dict[str, Any]: + return {'cutoff': self.cutoff} + + @classmethod + def from_dict(cls, args: Mapping[str, Any]) -> 'ErrorsClearArgs': + value = args.get('cutoff') + if isinstance(value, bool) or not isinstance(value, (int, float)): + raise ProtocolError(ErrorCode.INVALID_ARGS, 'cutoff must be a number of seconds') + if not math.isfinite(value) or value < 0: + raise ProtocolError(ErrorCode.INVALID_ARGS, + 'cutoff must be a finite, non-negative number of seconds') + return cls(cutoff=float(value)) + + CommandArgs = Union[HelloArgs, OnDemandStartArgs, OnDemandStopArgs, NoArgs, - BrightnessSetArgs, PluginReloadArgs, StateGetArgs, StateSubscribeArgs] + BrightnessSetArgs, PluginReloadArgs, StateGetArgs, StateSubscribeArgs, + ErrorsClearArgs] #: The arguments of a command that goes on the render thread's queue. QueuedArgs = Union[OnDemandStartArgs, OnDemandStopArgs, BrightnessSetArgs, PluginReloadArgs] @@ -581,6 +613,7 @@ _ARG_TYPES: Dict[str, Any] = { Command.PLUGIN_RELOAD: PluginReloadArgs, Command.STATE_GET: StateGetArgs, Command.STATE_SUBSCRIBE: StateSubscribeArgs, + Command.ERRORS_CLEAR: ErrorsClearArgs, } @@ -653,6 +686,13 @@ class PluginReloadResult(TypedDict): modes: List[str] +class ErrorsClearResult(TypedDict): + """``errors.clear``, once applied and the error snapshot republished.""" + request_id: str + cutoff: float + cleared: int + + class LoopState(TypedDict): """``loop``: is the render loop still going round? diff --git a/src/ipc/server.py b/src/ipc/server.py index e2b456f4..1b81ca04 100644 --- a/src/ipc/server.py +++ b/src/ipc/server.py @@ -6,7 +6,9 @@ rendering: a command that changes the panel is validated, put on a bounded queue and acknowledged, and the render thread drains that queue at the point where it reads the file mailbox (``DisplayController._poll_on_demand_requests``), handing each command to the same code. Queries (``on_demand.status``) are -answered from a snapshot callable the display provides. +answered from a snapshot callable the display provides, and the few commands +that touch nothing the render thread owns (``errors.clear``) by a handler the +display registers, on the connection thread. The queue also wakes the render thread: :meth:`ControlServer.wait_for_command` is what it waits on in place of a sleep, so a command lands within a frame on @@ -62,6 +64,7 @@ from src.ipc.contract import ( AWAITED_COMMANDS, COMMANDS, DEFAULT_SOCKET_DIR, + DIRECT_COMMANDS, DEFAULT_SOCKET_PATH, MAX_MESSAGE_BYTES, MAX_SUBSCRIBERS, @@ -107,7 +110,8 @@ MAX_CLIENTS = 8 #: Commands waiting for the render thread. It drains them at least every #: 0.25 s, so a full queue means the render thread is stuck, and the client -#: is told ``busy`` (and falls back to the mailbox) instead of piling up work. +#: is told ``busy`` instead of piling up work. The mailbox would not be read +#: either, so the web interface reports the failure rather than fall back. QUEUE_SIZE = 16 #: Timeout for one recv()/send() on a connection. @@ -552,6 +556,12 @@ def server_socket_path(environ: Optional[Mapping[str, str]] = None) -> Optional[ StatusProvider = Callable[[], Dict[str, Any]] +#: A handler for one of DIRECT_COMMANDS, ``(request_id, args) -> result``. It +#: runs on the connection thread, so it must not touch what the render thread +#: owns. It may raise ProtocolError to answer with that error's code; any +#: other exception is answered ``internal``. +DirectHandler = Callable[[str, Any], Mapping[str, Any]] + class ControlServer: """Serves the control socket on background threads. @@ -569,9 +579,12 @@ class ControlServer: await_seconds: Optional[Mapping[str, float]] = None, state_hub: Optional[StateHub] = None, max_subscribers: int = MAX_SUBSCRIBERS, - keepalive: float = SUBSCRIBE_KEEPALIVE_SECONDS): + keepalive: float = SUBSCRIBE_KEEPALIVE_SECONDS, + handlers: Optional[Mapping[str, DirectHandler]] = None): self.path = path self.state_hub = state_hub + self._handlers: Dict[str, DirectHandler] = { + cmd: fn for cmd, fn in (handlers or {}).items() if cmd in DIRECT_COMMANDS} self._subscriber_slots = threading.BoundedSemaphore(max_subscribers) self._keepalive = keepalive self._await_seconds: Dict[str, float] = dict(AWAIT_SECONDS) @@ -1028,6 +1041,9 @@ class ControlServer: snap = hub.snapshot() return Response.success(request.id, fit_snapshot(snap), v=request.v) + if request.cmd in DIRECT_COMMANDS: + return self._direct(request, args) + if request.cmd in QUEUED_COMMANDS and isinstance(args, ( OnDemandStartArgs, OnDemandStopArgs, BrightnessSetArgs, PluginReloadArgs)): awaited = request.cmd in AWAITED_COMMANDS @@ -1055,6 +1071,25 @@ class ControlServer: return Response.failure(request.id, ErrorCode.INTERNAL, f'{request.cmd} is not implemented', v=request.v) + def _direct(self, request: Request, args: Any) -> Response: + """A command the display answers on this thread (DIRECT_COMMANDS).""" + handler = self._handlers.get(request.cmd) + if handler is None: + # Answered as an older display would, so the client falls back. + return Response.failure(request.id, ErrorCode.UNKNOWN_COMMAND, + f'{request.cmd} is not served by this display', + v=request.v) + try: + result = handler(request.id, args) + except ProtocolError as e: + return Response.failure(request.id, e.code, e.message, v=request.v) + except Exception: # pylint: disable=broad-except + logger.exception("Control socket: %s %s failed", request.cmd, request.id) + return Response.failure(request.id, ErrorCode.INTERNAL, + 'the display failed to apply it', v=request.v) + logger.info("Control socket applied %s %s", request.cmd, request.id) + return Response.success(request.id, dict(result), v=request.v) + def _await_outcome(self, request: Request, outcome: CommandOutcome) -> Response: """Answer an awaited command once the render thread has applied it. @@ -1079,7 +1114,9 @@ class ControlServer: def start_control_server(status_provider: Optional[StatusProvider] = None, cache_dir: Optional[str] = None, environ: Optional[Mapping[str, str]] = None, - state_hub: Optional[StateHub] = None) -> Optional[ControlServer]: + state_hub: Optional[StateHub] = None, + handlers: Optional[Mapping[str, DirectHandler]] = 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 @@ -1091,7 +1128,7 @@ def start_control_server(status_provider: Optional[StatusProvider] = 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), - state_hub=state_hub) + state_hub=state_hub, handlers=handlers) return server if server.start() else None diff --git a/test/_run_loop_harness.py b/test/_run_loop_harness.py index 80b04e5c..7bed2e94 100644 --- a/test/_run_loop_harness.py +++ b/test/_run_loop_harness.py @@ -143,12 +143,24 @@ class FakeCache: def __init__(self): self.data: Dict[str, Any] = {} self.cache_dir = "/nonexistent/run-loop-harness" + self._writes = 0 + self._written: Dict[str, int] = {} def get(self, key, max_age=None, memory_ttl=None): return self.data.get(key) def set(self, key, data, ttl=None): self.data[key] = data + # Every write is a new file, as DiskCache's rename makes it. + self._writes += 1 + self._written[key] = self._writes + + def file_signature(self, key): + """CacheManager.file_signature: None without a file, else a value + that changes with every write.""" + if key not in self.data: + return None + return (self._written.get(key, 0), 0, 0) def delete(self, key): self.data.pop(key, None) diff --git a/test/test_api_v3_on_demand_socket.py b/test/test_api_v3_on_demand_socket.py index 8498b8bd..c93fadd0 100644 --- a/test/test_api_v3_on_demand_socket.py +++ b/test/test_api_v3_on_demand_socket.py @@ -1,10 +1,13 @@ """POST /display/on-demand/start and /stop: control socket first, mailbox fallback. The routes hand the request to the display over the control socket -(src/ipc) and get an acknowledgement. On any failure -- no socket (a stopped -display, or one older than the socket), a timeout, a refusal, a bug in the -client -- they write the file mailbox exactly as they did before the socket -existed. These tests pin both paths, that exactly one of them is used, that +(src/ipc) and get an acknowledgement. When the socket could not carry the +request -- no socket (a stopped display, or one older than the socket), a +refused or timed-out connect, a display too old to know the command, a bug +in the client -- they write the file mailbox exactly as they did before the +socket existed. When the display had the request and failed it (a full +queue, bad arguments, no answer in time) the route says so and writes +nothing. These tests pin each path, that at most one of them is used, that the response says which, and that the request id is the same either way (the display deduplicates on it). @@ -101,12 +104,15 @@ class TestSocketPath: class TestMailboxFallback: @pytest.mark.parametrize("reason", [ "no_socket", "refused", "timeout", "closed", "bad_response", "invalid_request", - "busy", "unknown_command", "unsupported_version", "disabled", "unsupported", + "busy", "forbidden", "unknown_command", "unsupported_version", "disabled", + "unsupported", ]) - def test_any_socket_failure_writes_the_mailbox_as_before( + def test_a_request_the_socket_never_carried_writes_the_mailbox( self, api_v3_client, service, reason): + # sent=False: the display never had it (no socket, a refused or + # timed-out connect, turned away at the door). with patch(f"{CLIENT}.on_demand_start", - side_effect=control_client.ControlError(reason, "x")): + side_effect=control_client.ControlError(reason, "x", sent=False)): resp = api_v3_client.post(START_URL, json={ "plugin_id": "weather", "mode": "weather_current", "duration": 60, "pinned": True}) @@ -168,6 +174,60 @@ class TestMailboxFallback: assert data["socket_error"] in ("disabled", "unsupported") # Linux, Windows +class TestTheDisplayHadIt: + """Once the display has the request, its answer stands: no mailbox copy. + + A busy queue, a refusal or silence after the request was sent mean the + display may have applied it, or would refuse the mailbox copy too, so + the route reports the failure instead of posting it a second time. + """ + + @pytest.mark.parametrize("reason,status", [ + ("busy", 503), ("internal", 503), ("timeout", 503), ("closed", 503), + ("bad_response", 503), ("invalid_args", 400), + ]) + def test_start_is_answered_with_the_failure(self, api_v3_client, service, reason, status): + with patch(f"{CLIENT}.on_demand_start", + side_effect=control_client.ControlError(reason, "x", sent=True)): + resp = api_v3_client.post(START_URL, json={"plugin_id": "weather"}) + assert resp.status_code == status + body = resp.get_json() + assert body["status"] == "error" + assert body["data"]["transport"] == "socket" + assert body["data"]["socket_error"] == reason + assert _mailbox_writes(service["cache"]) == [] + assert not [call for call in service["calls"] if call[0] == "systemctl"] + + def test_stop_is_answered_with_the_failure(self, api_v3_client, service): + with patch(f"{CLIENT}.on_demand_stop", + side_effect=control_client.ControlError("busy", "x", sent=True)): + resp = api_v3_client.post(STOP_URL, json={}) + assert resp.status_code == 503 + assert _mailbox_writes(service["cache"]) == [] + + def test_a_stop_with_stop_service_still_stops_the_service(self, api_v3_client, service): + with patch(f"{CLIENT}.on_demand_stop", + side_effect=control_client.ControlError("timeout", "x", sent=True)), \ + patch("web_interface.blueprints.api_v3.display._stop_display_service", + return_value={"active": False}) as stop: + resp = api_v3_client.post(STOP_URL, json={"stop_service": True}) + assert resp.status_code == 200 + data = resp.get_json()["data"] + assert data["transport"] == "socket" and data["socket_error"] == "timeout" + stop.assert_called_once() + assert _mailbox_writes(service["cache"]) == [] + + @pytest.mark.parametrize("reason", ["unknown_command", "unsupported_version"]) + def test_an_older_display_that_does_not_speak_it_gets_the_mailbox( + self, api_v3_client, service, reason): + # The upgrade case: new web interface, display still on an old build. + with patch(f"{CLIENT}.on_demand_start", + side_effect=control_client.ControlError(reason, "x", sent=True)): + data = api_v3_client.post(START_URL, json={"plugin_id": "weather"}).get_json()["data"] + assert data["transport"] == "mailbox" and data["socket_error"] == reason + assert len(_mailbox_writes(service["cache"])) == 1 + + @pytest.mark.skipif(not c.socket_supported(), reason="AF_UNIX sockets are Linux/macOS only") class TestRealSocket: @pytest.fixture @@ -204,3 +264,23 @@ class TestRealSocket: data = api_v3_client.post(START_URL, json={"plugin_id": "weather"}).get_json()["data"] assert data["transport"] == "mailbox" and data["socket_error"] == "no_socket" assert len(_mailbox_writes(service["cache"])) == 1 + + def test_a_full_queue_is_reported_not_mailed(self, api_v3_client, service, monkeypatch): + import shutil + import tempfile + from src.ipc.server import ControlServer + d = tempfile.mkdtemp(prefix="lmipc-") + path = os.path.join(d, "control.sock") + server = ControlServer(path, status_provider=dict, queue_size=1) + assert server.start() + monkeypatch.setenv(c.SOCKET_PATH_ENV, path) + try: + first = api_v3_client.post(START_URL, json={"plugin_id": "weather"}) + assert first.get_json()["data"]["transport"] == "socket" + second = api_v3_client.post(START_URL, json={"plugin_id": "weather"}) + assert second.status_code == 503 + assert second.get_json()["data"]["socket_error"] == "busy" + assert _mailbox_writes(service["cache"]) == [] + finally: + server.close() + shutil.rmtree(d, ignore_errors=True) diff --git a/test/test_error_snapshot_cross_process.py b/test/test_error_snapshot_cross_process.py index 3cc48614..f7b8ec3e 100644 --- a/test/test_error_snapshot_cross_process.py +++ b/test/test_error_snapshot_cross_process.py @@ -430,6 +430,157 @@ class TestClear: assert "clear request" in response.get_json()["message"] +CLIENT = "web_interface.blueprints.api_v3.control_client" + + +def _mailbox_file(shared_cache): + _, _, directory = shared_cache + return directory / f"{ERROR_CLEAR_REQUEST_KEY}.json" + + +class TestClearOverTheSocket: + """``errors.clear``: the display applies the clear before it answers, and + the mailbox is written only when the socket could not carry it.""" + + @pytest.fixture + def socket_up(self, display, monkeypatch): + """The control socket, as the display serves it: errors_clear runs + the display's own handler against its publisher.""" + from src.ipc import client as control_client + from src.ipc.contract import ErrorsClearArgs + _, publisher, _ = display + monkeypatch.setattr(errors, "_snapshot_publisher", publisher) + calls = [] + + def errors_clear(request_id, cutoff, **kw): + calls.append((request_id, cutoff)) + return errors.apply_error_clear(request_id, ErrorsClearArgs(cutoff=cutoff)) + + monkeypatch.setattr(f"{CLIENT}.errors_clear", errors_clear) + assert control_client.errors_clear is errors_clear + return calls + + def test_socket_clear_is_applied_before_the_answer(self, web, display, socket_up, + shared_cache): + aggregator, publisher, _ = display + for _ in range(3): + _fail(aggregator) + publisher.tick() + response = web.post("/api/v3/errors/clear", json={"all": True}) + assert response.status_code == 200 + body = response.get_json() + data = body["data"] + assert data["transport"] == "socket" and data["applied"] is True + assert data["cleared_count"] == 3 + assert body["message"] == "Cleared all errors" + [(request_id, _)] = socket_up + assert data["request_id"] == request_id + # Applied already: no display tick needed, nothing pending. + assert aggregator.get_error_summary()["total_errors"] == 0 + summary = _summary(web) + assert summary["total_errors"] == 0 and summary["clear_pending"] is False + # And no mailbox file. + assert not _mailbox_file(shared_cache).exists() + + def test_an_older_mailbox_request_does_not_read_as_pending(self, web, display, socket_up, + shared_cache, monkeypatch): + # A clear that went to the mailbox while the socket was down, then a + # wider one over the socket: the old request has nothing left to hide. + from unittest.mock import patch + from src.ipc import client as control_client + aggregator, publisher, _ = display + _fail(aggregator) + publisher.tick() + with patch(f"{CLIENT}.errors_clear", + side_effect=control_client.ControlError("no_socket")): + data = web.post("/api/v3/errors/clear", json={"max_age_hours": 1}).get_json()["data"] + assert data["transport"] == "mailbox" + assert _mailbox_file(shared_cache).exists() + data = web.post("/api/v3/errors/clear", json={"all": True}).get_json()["data"] + assert data["transport"] == "socket" + assert _summary(web)["clear_pending"] is False + + @pytest.mark.parametrize("reason", ["no_socket", "refused", "disabled", "unsupported"]) + def test_no_socket_writes_the_mailbox(self, web, display, shared_cache, monkeypatch, reason): + from src.ipc import client as control_client + monkeypatch.setattr(f"{CLIENT}.errors_clear", MagicMock( + side_effect=control_client.ControlError(reason, sent=False))) + aggregator, publisher, _ = display + _fail(aggregator) + publisher.tick() + data = web.post("/api/v3/errors/clear", json={"all": True}).get_json()["data"] + assert data["transport"] == "mailbox" and data["applied"] is False + assert _mailbox_file(shared_cache).exists() + assert _summary(web)["clear_pending"] is True + publisher.tick() + assert aggregator.get_error_summary()["total_errors"] == 0 + + def test_an_older_display_gets_the_mailbox(self, web, display, shared_cache, monkeypatch): + # The upgrade case: new web interface, a display from before errors.clear. + from src.ipc import client as control_client + monkeypatch.setattr(f"{CLIENT}.errors_clear", MagicMock( + side_effect=control_client.ControlError("unknown_command", sent=True))) + aggregator, publisher, _ = display + _fail(aggregator) + publisher.tick() + data = web.post("/api/v3/errors/clear", json={"all": True}).get_json()["data"] + assert data["transport"] == "mailbox" + publisher.tick() + assert aggregator.get_error_summary()["total_errors"] == 0 + + @pytest.mark.parametrize("reason", ["internal", "timeout", "busy", "invalid_args"]) + def test_a_display_that_had_it_and_failed_is_an_error(self, web, display, shared_cache, + monkeypatch, reason): + from src.ipc import client as control_client + monkeypatch.setattr(f"{CLIENT}.errors_clear", MagicMock( + side_effect=control_client.ControlError(reason, sent=True))) + response = web.post("/api/v3/errors/clear", json={"all": True}) + assert response.status_code == 503 + assert response.get_json()["context"]["socket_error"] == reason + assert not _mailbox_file(shared_cache).exists() + + +class TestPublisherMailboxPoll: + def test_the_mailbox_is_read_only_when_its_file_changed(self, display, shared_cache): + _, publisher, _ = display + _, web_cache, _ = shared_cache + publisher.tick() + publisher.cache_manager = MagicMock(wraps=publisher.cache_manager) + + def reads(): + return [c for c in publisher.cache_manager.get.call_args_list + if c.args[0] == ERROR_CLEAR_REQUEST_KEY] + + for _ in range(5): + publisher.tick() + assert reads() == [] # no file: a stat per tick, no read + errors.request_error_clear(web_cache, 1.0) + publisher.tick() + assert len(reads()) == 1 + for _ in range(5): + publisher.tick() + assert len(reads()) == 1 # unchanged file: not read again + errors.request_error_clear(web_cache, 2.0) + publisher.tick() + assert len(reads()) == 2 + + def test_clear_now_publishes_what_it_applied(self, display, shared_cache): + aggregator, publisher, _ = display + _, web_cache, _ = shared_cache + _fail(aggregator) + assert publisher.clear_now("sock-1", datetime.now().timestamp() + 1) == 1 + snapshot = web_cache.get(ERROR_SNAPSHOT_KEY, max_age=None, memory_ttl=0) + assert snapshot["applied_clear_id"] == "sock-1" + assert snapshot["total_errors"] == 0 + assert snapshot["applied_clear_cutoff"] is not None + + def test_the_handler_needs_a_running_publisher(self, monkeypatch): + from src.ipc.contract import ErrorsClearArgs + monkeypatch.setattr(errors, "_snapshot_publisher", None) + with pytest.raises(RuntimeError): + errors.apply_error_clear("x", ErrorsClearArgs(cutoff=1.0)) + + @pytest.mark.skipif(not hasattr(os, "fchmod") or os.name == "nt", reason="POSIX file modes") def test_both_files_are_group_readable(web, display, shared_cache): diff --git a/test/test_ipc_contract.py b/test/test_ipc_contract.py index 3d7d3cd8..1847f039 100644 --- a/test/test_ipc_contract.py +++ b/test/test_ipc_contract.py @@ -167,7 +167,8 @@ class TestOnDemandArgs: def test_every_command_has_an_argument_type(self, cmd): args = {Command.ON_DEMAND_START: {'plugin_id': 'p'}, Command.PLUGIN_RELOAD: {'plugin_id': 'p'}, - Command.BRIGHTNESS_SET: {'brightness': 50}}.get(cmd, {}) + Command.BRIGHTNESS_SET: {'brightness': 50}, + Command.ERRORS_CLEAR: {'cutoff': 1790000000.0}}.get(cmd, {}) c.parse_args(cmd, args) def test_hello_versions(self): diff --git a/test/test_ipc_stage4.py b/test/test_ipc_stage4.py new file mode 100644 index 00000000..175c8986 --- /dev/null +++ b/test/test_ipc_stage4.py @@ -0,0 +1,509 @@ +"""Stage 4 of the control socket: the file mailboxes are only a fallback. + +* The client knows whether the display had the request (``ControlError.sent``) + and ``should_fall_back`` allows a mailbox write only when it did not, or + when the display is too old to know the command (the upgrade case). +* ``errors.clear`` is answered on the connection thread by a handler the + display registers; a display without one answers like an older display. +* The display looks at the on-demand mailbox once a second while the socket + is up (0.25 s without it), reads it only when its file changed, never + touches it for a socket command, and logs who still writes it. +* ``CacheManager.file_signature`` / ``MailboxWatch`` make a look one stat(). + +The web routes are covered in test_api_v3_on_demand_socket.py and +test_error_snapshot_cross_process.py. +""" + +import json +import logging +import os +import socket +import time +from unittest.mock import MagicMock, patch + +import pytest + +from src.cache_manager import CacheManager, MailboxWatch +from src.ipc import client +from src.ipc import contract as c +from src.ipc.contract import Command, ErrorsClearArgs, OnDemandStartArgs, ProtocolError +from src.ipc.server import ControlServer, QueuedCommand + +MAILBOX = 'display_on_demand_request' + + +# -- the client: was the request sent? --------------------------------------------- + +class FakeSock: + """Stands in for a connected socket in client._exchange.""" + + def __init__(self, replies=(), send_error=None, recv_error=None): + self.replies = list(replies) + self.send_error = send_error + self.recv_error = recv_error + self.sent = b'' + + def settimeout(self, _t): + pass + + def sendall(self, data): + if self.send_error is not None: + raise self.send_error + self.sent += data + + def recv(self, _n): + if self.recv_error is not None: + raise self.recv_error + return self.replies.pop(0) if self.replies else b'' + + def close(self): + pass + + +def _reply(request_id, **body): + return (json.dumps(dict({'v': 1, 'id': request_id}, **body)) + '\n').encode() + + +def _call(sock=None, connect_error=None, request_id='rid-1'): + with patch.object(client, 'socket_supported', return_value=True), \ + patch.object(client, '_connect', + side_effect=connect_error, return_value=sock): + return client.request(Command.PING, {}, request_id=request_id, paths=['/x.sock']) + + +def _error(**kw): + with pytest.raises(client.ControlError) as e: + _call(**kw) + return e.value + + +class TestSent: + def test_an_answer_is_returned(self): + assert _call(FakeSock([_reply('rid-1', ok=True, result={'pong': True})])) == {'pong': True} + + @pytest.mark.parametrize('reason', ['no_socket', 'refused', 'timeout', 'busy']) + def test_a_failed_connect_was_not_sent(self, reason): + e = _error(connect_error=client.ControlError(reason)) + assert e.reason == reason and e.sent is False + assert client.should_fall_back(e) + + def test_a_send_that_timed_out_was_not_sent(self): + e = _error(sock=FakeSock(send_error=socket.timeout())) + assert e.reason == 'timeout' and e.sent is False + assert client.should_fall_back(e) + + def test_silence_after_the_request_was_sent(self): + e = _error(sock=FakeSock(recv_error=socket.timeout())) + assert e.reason == 'timeout' and e.sent is True + assert not client.should_fall_back(e) + + def test_a_hang_up_after_the_request_was_sent(self): + e = _error(sock=FakeSock([])) + assert e.reason == 'closed' and e.sent is True + assert not client.should_fall_back(e) + + def test_a_garbled_reply(self): + e = _error(sock=FakeSock([b'not json\n'])) + assert e.reason == 'bad_response' and e.sent is True + assert not client.should_fall_back(e) + + @pytest.mark.parametrize('code', ['busy', 'invalid_args', 'internal', 'pending', 'failed']) + def test_a_display_error_with_an_id_was_sent(self, code): + sock = FakeSock([_reply('rid-1', ok=False, error={'code': code, 'message': 'x'})]) + e = _error(sock=sock) + assert e.reason == code and e.sent is True + assert not client.should_fall_back(e) + + @pytest.mark.parametrize('code', ['forbidden', 'busy']) + def test_a_refusal_at_the_door_was_not_sent(self, code): + # forbidden, or too many connections: answered before the request + # was read, so with no id. + sock = FakeSock([(json.dumps({'v': 1, 'id': None, 'ok': False, + 'error': {'code': code, 'message': 'x'}}) + '\n').encode()]) + e = _error(sock=sock) + assert e.reason == code and e.sent is False + assert client.should_fall_back(e) + + @pytest.mark.parametrize('code', ['unknown_command', 'unsupported_version']) + def test_an_older_display_is_fallen_back_from(self, code): + sock = FakeSock([_reply('rid-1', ok=False, error={'code': code, 'message': 'x'})]) + e = _error(sock=sock) + assert e.sent is True + assert client.should_fall_back(e) + + def test_a_request_refused_locally_never_left(self): + with pytest.raises(client.ControlError) as e: + client.request(Command.ERRORS_CLEAR, {'cutoff': 'soon'}, paths=['/x.sock']) + assert e.value.reason == 'invalid_request' and e.value.sent is False + assert client.should_fall_back(e.value) + + def test_a_client_bug_falls_back(self): + assert client.should_fall_back(RuntimeError('boom')) + + +# -- errors.clear on the server ------------------------------------------------------ + +def _line(cmd, args, rid='r1'): + return c.encode_message({'v': 1, 'id': rid, 'cmd': cmd, 'args': args}) + + +class TestErrorsClearOnTheServer: + def test_the_handler_answers_it(self, tmp_path): + seen = [] + + def handler(request_id, args): + seen.append((request_id, args)) + return {'request_id': request_id, 'cutoff': args.cutoff, 'cleared': 4} + + server = ControlServer(str(tmp_path / 's.sock'), + handlers={Command.ERRORS_CLEAR: handler}) + response = server.handle_line(_line(Command.ERRORS_CLEAR, {'cutoff': 123})) + assert response.ok and response.result == {'request_id': 'r1', 'cutoff': 123.0, + 'cleared': 4} + assert seen == [('r1', ErrorsClearArgs(cutoff=123.0))] + assert not server.has_pending # not queued for the render thread + + def test_a_display_without_a_handler_answers_like_an_older_one(self, tmp_path): + server = ControlServer(str(tmp_path / 's.sock')) + response = server.handle_line(_line(Command.ERRORS_CLEAR, {'cutoff': 1})) + assert not response.ok and response.error.code == c.ErrorCode.UNKNOWN_COMMAND + + def test_only_direct_commands_take_a_handler(self, tmp_path): + server = ControlServer(str(tmp_path / 's.sock'), + handlers={Command.ON_DEMAND_START: lambda *a: {}}) + response = server.handle_line(_line(Command.ON_DEMAND_START, {'plugin_id': 'p'})) + assert response.ok and response.result['accepted'] is True # still queued + assert server.has_pending + + def test_a_handler_error_is_contained(self, tmp_path): + def boom(*_a): + raise ValueError('disk gone') + + server = ControlServer(str(tmp_path / 's.sock'), handlers={Command.ERRORS_CLEAR: boom}) + response = server.handle_line(_line(Command.ERRORS_CLEAR, {'cutoff': 1})) + assert response.error.code == c.ErrorCode.INTERNAL + assert 'disk gone' not in response.error.message + + def test_a_handler_can_refuse_with_a_code(self, tmp_path): + def refuse(*_a): + raise ProtocolError(c.ErrorCode.BUSY, 'later') + + server = ControlServer(str(tmp_path / 's.sock'), handlers={Command.ERRORS_CLEAR: refuse}) + assert server.handle_line( + _line(Command.ERRORS_CLEAR, {'cutoff': 1})).error.code == c.ErrorCode.BUSY + + @pytest.mark.parametrize('cutoff', ['1', None, True, float('inf'), -1]) + def test_bad_cutoffs_are_refused(self, cutoff): + with pytest.raises(ProtocolError) as e: + ErrorsClearArgs.from_dict({'cutoff': cutoff}) + assert e.value.code == c.ErrorCode.INVALID_ARGS + + def test_hello_lists_it(self, tmp_path): + server = ControlServer(str(tmp_path / 's.sock')) + result = server.handle_line(_line(Command.HELLO, {'versions': [1]})).result + assert Command.ERRORS_CLEAR in result['commands'] + + +# -- the display's mailbox poll ------------------------------------------------------ + +class SignedCache: + """The slice of CacheManager the poll uses, counting what it costs.""" + + def __init__(self): + self.data = {} + self.writes = 0 + self.version = {} + self.reads = [] + self.stats = 0 + self.deletes = [] + self.sets = [] + + def file_signature(self, key): + self.stats += 1 + return (self.version[key], 0, 0) if key in self.data else None + + def get(self, key, *a, **kw): + self.reads.append(key) + return self.data.get(key) + + def set(self, key, value, *a, **kw): + self.sets.append(key) + self.data[key] = value + self.writes += 1 + self.version[key] = self.writes + + def delete(self, key): + self.deletes.append(key) + self.data.pop(key, None) + + +class FakeServer: + def __init__(self): + self.commands = [] + + @property + def has_pending(self): + return bool(self.commands) + + def drain(self): + out, self.commands = self.commands, [] + return out + + +class Clock: + def __init__(self): + self.t = 1000.0 + + def __call__(self): + return self.t + + +@pytest.fixture +def controller(test_display_controller, monkeypatch): + dc = test_display_controller + dc.cache_manager = SignedCache() + dc._activate_on_demand = MagicMock() + dc.on_demand_active = False + dc.on_demand_request_id = None + dc._last_on_demand_poll = None + dc._on_demand_mailbox = None + dc._mailbox_writers_logged = frozenset() + clock = Clock() + monkeypatch.setattr('src.display_controller.time.monotonic', clock) + dc.clock = clock + return dc + + +def _post(dc, rid, action='start', **fields): + dc.cache_manager.set(MAILBOX, dict({'request_id': rid, 'action': action}, **fields)) + + +def _poll_for(dc, seconds, step=1 / 16): # exact in binary: no drift past a floor + end = dc.clock.t + seconds + while dc.clock.t < end: + dc._poll_on_demand_requests() + dc.clock.t += step + + +class TestMailboxCadence: + def test_without_a_socket_it_is_looked_at_every_quarter_second(self, controller): + controller._control_server = None + _poll_for(controller, 10.0) + assert 38 <= controller.cache_manager.stats <= 42 + + def test_with_a_socket_it_is_looked_at_once_a_second(self, controller): + controller._control_server = FakeServer() + _poll_for(controller, 10.0) + assert 9 <= controller.cache_manager.stats <= 11 + + def test_a_look_that_finds_nothing_reads_nothing(self, controller): + controller._control_server = FakeServer() + _poll_for(controller, 10.0) + assert controller.cache_manager.reads == [] + + def test_an_unchanged_mailbox_is_not_read_again(self, controller): + # An already-processed start the delete could not remove, say. + controller._control_server = FakeServer() + controller.cache_manager.delete = MagicMock() # the file stays + _post(controller, 'once', plugin_id='clock') + _poll_for(controller, 10.0) + assert controller.cache_manager.reads.count(MAILBOX) <= 2 # the read + the re-check + controller._activate_on_demand.assert_called_once() + + def test_a_mailbox_request_lands_within_a_second_with_the_socket_up(self, controller): + # The upgrade case the other way round: a new display, and a web + # interface (or a plugin) that still writes the mailbox. + controller._control_server = FakeServer() + controller._poll_on_demand_requests() + controller.clock.t += 0.1 + _post(controller, 'old-web', plugin_id='clock') + posted = controller.clock.t + while not controller._activate_on_demand.called: + controller._poll_on_demand_requests() + controller.clock.t += 0.05 + assert controller.clock.t - posted < 1.5 + assert controller.clock.t - posted <= controller.MAILBOX_POLL_INTERVAL_WITH_SOCKET + 0.06 + assert MAILBOX in controller.cache_manager.deletes # consumed + + def test_socket_commands_still_land_at_once(self, controller): + server = controller._control_server = FakeServer() + controller._poll_on_demand_requests() + server.commands.append(QueuedCommand('sock', Command.ON_DEMAND_START, + OnDemandStartArgs(plugin_id='clock'), time.time())) + controller._poll_on_demand_requests() # inside the mailbox interval + controller._activate_on_demand.assert_called_once() + + +class TestSocketCommandsLeaveTheMailboxAlone: + def test_a_socket_start_reads_and_deletes_no_mailbox(self, controller): + server = controller._control_server = FakeServer() + controller._poll_on_demand_requests() + before = list(controller.cache_manager.reads) + server.commands.append(QueuedCommand('s1', Command.ON_DEMAND_START, + OnDemandStartArgs(plugin_id='clock'), time.time())) + controller._poll_on_demand_requests() + controller._activate_on_demand.assert_called_once() + assert MAILBOX not in controller.cache_manager.reads[len(before):] + assert controller.cache_manager.deletes == [] + + def test_a_socket_stop_reads_and_deletes_no_mailbox(self, controller): + from src.ipc.contract import OnDemandStopArgs + controller.on_demand_active = True + controller._clear_on_demand = MagicMock() + server = controller._control_server = FakeServer() + server.commands.append(QueuedCommand('s2', Command.ON_DEMAND_STOP, + OnDemandStopArgs(), time.time())) + controller.clock.t += 5 + controller._drain_control_commands() + controller._clear_on_demand.assert_called_once() + assert MAILBOX not in controller.cache_manager.reads + assert controller.cache_manager.deletes == [] + + def test_a_mailbox_copy_of_a_socket_command_is_dropped(self, controller): + # An older web interface timed out after the display queued the + # command, then wrote the mailbox too. + server = controller._control_server = FakeServer() + server.commands.append(QueuedCommand('both', Command.ON_DEMAND_START, + OnDemandStartArgs(plugin_id='clock'), time.time())) + controller._poll_on_demand_requests() + _post(controller, 'both', plugin_id='clock') + _poll_for(controller, 2.0) + controller._activate_on_demand.assert_called_once() + assert MAILBOX not in controller.cache_manager.data + + +class TestDeprecationLog: + def test_each_mailbox_writer_is_logged_once(self, controller, caplog): + controller._control_server = FakeServer() + caplog.set_level(logging.INFO, logger='src.display_controller') + for i, plugin in enumerate(['on-air', 'on-air', 'pomodoro-timer']): + _post(controller, f'r{i}', plugin_id=plugin) + _poll_for(controller, 1.2) + lines = [r.getMessage() for r in caplog.records if 'file mailbox' in r.getMessage()] + assert len(lines) == 2 + assert 'on-air' in lines[0] and 'pomodoro-timer' in lines[1] + + def test_nothing_is_logged_without_a_socket(self, controller, caplog): + controller._control_server = None + caplog.set_level(logging.INFO, logger='src.display_controller') + _post(controller, 'r', plugin_id='on-air') + _poll_for(controller, 1.0) + controller._activate_on_demand.assert_called_once() + assert not [r for r in caplog.records if 'file mailbox' in r.getMessage()] + + +# -- file_signature and MailboxWatch ------------------------------------------------- + +@pytest.fixture +def real_cache(tmp_path, monkeypatch): + monkeypatch.setattr(CacheManager, '_get_writable_cache_dir', lambda self: str(tmp_path)) + cache = CacheManager() + yield cache + cache.stop_cleanup_thread() + + +class TestFileSignature: + def test_absent_key(self, real_cache): + assert real_cache.file_signature('nothing') is None + + def test_every_write_is_a_new_signature(self, real_cache): + seen = set() + for i in range(20): + # Same size each time, written as fast as possible. + real_cache.set(MAILBOX, {'request_id': f'r{i:02d}'}) + sig = real_cache.file_signature(MAILBOX) + assert isinstance(sig, tuple) + seen.add(sig) + assert len(seen) == 20 + + def test_gone_after_a_delete(self, real_cache): + real_cache.set(MAILBOX, {'a': 1}) + real_cache.delete(MAILBOX) + assert real_cache.file_signature(MAILBOX) is None + + +class TestMailboxWatch: + def test_reads_once_per_write(self, real_cache): + watch = MailboxWatch(MAILBOX) + assert watch.changed(real_cache) is False # no file + real_cache.set(MAILBOX, {'request_id': 'a'}) + assert watch.changed(real_cache) is True + assert watch.changed(real_cache) is False + real_cache.set(MAILBOX, {'request_id': 'b'}) + assert watch.changed(real_cache) is True + + def test_forget_reads_again(self, real_cache): + watch = MailboxWatch(MAILBOX) + real_cache.set(MAILBOX, {'request_id': 'a'}) + assert watch.changed(real_cache) is True + watch.forget() + assert watch.changed(real_cache) is True + + def test_a_rewrite_after_a_delete_is_seen(self, real_cache): + watch = MailboxWatch(MAILBOX) + real_cache.set(MAILBOX, {'request_id': 'a'}) + assert watch.changed(real_cache) + real_cache.delete(MAILBOX) + assert watch.changed(real_cache) is False + real_cache.set(MAILBOX, {'request_id': 'a'}) + assert watch.changed(real_cache) is True + + def test_a_cache_that_cannot_tell_is_read_every_time(self): + watch = MailboxWatch(MAILBOX) + assert watch.changed(MagicMock()) is True + assert watch.changed(MagicMock()) is True + assert watch.changed(object()) is True + + +# -- end to end over a real socket --------------------------------------------------- + +@pytest.mark.skipif(not c.socket_supported(), reason='AF_UNIX sockets are Linux/macOS only') +class TestOverTheSocket: + @pytest.fixture + def sock_path(self): + import shutil + import tempfile + d = tempfile.mkdtemp(prefix='lmipc-') + yield os.path.join(d, 'control.sock') + shutil.rmtree(d, ignore_errors=True) + + def test_errors_clear_round_trip(self, sock_path): + def handler(request_id, args): + return {'request_id': request_id, 'cutoff': args.cutoff, 'cleared': 2} + + server = ControlServer(sock_path, handlers={Command.ERRORS_CLEAR: handler}) + assert server.start() + try: + result = client.errors_clear('clr-1', 1790000000.0, paths=[sock_path]) + assert result == {'request_id': 'clr-1', 'cutoff': 1790000000.0, 'cleared': 2} + finally: + server.close() + + def test_an_older_display_is_an_upgrade_fallback(self, sock_path): + server = ControlServer(sock_path) # no errors.clear handler + assert server.start() + try: + with pytest.raises(client.ControlError) as e: + client.errors_clear('clr-2', 1.0, paths=[sock_path]) + assert e.value.reason == 'unknown_command' and e.value.sent is True + assert client.should_fall_back(e.value) + finally: + server.close() + + def test_no_display_is_a_fallback(self, sock_path): + with pytest.raises(client.ControlError) as e: + client.errors_clear('clr-3', 1.0, paths=[sock_path]) + assert e.value.reason == 'no_socket' and e.value.sent is False + assert client.should_fall_back(e.value) + + def test_a_full_queue_is_not_a_fallback(self, sock_path): + server = ControlServer(sock_path, queue_size=1) + assert server.start() + try: + client.on_demand_start('q1', 'clock', None, paths=[sock_path]) + with pytest.raises(client.ControlError) as e: + client.on_demand_start('q2', 'clock', None, paths=[sock_path]) + assert e.value.reason == 'busy' and e.value.sent is True + assert not client.should_fall_back(e.value) + finally: + server.close() diff --git a/web_interface/blueprints/api_v3/display.py b/web_interface/blueprints/api_v3/display.py index a3da2acb..3cbe6c60 100644 --- a/web_interface/blueprints/api_v3/display.py +++ b/web_interface/blueprints/api_v3/display.py @@ -33,16 +33,28 @@ def _cache_manager(): +class _NotDelivered(Exception): + """The display took an on-demand request over the socket and did not + accept it (``busy``, ``invalid_args``, ...) or did not answer in time.""" + + def __init__(self, reason): + super().__init__(reason) + self.reason = reason + + def _deliver_on_demand(payload): """Hand an on-demand request to the display: control socket, else mailbox. The socket (src/ipc) answers with an acknowledgement as soon as the - display has the command queued for its render thread. Any failure -- no - socket (the display is stopped or predates it), a timeout, a refusal -- - writes the file mailbox instead, exactly as before the socket existed; - the display reads it within ON_DEMAND_POLL_INTERVAL. Both carry the same - request_id, so a request that reached the display both ways (a reply - that timed out after the command was queued) is still processed once. + display has the command queued for its render thread. The file mailbox + is written only when the socket could not carry the request at all + (``control_client.should_fall_back``): no socket (the display is stopped + or predates it), a refused connection, or a display too old to know the + command. The display reads it within its mailbox poll interval. + + A display that had the request and refused it, or did not answer in + time, raises :class:`_NotDelivered`: a mailbox copy would be refused + the same way, or hide a stuck display behind a "success". Returns ``(transport, socket_error)``: ``'socket'`` and None, or ``'mailbox'`` and the socket failure's reason code. @@ -57,6 +69,10 @@ def _deliver_on_demand(payload): return 'socket', None except control_client.ControlError as e: reason = _socket_reason_code(e.reason) + if not control_client.should_fall_back(e): + logger.warning("The display did not accept on-demand %s %s (%s)", + payload['action'], payload['request_id'], e) + raise _NotDelivered(reason) from None if reason in _QUIET_SOCKET_REASONS: logger.debug("On-demand %s via the mailbox: %s", payload['action'], e) else: @@ -69,6 +85,17 @@ def _deliver_on_demand(payload): return 'mailbox', reason +def _not_delivered_response(request_id, action, reason): + """The answer when the display had the request and did not accept it.""" + status = 400 if reason == 'invalid_args' else 503 + return jsonify({ + 'status': 'error', + 'message': (f'The display service did not accept the on-demand {action} ' + f'request ({reason})'), + 'data': {'request_id': request_id, 'transport': 'socket', 'socket_error': reason}, + }), status + + def _withdraw_on_demand(request_id): """Take a start request the route has refused back out of the mailbox. @@ -277,7 +304,10 @@ def start_on_demand_display(): 'pinned': pinned, 'timestamp': _pkg.time.time() } - transport, socket_error = _deliver_on_demand(request_payload) + try: + transport, socket_error = _deliver_on_demand(request_payload) + except _NotDelivered as e: + return _not_delivered_response(request_id, 'start', e.reason) # A socket acknowledgement is the display itself answering: it is # running and has the request queued, whatever systemd says (a display @@ -312,9 +342,10 @@ def start_on_demand_display(): # MQTT on-demand command, which posts here with the default -- cold- # restarted the display process: every plugin reloaded and the panel was # blank for seconds. The restart bought nothing. The running process - # reads this mailbox every ON_DEMAND_POLL_INTERVAL (0.25s), from its - # dwell sleep, its render loops and Vegas's interrupt check as well as - # the main loop, and a restarted one got the request the same way: the + # looks at this mailbox at least once a second, from its dwell sleep, + # its render loops and Vegas's interrupt check as well as the main loop + # (DisplayController._mailbox_poll_interval), and a restarted one got + # the request the same way: the # startup path only restores a session the display itself saved # (display_on_demand_config), so it loaded nothing it would not have had. service_result = None @@ -365,7 +396,14 @@ def stop_on_demand_display(): 'action': 'stop', 'timestamp': _pkg.time.time() } - transport, socket_error = _deliver_on_demand(request_payload) + try: + transport, socket_error = _deliver_on_demand(request_payload) + except _NotDelivered as e: + if not stop_service: + return _not_delivered_response(request_id, 'stop', e.reason) + # Stopping the service ends on-demand too, whatever the display did + # with the request. + transport, socket_error = 'socket', e.reason service_result = None if stop_service: diff --git a/web_interface/blueprints/api_v3/misc.py b/web_interface/blueprints/api_v3/misc.py index dd089a19..60db7e08 100644 --- a/web_interface/blueprints/api_v3/misc.py +++ b/web_interface/blueprints/api_v3/misc.py @@ -374,6 +374,23 @@ def _read_errors(): return snapshot, clear_request +def _send_error_clear(request_id, cutoff): + """``errors.clear`` over the control socket: the display's answer, or None + when the socket could not carry it (no socket, or a display older than + the command) and the clear goes to the mailbox instead. A display that + took the request and failed raises ControlError (no second copy).""" + client = _pkg.control_client + try: + return client.errors_clear(request_id, cutoff) + except client.ControlError as e: + if not client.should_fall_back(e): + raise + _pkg._log_socket_failure('errors.clear', e, _pkg._socket_reason_code(e.reason)) + except Exception: # never let the socket path break the route + logger.exception("Control socket client failed clearing errors; using the mailbox") + return None + + @api_v3.route('/errors/summary', methods=['GET']) def get_error_summary(): """ @@ -471,7 +488,17 @@ def clear_old_errors(): now = _pkg.time.time() cutoff = now if clear_all else now - max_age_hours * 3600 try: - result = _errors.request_error_clear(_errors_cache(), cutoff) + result = _errors.request_error_clear(_errors_cache(), cutoff, + send=_send_error_clear) + except _pkg.control_client.ControlError as e: + reason = _pkg._socket_reason_code(e.reason) + logger.warning("The display did not apply the error clear: %s", e) + return error_response( + error_code=ErrorCode.SYSTEM_ERROR, + message="The display service did not apply the clear", + context={'socket_error': reason}, + status_code=503 + ) except OSError as e: logger.error("Could not record an error clear request: %s", e) return error_response( @@ -481,11 +508,12 @@ def clear_old_errors(): ) scope = "all errors" if clear_all else f"errors older than {max_age_hours} hours" - return success_response( - data=result, - message=(f"Clear of {scope} requested; the display service applies it " - f"within about {int(_errors.SNAPSHOT_TICK_INTERVAL)} seconds") - ) + if result.get('applied'): + message = f"Cleared {scope}" + else: + message = (f"Clear of {scope} requested; the display service applies it " + f"within about {int(_errors.SNAPSHOT_TICK_INTERVAL)} seconds") + return success_response(data=result, message=message) except Exception as e: logger.error(f"Error clearing old errors: {e}", exc_info=True) return error_response(