feat(ipc): the control socket carries every web command; mailboxes are a fallback (stage 4) (#765)

- client: ControlError.sent says whether the display had the request;
  should_fall_back() allows a mailbox write only when it did not, or when
  the display is too old to know the command (upgrade case)
- on-demand start/stop: a display that had the request and failed it is
  answered 503 (400 for invalid_args), no mailbox copy
- errors.clear: new socket command, answered on the connection thread by
  a handler the display registers; applied and republished before the
  answer; plugin_error_clear_request only on fallback
- display: on-demand mailbox looked at once a second while the socket is
  up (0.25 s without), read only when its file changed (one stat via
  CacheManager.file_signature / MailboxWatch); socket commands no longer
  touch the mailbox; a processed duplicate is consumed; writers logged once
- error publisher: mailbox read only when changed; snapshot carries
  applied_clear_cutoff so an older mailbox request is not shown pending
- docs and CHANGELOG (mailboxes kept for one release)

Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
Chuck
2026-10-04 22:44:17 -04:00
committed by GitHub
co-authored by Claude Opus 5.5
parent 26cae3e5d6
commit 3866aa4519
18 changed files with 1498 additions and 149 deletions
+57 -2
View File
@@ -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.
+101 -23
View File
@@ -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'))
+125 -31
View File
@@ -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))
+71 -10
View File
@@ -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)
+44 -4
View File
@@ -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?
+42 -5
View File
@@ -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