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:
Chuck
2026-10-05 09:32:22 -04:00
co-authored by Claude Opus 5.5
parent bb475a79ea
commit c461f4efdb
33 changed files with 1307 additions and 1756 deletions
+54 -175
View File
@@ -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
@@ -487,33 +487,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
@@ -582,8 +577,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.
@@ -601,11 +596,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
@@ -617,39 +607,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.
@@ -667,23 +627,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
@@ -754,16 +711,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]:
@@ -776,36 +730,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,
@@ -818,11 +742,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():
@@ -831,32 +755,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",
@@ -865,19 +774,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
@@ -899,61 +804,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))