mirror of
https://github.com/ChuckBuilds/LEDMatrix.git
synced 2026-10-10 17:16:36 +00:00
feat(ipc)!: remove the cache-key mailboxes (control socket stage 5)
The control socket is now the only way the web interface sends the display a command. The display stops reading display_on_demand_request and plugin_error_clear_request, and the web interface stops writing them. - Display: no mailbox poll (MailboxWatch, the 1 s / 0.25 s cadence, _consume_on_demand_request, the deprecation log) and no persisted display_on_demand_processed_id guard; the error publisher reads no clear request. CacheManager.file_signature and MailboxWatch are removed. - A write to either retired key is dropped by CacheManager.save_cache and logged once per writer, naming the plugin from the call stack (or the request's plugin_id), with the API to move to. - Web: on-demand start with no display listening starts the service (when start_service) and sends the request again once the socket answers (45 s, 10 s for a running service without a socket yet); every other failure is a 503 (400 for invalid_args). Stop answers 503 when no display listens, unless stop_service. errors/clear answers 503 with a reason-specific message instead of writing a request; clear_pending is always false. src.ipc.client.should_fall_back is replaced by display_not_listening. - Kept: display_current_state, display_on_demand_state, plugin_runtime_snapshot and the heartbeat (read whenever the socket cannot answer), and display_on_demand_config (the display's resume record). Tests: mailbox-only tests removed (test_on_demand_mailbox.py, the mailbox cadence, file_signature and MailboxWatch tests); tests that injected requests through the mailbox now use the socket queue or a plugin's in-process request. The run-loop harness sends on-demand requests over its fake control socket, so four golden traces change: on-demand starts and stops land at the request instant instead of the next 0.25 s mailbox look (one frame fewer on the screen they end), and in vegas.json within one frame instead of 263 ms, which shifts the later 1 s-throttled WiFi-notice check by under a second. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
+54
-48
@@ -25,10 +25,11 @@ Typical plugin usage::
|
||||
|
||||
import json
|
||||
import os
|
||||
import sys
|
||||
import time
|
||||
from datetime import datetime
|
||||
import pytz
|
||||
from typing import Any, Dict, List, Optional, Tuple
|
||||
from typing import Any, Dict, List, Optional
|
||||
import logging
|
||||
import threading
|
||||
import tempfile
|
||||
@@ -72,41 +73,60 @@ def _outlived(record: Any, max_age: Optional[float], now: float) -> bool:
|
||||
return False
|
||||
|
||||
|
||||
_NOT_SEEN: Any = object()
|
||||
#: Cache keys that were file "mailboxes" from the web interface (and some
|
||||
#: plugins) to the display. The control socket replaced them, and nothing
|
||||
#: reads them any more, so a write is refused rather than left on the SD card
|
||||
#: for nobody: see :func:`_refuse_retired_mailbox_write`.
|
||||
RETIRED_MAILBOX_KEYS = frozenset({'display_on_demand_request', 'plugin_error_clear_request'})
|
||||
|
||||
#: (key, writer) pairs already warned about, so a plugin that writes on every
|
||||
#: event logs once per process, not once per write.
|
||||
_retired_writers_warned: set = set()
|
||||
_retired_writers_lock = threading.Lock()
|
||||
|
||||
|
||||
class MailboxWatch:
|
||||
"""Tells the poller of a mailbox key whether its file changed since the
|
||||
last look, from one stat() (:meth:`CacheManager.file_signature`).
|
||||
def _retired_mailbox_writer(data: Any) -> str:
|
||||
"""Name whoever is writing a retired mailbox key, as well as can be told.
|
||||
|
||||
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.
|
||||
The plugin instance on the call stack when there is one (a ``self`` with
|
||||
a string ``plugin_id`` and a ``cache_manager``: what BasePlugin gives
|
||||
every plugin), else the ``plugin_id`` the request itself names, else
|
||||
``'unknown'``.
|
||||
"""
|
||||
frame = sys._getframe(2) # pylint: disable=protected-access
|
||||
depth = 0
|
||||
while frame is not None and depth < 25:
|
||||
owner = frame.f_locals.get('self')
|
||||
plugin_id = getattr(owner, 'plugin_id', None) if owner is not None else None
|
||||
if isinstance(plugin_id, str) and plugin_id and hasattr(owner, 'cache_manager'):
|
||||
return f"plugin '{plugin_id}'"
|
||||
frame = frame.f_back
|
||||
depth += 1
|
||||
payload = data.get('data', data) if isinstance(data, dict) else None
|
||||
named = payload.get('plugin_id') if isinstance(payload, dict) else None
|
||||
if isinstance(named, str) and named:
|
||||
return f"plugin '{named}' (named in the request)"
|
||||
return 'unknown'
|
||||
|
||||
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
|
||||
def _refuse_retired_mailbox_write(logger: logging.Logger, key: str, data: Any) -> None:
|
||||
"""Warn, once per writer, that a write to a retired mailbox key was dropped."""
|
||||
try:
|
||||
writer = _retired_mailbox_writer(data)
|
||||
except Exception: # pylint: disable=broad-except
|
||||
writer = 'unknown'
|
||||
with _retired_writers_lock:
|
||||
if (key, writer) in _retired_writers_warned:
|
||||
return
|
||||
_retired_writers_warned.add((key, writer))
|
||||
if key == 'display_on_demand_request':
|
||||
hint = ("call self.request_on_demand() / self.end_on_demand() instead "
|
||||
"(BasePlugin, LEDMatrix 3.8.1 and later)")
|
||||
else:
|
||||
hint = "clear errors through POST /api/v3/errors/clear instead"
|
||||
logger.warning("Ignored a write to the retired '%s' cache key by %s: the display no "
|
||||
"longer reads this file mailbox. Update it to %s. (Logged once per writer.)",
|
||||
key, writer, hint)
|
||||
|
||||
|
||||
class CacheManager:
|
||||
@@ -330,24 +350,6 @@ class CacheManager:
|
||||
"""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.
|
||||
|
||||
@@ -385,6 +387,10 @@ class CacheManager:
|
||||
key: Cache key
|
||||
data: Data to cache
|
||||
"""
|
||||
if key in RETIRED_MAILBOX_KEYS:
|
||||
_refuse_retired_mailbox_write(self.logger, key, data)
|
||||
return
|
||||
|
||||
# Periodic cleanup before adding new entries
|
||||
self._cleanup_memory_cache()
|
||||
|
||||
|
||||
+30
-182
@@ -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, MailboxWatch
|
||||
from src.cache_manager import CacheManager
|
||||
from src.font_manager import FontManager
|
||||
from src.logging_config import get_logger
|
||||
from src.exceptions import PluginError
|
||||
@@ -69,10 +69,6 @@ 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
|
||||
@@ -412,19 +408,15 @@ class DisplayController:
|
||||
# coordinator exists (Vegas was off at startup). The render thread
|
||||
# creates it in _is_vegas_mode_active(), never the watcher thread.
|
||||
self._pending_vegas_init = False
|
||||
# Monotonic stamp of the last mailbox disk read; see
|
||||
# _poll_on_demand_requests. None means "never polled", so the first
|
||||
# call always goes through.
|
||||
self._last_on_demand_poll: Optional[float] = None
|
||||
# Monotonic stamp of the last _service_pending_changes pass; same
|
||||
# "None means never" convention as _last_on_demand_poll.
|
||||
# Monotonic stamp of the last _service_pending_changes pass. None
|
||||
# means "never", so the first call always goes through.
|
||||
self._last_pending_service: Optional[float] = None
|
||||
# Monotonic stamp of the last scheduled-update pass; see
|
||||
# _tick_plugin_updates_if_due. Same "None means never" convention.
|
||||
self._last_plugin_update_tick: Optional[float] = None
|
||||
# The control socket (src/ipc), started by run(). None when it is not
|
||||
# served (Windows, LEDMATRIX_CONTROL_SOCKET=off, a bind failure);
|
||||
# the file mailbox works either way.
|
||||
# then only plugins in this process can start on-demand sessions.
|
||||
self._control_server = None
|
||||
# A brightness set_brightness() refused, so the periodic service pass
|
||||
# doesn't retry (and log) the same failure several times a second.
|
||||
@@ -1859,34 +1851,14 @@ class DisplayController:
|
||||
self.force_change = True
|
||||
self._publish_on_demand_state()
|
||||
|
||||
#: Shortest gap between mailbox disk reads. This is called after every
|
||||
#: frame -- about 125 times a second on a scrolling mode -- and the read
|
||||
#: below is deliberately uncached, so without a floor it was 125 disk reads
|
||||
#: per second to find nothing. An on-demand request comes from a person
|
||||
#: clicking in the web UI, so a quarter second of latency is not
|
||||
#: perceptible, and it cuts the read rate by 30x.
|
||||
ON_DEMAND_POLL_INTERVAL = 0.25
|
||||
|
||||
#: 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
|
||||
|
||||
#: 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,
|
||||
#: Shortest gap between _service_pending_changes passes. 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()
|
||||
#: Most plugin on-demand requests waiting for the render thread at once.
|
||||
#: A plugin that asks faster than the display drains (four times a
|
||||
#: second at worst) is refused, not queued without end.
|
||||
@@ -2029,44 +2001,11 @@ class DisplayController:
|
||||
on_demand_plugin_id, len(enabled_plugins))
|
||||
return enabled_plugins
|
||||
|
||||
def _consume_on_demand_request(self, request_id: str) -> None:
|
||||
"""Remove the request we just handled from the mailbox.
|
||||
|
||||
Leaving it on disk meant a restart replayed the previous request: the
|
||||
fresh controller read it, activated it and cached it, so the request
|
||||
the caller had just made was ignored and the panel silently showed the
|
||||
earlier plugin.
|
||||
|
||||
Compare before deleting. The web process can post a newer request
|
||||
between the read and this delete; an unconditional delete threw that
|
||||
one away and it was never processed -- the user's second click did
|
||||
nothing. Re-reading uncached and only deleting our own request_id
|
||||
leaves a newer request in the mailbox for the next poll instead.
|
||||
|
||||
This narrows the window rather than closing it: a request landing
|
||||
between the re-read and the delete is still lost. Closing it properly
|
||||
needs an atomic claim (a rename, or a compare-and-delete primitive)
|
||||
that the cache layer does not currently offer, so the honest fix is a
|
||||
smaller window plus this note, not a bigger lock. For start requests
|
||||
processed_id still guards against reprocessing if the delete fails.
|
||||
"""
|
||||
try:
|
||||
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(ON_DEMAND_MAILBOX_KEY)
|
||||
else:
|
||||
logger.debug("Newer on-demand request %s arrived while processing "
|
||||
"%s; leaving it in the mailbox",
|
||||
current.get('request_id'), request_id)
|
||||
except (OSError, AttributeError, KeyError) as err:
|
||||
logger.debug("Could not clear the on-demand request mailbox: %s", err)
|
||||
|
||||
def _start_control_server(self) -> None:
|
||||
"""Serve the control socket (src/ipc/server.py). Never raises.
|
||||
|
||||
Its handlers only queue commands; _drain_control_commands applies
|
||||
them on the render thread, where the mailbox is read.
|
||||
them on the render thread.
|
||||
"""
|
||||
if self._control_server is not None:
|
||||
return
|
||||
@@ -2079,7 +2018,8 @@ class DisplayController:
|
||||
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")
|
||||
logger.exception("Control socket not started; the web interface cannot "
|
||||
"send this display commands")
|
||||
if self._control_server is not None:
|
||||
self._start_state_stream(hub)
|
||||
|
||||
@@ -2107,10 +2047,8 @@ class DisplayController:
|
||||
def _drain_control_commands(self) -> None:
|
||||
"""Apply the commands that arrived over the control socket.
|
||||
|
||||
On-demand commands go through _handle_on_demand_request, the
|
||||
mailbox's own handler, so both ways in behave the same, and a
|
||||
request that came both ways (a client that timed out and fell back)
|
||||
has one request id and is processed once. A brightness is applied
|
||||
On-demand commands go through _handle_on_demand_request, as plugins'
|
||||
own requests do, so both ways in behave the same. A brightness is applied
|
||||
here. A plugin reload waits for the top of the next loop pass, where
|
||||
no plugin is on the stack (_apply_pending_plugin_reloads); until
|
||||
then the current screen ends early (_plugin_reload_pending).
|
||||
@@ -2143,7 +2081,7 @@ class DisplayController:
|
||||
|
||||
``PluginManager.request_on_demand`` / ``end_on_demand`` (which
|
||||
BasePlugin's methods of the same names call) build ``request``: the
|
||||
mailbox's shape, with ``source: 'plugin'`` and the asking plugin's
|
||||
on-demand request shape, with ``source: 'plugin'`` and the asking plugin's
|
||||
id. Nothing here touches the panel or the on-demand state; the render
|
||||
thread applies the request where it applies a socket command
|
||||
(_drain_control_commands), through _handle_on_demand_request, and is
|
||||
@@ -2438,105 +2376,35 @@ 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:
|
||||
"""Apply on-demand requests: the control socket's, then the mailbox's.
|
||||
"""Apply on-demand requests: the control socket's, then plugins' own.
|
||||
|
||||
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().
|
||||
Both are in memory (no disk read), so this has no floor and is
|
||||
cheap to call every frame. The file mailbox
|
||||
(``display_on_demand_request``) is gone: nothing reads it, and a
|
||||
write to it is dropped with a warning (CacheManager).
|
||||
"""
|
||||
# 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._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.
|
||||
#
|
||||
# memory_ttl=0 is required, not optional: this key is a mailbox the
|
||||
# web process writes and this process reads. get() defaults the
|
||||
# in-memory TTL to max_age, so without it the first request read was
|
||||
# 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(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 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 control socket or a plugin.
|
||||
|
||||
A socket command carries ``source: 'socket'``, and a plugin's own
|
||||
request (submit_plugin_on_demand) ``source: 'plugin'``. Only a
|
||||
mailbox request is removed from the mailbox afterwards: the others
|
||||
never put anything there, so that would be a disk read and maybe a
|
||||
delete for nothing.
|
||||
request (submit_plugin_on_demand) ``source: 'plugin'``.
|
||||
|
||||
A plugin's stop ends only that plugin's own session: a plugin
|
||||
releasing the screen must not end one the user started for
|
||||
another plugin. (A stop through the mailbox ends any session, as it
|
||||
always has.)
|
||||
another plugin. A stop from the socket (the web interface) ends any
|
||||
session.
|
||||
"""
|
||||
request_id = request.get('request_id')
|
||||
if not request_id:
|
||||
return
|
||||
source = request.get('source')
|
||||
from_mailbox = source not in ('socket', 'plugin')
|
||||
|
||||
action = request.get('action')
|
||||
|
||||
# For stop requests, always process them (don't check processed_id)
|
||||
# For stop requests, always process them (no request-id guard)
|
||||
# This allows stopping even if the same stop request was sent before
|
||||
if action == 'stop':
|
||||
if source == 'plugin' and not (
|
||||
@@ -2562,42 +2430,22 @@ class DisplayController:
|
||||
# without this the status route kept reporting it until
|
||||
# the state aged out (120s) or another request came in.
|
||||
self._clear_on_demand(reason='requested-stop')
|
||||
# Stop requests are deliberately exempt from the request_id/
|
||||
# processed_id guards above, so that a second click stops a mode
|
||||
# that a race left running. Consuming the mailbox is therefore the
|
||||
# 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.
|
||||
if from_mailbox:
|
||||
self._consume_on_demand_request(request_id)
|
||||
# Stop requests are deliberately exempt from the request_id
|
||||
# guard below, so that a second click stops a mode that a race
|
||||
# left running.
|
||||
return
|
||||
|
||||
# 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.
|
||||
# A start already processed (a client that sent the same request
|
||||
# id twice) is not applied a second time.
|
||||
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)
|
||||
logger.debug("On-demand start request %s already processed", request_id)
|
||||
return
|
||||
|
||||
logger.info("Received on-demand request %s: %s (plugin_id=%s, mode=%s, via %s)",
|
||||
request_id, action, request.get('plugin_id'), request.get('mode'),
|
||||
'mailbox' if from_mailbox else source)
|
||||
source or 'unknown')
|
||||
|
||||
# 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
|
||||
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'))
|
||||
|
||||
+54
-175
@@ -19,7 +19,7 @@ import uuid
|
||||
from collections import defaultdict
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import datetime, timedelta
|
||||
from typing import Dict, List, Optional, Any, Callable, Tuple
|
||||
from typing import Dict, List, Optional, Any, Callable
|
||||
import logging
|
||||
|
||||
from src.exceptions import LEDMatrixError
|
||||
@@ -449,33 +449,28 @@ def record_error(
|
||||
# service publishes to the shared cache directory -- the same channel, and the
|
||||
# same file permissions, as display_current_state and plugin_metrics_snapshot: files
|
||||
# are 0660 and carry the cache directory's group, so root writes and the web
|
||||
# user reads, and the other way round for the clear request.
|
||||
# user reads.
|
||||
#
|
||||
# ERROR_SNAPSHOT_KEY written by the display service only
|
||||
# ERROR_CLEAR_REQUEST_KEY written by the web interface only, as a fallback
|
||||
#
|
||||
# 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.
|
||||
# The web interface never writes the snapshot itself: two writers would race,
|
||||
# and a snapshot owned by the web user is one more file root's write has to
|
||||
# replace.
|
||||
# it (clear_before) and republishes the snapshot before it answers. When the
|
||||
# socket cannot carry it, the clear fails and the route says so: the
|
||||
# ``plugin_error_clear_request`` file mailbox it used to fall back to is gone.
|
||||
# A stopped display's errors go anyway: its next run publishes an empty
|
||||
# snapshot over the old one. The web interface never writes the snapshot
|
||||
# itself: two writers would race, and a snapshot owned by the web user is one
|
||||
# more file root's write has to replace.
|
||||
|
||||
ERROR_SNAPSHOT_KEY = "plugin_error_snapshot"
|
||||
ERROR_CLEAR_REQUEST_KEY = "plugin_error_clear_request"
|
||||
|
||||
#: Shortest gap between two snapshot writes, in seconds. A plugin failing in
|
||||
#: a tight loop changes the aggregator many times a second; the snapshot is
|
||||
#: rewritten at most this often, and only when something changed.
|
||||
SNAPSHOT_MIN_INTERVAL = 10.0
|
||||
|
||||
#: How often the display service checks for changes and clear requests. A
|
||||
#: check is an in-memory comparison plus reading one small file.
|
||||
#: How often the display service checks for changes. A check is an
|
||||
#: in-memory comparison.
|
||||
SNAPSHOT_TICK_INTERVAL = 5.0
|
||||
|
||||
_SNAPSHOT_RECENT_ERRORS = 20
|
||||
@@ -544,8 +539,8 @@ class ErrorSnapshotPublisher:
|
||||
Runs in the display service only. tick() is the whole job; start() just
|
||||
calls it from a daemon thread every SNAPSHOT_TICK_INTERVAL seconds, which
|
||||
also means errors recorded while a write was being throttled still reach
|
||||
the cache once the interval has passed, and a clear request is applied
|
||||
even when no new error arrives to trigger a publish.
|
||||
the cache once the interval has passed. A clear (``errors.clear`` over
|
||||
the control socket) is applied by :meth:`clear_now`.
|
||||
|
||||
Nothing here raises: a failure to read or write the cache is logged at
|
||||
debug and retried on a later tick.
|
||||
@@ -563,11 +558,6 @@ 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
|
||||
@@ -579,39 +569,9 @@ class ErrorSnapshotPublisher:
|
||||
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 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")
|
||||
if not isinstance(request_id, str) or not request_id or request_id == self._applied_clear_id:
|
||||
return False
|
||||
try:
|
||||
cutoff = float(request.get("cutoff"))
|
||||
except (TypeError, ValueError):
|
||||
cutoff = float("nan")
|
||||
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.
|
||||
@@ -629,23 +589,20 @@ class ErrorSnapshotPublisher:
|
||||
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."""
|
||||
"""Publish if due. True if a snapshot was written."""
|
||||
with self._tick_lock:
|
||||
try:
|
||||
cleared = self._apply_clear_request()
|
||||
version = self.aggregator.version
|
||||
now = self._clock()
|
||||
if not cleared:
|
||||
if version == self._published_version:
|
||||
return False
|
||||
if (self._last_attempt is not None
|
||||
and now - self._last_attempt < self.min_interval):
|
||||
return False
|
||||
if version == self._published_version:
|
||||
return False
|
||||
if (self._last_attempt is not None
|
||||
and now - self._last_attempt < self.min_interval):
|
||||
return False
|
||||
self._publish(version, now)
|
||||
return True
|
||||
except Exception as err: # never let reporting break the display
|
||||
@@ -716,16 +673,13 @@ def apply_error_clear(request_id: str, args: Any) -> Dict[str, Any]:
|
||||
|
||||
# --- Reading side (web interface) -------------------------------------------
|
||||
|
||||
def read_error_report(cache_manager: Any) -> Tuple[Optional[Dict[str, Any]], Optional[Dict[str, Any]]]:
|
||||
"""The display service's latest snapshot and the latest clear request.
|
||||
def read_error_report(cache_manager: Any) -> Optional[Dict[str, Any]]:
|
||||
"""The display service's latest snapshot, or None before it has published.
|
||||
|
||||
memory_ttl=0: both keys are written by the other process, so only the
|
||||
file is current.
|
||||
memory_ttl=0: the display writes the key, so only the file is current.
|
||||
"""
|
||||
snapshot = cache_manager.get(ERROR_SNAPSHOT_KEY, max_age=None, memory_ttl=0)
|
||||
clear_request = cache_manager.get(ERROR_CLEAR_REQUEST_KEY, max_age=None, memory_ttl=0)
|
||||
return (snapshot if isinstance(snapshot, dict) else None,
|
||||
clear_request if isinstance(clear_request, dict) else None)
|
||||
return snapshot if isinstance(snapshot, dict) else None
|
||||
|
||||
|
||||
def _epoch(iso: Any) -> Optional[float]:
|
||||
@@ -738,36 +692,6 @@ def _epoch(iso: Any) -> Optional[float]:
|
||||
return None
|
||||
|
||||
|
||||
def _pending_cutoff(snapshot: Optional[Dict[str, Any]],
|
||||
clear_request: Optional[Dict[str, Any]]) -> Optional[float]:
|
||||
"""The cutoff of a clear the snapshot has not applied yet, if any."""
|
||||
if not clear_request:
|
||||
return None
|
||||
request_id = clear_request.get("request_id")
|
||||
if not request_id:
|
||||
return None
|
||||
if snapshot is not None and snapshot.get("applied_clear_id") == request_id:
|
||||
return None
|
||||
try:
|
||||
cutoff = float(clear_request.get("cutoff"))
|
||||
except (TypeError, ValueError):
|
||||
return 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:
|
||||
when = _epoch(item.get(field_name)) if isinstance(item, dict) else None
|
||||
return when is not None and when > cutoff
|
||||
|
||||
|
||||
def _empty_summary(snapshot: Optional[Dict[str, Any]]) -> Dict[str, Any]:
|
||||
return {
|
||||
"session_start": snapshot.get("session_start") if snapshot else None,
|
||||
@@ -780,11 +704,11 @@ def _empty_summary(snapshot: Optional[Dict[str, Any]]) -> Dict[str, Any]:
|
||||
}
|
||||
|
||||
|
||||
def error_summary_from_report(snapshot: Optional[Dict[str, Any]],
|
||||
clear_request: Optional[Dict[str, Any]]) -> Dict[str, Any]:
|
||||
def error_summary_from_report(snapshot: Optional[Dict[str, Any]]) -> Dict[str, Any]:
|
||||
"""The /errors/summary payload: get_error_summary()'s shape plus
|
||||
``generated_at``, ``snapshot_available`` and ``clear_pending``."""
|
||||
cutoff = _pending_cutoff(snapshot, clear_request)
|
||||
``generated_at``, ``snapshot_available`` and ``clear_pending`` (always
|
||||
False now: a clear is applied before its route answers; kept for API
|
||||
compatibility)."""
|
||||
summary = _empty_summary(snapshot)
|
||||
if snapshot is not None:
|
||||
for name, default in summary.items():
|
||||
@@ -793,32 +717,17 @@ def error_summary_from_report(snapshot: Optional[Dict[str, Any]],
|
||||
isinstance(default, float) and isinstance(value, int)) or (
|
||||
name == "session_start" and isinstance(value, str)):
|
||||
summary[name] = value
|
||||
if cutoff is not None:
|
||||
recent = summary["recent_errors"]
|
||||
newest = _epoch(recent[-1].get("timestamp")) if recent and isinstance(recent[-1], dict) else None
|
||||
if newest is None or newest <= cutoff:
|
||||
# Everything the display has reported predates the clear.
|
||||
summary = _empty_summary(snapshot)
|
||||
else:
|
||||
# Only part of it does. The lists can be filtered exactly; the
|
||||
# counts cannot, and stay as reported until the display
|
||||
# applies the clear (clear_pending says so).
|
||||
summary["recent_errors"] = [r for r in recent if _is_after(r, "timestamp", cutoff)]
|
||||
summary["active_patterns"] = {
|
||||
k: p for k, p in summary["active_patterns"].items()
|
||||
if _is_after(p, "last_seen", cutoff)
|
||||
}
|
||||
summary["generated_at"] = snapshot.get("generated_at") if snapshot else None
|
||||
summary["snapshot_available"] = snapshot is not None
|
||||
summary["clear_pending"] = cutoff is not None
|
||||
summary["clear_pending"] = False
|
||||
return summary
|
||||
|
||||
|
||||
def plugin_health_from_report(snapshot: Optional[Dict[str, Any]],
|
||||
clear_request: Optional[Dict[str, Any]],
|
||||
plugin_id: str) -> Dict[str, Any]:
|
||||
"""The /errors/plugin/<id> payload: get_plugin_health()'s shape plus
|
||||
``generated_at``, ``snapshot_available`` and ``clear_pending``."""
|
||||
``generated_at``, ``snapshot_available`` and ``clear_pending`` (always
|
||||
False, as in error_summary_from_report)."""
|
||||
health: Dict[str, Any] = {
|
||||
"plugin_id": plugin_id,
|
||||
"status": "healthy",
|
||||
@@ -827,19 +736,15 @@ def plugin_health_from_report(snapshot: Optional[Dict[str, Any]],
|
||||
"recent_error_count": 0,
|
||||
"last_error": None,
|
||||
}
|
||||
cutoff = _pending_cutoff(snapshot, clear_request)
|
||||
table = snapshot.get("plugin_health") if snapshot else None
|
||||
entry = table.get(plugin_id) if isinstance(table, dict) else None
|
||||
if isinstance(entry, dict):
|
||||
# last_error is the plugin's newest error: if even that predates a
|
||||
# pending clear, so does everything else the display reported for it.
|
||||
if cutoff is None or _is_after(entry.get("last_error"), "timestamp", cutoff):
|
||||
for name in ("status", "total_errors", "error_types", "recent_error_count", "last_error"):
|
||||
if name in entry:
|
||||
health[name] = entry[name]
|
||||
for name in ("status", "total_errors", "error_types", "recent_error_count", "last_error"):
|
||||
if name in entry:
|
||||
health[name] = entry[name]
|
||||
health["generated_at"] = snapshot.get("generated_at") if snapshot else None
|
||||
health["snapshot_available"] = snapshot is not None
|
||||
health["clear_pending"] = cutoff is not None
|
||||
health["clear_pending"] = False
|
||||
return health
|
||||
|
||||
|
||||
@@ -861,61 +766,35 @@ def _count_cleared(summary: Dict[str, Any], cutoff: float) -> Optional[int]:
|
||||
|
||||
|
||||
#: ``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]]]
|
||||
#: socket and returns its ErrorsClearResult. It raises when the socket could
|
||||
#: not carry it or the display failed it (``src.ipc.client.ControlError``).
|
||||
ClearSender = Callable[[str, float], Dict[str, Any]]
|
||||
|
||||
|
||||
def request_error_clear(cache_manager: Any, cutoff: float,
|
||||
send: Optional[ClearSender] = None) -> Dict[str, Any]:
|
||||
send: ClearSender) -> Dict[str, Any]:
|
||||
"""Ask the display service to forget errors recorded at or before ``cutoff``.
|
||||
|
||||
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.
|
||||
Over the control socket: the display applies the clear and republishes
|
||||
its snapshot before it answers, so nothing is written here. Whatever
|
||||
``send`` raises reaches the caller.
|
||||
|
||||
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
|
||||
would otherwise bring back the errors the first one hid.
|
||||
Returns ``request_id``, ``cutoff`` (ISO, local time), ``cleared_count``
|
||||
(the display's own count, else an estimate from the snapshot; see
|
||||
_count_cleared), and ``clear_requested``, ``applied`` and ``transport``,
|
||||
which are always True, True and ``"socket"`` now and are kept for API
|
||||
compatibility.
|
||||
"""
|
||||
snapshot, clear_request = read_error_report(cache_manager)
|
||||
pending = _pending_cutoff(snapshot, clear_request)
|
||||
if pending is not None:
|
||||
cutoff = max(cutoff, pending)
|
||||
before = error_summary_from_report(snapshot, clear_request)
|
||||
before = error_summary_from_report(read_error_report(cache_manager))
|
||||
request_id = uuid.uuid4().hex
|
||||
answer = {
|
||||
result = send(request_id, cutoff)
|
||||
count = result.get("cleared") if isinstance(result, dict) else None
|
||||
return {
|
||||
"clear_requested": True,
|
||||
"request_id": request_id,
|
||||
"cutoff": datetime.fromtimestamp(cutoff).isoformat(),
|
||||
"applied": True,
|
||||
"transport": "socket",
|
||||
"cleared_count": (count if isinstance(count, int) and not isinstance(count, bool)
|
||||
else _count_cleared(before, cutoff)),
|
||||
}
|
||||
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": request_id,
|
||||
"cutoff": cutoff,
|
||||
"requested_at": time.time(),
|
||||
}
|
||||
cache_manager.set(ERROR_CLEAR_REQUEST_KEY, request)
|
||||
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 dict(answer, applied=False, transport="mailbox",
|
||||
cleared_count=_count_cleared(before, cutoff))
|
||||
|
||||
+23
-28
@@ -5,11 +5,11 @@ 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``. 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.
|
||||
There is no other way to reach the display: the file mailboxes the web
|
||||
interface used to fall back to are gone. :func:`display_not_listening` tells
|
||||
a caller when no display is listening yet (it is stopped, or still
|
||||
starting), the one case where sending the same request again later, once a
|
||||
display is up, can work.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -43,8 +43,7 @@ from src.ipc.contract import (
|
||||
|
||||
#: Total budget for one request: connect, send and the reply. The display
|
||||
#: answers from a thread that does no rendering, normally within a few
|
||||
#: milliseconds; this only bounds a wedged one. The web route then falls back
|
||||
#: to the mailbox, so a timeout costs this much latency and nothing else.
|
||||
#: milliseconds; this only bounds a wedged one.
|
||||
DEFAULT_TIMEOUT_SECONDS = 1.0
|
||||
|
||||
|
||||
@@ -74,28 +73,25 @@ class ControlError(Exception):
|
||||
|
||||
#: 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})
|
||||
|
||||
#: Transport reasons that mean nothing is listening at the socket: no socket
|
||||
#: file (the display is stopped, or has not reached its run loop), or a file
|
||||
#: nobody accepts on (a stale socket) or that this user may not open.
|
||||
NOT_LISTENING_REASONS = frozenset({'no_socket', 'refused'})
|
||||
|
||||
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.
|
||||
def display_not_listening(error: BaseException) -> bool:
|
||||
"""True when ``error`` says no display took the request because none is
|
||||
listening: it never reached one (``sent`` is False) and the reason is in
|
||||
:data:`NOT_LISTENING_REASONS`. A display that is started, or finishes
|
||||
starting, may take the same request later. False for everything else:
|
||||
a display that had the request and failed it, one too old to know the
|
||||
command, a client that cannot use the socket at all (``disabled``,
|
||||
``unsupported``), and an exception that is not a :class:`ControlError`.
|
||||
"""
|
||||
if not isinstance(error, ControlError):
|
||||
return True
|
||||
return not error.sent or error.reason in UPGRADE_REASONS
|
||||
return (isinstance(error, ControlError) and not error.sent
|
||||
and error.reason in NOT_LISTENING_REASONS)
|
||||
|
||||
|
||||
def request(cmd: str, args: Optional[Mapping[str, Any]] = None, *,
|
||||
@@ -218,8 +214,8 @@ def on_demand_start(request_id: str, plugin_id: Optional[str], mode: Optional[st
|
||||
paths: Optional[Sequence[str]] = None) -> Dict[str, Any]:
|
||||
"""Ask the display to show a plugin now. Returns the ack; raises :class:`ControlError`.
|
||||
|
||||
``request_id`` doubles as the on-demand request id, so a request that a
|
||||
timed-out caller then also writes to the mailbox is processed only once.
|
||||
``request_id`` doubles as the on-demand request id, so a request sent
|
||||
twice with the same id is processed only once.
|
||||
"""
|
||||
args = {'plugin_id': plugin_id, 'mode': mode, 'duration': duration, 'pinned': pinned}
|
||||
return request(Command.ON_DEMAND_START, args, request_id=request_id,
|
||||
@@ -281,8 +277,7 @@ def errors_clear(request_id: str, cutoff: float, *,
|
||||
|
||||
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.
|
||||
than the command.
|
||||
"""
|
||||
return request(Command.ERRORS_CLEAR, {'cutoff': cutoff}, request_id=request_id,
|
||||
timeout=timeout, paths=paths)
|
||||
|
||||
+10
-11
@@ -85,8 +85,8 @@ DEFAULT_SOCKET_PATH = DEFAULT_SOCKET_DIR + '/' + SOCKET_NAME
|
||||
|
||||
#: Overrides the socket path for both processes (a dev checkout, a second
|
||||
#: instance, tests). One of :data:`DISABLED_VALUES` turns the socket off: the
|
||||
#: display does not serve it and the web interface goes straight to the
|
||||
#: file mailbox.
|
||||
#: display does not serve it, and the web interface cannot send it commands
|
||||
#: (it still reads the state the display writes to the cache).
|
||||
SOCKET_PATH_ENV = 'LEDMATRIX_CONTROL_SOCKET'
|
||||
DISABLED_VALUES = frozenset({'off', '0', 'false', 'no', 'none', 'disabled'})
|
||||
|
||||
@@ -375,8 +375,8 @@ def _optional_name(args: Mapping[str, Any], key: str) -> Optional[str]:
|
||||
def _optional_duration(value: Any) -> Optional[float]:
|
||||
"""Seconds, or None for "until stopped". 0 means the same as None.
|
||||
|
||||
Numbers and numeric strings are accepted, the same as the REST route and
|
||||
the file mailbox take them; anything else is refused rather than guessed.
|
||||
Numbers and numeric strings are accepted, the same as the REST route
|
||||
takes them; anything else is refused rather than guessed.
|
||||
"""
|
||||
if value is None or value == '':
|
||||
return None
|
||||
@@ -415,7 +415,7 @@ class HelloArgs:
|
||||
class OnDemandStartArgs:
|
||||
"""``on_demand.start``: show a plugin (or one of its modes) now.
|
||||
|
||||
The same fields the file mailbox carries. At least one of ``plugin_id``
|
||||
The same fields as the REST route's body. At least one of ``plugin_id``
|
||||
and ``mode`` is required; the display resolves the other.
|
||||
"""
|
||||
plugin_id: Optional[str] = None
|
||||
@@ -600,13 +600,12 @@ def parse_args(cmd: str, args: Mapping[str, Any]) -> CommandArgs:
|
||||
|
||||
def on_demand_request(request_id: str, args: Union[OnDemandStartArgs, OnDemandStopArgs],
|
||||
timestamp: float) -> Dict[str, Any]:
|
||||
"""The file-mailbox payload for a queued on-demand command.
|
||||
"""The on-demand request dict for a queued on-demand command.
|
||||
|
||||
The display hands socket commands to the same code that handles the
|
||||
mailbox (``DisplayController._handle_on_demand_request``), so a command
|
||||
behaves identically whichever way it arrived, and a request that came
|
||||
both ways (a client that timed out and fell back) is processed once: the
|
||||
request id is the same.
|
||||
The display hands socket commands to the same code that handles
|
||||
plugins' own requests (``DisplayController._handle_on_demand_request``),
|
||||
so a command behaves identically whichever way it arrived. (This was
|
||||
the file mailbox's payload, which the display no longer reads.)
|
||||
"""
|
||||
if isinstance(args, OnDemandStartArgs):
|
||||
return {'request_id': request_id, 'action': 'start', 'plugin_id': args.plugin_id,
|
||||
|
||||
+13
-12
@@ -3,9 +3,9 @@
|
||||
A small threaded server on a Unix stream socket (``/run/ledmatrix/control.sock``
|
||||
by default; see :mod:`src.ipc.contract` for the protocol). It never touches
|
||||
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
|
||||
queue and acknowledged, and the render thread drains that queue
|
||||
(``DisplayController._poll_on_demand_requests``), handing each on-demand
|
||||
command to the code that handles plugins' own requests. Queries (``on_demand.status``) are
|
||||
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.
|
||||
@@ -110,8 +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`` instead of piling up work. The mailbox would not be read
|
||||
#: either, so the web interface reports the failure rather than fall back.
|
||||
#: is told ``busy`` instead of piling up work, and the web interface reports
|
||||
#: the failure.
|
||||
QUEUE_SIZE = 16
|
||||
|
||||
#: Timeout for one recv()/send() on a connection.
|
||||
@@ -185,7 +185,7 @@ class QueuedCommand:
|
||||
outcome: Optional[CommandOutcome] = field(default=None, compare=False, repr=False)
|
||||
|
||||
def as_on_demand_request(self) -> Dict[str, Any]:
|
||||
"""The mailbox-shaped payload the display's on-demand handler takes."""
|
||||
"""The on-demand request dict the display's on-demand handler takes."""
|
||||
if not isinstance(self.args, (OnDemandStartArgs, OnDemandStopArgs)):
|
||||
raise TypeError(f'{self.cmd} is not an on-demand command')
|
||||
return on_demand_request(self.request_id, self.args, self.received_at)
|
||||
@@ -619,8 +619,8 @@ class ControlServer:
|
||||
def start(self) -> bool:
|
||||
"""Bind and start serving. False (logged) when the socket cannot be served.
|
||||
|
||||
Never raises: without the socket the web interface uses the file
|
||||
mailbox, exactly as before.
|
||||
Never raises: without the socket the display runs, but the web
|
||||
interface cannot send it commands.
|
||||
"""
|
||||
if not socket_supported():
|
||||
logger.debug("Control socket not started: no Unix sockets on this platform")
|
||||
@@ -632,7 +632,7 @@ class ControlServer:
|
||||
self._bind()
|
||||
except OSError as e:
|
||||
logger.warning("Control socket not started at %s (%s); the web interface "
|
||||
"will use the file mailbox", self.path, e)
|
||||
"cannot send this display commands", self.path, e)
|
||||
self._close_socket()
|
||||
return False
|
||||
self._stopping.clear()
|
||||
@@ -1134,12 +1134,13 @@ def start_control_server(status_provider: Optional[StatusProvider] = None,
|
||||
"""Start the display's control socket, or return None when it can't run.
|
||||
|
||||
None covers Windows, ``LEDMATRIX_CONTROL_SOCKET=off`` and any failure to
|
||||
bind; in every case the web interface falls back to the file mailbox
|
||||
and to the cache keys the display still writes.
|
||||
bind; in every case the web interface cannot send the display commands,
|
||||
and reads the cache keys the display still writes.
|
||||
"""
|
||||
path = server_socket_path(environ)
|
||||
if path is None:
|
||||
logger.debug("Control socket disabled or unsupported here; using the file mailbox only")
|
||||
logger.debug("Control socket disabled or unsupported here; the web interface "
|
||||
"cannot send this display commands")
|
||||
return None
|
||||
server = ControlServer(path, status_provider, resolve_socket_group(cache_dir),
|
||||
state_hub=state_hub, handlers=handlers)
|
||||
|
||||
@@ -1114,16 +1114,16 @@ class BasePlugin(ABC):
|
||||
Returns:
|
||||
The request id once the display has queued it, or None when
|
||||
there is no display in this process to ask (the web interface,
|
||||
scripts/check_plugin.py) or its queue is full. A plugin that
|
||||
also runs on cores without this method writes the
|
||||
``display_on_demand_request`` mailbox on None, as before; see
|
||||
"On-demand display" in docs/PLUGIN_API_REFERENCE.md.
|
||||
scripts/check_plugin.py) or its queue is full. Nothing else can
|
||||
take the request then: the ``display_on_demand_request`` file
|
||||
mailbox older cores read is gone, and a write to it is dropped
|
||||
with a warning. See "On-demand display" in
|
||||
docs/PLUGIN_API_REFERENCE.md.
|
||||
|
||||
Example::
|
||||
|
||||
if not (hasattr(self, 'request_on_demand')
|
||||
and self.request_on_demand(mode='my_alert', duration=15)):
|
||||
self._write_on_demand_mailbox(...) # older cores
|
||||
if self.request_on_demand(mode='my_alert', duration=15) is None:
|
||||
self.logger.info("No display to show the alert on")
|
||||
"""
|
||||
request = getattr(getattr(self, 'plugin_manager', None), 'request_on_demand', None)
|
||||
if not callable(request):
|
||||
|
||||
@@ -1869,7 +1869,7 @@ class PluginManager:
|
||||
"""Route plugins' on-demand requests to ``handler`` (None: nowhere).
|
||||
|
||||
The display controller sets its ``submit_plugin_on_demand`` here
|
||||
before any plugin loads. The handler takes a mailbox-shaped request
|
||||
before any plugin loads. The handler takes an on-demand request dict
|
||||
from any thread, queues it for the render thread and returns True,
|
||||
or False when it could not. A plugin manager with no handler (the
|
||||
web interface's, a test's, scripts/check_plugin.py's) has no screen
|
||||
|
||||
Reference in New Issue
Block a user