"""The display's plugin runtime snapshot, shared with the web interface. Only the display process runs plugins, so only it knows which ones it has loaded, where each is in its lifecycle (``plugin_state.PluginStateManager``), why one failed and which version it is running. It publishes that to the shared cache directory -- the channel, and the file permissions, that the error snapshot, plugin health and ``display_current_state`` already use -- and the web interface reads it back for ``/api/v3/plugins/installed``, ``/api/v3/plugins/state`` and state reconciliation. PLUGIN_RUNTIME_KEY written by the display service only Writes. The cache lives on disk, usually the SD card, so the snapshot is written when something a reader would see changes, at most once every ``MIN_INTERVAL`` seconds, and otherwise once every ``REFRESH_INTERVAL`` seconds as a heartbeat. An ordinary plugin update is not a change: the RUNNING state it passes through is published as ENABLED (``plugin_state.published_state``). A display with nothing changing writes this one small file once a minute. Staleness. Every snapshot carries ``published_at`` (wall clock) and ``stale_after``. A reader treats a snapshot older than that as unknown, not as the truth: a display that died without cleaning up leaves its last snapshot behind. A display that stops cleanly publishes ``running: false`` on the way out, so readers see "stopped" at once rather than after the stale window. Nothing on the reading side reports a runtime fact from a snapshot that is not live. """ import math import os import threading import time from dataclasses import dataclass, field from typing import Any, Callable, Dict, Optional from src.logging_config import get_logger from src.redaction import redact_credentials logger = get_logger(__name__) PLUGIN_RUNTIME_KEY = "plugin_runtime_snapshot" SNAPSHOT_SCHEMA = 1 #: Shortest gap, in seconds, between two change-driven writes. Startup loads #: every plugin in a burst, and a plugin failing each update cycle changes its #: error info each time; either is written at most this often. MIN_INTERVAL = 10.0 #: An unchanged snapshot is rewritten this often so readers can tell a quiet #: display from a dead one. REFRESH_INTERVAL = 60.0 #: How often the publisher thread looks for changes: an in-memory comparison. TICK_INTERVAL = 5.0 #: A snapshot older than this is stale: three missed refreshes. STALE_AFTER = 3 * REFRESH_INTERVAL #: Bounds on a published ``stale_after``, so a corrupt value can make a #: reader neither trust a dead display for hours nor distrust a live one. _STALE_AFTER_MIN = 30.0 _STALE_AFTER_MAX = 3600.0 _ERROR_MESSAGE_CHARS = 200 _ERROR_TYPE_CHARS = 80 _ID_CHARS = 100 _VERSION_CHARS = 40 #: Reader statuses. Only LIVE carries runtime facts. LIVE = "live" STALE = "stale" STOPPED = "stopped" UNKNOWN = "unknown" def _clip(value: Any, limit: int) -> str: text = value if isinstance(value, str) else str(value) return text if len(text) <= limit else text[:limit - 3] + "..." def _epoch(value: Any) -> Optional[float]: """Seconds since the epoch for a float or a datetime; None otherwise.""" if isinstance(value, bool): return None if isinstance(value, (int, float)): number = float(value) return number if math.isfinite(number) else None timestamp = getattr(value, "timestamp", None) if callable(timestamp): try: number = float(timestamp()) except (TypeError, ValueError, OverflowError, OSError): return None return number if math.isfinite(number) else None return None def summarize_error(error_info: Optional[Dict[str, Any]]) -> Optional[Dict[str, Any]]: """A short, redacted summary of the state machine's error info. ``message`` is redacted before it is clipped: clipping first could cut a ``token=`` marker off and keep the secret after it. No stack trace: the full error, with its trace, is in the error snapshot (/api/v3/errors). """ if not isinstance(error_info, dict): return None message = error_info.get("error") error_type = error_info.get("error_type") return { "type": _clip(error_type, _ERROR_TYPE_CHARS) if error_type else None, "message": _clip(redact_credentials(message if isinstance(message, str) else str(message or "")), _ERROR_MESSAGE_CHARS), "at": _epoch(error_info.get("timestamp")), "recoverable": bool(error_info.get("recoverable", False)), } def build_runtime_snapshot(state_manager: Any, *, started_at: float, now: Optional[float] = None, running: bool = True) -> Dict[str, Any]: """The snapshot for ``state_manager`` (a plugin_state.PluginStateManager). A stopped snapshot (``running=False``) lists no plugins: nothing is loaded once the display has gone. """ plugins: Dict[str, Dict[str, Any]] = {} if running: for plugin_id, record in state_manager.runtime_records().items(): version = record.get("version") plugins[_clip(plugin_id, _ID_CHARS)] = { "loaded": bool(record.get("loaded")), "state": record.get("state"), "error": summarize_error(record.get("error_info")), "version": _clip(version, _VERSION_CHARS) if version else None, "loaded_at": _epoch(record.get("loaded_at")), } return { "schema": SNAPSHOT_SCHEMA, "running": running, "published_at": time.time() if now is None else now, "started_at": started_at, "refresh_interval": REFRESH_INTERVAL, "stale_after": STALE_AFTER, "pid": os.getpid(), "plugins": plugins, } class PluginRuntimePublisher: """Publishes the display's plugin state machine to the shared cache. Runs in the display service only. tick() is the whole job; start() calls it from a daemon thread every TICK_INTERVAL seconds. Nothing here raises: a failed write is logged at debug and retried on a later tick, at the throttled rate. """ def __init__(self, cache_manager: Any, state_manager: Any, min_interval: float = MIN_INTERVAL, refresh_interval: float = REFRESH_INTERVAL, clock: Callable[[], float] = time.monotonic, wall_clock: Callable[[], float] = time.time) -> None: self.cache_manager = cache_manager self.state_manager = state_manager self.min_interval = min_interval self.refresh_interval = refresh_interval self._clock = clock self._wall_clock = wall_clock self.started_at = wall_clock() # None forces a first publish, which replaces whatever a previous run # of the service left behind. self._published_change: Optional[int] = None self._last_attempt: Optional[float] = None self._tick_lock = threading.Lock() self._stop = threading.Event() self._thread: Optional[threading.Thread] = None def _write(self, running: bool) -> None: snapshot = build_runtime_snapshot(self.state_manager, started_at=self.started_at, now=self._wall_clock(), running=running) self.cache_manager.set(PLUGIN_RUNTIME_KEY, snapshot) def tick(self) -> bool: """Publish if something changed (throttled) or the refresh is due. True if a snapshot was written.""" with self._tick_lock: try: change = self.state_manager.change_count now = self._clock() since = None if self._last_attempt is None else now - self._last_attempt if since is not None: if change == self._published_change: if since < self.refresh_interval: return False elif since < self.min_interval: return False # Stamp the attempt before writing: a cache that keeps failing # is retried at the throttled rate, not on every tick. self._last_attempt = now self._write(running=True) self._published_change = change return True except Exception as err: # never let reporting break the display logger.debug("Could not publish the plugin runtime snapshot: %s", err, exc_info=True) return False def start(self, interval: float = TICK_INTERVAL) -> None: """Tick from a daemon thread until stop(). A no-op while running.""" if self._thread is not None and self._thread.is_alive(): return self._stop.clear() def run() -> None: self.tick() while not self._stop.wait(interval): self.tick() self._thread = threading.Thread(target=run, name="plugin-runtime-publisher", daemon=True) self._thread.start() def stop(self, publish_stopped: bool = True) -> None: """Stop ticking and, by default, publish ``running: false`` so readers see the display as stopped now rather than after the stale window.""" self._stop.set() if self._thread is not None: self._thread.join(timeout=2) self._thread = None if publish_stopped: with self._tick_lock: try: self._write(running=False) except Exception as err: logger.debug("Could not publish the stopped plugin runtime snapshot: %s", err, exc_info=True) def start_plugin_runtime_publisher(cache_manager: Any, state_manager: Any) -> Optional[PluginRuntimePublisher]: """Start publishing the display's plugin runtime state. Display service only: whichever process calls it becomes the source readers trust. Never raises.""" try: publisher = PluginRuntimePublisher(cache_manager, state_manager) publisher.start() return publisher except Exception as err: logger.warning("Plugin runtime reporting to the web interface is unavailable: %s", err) return None # --- Reading side (web interface) ------------------------------------------- #: What a reader reports for a plugin when it does not know. _UNKNOWN_PLUGIN: Dict[str, Any] = { "loaded": None, "state": None, "error_info": None, "loaded_version": None, "loaded_at": None, } #: A plugin a live snapshot does not list: the display has not loaded it #: (never enabled, or unloaded since), which is what its state machine #: reports for an id it has no record of. _NOT_LOADED_PLUGIN: Dict[str, Any] = { "loaded": False, "state": "unloaded", "error_info": None, "loaded_version": None, "loaded_at": None, } @dataclass(frozen=True) class PluginRuntimeView: """What a reader may say about the display's plugins right now. ``status``: ``live`` (a fresh snapshot from a running display), ``stale`` (the last snapshot is older than its ``stale_after``: the display is hung or died without cleaning up), ``stopped`` (the display said so on its way out) or ``unknown`` (no readable snapshot). Only a live view reports per-plugin facts; every other status answers None for them, so a caller cannot pass stale truth on by accident. """ status: str published_at: Optional[float] = None age_seconds: Optional[float] = None stale_after: float = STALE_AFTER plugins: Dict[str, Dict[str, Any]] = field(default_factory=dict) @property def live(self) -> bool: return self.status == LIVE def plugin(self, plugin_id: str) -> Dict[str, Any]: """``loaded``, ``state``, ``error_info``, ``loaded_version`` and ``loaded_at`` for one plugin; all None unless the view is live.""" if not self.live: return dict(_UNKNOWN_PLUGIN) record = self.plugins.get(plugin_id) if not isinstance(record, dict): return dict(_NOT_LOADED_PLUGIN) error = record.get("error") return { "loaded": bool(record.get("loaded")), "state": record.get("state") if isinstance(record.get("state"), str) else None, "error_info": dict(error) if isinstance(error, dict) else None, "loaded_version": record.get("version"), "loaded_at": record.get("loaded_at"), } def describe(self) -> Dict[str, Any]: """The view's own status, for a response to carry beside the facts.""" return { "status": self.status, "published_at": self.published_at, "age_seconds": None if self.age_seconds is None else round(self.age_seconds, 1), "stale_after": self.stale_after, } def _stale_after_of(snapshot: Dict[str, Any]) -> float: value = snapshot.get("stale_after") if isinstance(value, bool) or not isinstance(value, (int, float)): return STALE_AFTER number = float(value) if not math.isfinite(number): return STALE_AFTER return min(max(number, _STALE_AFTER_MIN), _STALE_AFTER_MAX) def view_from_snapshot(snapshot: Any, now: Optional[float] = None) -> PluginRuntimeView: """Judge a snapshot read from the cache; never raises.""" if not isinstance(snapshot, dict) or snapshot.get("schema") != SNAPSHOT_SCHEMA: return PluginRuntimeView(status=UNKNOWN) published_at = _epoch(snapshot.get("published_at")) if published_at is None: return PluginRuntimeView(status=UNKNOWN) stale_after = _stale_after_of(snapshot) age = (time.time() if now is None else now) - published_at if snapshot.get("running") is not True: return PluginRuntimeView(status=STOPPED, published_at=published_at, age_seconds=max(age, 0.0), stale_after=stale_after) # A snapshot from the future is trusted a little: the Pi has no RTC and # its clock steps when NTP syncs. Far in the future, it cannot be dated. if age > stale_after or age < -stale_after: return PluginRuntimeView(status=STALE, published_at=published_at, age_seconds=age, stale_after=stale_after) plugins = snapshot.get("plugins") return PluginRuntimeView( status=LIVE, published_at=published_at, age_seconds=max(age, 0.0), stale_after=stale_after, plugins={k: v for k, v in plugins.items() if isinstance(v, dict)} if isinstance(plugins, dict) else {}, ) def read_plugin_runtime(cache_manager: Any, now: Optional[float] = None) -> PluginRuntimeView: """The display's latest snapshot, judged for staleness. Never raises; a missing cache manager or an unreadable snapshot is ``unknown``. memory_ttl=0: the key is written by the other process, so only the file is current. """ if cache_manager is None: return PluginRuntimeView(status=UNKNOWN) try: snapshot = cache_manager.get(PLUGIN_RUNTIME_KEY, max_age=None, memory_ttl=0) except Exception as err: logger.debug("Could not read the plugin runtime snapshot: %s", err, exc_info=True) return PluginRuntimeView(status=UNKNOWN) return view_from_snapshot(snapshot, now=now)