mirror of
https://github.com/ChuckBuilds/LEDMatrix.git
synced 2026-10-04 14:25:08 +00:00
refactor(plugins): the display publishes plugin runtime state; retire plugin_state.json (#690)
Stage 2 of the web plugin catalog, after #688. - The display publishes a plugin runtime snapshot (plugin_runtime.py) to the shared cache: per plugin loaded, lifecycle state, a short redacted error summary, the version it loaded and when, plus published_at / stale_after / running. Written on change (throttled to 10 s; the RUNNING/ENABLED flip of an ordinary update is not a change) and once a minute otherwise; cleanup() publishes running: false. - The web reads it back and restores loaded / state / error_info in /api/v3/plugins/installed (plus loaded_version, loaded_at and data.runtime). Only a live snapshot counts; stale, stopped or missing answers null and says which. - data/plugin_state.json is retired: every reader and writer moved to config + disk (desired) or the snapshot (observed). Nothing in it was non-derivable, so nothing is migrated and an existing file is left unread. The web-side PluginStateManager (state_manager.py) is removed; the display's plugin_state.PluginStateManager is the only state machine. - StateReconciliation compares config + disk with the snapshot, reporting enabled-but-not-loaded and older-version-loaded as no_action findings. - Backups list installed manifests with enabled from config.json. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,377 @@
|
||||
"""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)
|
||||
Reference in New Issue
Block a user